Skip to main content

AIND Code Ocean pipeline utils

CI PyPI - Version semantic-release: angular License ruff uv Copier

Utilities for AIND Code Ocean capsules. They assemble a pipeline's processing.json, lay a capsule out as launcher, workers, and aggregator, and guard long runs against termination, flaky I/O, stale caches, logs lost under progress bars, and context lost across thread pools.

Installation

The core is stdlib-only; two extras add dependencies.

pip install aind-code-ocean-pipeline-utils
pip install "aind-code-ocean-pipeline-utils[rich]"      # log: rich handler, progress bars
pip install "aind-code-ocean-pipeline-utils[metadata]"  # metadata: aind-data-schema

A capsule, end to end

from pathlib import Path

from aind_code_ocean_pipeline_utils import (
    check_shutdown,
    forward_metadata,
    log_data_tree,
    package_version,
    record_step,
    shutdown_handler,
)
from aind_code_ocean_pipeline_utils.log import install_rich_handler
from aind_code_ocean_pipeline_utils.metadata import write_assembled_processing

install_rich_handler()
log_data_tree(Path("/data"))  # what Code Ocean mounted
with shutdown_handler():  # SIGTERM exits at the next check_shutdown()
    with record_step("sort", process_type="Spike sorting", version=package_version("my-package")):
        for shard in shards:
            check_shutdown()
            process(shard)
forward_metadata("/data", "/results")  # subject.json, procedures.json, ...
write_assembled_processing("/data", "/results")  # in the pipeline's final capsule only

What's in the package

To Use
write a pipeline's processing.json metadata.write_assembled_processing
record this capsule's step record_step
carry provenance across a fan-out step.fanout_shards()
copy or author the other metadata files forward_metadata, metadata.make_derived_data_description
stamp a commit or package version capsule_commit, package_version
fan work out and merge the results write_stream_configs, find_stream_config, merge_manifests
read a boolean App Panel parameter parse_truthy
stop cleanly on SIGTERM shutdown_handler, check_shutdown
retry flaky I/O, write files atomically retry_on_oserror, atomic_json_write
tell whether cached output is stale input_fingerprint
keep contextvars in thread-pool workers submit_with_context
log under progress bars, keep a log file log.install_rich_handler, log.build_progress, attach_file_log
see the mounts and memory use log_data_tree, start_memory_reporter

Everything imports from aind_code_ocean_pipeline_utils except the log. and metadata. names, which live in submodules so the package root needs neither extra.

Recording provenance and processing.json

Every AIND asset should carry a processing.json: the steps that produced it, and a dependency_graph saying which fed which. No capsule sees the whole pipeline, but each sees its inputs, which are its parents. So each capsule that opts in records its own step, the records travel with the data, and the final capsule assembles them. A capsule that never opts in still appears if it writes its own processing.json.

The final capsule writes processing.json

from aind_code_ocean_pipeline_utils.metadata import write_assembled_processing

write_assembled_processing("/data", "/results")  # -> /results/processing.json

It reads the records and upstream processing.json files under /data, plus this capsule's own records in /results, so call it after this capsule's record_step block. A step it cannot recover in full, such as a parent that left no record or a step recorded under a schema version the installed aind-data-schema rejects, becomes a placeholder whose notes say why. The call logs a failure and returns None rather than raising; assemble_processing returns the Processing without writing it.

Opted-in capsules record their own step

with record_step(
    "mri-registration",  # unique node id in the pipeline
    process_type="Image atlas alignment",  # a ProcessName value; any other becomes "Other"
    version=package_version("my-package"),
) as step:
    step.parameters = {"mask_dilate": 4}  # values known only at runtime
    ...

On a clean exit, record_step writes /results/provenance/mri-registration.json beside copies of every upstream record; if the block raises, it writes nothing. Parents are inferred from what arrives in /data, so a capsule never states its position in the DAG. code_url and commit_hash come from the /code checkout, but version must be passed. Without the [metadata] extra, the record keeps the step's place in the graph and loses its details.

Two wiring rules hold. Every pipeline edge must carry /results/provenance/, and a fan-out worker's node id must include its unit, as in f"sort-{probe}".

A launcher hands its record to fan-out workers

A Flatten fan-out gives each worker only its own stream_<name>/ directory, so the launcher writes its records into each one:

with record_step("discover", process_type="Other", run_experimenters=["Jane Doe"]) as step:
    write_stream_configs(
        items,
        results_dir=Path("/results"),
        schema_marker=MARKER,
        provenance=step.fanout_shards(),
    )

Each worker then infers discover as its parent. At assembly, run_experimenters fills every step of the run that names none, and a pipeline={"name": ..., "code": {"url": ...}} block fills Processing.pipelines.

The other metadata files are forwarded or authored

forward_metadata("/data", "/results") copies subject.json, procedures.json, instrument.json, and acquisition.json verbatim. The derived asset's data_description.json describes a new asset, so it is authored instead, with metadata.make_derived_data_description or from the fields read_data_description_fields extracts from any schema version.

Commit and version stamps

capsule_commit() reads CO_COMMIT, GIT_COMMIT, or COMMIT_ID, then falls back to git -C /code rev-parse HEAD. package_version(name) wraps importlib.metadata.version. Both return None rather than raise, so a manifest can be stamped unconditionally.

Structuring a pipeline capsule

Launcher, workers, aggregator

Most parallel AIND capsules take three roles. The launcher writes one config.json per item under /results/stream_<name>/, Code Ocean's Flatten hands each directory to its own worker, and the aggregator merges the workers' manifests. Role names these three, plus MONOLITH for a capsule that runs all of them in one process.

MARKER = "_mycapsule_stream_config"

# Launcher
write_stream_configs(items, results_dir=Path("/results"), schema_marker=MARKER)

# Worker: exactly one staged config anywhere under /data
cfg_path, cfg = find_stream_config(Path("/data"), schema_marker=MARKER)

# Aggregator
workers = find_worker_manifests(Path("/data"))
launcher = find_launcher_manifest(Path("/data"))
merged = merge_manifests(m for _, m in workers)  # {"built": [...], "skipped": [...]}

A worker finds its config by the marker key in the JSON, not by path, because Flatten and Target Map Path nest inputs unpredictably. find_stream_config raises StreamConfigError, listing the candidate paths, on zero or several matches.

App Panel parameters

With named_parameters: true, the App Panel passes every parameter as a string, so a boolean flag does not survive. parse_truthy(args.flag) reads one back: true, yes, y, t (any case) and non-zero numbers are True; everything else, including "0", "false", and "", is False.

Surviving long runs

Graceful shutdown

with shutdown_handler():
    for shard in shards:
        check_shutdown()  # raises GracefulExit after SIGINT/SIGTERM
        process(shard)

A signal only sets a flag, so work stops at the next check_shutdown() rather than mid-write, and the process exits with 128 + signum (143 for SIGTERM). GracefulExit inherits from BaseException, so except Exception: cannot swallow it. A second signal exits at once through os._exit.

Retries and atomic writes

download = retry_on_oserror(_raw_download, retries=5)
atomic_json_write(out_path, download(url))

retry_on_oserror retries only TRANSIENT_ERRNOS (EIO, EAGAIN, EBUSY, network errnos); pass transient_errnos= to widen it. A permanent error such as ENOENT raises at once instead of hiding a config mistake behind minutes of backoff. atomic_json_write and atomic_write_text write a temp file, fsync it, and os.replace it over the destination, so no reader sees a half-written file.

Input fingerprints

input_fingerprint({"window": 0.01, "channels": [0, 1, 2]}) returns "sha256:…", equal for equal inputs in any key order. Store it beside cached output and compare on resume. A non-JSON value raises TypeError naming its key path, so coerce Path or ndarray at the call site.

Context across thread pools

ThreadPoolExecutor.submit runs its function in an empty context, so a ContextVar setting such as scipy.fft.set_workers silently lapses in the worker. submit_with_context(pool, fn, *args) copies the caller's context for each submit; one shared copy would raise RuntimeError once two threads ran it.

Seeing what happened

Logging under progress bars

A rich Progress repaints several times a second, painting over any log line that bypassed rich, often the tail of a traceback. install_rich_handler ([rich] extra) routes logging through rich and returns its Console; pass that console to every Progress.

from aind_code_ocean_pipeline_utils.log import build_progress, install_rich_handler, make_progress_callback

install_rich_handler()
with build_progress(len(items)) as (progress, overall, item):  # uses the installed console
    for it in items:
        progress.reset(item, total=it.size, description=it.name, visible=True)
        do_work(it, on_progress=make_progress_callback(progress, item))
        progress.advance(overall)

attach_file_log(Path("/results/run.log")) adds a file handler beside the console one, so the log survives the run in /results. It needs no extra.

Mounts and memory

log_data_tree(Path("/data"))
reporter = start_memory_reporter()  # logs RSS against the cgroup limit every 15 s
...
reporter.stop()

log_data_tree follows Code Ocean's symlinked mounts to a bounded depth. An OOM kill runs no except block and flushes no output, so the reporter's last line is what tells an OOM apart from a spot reclamation afterwards.

Development

# Set up environment
uv sync

# Run the full check suite
./scripts/run_linters_and_checks.sh -c

# Individual tools
uv run pytest
uv run ruff format
uv run ruff check
uv run mypy

See CLAUDE.md for the design invariants each module is required to preserve.

License

MIT; see LICENSE.

Metadata

Release files for aind-code-ocean-pipeline-utils 0.7.1

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

Source distribution (sdist)

Source distribution for aind-code-ocean-pipeline-utils 0.7.1
File Size Uploaded
aind_code_ocean_pipeline_utils-0.7.1.tar.gz 184.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for aind-code-ocean-pipeline-utils 0.7.1
File Interpreter ABI Platform
aind_code_ocean_pipeline_utils-0.7.1-py3-none-any.whl Python 3 none any Details

Total release size: 241.0 kB

Release files / aind_code_ocean_pipeline_utils-0.7.1.tar.gz

Download URL aind_code_ocean_pipeline_utils-0.7.1.tar.gz
Size 184.6 kB
Tags Source
SHA-256 checksum
How to use checksums
342f2126280f0f0e7cfcddeff9f12eb855d17ad56443556dc4c8b40b74599131
BLAKE2b-256 checksum
How to use checksums
f3b2cf81576f3023853474f419c622dfddc2e695ab7a62b9ea960a77909ee578
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.12.20 {"installer":{"name":"uv","version":"0.12.20","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release files / aind_code_ocean_pipeline_utils-0.7.1-py3-none-any.whl

Download URL aind_code_ocean_pipeline_utils-0.7.1-py3-none-any.whl
Size 56.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
02839f0034c0da9ca39fcd40cc3a457eb2cf3d61cdce5d78fc4341a1551eef08
BLAKE2b-256 checksum
How to use checksums
8fe6a468a8f20f20325b54d4633fcec05b9cb6fc7e00a493c041868c1831b743
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.12.20 {"installer":{"name":"uv","version":"0.12.20","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release history Release notifications | RSS feed

This release

0.7.1 This release

2 release files

0.7.0

2 release files

0.6.1

2 release files

0.6.0

2 release files

0.5.0

2 release files

0.4.2

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