Skip to main content

CI CI platforms CI Python PyPI DOI Python License: MIT

necroflow

necroflow logo

Python pipeline framework inspired by Snakemake. Define rules, wire them into pipelines, run with automatic parallelism and caching. All in Python. All safe. All readable.

Compact overview of the current software surface, see features.txt.

See COMPARISON.md for a detailed comparison with Snakemake, Nextflow, Luigi, CWL/WDL, and Prefect/Airflow across 20 axes.

Define a pipeline

A command-line run points at a Python workflow. Rules describe typed outputs and shell commands; the workflow wires rule calls into a pipeline.

# pipeline.py
from necroflow import DAG, NodeType, Pipeline, command, symlink_file, output, workflow

class Fastq(NodeType):
    filename = "reads.fastq.gz"

class Bam(NodeType):
    filename = "aligned.bam"

class Counts(NodeType):
    filename = "counts.txt"

@symlink_file
def raw_fastq(path: str):
    fastq = output(Fastq)
    return fastq

@command("bwa mem {ref} {fastq} > {bam}", threads=4)
def align(fastq: Fastq, ref: str):
    bam = output(Bam)
    return bam

@command("featureCounts -a {gene_model} {bam} -o {counts}")
def count(bam: Bam, gene_model: str):
    counts = output(Counts)
    return counts

@workflow
def rna_pipeline(P: Pipeline, config: dict) -> None:
    P.fastq = raw_fastq(path=config["path"])
    P.bam = align(P.fastq, ref=config["ref"])
    P.counts = count(P.bam, gene_model=config["gene_model"])

Request pipeline results

Assignments such as P.counts = ... give Nodes public Pipeline labels. Select which labelled outputs to produce with .requests in the job TOML:

".pipeline" = "pipeline.py:rna_pipeline"
".requests" = ["counts"]

Running necroflow job.toml executes counts and all its ancestors, then copies the requested output to results/job/counts/counts.txt. If .requests is omitted, necroflow requests every labelled sink Node. More on jobs in Job TOML and parameter grids.

Reusable subpipelines

P.subpipeline(prefix) returns a view over the same Pipeline. Assignments through the view are registered on the root with the prefix, while rule identity and DAG deduplication remain unchanged:

from necroflow import workflow

@workflow
def sample_pipeline(P: Pipeline, reference, sample: dict) -> None:
    P.fastq = raw_fastq(path=sample["reads"])
    P.bam = align(P.fastq, ref=reference)
    P.counts = count(P.bam, gene_model=sample["gene_model"])

@workflow
def cohort_pipeline(P: Pipeline, config: dict) -> None:
    for sample in config["samples"]:
        sample_pipeline(
            P.subpipeline(f"samples/{sample['name']}"),
            config["reference"],
            sample,
        )

The resulting labels include samples/A/bam and samples/A/counts. Prefixes are request/result names only and never enter fingerprints. The CLI calls P.finish() after a successful workflow return. Direct Python callers must finish the root before selecting P.sinks(); finishing freezes the root and every subpipeline view.

Core ideas

  • Rules describe how to produce outputs from inputs — shell command templates with typed I/O and lint-clean name = output(NodeType) declarations.
  • Workflows wire rule calls together for a single config using a Pipeline for labels; prefixed subpipeline views make reusable loop-generated outputs requestable.
  • DAG runs many pipelines at once, deduplicating shared upstream work across samples automatically.
  • Paths are derived from a lineage-derived fingerprint of the full input chain — same inputs always produce the same path, different inputs produce different paths. The filesystem is the cache.

Install

cd necroflow
make venv
source .venv/bin/activate

Platform support

necroflow supports POSIX systems (Linux and macOS). We do not offer native Windows support because POSIX commands are the reproducible execution target for workflows. On Windows, use Windows Subsystem for Linux (WSL) to run necroflow in a POSIX environment.

Compose pipeline fragments

A command-line workflow has the signature build_workflow(P, config) -> None. The CLI creates the shared DAG and an open Pipeline, then calls the workflow. @workflow activates its first positional Pipeline argument while the function runs. Rules retrieve that context automatically; their Nodes already have final paths and fingerprints and are interned in the DAG when the calls return.

The decorator restores the previous context on return or exception. Decorated subworkflows can therefore activate P.subpipeline(prefix) temporarily. Ordinary helpers inherit the active context; helpers switching views must be decorated or pass their Pipeline explicitly to rules. Labels still require P.name = node or P[label] = node; returning a Node does not publish it.

Explicit calls such as align(P, reads, ref="hg38") remain supported, including outside workflows. An explicit Pipeline takes precedence for that call without changing the active context. Calls with neither an explicit Pipeline nor an active workflow raise RuntimeError explaining both options.

The first workflow argument must be an open Pipeline: missing/wrong owners raise TypeError, and finished owners raise RuntimeError. Each message identifies the workflow and states the requirement. Coroutine, generator, and async-generator functions raise TypeError at decoration because their bodies defer execution beyond the synchronous context.

For reusable internal fragments, pass an existing pipeline to a helper that adds its named nodes. This lets several fragments contribute to one public workflow without changing the CLI workflow signature:

from necroflow import workflow

@workflow
def add_alignment(P, config):
    P.fastq = raw_fastq(path=config["path"])
    P.bam = align(P.fastq, ref=config["ref"])

@workflow
def rna_pipeline(P, config):
    add_alignment(P, config)
    P.counts = count(P.bam, gene_model=config["gene_model"])

An assembler mutates the supplied pipeline, so its labels must not conflict with labels added by another fragment. Use this form for components that belong to one pipeline. The caller creates a fresh Pipeline(dag) for each independent config. Equivalent upstream calls are canonicalized immediately in the shared DAG; after each workflow, the caller marks its sinks or explicit outputs with dag.require(...).

Attribute and item labels share one namespace. Use P.counts for ordinary Python identifiers and item syntax for generated paths:

for dataset, config in combinations:
    P[f"{dataset}/{config}"] = count(P, inputs[dataset], config=config)

Labels are canonical relative POSIX paths, so P["dataset/config"] creates a nested result at results/<job>/dataset/config/<filename> and is requested with the same string in .requests. Absolute paths, empty or dot-prefixed components, ., .., repeated/trailing separators, and paths exceeding Linux NAME_MAX/PATH_MAX byte limits are rejected at assignment. Labels that collide with Pipeline API attributes such as nodes are item-only.

Run from the CLI

Create a job TOML that references the workflow and carries the concrete parameters for one run.

# job.toml
".pipeline" = "pipeline.py:rna_pipeline" # from pipeline import rna_pipeline

path = "/data/s1.fastq.gz"
ref = "hg38"
gene_model = "gencode_v44"

Run it with the necroflow command:

necroflow job.toml

By default, real cached node outputs go under nodes/, while user-facing results and manifest.toml go under results/; above, simply results/job. Use explicit roots when you want them elsewhere:

necroflow --nodes-dir nodes --results-dir results job.toml

For many runs, use multiple job TOMLs or __grid values inside one job TOML:

".pipeline" = "pipeline.py:rna_pipeline"

path__grid = ["/data/s1.fastq.gz", "/data/s2.fastq.gz"]
ref = "hg38"
gene_model = "gencode_v44"

The same pipeline can also be assembled and executed from Python directly; see Rules and typed outputs and Executor, classification, scheduling, and cleanup. See Command-line interface and Job TOML and parameter grids for the full CLI format.

Where outputs live

DAG("some-dir") writes real lineage-addressed node outputs directly under that directory. The CLI defaults to a split layout: canonical cached outputs under nodes/, plus per-job copies and manifest.toml files under results/. Linux and macOS opportunistically use filesystem copy-on-write cloning, so only requested outputs can require additional physical storage. See Where outputs live and caching for the full layout.

Manual

Start with the canonical workflow in examples/canonical, or copy it with necroflow init my-workflow.

CLI subcommands

The default command form is kept for convenience, but the same run can be written explicitly:

necroflow run job.toml

This executes the requested pipeline and creates cached outputs under nodes/ plus job-facing copies and a manifest under results/job/.

Create a starter workflow from the canonical template:

necroflow init my-workflow

Example output:

created my-workflow

Render the requested DAG without executing commands:

necroflow graph job.toml

Example output, abridged:

DAG  4 nodes  (1 required)

import_text[RawText:raw_text] (path='input.txt')
write_tool_config[ToolConfig:tool_config] (text='{\n  "mode": "uppercase"\n}\n')
process_text[ProcessedText:processed_text]
summarize[Summary:summary] *

List requested output paths without executing commands:

necroflow outputs job.toml

Example output:

[job]
summary	node=nodes/summarize/d18e6af2070f14be/summary.txt	result=results/job/summary/summary.txt

Inspect stored metadata for an existing cached output:

necroflow provenance nodes/summarize/<provenance_hash>/summary.txt

Example output:

path = nodes/summarize/<provenance_hash>/summary.txt
rule = summarize
rule_hash = <64 hex characters>
provenance_hash = <64 hex characters>
[config]
path = 'input.txt'
text = '{\n  "mode": "uppercase"\n}\n'

What is not yet implemented

  • Cluster / cloud backends

Metadata

Release files for necroflow 0.0.8

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for necroflow 0.0.8
File Size Uploaded
necroflow-0.0.8.tar.gz 140.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for necroflow 0.0.8
File Interpreter ABI Platform
necroflow-0.0.8-py3-none-any.whl Python 3 none any Details

Total release size: 219.7 kB

Release files / necroflow-0.0.8.tar.gz

Download URL necroflow-0.0.8.tar.gz
Size 140.2 kB
Tags Source
SHA-256 checksum
How to use checksums
f04caab29b469ce2cde2bf5f4f426fad477373089163a6bd0e271d2c9e55f48f
BLAKE2b-256 checksum
How to use checksums
7dae11e8aaba90620517e4e3cf6f060d40cfbf2beccb6bc79c82651b3ddacc11
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.14.3

Release files / necroflow-0.0.8-py3-none-any.whl

Download URL necroflow-0.0.8-py3-none-any.whl
Size 79.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
d101eab0bc269493e772d0f73a5d29c431598b8123e4971bac5e69871bc730c1
BLAKE2b-256 checksum
How to use checksums
504e52dbf53e74345824ef8a08b3b69a1ec3ed973ac6fd6db84ad1c77edcada0
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.14.3

Release history Release notifications | RSS feed

This release

0.0.8 This release

2 release files

0.0.7

2 release files

0.0.6

2 release files

0.0.5

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

2 release files

0.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page