spacr.qt.bridge

Run spaCR pipelines in Qt workers and relay their progress to the GUI.

PipelineWorker executes a pipeline call in a QThread and routes captured standard output and errors through line_ready. The process-wide registry tracks active jobs for shared surfaces such as Home, while make_thread centralizes worker registration and lifecycle setup.

Stopping and pausing are cooperative. Every worker receives a spacr.cancellation.CancellationToken, which long-running pipelines poll at safe boundaries. PauseGate and checkpoint() provide the pause protocol, but pipeline entries that do not call checkpoint report supports_pause=False and must expose Pause as unavailable.

Classes

PauseGate

A latch a worker thread waits on, so a pause is a pause.

PipelineWorker

Runs one pipeline function in its own thread.

RunHandle

One in-flight job: what it is, how far along, and its pause gate.

RunRegistry

Every job started through make_thread(), while it runs.

Functions

apply_worker_budget(→ int)

Cap pool settings in-place and return this run's worker allocation.

available_worker_count(→ int)

Workers a newly-started run may use, never fewer than one.

checkpoint(→ None)

Pipeline-side cancellation/pause point. Call only where stopping is safe.

current_gate(→ Optional[PauseGate])

The PauseGate for the calling thread, or None.

drain_thread(→ bool)

Ask thread to stop, wait for it, and never terminate it.

emit_safely(→ bool)

Emit a Qt signal unless its C++ receiver has been destroyed.

make_thread(→ tuple[PySide6.QtCore.QThread, ...)

Return (thread, worker) — the caller connects the worker's signals

parked_thread_count(→ int)

How many stubborn threads are currently parked. For diagnostics.

pausable(→ Callable)

Mark an entry point as honouring checkpoint().

prune_job_pairs(→ List[tuple])

Return the (thread, worker) pairs a screen must still hold.

prune_parked_threads(→ int)

Release parked (thread, worker) pairs whose thread has exited.

registry(→ RunRegistry)

The process-wide RunRegistry (created on first use).

resolve_pipeline_entry(→ Callable[[Dict[str, Any]], ...)

Return the pipeline function that runs a given app, or None if the

thread_has_stopped(→ bool)

True when thread is finished, never started, or already deleted.

wait_for_parked_threads(→ int)

Wait for parked Qt threads within a shared time limit.

worker_capacity(→ int)

Logical CPU capacity used by the cooperative run allocator.

Module Contents

class spacr.qt.bridge.PauseGate[source]

A latch a worker thread waits on, so a pause is a pause.

Why this is not connected to every pipeline. Pausing a running job can only mean one of two things:

  • stop the thread wherever it happens to be — which, for spaCR, means possibly mid-np.save of a mask (spacr.object writes .npy without a tmp+rename) or between the several INSERTs that make up one measured field. Both leave exactly the half-written artefact spacr.resume exists to clean up. That is not a pause, it is corruption with a friendly label.

  • let the pipeline reach a boundary where nothing is half-written, and hold it there. That requires the pipeline to ask, which is what checkpoint() is for.

Only the second is honest, and it cannot be bolted on from outside: a gate can be set from the GUI thread, but if nothing ever reads it the button does nothing. Pause therefore remains disabled until an entry point opts in via pausable(). Stop is a separate contract: spacr.cancellation is polled by shipped workflows at durable boundaries and never claims to suspend an in-progress write.

Thread-safety: pause() / resume() are called from the GUI thread; wait_if_paused() blocks the worker thread. That is the whole contract — a threading.Event does the rest.

Create the gate open, with nothing paused.

is_paused() → bool[source]

True once pause() has been called and not yet released.

pause() → None[source]

Ask the worker to stop at its next checkpoint.

paused_for() → float[source]

Seconds spent waiting, or 0.0 when not paused.

resume() → None[source]

Release a paused worker.

wait_if_paused(timeout: float | None = None) → bool[source]

Block while paused. Returns True once running again.

Called by pipeline code through checkpoint(); safe to call when not paused, where it returns immediately.

class spacr.qt.bridge.PipelineWorker(fn: Callable[..., Any], settings: Dict[str, Any], worker_count: int = 1, app_key: str = '', journal: bool = True, capture_figures: bool = True)[source]

Bases: PySide6.QtCore.QObject

Runs one pipeline function in its own thread.

Signals:

  • line_ready(str) — a chunk of stdout/stderr text

  • finished(bool) — True if the function returned without an unhandled exception

  • error(str) — traceback string on failure

  • figure_ready(object) — a matplotlib Figure that the pipeline asked to show(); emitted from the worker thread so the UI slot can attach it.

Parameters:
  • fn – the pipeline function to run. It is called on the worker thread, so anything it constructs belongs to that thread.

  • settings – the settings dict handed to fn.

  • worker_count – how many workers the run is allowed, passed through to the pipeline rather than used here.

  • app_key – which module is running, for the journal and for routing output back to the right screen.

  • journal – whether to record the run in the journal.

  • capture_figures – whether a figure the pipeline shows is captured and emitted through figure_ready. False leaves matplotlib alone, which is what a run wants when nobody is watching it.

Prepare to run fn(settings) in a worker thread.

Parameters:
  • fn – pipeline entry point (see resolve_pipeline_entry()).

  • settings – keyword-style dict passed as the sole argument.

  • worker_count – worker allocation reserved in the run registry.

  • app_key – optional explicit module name for run-history records.

  • journal – create a reproducibility manifest. Set false only for read-only background UI maintenance such as refreshing history.

  • capture_figures – intercept Matplotlib output for an analysis run. Read-only UI jobs disable this so polling cannot load the plotting stack merely by opening a screen.

request_cancel(reason: str = 'cancelled by the user') → bool[source]

Request a stop at the pipeline’s next declared safe boundary.

run() → None[source]

Invoked by QThread.started; runs the pipeline function to completion.

property app_key: str[source]

App key this job belongs to, or "" for an ad-hoc job.

property blocks_shutdown: bool[source]

Whether this job still running is a reason not to close the app.

True for analysis runs. They write masks, measurements and model files, so stopping one anywhere other than a declared safe boundary risks the half-written artefact spacr.resume exists to clean up, and the honest response to “close now” is to wait.

False for the read-only background jobs the UI runs on its own behalf — refreshing run history, scanning a model folder, polling a job queue. journal=False already means exactly that (see __init__()): nothing is produced, so nothing can be left half-produced, and a caller that lets one veto shutdown converts a list refresh into an application that will not quit.

property supports_pause: bool[source]

True only when the entry point declares pausable().

No shipped pipeline does, so this is False everywhere today — see PauseGate for why that is the honest answer rather than a missing feature.

class spacr.qt.bridge.RunHandle(app_key: str, worker: PipelineWorker, thread: PySide6.QtCore.QThread, parent=None)[source]

Bases: PySide6.QtCore.QObject

One in-flight job: what it is, how far along, and its pause gate.

Lives on the GUI thread. progress is scraped from the worker’s stdout (see _PROGRESS_RE) rather than reported by the pipeline, because the pipelines have no reporting channel.

Parameters:
  • app_key – which module is running. Falls back to "job" so a handle always has a name to show.

  • worker – the PipelineWorker doing the work. Its worker_count is read once here rather than on every update.

  • thread – the thread the worker was moved to. Held so the handle can wait on it, not so it can be restarted.

  • parent – parent object.

Wrap one running job for the run registry.

blocks_shutdown and user_visible are read from the worker at construction rather than on demand: retire() drops the worker reference, and the answers are still needed after that.

Parameters:
  • app_key – what the job is called.

  • worker – the pipeline worker doing the work.

  • thread – the thread it runs on.

  • parent – parent object, or None.

elapsed() → float[source]

How long this run has been going, in seconds.

Floored at zero: a clock adjusted backwards mid-run would otherwise report a negative age, which reads as a bug in the run rather than in the clock.

Returns:

the elapsed seconds.

fraction() → float | None[source]

Completed fraction in 0..1, or None when unknown.

is_running() → bool[source]

Whether this handle still owns a live QThread.

request_cancel(reason: str = 'cancelled by the user') → None[source]

Cooperatively cancel this job and ask its thread to retire.

retire() → None[source]

The job’s thread has stopped — drop out of the registry.

Wired to QThread.finished rather than worker.finished: the latter is emitted inside run(), before the thread’s event loop has actually stopped, and releasing the last reference at that moment is how this module segfaulted once already.

property gate: PauseGate[source]

The pause gate this run checks between steps.

Returns:

the gate.

property supports_pause: bool[source]

Whether pausing this job would actually pause it.

class spacr.qt.bridge.RunRegistry(parent=None)[source]

Bases: PySide6.QtCore.QObject

Every job started through make_thread(), while it runs.

Deliberately tiny: a list plus a changed signal. It exists so that a surface which is not the screen that started the job — the Home page — can show what spaCR is doing, without teaching every screen to report in.

Parameters:

parent – parent widget.

Create the empty run registry.

Parameters:

parent – parent object, or None.

active() → List[RunHandle][source]

Handles for the jobs running right now, oldest first.

cancel_all(timeout_ms: int = 5000, reason: str = 'application shutdown') → List[RunHandle][source]

Cancel every active job and wait up to one shared deadline.

The return value contains workers that are still live and whose survival is a reason not to close — see RunHandle.blocks_shutdown. They are kept registered and strongly referenced; callers must refuse to destroy the GUI rather than use QThread.terminate() mid-write.

Read-only UI housekeeping (journal=False — a history refresh, a model scan) is cancelled and waited for exactly like everything else, but it never appears in the return value. It writes nothing, so there is no half-written artefact to protect, and a caller that treats it as a veto turns “the run-history list is still loading” into “the application cannot be closed”. MainWindow.closeEvent answers a non-empty return by refusing to close and showing a modal: on a headless run — CI, a test shard — nothing can ever dismiss that modal, and the process spins in its nested event loop with no forward progress. A housekeeping job that outlived the screen which started it used to be enough to trigger it.

A handle whose stop_on_quit is false (an install that must not be killed half way) is left running and is not waited for.

Parameters:
  • timeout_ms – total wait budget across all active threads.

  • reason – cancellation reason recorded by each worker.

Returns:

handles that block shutdown and did not stop in the budget.

clear() → None[source]

Drop every handle. For tests — never call this on a live app.

is_busy() → bool[source]

Whether any run is still registered.

What the window asks before quitting.

Returns:

True while a run is tracked.

register(handle: RunHandle) → RunHandle[source]

Take ownership of a run and start reporting it.

Parameters:

handle – the run to track.

Returns:

the same handle, for chaining.

unregister(handle: RunHandle) → None[source]

Stop tracking a run and hand ownership back to Python.

Parameters:

handle – the run to drop.

spacr.qt.bridge.apply_worker_budget(settings: Dict[str, Any], total: int | None = None) → int[source]

Cap pool settings in-place and return this run’s worker allocation.

-1 and None mean “all available” in the libraries used by spaCR, so they resolve to the current remaining budget. Explicit smaller values are preserved. The return value is stored on the run handle so the next run sees what this one actually reserved.

Parameters:

settings – the run settings dict, modified in place; each WORKER_SETTING_KEYS entry present is clamped to between 1 and the available workers, with None, -1 and other non-positive values taking all of them and unparseable values left alone.

spacr.qt.bridge.available_worker_count(total: int | None = None) → int[source]

Workers a newly-started run may use, never fewer than one.

spacr.qt.bridge.checkpoint() → None[source]

Pipeline-side cancellation/pause point. Call only where stopping is safe.

“Safe” means: no file is half-written, no multi-table insert is half-done, and no child process is still working. The top of a per-field or per-batch loop, before that unit’s first write, qualifies; anywhere inside the unit does not.

A no-op when the calling thread has no gate (i.e. outside a PipelineWorker), so pipeline code can call it unconditionally and stay importable from a plain script.

spacr.qt.bridge.current_gate() → PauseGate | None[source]

The PauseGate for the calling thread, or None.

spacr.qt.bridge.drain_thread(thread, worker=None, timeout_ms: int = 3000) → bool[source]

Ask thread to stop, wait for it, and never terminate it.

QThread.terminate() is pthread_cancel under the covers, and every thread in this application runs Python. A cancelled thread that happened to hold the GIL never releases it and the whole process stops making progress; one cancelled inside a Qt or PySide internal leaves a corrupt heap and the process dies later, somewhere unrelated. Neither failure names the line that caused it, which is why RunRegistry.cancel_all() documents “refuse to destroy the GUI rather than use QThread.terminate()” — this function is how a caller does that.

A thread that will not stop is parked (see _PARKED_THREADS) so that nothing drops the last reference to a running QThread — which is the abort this whole module is arranged around.

Parameters:
  • thread – the QThread to drain. None is accepted and is a no-op.

  • worker – object to keep alive alongside a parked thread.

  • timeout_ms – how long to wait for the thread to exit.

Returns:

True when the thread has stopped (or was already gone).

spacr.qt.bridge.emit_safely(signal, *args) → bool[source]

Emit a Qt signal unless its C++ receiver has been destroyed.

Parameters:
  • signal (PySide6.QtCore.SignalInstance) – Bound signal to emit.

  • *args – Values passed to signal.emit.

Returns:

bool – True when emission succeeds; False when PySide raises RuntimeError because the receiving object no longer exists.

spacr.qt.bridge.make_thread(fn: Callable[[Dict[str, Any]], Any], settings: Dict[str, Any], app_key: str = '', *, journal: bool = True, user_visible: bool = True, capture_figures: bool = True) → tuple[PySide6.QtCore.QThread, PipelineWorker][source]

Return (thread, worker) — the caller connects the worker’s signals and calls thread.start().

The worker remains Python-owned and is not connected to deleteLater; its last strong reference releases it. The QThread itself keeps deleteLater because it belongs to the caller’s running event loop.

A caller must hold a strong reference to both until thread.finished: a QThread garbage-collected while running takes the process down.

The job is also added to registry() for as long as it runs, so surfaces that did not start it (the Home screen) can show what spaCR is doing. Deregistration hangs off thread.finished — which has the caller’s affinity — rather than worker.finished, which is emitted on the worker thread and would take the registry (a GUI-thread QObject) with it.

Parameters:
  • fn – callable invoked as fn(settings) on the worker thread.

  • settings – the single argument handed to fn.

  • app_key – overrides the key stamped on fn by resolve_pipeline_entry(). Only needed for ad-hoc jobs that want to show up on Home under a name of their own.

  • journal – create a reproducibility record. Disable only for read-only UI housekeeping that is not an analysis run.

  • capture_figures – prepare Matplotlib figure interception. Keep true for analysis runs; read-only UI jobs which cannot emit figures set it false to preserve the operation import boundary.

Returns:

an unstarted (QThread, PipelineWorker) pair.

spacr.qt.bridge.parked_thread_count() → int[source]

How many stubborn threads are currently parked. For diagnostics.

spacr.qt.bridge.pausable(fn: Callable) → Callable[source]

Mark an entry point as honouring checkpoint().

Setting this on a function that does not actually call checkpoint() is how you ship a Pause button that lies, so the marker is deliberately explicit rather than inferred.

Parameters:

fn – the entry-point callable to mark; it is returned, and an object that cannot take attributes is returned unmarked.

spacr.qt.bridge.prune_job_pairs(pairs, finished=None) → List[tuple][source]

Return the (thread, worker) pairs a screen must still hold.

The idiom this replaces — [p for p in jobs if p[0].isRunning()] wired to thread.finished — is wrong twice over, and the two errors hid each other:

  • The slot is queued onto the GUI thread, so it runs after the OS thread has gone. make_thread also wires thread.finished -> thread.deleteLater, and the same pump that delivers this call flushes that deferred delete first — so isRunning() raises RuntimeError: Internal C++ object already deleted from inside a Qt slot. The list comprehension never completes and every pair is retained, forever, along with its worker.

  • self.sender() is not a way out: Qt nulls the sender of a queued call whose emitter has since been destroyed, which is exactly what deleteLater just did. A screen keyed on it silently retires nothing.

So: a wrapper whose C++ half is gone is proof of retirement ( deleteLater is only ever posted from thread.finished), and a finished sender is used when Qt still offers one.

Parameters:
  • pairs – the screen’s (thread, worker) ownership list.

  • finished – the QThread that just finished, when known.

Returns:

the pairs still worth holding a reference to.

spacr.qt.bridge.prune_parked_threads() → int[source]

Release parked (thread, worker) pairs whose thread has exited.

Returns:

how many pairs are still parked afterwards.

spacr.qt.bridge.registry() → RunRegistry[source]

The process-wide RunRegistry (created on first use).

spacr.qt.bridge.resolve_pipeline_entry(app_key: str) → Callable[[Dict[str, Any]], Any] | None[source]

Return the pipeline function that runs a given app, or None if the app is interactive-only (annotate / make_masks) or unknown.

Each returned entry point is wrapped with spacr.qt.verbose_logger.log_call() so that when the user has “Verbose logging” enabled, every pipeline invocation emits an entry-and-return trace in the console. Zero cost when verbose is off (the wrapper is a single attribute check).

The result is also stamped with the app key (_tag()) so the run registry can say which module is running. Note that none of these are stamped pausable() — see PauseGate.

Parameters:

app_key – the app key, e.g. 'mask' or 'measure'; keys outside the built-in chain fall back to the entry= an app registered, then to a plugin app of that key.

spacr.qt.bridge.thread_has_stopped(thread) → bool[source]

True when thread is finished, never started, or already deleted.

None counts as stopped, and so does a QThread whose C++ half PySide6 has already taken away: asking it anything raises RuntimeError, and an object that no longer exists is certainly not still running. That case is not exotic — make_thread wires thread.finished -> thread.deleteLater, so any GUI pump that delivers a retirement slot has usually flushed the deferred delete on the way.

Parameters:

thread – a QThread, or None.

spacr.qt.bridge.wait_for_parked_threads(timeout_ms: int = PARKED_EXIT_WAIT_MS) → int[source]

Wait for parked Qt threads within a shared time limit.

Parked threads remain referenced until native work finishes, preventing Qt from destroying a running QThread during application shutdown.

Parameters:

timeout_ms – Total waiting time, in milliseconds, shared by all parked threads.

Returns:

Number of threads still running when the time limit expires.

spacr.qt.bridge.worker_capacity(total: int | None = None) → int[source]

Logical CPU capacity used by the cooperative run allocator.

Nested helpers

PipelineWorker.run._already_emitted(fig)

Return whether this figure was emitted during this run.

Both the object identity recorded for the run and the figure marker must match. Requiring both avoids false matches when Python reuses an object ID or a figure persists across runs.

spacr/qt/bridge.py:1441

PipelineWorker.run._capture_show(*args, **kwargs)

Emit each new figure once, rendering it HERE on the worker thread.

Agg’s savefig touches no Qt, so the expensive part happens off the GUI thread and the GUI only does a file move and a pixmap load – which is what stops it hanging while figures stream in.

Figures marked _spacr_live_update are re-emitted in place instead, which is how the training monitor refreshes without filling the gallery with one snapshot per epoch.

spacr/qt/bridge.py:1463

PipelineWorker.run._forward_output(text: str) → None

Emit worker text and retain warning lines in its manifest.

spacr/qt/bridge.py:1398

PipelineWorker.run._mark_emitted(fig)

Record that a figure has been emitted, by id and on the figure.

The attribute is the cross-route half of the guard and a figure that refuses one still gets its tile – it just loses that half.

spacr/qt/bridge.py:1451

PipelineWorker.run._publish_figure(fig, path='')

Publish one figure, unless it has already been shown.

A picture that was saved and then plotted is ONE picture: without this the gallery held two tiles for one file, which generate_ml_scores produces as an ordinary sequence.

spacr/qt/bridge.py:1501

PipelineWorker.run._registered_figures()

Yield existing pyplot figures without creating one.

pyplot.figure(number) is both a lookup and a constructor. The numbers below come from pyplot’s process-wide registry, but another thread can close one between the two calls. Read its manager directly so capture never recreates a figure and accidentally assigns it to this run.

spacr/qt/bridge.py:1424

_say_what_is_wrong_with_the_settings.run(settings=None, *args, **kwargs)

Run the entry point, turning a settings problem into a message.

spacr/qt/bridge.py:1709

resolve_pipeline_entry._ret(fn)

Tag the entry point with its app key and its settings check.

spacr/qt/bridge.py:1755