spacr.runctx¶
The run context: one id, one seed, one error policy, for a whole run.
Three things every spaCR pipeline needs and none of them had:
One run id, on every log line and every output. A run used to be
untraceable. The log said a field failed; the registry said an artifact
existed; nothing connected them, so “show me everything from the run that
produced this file” had no answer. run_context() mints one id, stamps
it onto every logging.LogRecord created anywhere in the process,
writes a per-run JSONL log that read_run_log() reads back, and hands
the same id to spacr.artifacts.register_run_outputs(). A log line
and an output can therefore be joined on run_id, which is the whole
point.
One seed that reaches everything. random_seed used to be read by
spacr.deep_spacr and spacr.sim and nowhere else, so a
“reproducible” run still shuffled its fields differently, split its folds
differently and initialised Cellpose differently every time.
seed_everything() seeds Python, NumPy (legacy global and the
Generator stream spacr_rng() hands out), Torch
on CPU and every CUDA device, and — through those two — Cellpose. sklearn
has no global seed at all; random_state() is what estimator
construction sites pass. What cannot be made deterministic is listed in
SeedReport.caveats and in seed_everything()’s docstring
rather than papered over.
One error policy, honoured at every batch boundary. on_error is a
tri-state, default "stop":
stopthe first failed unit aborts the run. The default, because a pipeline that quietly drops a third of its plates produces a number that looks exactly like a good one.
skipthe unit is recorded — on the
RunLedgerand as aSkipRecordnaming the unit, the stage and why — and the run carries on. Never a silent drop:ErrorPolicy.skipsis what was lost.retrythe unit is attempted
ErrorPolicy.attemptstimes with an exponential backoff, and if the budget runs out it behaves exactly likestop.
Usage¶
from spacr.runctx import run_context
with run_context("mask", settings) as run:
for plate in plates:
for attempt in run.policy.attempts_for(plate, stage="plate"):
with attempt:
process(plate)
run.register_outputs(roots=plates)
Public API¶
run_context,RunContext,current_run_context,current_run_idThe run itself, and the ambient lookup a worker or a library call uses.
new_run_id,install_run_id_logging,uninstall_run_id_logging,RunIdFilter,runs_log_dir,run_log_path,read_run_log,run_resource_path,read_run_resourcesThe S7 machinery: minting, stamping and querying by run id.
seed_everything,SeedReport,resolve_seed,random_state,spacr_rng,torch_generator,seed_worker,DEFAULT_SEEDThe S5 machinery: one call that seeds them all, plus the per-library handles for the places a global seed cannot reach.
ErrorPolicy,resolve_error_policy,SkipRecord,SKIPPED,ON_ERROR_STOP,ON_ERROR_SKIP,ON_ERROR_RETRY,ON_ERROR_MODESThe S9 machinery.
apply_defaults,RUN_SETTING_KEYSThe settings seam: the keys this module owns, applied to any settings dict.
Classes¶
What a run does when one unit of work fails: stop, skip, or retry. |
|
One pipeline run: its id, its seed, its error policy, its ledger. |
|
Give every record a |
|
What |
|
One unit of work that |
Functions¶
|
Fill this module's keys into |
|
Return the |
|
Return the active run id, or |
|
Stamp |
|
Mint a fresh run id. |
|
Return the seed to hand an estimator's |
|
Return every log line this run emitted — "show me everything from run X". |
|
Read one run's process-tree resource document. |
|
Build the |
|
Return the seed a run should use. |
|
Open a run: mint an id, seed the world, arm the error policy. |
|
Return the JSONL log path for |
|
Return the persisted process-tree resource path for |
|
Return the folder holding per-run logs, creating it when needed. |
|
Seed every RNG spaCR can reach, and report what that does not cover. |
|
Seed one DataLoader worker. Pass as |
|
Return an independent |
|
Return a seeded |
|
Restore the record factory that was installed before us. |
Module Contents¶
- class spacr.runctx.ErrorPolicy(mode: str = DEFAULT_ON_ERROR, *, attempts: int = DEFAULT_RETRIES, backoff: float = DEFAULT_BACKOFF, ledger: spacr.errors.RunLedger | None = None, logger: logging.Logger | None = None, run_id: str = '', record: bool = True, sleep: Callable[[float], None] | None = None)[source]¶
What a run does when one unit of work fails: stop, skip, or retry.
Applied at a batch boundary — per plate, per well, per field — wherever the pipeline has a unit it can name. Every mode records the failure on the
RunLedger, so the artifact stamp tells the truth in all three cases; the mode only decides whether the run survives it.- Parameters:
mode –
ON_ERROR_STOP(default),ON_ERROR_SKIPorON_ERROR_RETRY.attempts – total tries per unit in
retrymode, including the first. Bounded on purpose: an unbounded retry against a dead NAS is an infinite loop with a progress bar.backoff – seconds before the second attempt. Doubled before each subsequent one, capped at
MAX_BACKOFF.ledger – the ledger to record on. One is created when omitted.
logger – where to log; defaults to
spacr.runctx.run_id – stamped onto every
SkipRecord.record – write successes and failures to
ledger. Pass False at a boundary whose call site already records them — the Measure pool does, through its own job and error callbacks — so the ledger counts each field once rather than twice.sleep – the sleep function, injectable so a test can assert the backoff schedule without waiting for it.
- Raises:
ValueError – on an unknown mode, or a non-positive attempt count — a “retry” that never retries is a silent stop.
Validate and store the error policy’s mode and retry controls.
- attempts_for(unit: Any, stage: str | None = None) Iterator[_Attempt][source]¶
Yield attempts at one unit, honouring the mode. The core of S9.
Drive it with an inner
with:for attempt in policy.attempts_for(plate, stage='plate'): with attempt: process(plate)
In
stopandskipmode exactly one attempt is yielded; inretrymode up toattempts, with a sleep in between. The generator raises the last exception after the final attempt unless the mode isskip, so theforstatement is where astoprun aborts.A body that omits the inner
withis a bug this cannot detect — the failure would propagate straight out of thefor, which isstopbehaviour whatever the mode says.- Parameters:
unit – the thing being processed; stringified for the record.
stage – pipeline stage, recorded on the ledger and the skip.
- Raises:
Exception – the unit’s own exception, in
stopmode and inretrymode once the budget is spent.
- bind(ledger: spacr.errors.RunLedger | None = None, record: bool | None = None) ErrorPolicy[source]¶
Point this policy at another ledger, and return it.
A run with several ledgers — Measure keeps one per source folder — needs the failures recorded against the right one, while the skip list stays with the run. Mutates and returns
selfrather than copying, soskipsremains the single account of what the run did not cover.- Parameters:
ledger – the ledger to record on from now on.
record – whether to record at all; see the constructor.
- run(unit: Any, fn: Callable[..., Any], *args: Any, stage: str | None = None, **kwargs: Any) Any[source]¶
Call
fn(*args, **kwargs)for one unit under this policy.The ergonomic form of
attempts_for(), for a boundary whose body is already a callable.- Parameters:
unit – the thing being processed.
fn – the work.
stage – pipeline stage.
- Returns:
whatever
fnreturned, orSKIPPEDwhen the unit was skipped.- Raises:
Exception – as
attempts_for().
- property retries: List[Tuple[str, int]][source]¶
(unit, attempts_made)for every unit that had to be retried.
- property skipped_units: List[str][source]¶
Just the names — the answer to “what did this run not cover?”.
- property skips: List[SkipRecord][source]¶
Every unit
skipdropped, in the order it dropped them.
- class spacr.runctx.RunContext[source]¶
One pipeline run: its id, its seed, its error policy, its ledger.
Created by
run_context(); read anywhere bycurrent_run_context(). The id on this object is the id on every log line the run emits (read_run_log()), on everyRunLedgerit hands out, and on everyspacr.artifacts.Artifactit registers — which is what lets a log line and an output be joined.- Parameters:
run_id – the id.
module – the producing module key —
"mask","measure", the same keysspacr.portsuses.seed – the seed applied, or None when the run is unseeded.
policy – the
ErrorPolicyin force.ledger – the run’s
RunLedger, whoserun_idhas been set to this run’s.settings – the settings the run was started with.
seed_report – what
seed_everything()managed to seed.started_utc – when the run opened.
log_path – the run’s JSONL log, or
""when logging is off.resource_log_path – path to the run’s process-tree resource JSON document, or
""when resource accounting is off or unavailable. Set when sampling starts and used to register the document.resource_artifact_id – artifact-registry identifier assigned after successful resource-document registration, or
""when no record was created.
- adopt(ledger: spacr.errors.RunLedger) spacr.errors.RunLedger[source]¶
Re-stamp an existing ledger with this run’s id, and return it.
- Parameters:
ledger – existing run ledger to associate with this context.
For a call site that already builds its own ledger and should not have to change how.
- new_ledger(name: str) spacr.errors.RunLedger[source]¶
Return a
RunLedgerstamped with this run id.A ledger mints its own uuid, which would put a second id on the run’s
run_statusrows and break the join to the artifact registry. This overwrites it, so every stamp the run leaves — ledger row, log line, artifact — carries one id.- Parameters:
name – the ledger name, i.e. the pipeline stage.
- random_state(default: int | None = None) int | None[source]¶
This run’s seed, for an estimator’s
random_state=.
- register_outputs(module: str | None = None, settings: Mapping[str, Any] | None = None, **kwargs: Any) Tuple[Any, ...][source]¶
Register this run’s outputs, stamped with this run id.
Thin wrapper over
spacr.artifacts.register_run_outputs()that suppliesrun_idand defaultsstrictto False, so a registry that cannot be written costs one printed line and never the run.- Parameters:
module – override the module key.
settings – override the settings hashed into the artifacts.
kwargs – passed through, e.g.
roots=[...],status=.
- Returns:
the registered artifacts.
- register_worker(worker_kind: Any, worker_id: Any = None, *, pid: int | None = None, create_time: float | None = None) str[source]¶
Give one sampled child process its run-specific identity.
The sampler can always report a PID and process name. This method is the seam a parameter sweep or sequencing parent uses to add the fact that the PID is, for example,
trial 17orFASTQ saver.worker_kindmay also be the complete stamp returned byspacr.fit_resources._worker_stamp(); that is how a spawned worker reports its own creation time without a PID-reuse race.- Parameters:
worker_kind – worker category, or a complete stamp mapping.
worker_id – identity within that category.
pid – process id; defaults to the calling process.
create_time – psutil process creation time, resolved when omitted.
- Returns:
the sampler’s stable process identity, or
""when resource accounting is off or unavailable.
- rng(stream: str = '') numpy.random.Generator[source]¶
An independent seeded Generator; see
spacr_rng().
- property log: logging.Logger[source]¶
The run’s logger. Every record it makes carries
run_id.
- property skips: List[SkipRecord][source]¶
Units the policy skipped — shorthand for
policy.skips.
- class spacr.runctx.RunIdFilter(run_id: str | None = None)[source]¶
Bases:
logging.FilterGive every record a
run_idattribute, so a formatter can use it.Belt and braces for the record factory installed by
install_run_id_logging(): a handler whose format string contains%(run_id)smust never raise on a record that came from somewhere the factory did not reach (a record unpickled from a worker, say).- Parameters:
run_id – stamp this id instead of the ambient one. Used by the per-run log so a nested run’s records are not misattributed.
Initialise the filter with an optional fixed run identifier.
- filter(record: logging.LogRecord) bool[source]¶
Set
record.run_idwhen it is missing. Never drops a record.- Parameters:
record – log record to stamp with a run id.
- class spacr.runctx.SeedReport[source]¶
What
seed_everything()actually managed to seed.- Parameters:
seed – the seed applied.
seeded – library handles that were seeded, e.g.
"python","numpy","torch","torch.cuda".unavailable – handles that could not be seeded because the library is not installed. Not an error — a headless analysis box without Torch is a supported install.
caveats – the honest part. Every place this seed does not buy determinism, in plain sentences, so a caller can quote them rather than assume a guarantee that does not exist.
deterministic – whether the deterministic-kernel switches were requested as well.
- class spacr.runctx.SkipRecord[source]¶
One unit of work that
on_error='skip'dropped, and why.The whole point of
skipover a bareexcept: pass: what was lost is named, counted and persisted, so a run that covered 97 of 100 plates cannot be read as one that covered 100.- Parameters:
unit – the unit skipped — a plate folder, a well, a field file.
stage – the pipeline stage it was skipped at.
reason – why, in a sentence.
exc_type – the exception class name.
message –
str(exc).attempts – how many times it was tried before being given up on.
run_id – the run that skipped it.
utc – when.
traceback_str – the full traceback, kept for the ledger.
- spacr.runctx.apply_defaults(settings: Dict[str, Any] | None = None) Dict[str, Any][source]¶
Fill this module’s keys into
settings, in place when given a dict.For a caller that wants an explicit, complete settings dict — a batch runner writing a settings CSV, or a test. The pipelines do not need it:
resolve_seed()andresolve_error_policy()default every key, so a settings dict that names none of them still getson_error='stop'andDEFAULT_SEED.- Parameters:
settings – the dict to fill; a new one is made when None.
- Returns:
the same dict, with
RUN_SETTING_KEYSpresent.
- spacr.runctx.current_run_context() RunContext | None[source]¶
Return the
RunContextof the innermost active run, or None.Context-local, so two runs on two threads do not see each other’s.
- spacr.runctx.current_run_id() str[source]¶
Return the active run id, or
""when no run is open.Falls back to
RUN_ID_ENV, which is how aspawn-ed Measure worker — a fresh interpreter with an empty context variable — still logs under the run that started it.
- spacr.runctx.install_run_id_logging() None[source]¶
Stamp
run_idonto everylogging.LogRecordin this process.Done with
logging.setLogRecordFactory()rather than a filter on the root logger, because a filter attached to a logger is consulted only for records logged through that logger — records propagating up fromspacr.measurenever see it, which is every record that matters. The factory sees them all, whichever handler is attached and whenever it was attached.Idempotent, and chains onto whatever factory is already installed instead of replacing it.
- spacr.runctx.new_run_id() str[source]¶
Mint a fresh run id.
Twelve hex characters — deliberately the same shape as
spacr.errors.RunLedger.run_id, so ledger stamps,run_statusrows andspacr.artifacts.Artifactrows all join on one column of one format.
- spacr.runctx.random_state(default: int | None = None) int | None[source]¶
Return the seed to hand an estimator’s
random_state=.sklearn, XGBoost, LightGBM and CatBoost all take a
random_state(orseed) and all ignore the NumPy global stream once one is given, so a construction site that hard-codesrandom_state=42silently overrides the run’s seed. Call this instead:RandomForestClassifier(random_state=random_state(42))
- Parameters:
default – what to return when no run is open and nothing has been seeded.
- Returns:
the active run’s seed, or
default.
- spacr.runctx.read_run_log(run_id: str, *, level: int | str | None = None, logger: str | None = None, contains: str | None = None) List[Dict[str, Any]][source]¶
Return every log line this run emitted — “show me everything from run X”.
The query side of S7. Each record is a dict with
run_id,utc,level,logger,message,file,line,processandthread; a record carrying an exception also hastraceback.- Parameters:
run_id – the run to read.
level – minimum level, as a number or a name (
"WARNING").logger – only records from this logger or its children.
contains – only records whose message contains this substring.
- Returns:
the matching records in the order they were written. An empty list when the run wrote no log — never an exception, because “nothing was logged” is an ordinary answer.
- spacr.runctx.read_run_resources(run_id: str) Dict[str, Any][source]¶
Read one run’s process-tree resource document.
- Parameters:
run_id – the run to read.
- Returns:
the JSON object, or an empty dict when no readable checkpoint exists. Missing accounting is never represented as a zero.
- spacr.runctx.resolve_error_policy(settings: Mapping[str, Any] | None = None, *, ledger: spacr.errors.RunLedger | None = None, logger: logging.Logger | None = None, run_id: str = '', sleep: Callable[[float], None] | None = None, default: str = DEFAULT_ON_ERROR) ErrorPolicy[source]¶
Build the
ErrorPolicya settings dict asks for.Reads
on_error,on_error_attemptsandon_error_backoff. An unseton_errormeansDEFAULT_ON_ERROR.- Parameters:
settings – a settings dict, or None.
ledger – the ledger failures are recorded on.
logger – where the policy logs.
run_id – stamped onto skip records.
sleep – injectable sleep, for tests.
default – mode to use when the settings name none.
- Returns:
an
ErrorPolicy.- Raises:
ValueError – when
on_erroris not one of the three modes. Loud on purpose: a typo likeon_error='continue'silently falling back tostopis how a user believes they asked for tolerance and did not get it.
- spacr.runctx.resolve_seed(settings: Mapping[str, Any] | None = None, default: Any = DEFAULT_SEED) int | None[source]¶
Return the seed a run should use.
Reads
random_seedfromsettings, thenSEED_ENV, thendefault. An explicitrandom_seed=Nonemeans “do not seed” and is honoured — a caller who deliberately wants a free-running RNG (a simulation sweep that must not produce the same draw twice) gets one.- Parameters:
settings – a settings dict, or None.
default – what to use when nothing names a seed. Pass
Nonefor “leave the RNGs alone unless asked”.
- Returns:
the seed, or None for “do not seed”.
- spacr.runctx.run_context(module: str = '', settings: Mapping[str, Any] | None = None, *, run_id: str | None = None, seed: Any = _UNSET, on_error: str | None = None, deterministic: bool | None = None, ledger: spacr.errors.RunLedger | None = None, log: bool = True, sleep: Callable[[float], None] | None = None) Iterator[RunContext][source]¶
Open a run: mint an id, seed the world, arm the error policy.
The one call a pipeline entry point makes. Inside the block,
current_run_id()answers everywhere in this process (and, viaRUN_ID_ENV, in its children), every log record carries the id, andread_run_log()can pull the run’s own log back out afterwards.- Parameters:
module – the producing module key —
"mask","measure", … — used for the artifact registry and the logger name.settings – the run’s settings.
random_seed,on_error,on_error_attemptsandon_error_backoffare read from here.run_id – use this id instead of minting one — for a resumed run, or a distributed worker continuing its parent’s run.
seed – override the seed. An explicit
Nonemeans “do not seed at all”; omit the argument to readsettings.on_error – override the mode.
deterministic – also request deterministic kernels; see
seed_everything(). Defaults to thedeterministicsetting.ledger – use this ledger rather than making one.
log – write the per-run JSONL log. False for a caller that only wants the id and the policy.
sleep – injectable sleep for the retry backoff, for tests.
- Yields:
the
RunContext.
- spacr.runctx.run_log_path(run_id: str) str[source]¶
Return the JSONL log path for
run_id. It need not exist yet.- Parameters:
run_id – run identifier used as the log filename.
- spacr.runctx.run_resource_path(run_id: str) str[source]¶
Return the persisted process-tree resource path for
run_id.The resource document lives beside the ordinary run log but has its own suffix and schema. It is written atomically by
spacr.fit_resources._ResourceSampler, so a reader sees either a complete checkpoint or the preceding complete checkpoint.- Parameters:
run_id – the run whose accounting record is wanted.
- Returns:
an absolute JSON path. The file need not exist yet.
- spacr.runctx.runs_log_dir() str[source]¶
Return the folder holding per-run logs, creating it when needed.
<log dir>/runs, where the log dir honoursSPACR_LOG_DIR— seespacr.logging_util.log_dir(). Pointing that variable at a scratch folder is how a test gets a private set of run logs.
- spacr.runctx.seed_everything(seed: int | None = None, *, deterministic: bool = False, quiet: bool = True) SeedReport[source]¶
Seed every RNG spaCR can reach, and report what that does not cover.
Seeds, in order:
random, NumPy’s legacy global (np.random), theGeneratorstreamspacr_rng()derives, Torch on CPU, and Torch on every visible CUDA device. Cellpose is seeded transitively — it has no seed API and draws from NumPy and Torch.PYTHONHASHSEEDis exported for child processes.What this does not buy you. A seed makes the draws reproducible. It does not make CUDA reductions associative, it does not stop a forked worker pool from sharing one stream, and it cannot re-seed this interpreter’s string hashing. Every such limit is named in
SeedReport.caveatsand inSEED_CAVEATS; do not promise a user more than that list allows. In particular, calling this and then reporting “the run is deterministic” is wrong on a GPU unlessdeterministic=Trueand the model avoids the kernels Torch has no deterministic implementation for.- Parameters:
seed – the seed.
NoneusesDEFAULT_SEED.deterministic – also ask for deterministic kernels —
cudnn.deterministic,cudnn.benchmark=False,torch.use_deterministic_algorithms(warn_only=True)andCUBLAS_WORKSPACE_CONFIG. Slower, sometimes much slower, and still not a guarantee; see the caveats.quiet – when False, log the report at INFO.
- Returns:
a
SeedReport.
- spacr.runctx.seed_worker(worker_id: int) None[source]¶
Seed one DataLoader worker. Pass as
worker_init_fn=seed_worker.A spawned worker inherits none of the parent’s RNG state and a forked one inherits all of it — so every worker augments identically, which is the classic silent bug where a batch of eight “random” crops is eight copies of the same transform. This derives a per-worker stream from Torch’s initial seed, which the DataLoader has already varied per worker and per epoch.
- Parameters:
worker_id – the worker’s index, supplied by the DataLoader.
- spacr.runctx.spacr_rng(stream: str = '', seed: int | None = None) numpy.random.Generator[source]¶
Return an independent
numpy.random.Generatorfor one stream.Derived from the run seed by
SeedSequencespawning rather than by re-seeding fromseed + 1: adjacent seeds produce correlated streams in some bit generators, and two workers drawing correlated “random” subsamples is a bug that looks like data.- Parameters:
stream – a name for this stream — a worker id, a stage, a fold. Different names give independent streams; the same name gives the same stream every run.
seed – override the run seed.
- Returns:
a fresh Generator.
- spacr.runctx.torch_generator(device: str = 'cpu', stream: str = '')[source]¶
Return a seeded
torch.Generatorfor a DataLoader or sampler.- Parameters:
device – the device the generator belongs to.
stream – a stream name, as for
spacr_rng().
- Returns:
a
torch.Generatorseeded from the run seed.- Raises:
RuntimeError – when Torch is not installed. Deliberate: a caller asking for a Torch generator cannot proceed without one, and handing back None would fail further away.
Nested helpers¶
- install_run_id_logging._factory(*args: Any, **kwargs: Any) logging.LogRecord¶
Create a record and stamp the ambient run id onto it.
spacr/runctx.py:284