spacr.batch¶
Workflow inputs and outputs¶
Batch Runner¶
Validate and execute module jobs in dependency order. Each job writes its own module outputs and log.
Open: the application’s Help/tools menus.
Inputs and outputs below include conditional alternatives. The guidance and handoff notes say which route applies.
Inputs
Run queue — Saved module/plate/settings job definitions and dependency order.
Outputs
Run history and artifacts — Project run records, settings, output paths, artifact provenance, status and logs.
spacr.batch — stack arbitrary module+settings jobs into one queue and
run them unattended.
The Plate Queue (spacr.qt.plate_queue) chains plates through one
pipeline with one settings dict.
This module is the other axis. It holds a queue of arbitrary
(module, settings) jobs in any order:
Mask → Measure → Classify (CV) → Classify (ML), then the same four again
with a different diameter, then a fifth plate’s Mask. That is what a night of
plate-scale work actually looks like.
A queued job is a spacr-run invocation. Nothing here re-implements
module dispatch or settings loading: spacr.cli owns the module
registry, the settings loader, the --set coercion and the exit-code
contract (0 ok, 1 the module raised, 2 bad arguments or settings), and this
module drives it. Likewise a queue is a ledger of ledgers: the queue’s own
verdict is a spacr.errors.RunLedger, stamped next to the queue file,
and each job’s spacr.errors.read_run_status() stamp is read back so a
job that exited 0 having silently skipped 40 fields is reported as partial
rather than as a success.
What makes a queue survive a night nobody is watching:
Every job is validated when it is added, not when it runs. Discovering at 3 a.m. that job 9’s
srcis misspelled wastes the whole night.validate_queue()reports all problems at once andrun_queue()refuses to start while any of them is an error.State is persisted after every transition, atomically. A machine that reboots mid-queue is resumed with
resume_queue(), not restarted. The file is written to a temp name andos.replaced, because a queue file truncated by a crash is worse than no queue file at all.A job whose dependency failed is SKIPPED, not run. Measure after a failed Mask produces a database that looks like a real result.
skipped(deliberately not run),failedandnot_run(the queue halted before reaching it) are three different things and stay that way.continue-on-error stops hiding a systematic failure. If the first three jobs died the same way the remaining nine will too;
max_consecutive_failureshalts the queue and says so.Per-job logs go to their own file — one interleaved log from twelve overnight jobs is unreadable — and the path is on the job record.
Concurrency: jobs run strictly one at a time. They compete for one GPU,
and two Cellpose jobs sharing a card is how an overnight run turns into an
overnight CUDA OOM. There is deliberately no max_workers; run two queues
in two processes if you really have two cards.
Import cost: this module imports only the standard library plus
spacr.cli, spacr.errors and spacr.validate, all of which
are torch-free. A queue file can be read, validated and planned on a login
node without a GPU stack — the pipeline itself is only imported inside the
subprocess that runs a job.
Typical use:
from spacr.batch import Job, Queue, run_queue, save_queue
q = Queue(name='overnight')
q.add(Job(module='mask', settings='/data/p1/settings/mask.csv'))
q.add(Job(module='measure', settings='/data/p1/settings/measure.csv',
depends_on=['mask-1']))
save_queue(q, '/data/overnight.queue.json')
result = run_queue(q, path='/data/overnight.queue.json')
print(result.summary())
Exceptions¶
The queue itself is wrong — a bad job, a cycle, a refused start. |
Classes¶
One |
|
One thing wrong with a job or with the queue. |
|
One incremental progress report, handed to |
|
An ordered list of |
|
What happened, and the summary that is the actual deliverable. |
Functions¶
|
Turn an exit code plus a log into |
|
Render a duration the way an overnight summary should read. |
|
Render every problem at once, errors first. |
|
Run one job in this interpreter, with its output tee'd to its log. |
|
Return the exact |
|
Read a queue file written by |
|
Describe what the queue would do, without doing any of it. |
|
Build the settings dict this job's module will actually receive. |
|
Pick a queue up where a crash, a reboot or a Stop left it. |
|
Run the queue, one job at a time, and report on all of it. |
|
Write |
|
Run one job as its own |
|
Check one job the way |
|
Validate every job, and the queue's own structure, all at once. |
Module Contents¶
- exception spacr.batch.QueueError[source]¶
Bases:
spacr.errors.SpacrErrorThe queue itself is wrong — a bad job, a cycle, a refused start.
A subclass of
spacr.errors.SpacrErrorsoexcept SpacrErrorcatches it alongside everything else spaCR raises deliberately.Initialize self. See help(type(self)) for accurate signature.
- class spacr.batch.Job[source]¶
One
spacr-runinvocation, with its place in the queue’s history.- Parameters:
module – module key or alias understood by
spacr.cli.resolve_module()—'mask','measure','ml_analyze', …settings – path to a settings CSV/JSON (what
spacr-run --settingstakes), or an inline settings dict for a job that has no file of its own.id – unique, human-typable identifier. Left empty,
Queue.add()assigns one like'mask-1'.label – what to show a human; defaults to
module @ src.overrides –
key=valuestrings applied on top ofsettings, exactly likespacr-run --set. This is how “the same four jobs again with a different diameter” is written without copying four settings files. A mapping is accepted and normalised.depends_on – ids of jobs that must succeed first. A job whose dependency failed is skipped, never run.
status – one of
ALL_STATUSES.started – ISO-8601 UTC start time, or
''.finished – ISO-8601 UTC end time, or
''.exit_code – the process exit code — 0 ok, 1 the module raised, 2 bad settings (
spacr.cli’s contract).error – one-line explanation of a failure or a skip.
log_path – this job’s own log file.
run_status – the job’s
spacr.errors.read_run_status()verdict, summarised —Nonewhen the job stamped nothing.
- copy(**changes: Any) Job[source]¶
Return a fresh, never-run copy of this job with
changesapplied.The GUI’s “Duplicate” button: the common way to build a queue is one job, then eleven variations of it.
- classmethod from_dict(data: Mapping[str, Any]) Job[source]¶
Rebuild a job from
to_dict(), tolerating a hand-edited file.Every field has a default and unknown keys are ignored, so a user can delete the bookkeeping (
status,started, …) from a queue file and still have it load — which is the whole point of a hand-editable format.- Parameters:
data – one entry of the queue file’s
jobslist.- Raises:
QueueError – when
datais not a mapping with amodule.
- to_dict() Dict[str, Any][source]¶
Return the job as a JSON-serialisable dict, in a stable key order.
- property duration_s: float | None[source]¶
Wall-clock seconds the job took, or None if it never finished.
- property elapsed_s: float | None[source]¶
Seconds the job has taken, counting up while it is still running.
- Returns:
duration_sonce it has finished, the time since it started while it is running, else None.
- property is_partial: bool[source]¶
True when the job exited 0 but its own ledger says items failed.
This is the case the queue exists to catch: a measure run that processed 344 of 384 wells exits 0 and looks like a success.
- class spacr.batch.Problem[source]¶
One thing wrong with a job or with the queue.
Deliberately the same shape as
spacr.validate.Problem— severity, message, fix — plus thejob_idneeded to say which of twelve jobs is at fault. Problems produced byspacr.validate.validate_settings()are wrapped into this, not re-invented.- Parameters:
job_id – id of the offending job;
''for a queue-level problem.severity –
'error'(the queue must not start) or'warning'.message – what is wrong, in the user’s terms.
fix – what to actually do about it.
setting – the settings key at fault, when there is one.
- class spacr.batch.Progress[source]¶
One incremental progress report, handed to
on_progress.A queue runs for hours; a GUI that only learns the outcome at the end is not showing progress. Every transition emits one of these.
- Parameters:
event –
'queue_started','job_started','job_finished','job_skipped','queue_stopped'or'queue_finished'.job_id – the job this is about, or
''for queue-level events.index – 1-based position of the job in the queue,
0when N/A.total – number of jobs in the queue.
status – the job’s status at the moment of the event.
message – one line fit to put in a status bar.
- class spacr.batch.Queue[source]¶
An ordered list of
Jobs, run one at a time, top to bottom.- Parameters:
jobs – the jobs, in the order they will run.
created – ISO-8601 UTC creation time.
name – what to call this queue in the summary and log folder.
- add(job: Job, validate: bool = True) Job[source]¶
Append
job, validating it now rather than at 3 a.m.- Parameters:
job – the job to add; its
idandlabelare filled in when empty.validate – set False only when deliberately building an invalid queue (loading a hand-edited file, for instance, which reports its problems through
validate_queue()instead of raising).
- Returns:
the job, now owned by this queue.
- Raises:
QueueError – when the job cannot run — an unknown or GUI-only module, an unreadable settings file, a bad override, a duplicate id, or a dependency that is not already in the queue.
- find(job_id: str) Job | None[source]¶
Return the job with
job_id, or None.- Parameters:
job_id – unique queue job identifier to find.
- classmethod from_dict(data: Mapping[str, Any]) Queue[source]¶
Rebuild a queue from
to_dict()or from a hand-written file.- Parameters:
data – serialized or hand-written queue mapping to rebuild.
- Raises:
QueueError – when the document is not a queue, is a format from the future, or holds a job entry that cannot be read.
- index(job_id: str) int[source]¶
Position of
job_idin run order, or-1.- Parameters:
job_id – unique queue job identifier to locate.
- mint_id(module: str) str[source]¶
Return an unused, human-typable id for a job of
module.- Parameters:
module – module name used as the identifier’s readable stem.
- move(job_id: str, offset: int) int[source]¶
Move
job_idoffsetplaces (negative is earlier).- Parameters:
job_id – unique identifier of the job to reposition.
offset – relative number of queue positions to move.
- Returns:
the job’s new index, or
-1when it is not in the queue.
- remove(job_id: str) bool[source]¶
Remove
job_idand drop it from every other job’sdepends_on.- Parameters:
job_id – unique identifier of the job to remove.
Leaving a dangling dependency behind would silently skip the jobs that referred to it, so the reference is cleaned up here.
- Returns:
True when a job was removed.
- class spacr.batch.QueueResult[source]¶
What happened, and the summary that is the actual deliverable.
- Parameters:
queue – the queue, with every job’s final status on it.
log_dir – folder holding the per-job logs.
started – ISO-8601 UTC start of the whole queue.
finished – ISO-8601 UTC end of the whole queue.
stopped_reason – why the queue halted early, or
''.ledger – the queue’s own
spacr.errors.RunLedger— one recorded item per job, which is what groups identical failures and what gets stamped next to the queue file.path – the queue file that was kept up to date, or
''.
- jobs_with(status: str) List[Job][source]¶
Every job that ended in
status.- Parameters:
status – job status value to select.
- summary() str[source]¶
Render the end-of-queue report.
What ran, what failed and why (identical failures grouped), what was skipped and because of which upstream job, what is only partial, and how long each took. This is the thing a user reads over coffee instead of scrolling four thousand lines of interleaved log.
- spacr.batch.classify_failure(exit_code: int, log_path: str) Tuple[str, str][source]¶
Turn an exit code plus a log into
(kind, one-line message).The kind is what groups identical failures: twelve jobs that all died on the same missing share are one problem worth one line, not twelve.
- Parameters:
exit_code – the runner’s exit code.
log_path – the job’s log file.
- Returns:
(kind, message).
- spacr.batch.fmt_duration(seconds: float | None) str[source]¶
Render a duration the way an overnight summary should read.
- Parameters:
seconds – elapsed seconds, or None when the job never finished.
- Returns:
'—','42.1s','7m 12s'or'7h 41m'.
- spacr.batch.format_problems(problems: Sequence[Problem], title: str = 'queue check') str[source]¶
Render every problem at once, errors first.
Reporting the first error and stopping is what makes a twelve-job queue take twelve rounds of fixing; this prints all of them.
- Parameters:
problems – what
validate_queue()returned.title – heading for the block.
- Returns:
the report as one string, no trailing newline.
- spacr.batch.inprocess_runner(job: Job, settings_path: str, log_path: str) int[source]¶
Run one job in this interpreter, with its output tee’d to its log.
- Parameters:
job – queue job whose module and overrides are executed.
settings_path – resolved settings file supplied to the job command.
log_path – destination file for redirected output.
Same argv, same exit-code contract as
subprocess_runner()— it callsspacr.cli.main()directly — but with none of the isolation: a segfault here takes the queue with it. Use it only where spawning is not possible (a frozen single-file build), and never from the GUI.- Returns:
the exit code
spacr.cli.main()returned.
- spacr.batch.job_command(job: Job, settings_path: str, python: str | None = None) List[str][source]¶
Return the exact
spacr-runcommand line forjob.Spelled as
<this python> -m spacr.cli ...rather than thespacr-runconsole script so the job runs in the interpreter the queue is running in — the venv the user actually installed spaCR into — even whenPATHsays otherwise.- Parameters:
job – the job.
settings_path – settings file the job will be given.
python – interpreter to use;
sys.executableby default.
- Returns:
the argv list.
- spacr.batch.load_queue(path: str | os.PathLike) Queue[source]¶
Read a queue file written by
save_queue(), or hand-written.- Parameters:
path – the queue file.
- Returns:
the
Queue.- Raises:
QueueError – when the file is missing or is not a queue. Never a traceback: an unattended runner should fail with a sentence.
- spacr.batch.plan(queue: Queue, detail: bool = False) str[source]¶
Describe what the queue would do, without doing any of it.
- Parameters:
queue – the queue.
detail – also render
spacr.validate.describe_plan()for every job — the full “here is what would actually happen” per job. Off by default because it lists directories, which is slow over NFS.
- Returns:
the plan as one string.
- spacr.batch.resolve_job_settings(job: Job) Dict[str, Any][source]¶
Build the settings dict this job’s module will actually receive.
Layered exactly as
spacr-runlayers them — module defaults, then the settings file (or the inline dict) read under today’s setting names, then the--setoverrides — usingspacr.cli’s own loader, migration and coercion so a value the CLI accepts is a value the queue accepts.- Parameters:
job – the job.
- Returns:
the resolved settings.
- Raises:
SettingsError – on an unknown module, an unreadable settings file, an unknown override key or an uncoercible override value.
- spacr.batch.resume_queue(path: str | os.PathLike, retry_failed: bool = False, **kwargs: Any) QueueResult[source]¶
Pick a queue up where a crash, a reboot or a Stop left it.
Jobs that already succeeded are left alone. Jobs that were
not_run(the queue halted before reaching them) or stillpendingare run. A job that wasrunningwhen the machine went down is reset and run again. For Mask, Measure, and Format Converter jobs the rerun receivesresume=Trueautomatically, so their verified field checkpoints are reused instead of repeating the whole plate. Failed and skipped jobs stay as they are unlessretry_failedis set.- Parameters:
path – the queue file
run_queue()was persisting to.retry_failed – also re-run jobs that failed, and un-skip the jobs that were skipped because of them.
kwargs – forwarded to
run_queue().
- Returns:
the
QueueResultfor the resumed run.
- spacr.batch.run_queue(queue: Queue, path: str | os.PathLike | None = None, on_error: str = 'continue', log_dir: str | os.PathLike | None = None, max_consecutive_failures: int | None = 3, on_progress: Callable[[Progress], None] | None = None, runner: Callable[[Job, str, str], int] | None = None, stop_flag: Callable[[], bool] | None = None, force: bool = False, echo: bool = True) QueueResult[source]¶
Run the queue, one job at a time, and report on all of it.
Jobs run sequentially and only sequentially — they compete for one GPU. The queue is validated in full before the first job starts, its state is written to
pathafter every transition, and a job whose dependency did not succeed is skipped rather than run against a half-written input.- Parameters:
queue – the queue to run. Mutated in place: each job’s status, timestamps, exit code, log path and run_status are filled in.
path – queue file kept up to date after every transition. Without it the run cannot be resumed, which is a real loss on a long night.
on_error –
'continue'(default) runs the jobs that do not depend on the failed one;'stop'halts the queue immediately.log_dir – folder for per-job logs; defaults to
<queue file>_logsor a temp folder.max_consecutive_failures – halt after this many failures in a row. Continue-on-error is meant to save a night from one bad plate, not to spend it repeating one systematic mistake twelve times.
Noneor0disables the check.on_progress – called with a
Progresson every transition, so a GUI can show the queue moving. Exceptions from it are logged and swallowed.runner –
(job, settings_path, log_path) -> exit_code;subprocess_runner()by default. Injected by the tests, and by anyone who wants to submit jobs to a scheduler instead.stop_flag – polled between jobs; return True to stop the queue after the running job finishes (the GUI’s Stop button).
force – run even though validation found errors. The escape hatch for a check that is wrong about your data — it is not the default for a reason.
echo – print the summary at the end. Always logged either way.
- Returns:
the
QueueResult.- Raises:
QueueError – when validation found errors and
forceis False — with every problem in the message, not just the first.
- spacr.batch.save_queue(queue: Queue, path: str | os.PathLike) pathlib.Path[source]¶
Write
queuetopathatomically, as indented JSON.JSON rather than a bespoke format because it is stdlib (no PyYAML on a compute node), because
spacr.errorsalready persists its stamps this way, and because indented JSON is genuinely hand-editable: a user fixing job 9’ssrcat 3 a.m. opens this file in an editor, not a GUI.- Parameters:
queue – the queue to persist.
path – destination file.
- Returns:
the written path.
- spacr.batch.subprocess_runner(job: Job, settings_path: str, log_path: str) int[source]¶
Run one job as its own
spacr-runprocess. The default runner.A separate process is the point: cellpose segfaulting or the CUDA driver wedging kills that job, not the other eleven and not the GUI. The exit code comes straight from
spacr.cli— 0 ok, 1 the module raised, 2 bad settings — and stdout+stderr go to this job’s own log file.- Parameters:
job – the job to run.
settings_path – settings file to pass to
--settings.log_path – file to write this job’s output to.
- Returns:
the process exit code.
- spacr.batch.validate_job(job: Job, queue: Queue | None = None) List[Problem][source]¶
Check one job the way
spacr-run --dry-runwould, plus queue rules.Every check that can be made without running anything is made here, so it is made when the job is added:
the module resolves, and is not one of the GUI-only apps (
spacr.cli.INTERACTIVE_ONLY) that has no headless callable;the settings file exists, parses, and the overrides name real settings with coercible values;
spacr.validate.validate_settings()agrees the settings are runnable against the data on disk — with data that an upstream job in this queue has not produced yet deferred to a warning (see_deferrable());dependencies exist and come earlier in the queue.
- Parameters:
job – the job to check.
queue – the queue it belongs to, needed for the dependency rules.
- Returns:
every problem found, errors and warnings mixed. Never raises.
- spacr.batch.validate_queue(queue: Queue) List[Problem][source]¶
Validate every job, and the queue’s own structure, all at once.
Reporting the first problem and stopping is how a twelve-job queue takes twelve rounds of fixing, so this returns everything: duplicate ids, cycles, unknown dependencies, and each job’s own settings problems.
- Parameters:
queue – the queue to check.
- Returns:
every problem found.
[p for p in problems if p.is_error]being empty is exactly the conditionrun_queue()requires.
Nested helpers¶
- _cycle_problems.walk(job_id: str, trail: List[str]) None¶
Depth-first search one captured dependency subgraph for cycles.
- Parameters:
job_id – queue job whose known dependencies should be visited.
trail – active ancestor IDs used to reconstruct a back-edge cycle.
- Returns:
None. Captured state records active and completed jobs, and a back edge appends one problem; missing dependency IDs are ignored here because ordinary dependency validation reports those typos.
spacr/batch.py:964