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¶
A latch a worker thread waits on, so a pause is a pause. |
|
Runs one pipeline function in its own thread. |
|
One in-flight job: what it is, how far along, and its pause gate. |
|
Every job started through |
Functions¶
|
Cap pool settings in-place and return this run's worker allocation. |
|
Workers a newly-started run may use, never fewer than one. |
|
Pipeline-side cancellation/pause point. Call only where stopping is safe. |
|
The |
|
Ask |
|
Emit a Qt signal unless its C++ receiver has been destroyed. |
|
Return |
|
How many stubborn threads are currently parked. For diagnostics. |
|
Mark an entry point as honouring |
|
Return the |
|
Release parked |
|
The process-wide |
|
Return the pipeline function that runs a given app, or None if the |
|
True when |
|
Wait for parked Qt threads within a shared time limit. |
|
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.saveof a mask (spacr.objectwrites.npywithout a tmp+rename) or between the severalINSERTs that make up one measured field. Both leave exactly the half-written artefactspacr.resumeexists 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.cancellationis 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 — athreading.Eventdoes the rest.Create the gate open, with nothing paused.
- 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.QObjectRuns one pipeline function in its own thread.
Signals:
line_ready(str)— a chunk of stdout/stderr textfinished(bool)— True if the function returned without an unhandled exceptionerror(str)— traceback string on failurefigure_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.
- 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.resumeexists 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=Falsealready 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
PauseGatefor 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.QObjectOne in-flight job: what it is, how far along, and its pause gate.
Lives on the GUI thread.
progressis 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
PipelineWorkerdoing the work. Itsworker_countis 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_shutdownanduser_visibleare 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.
- 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.finishedrather thanworker.finished: the latter is emitted insiderun(), before the thread’s event loop has actually stopped, and releasing the last reference at that moment is how this module segfaulted once already.
- class spacr.qt.bridge.RunRegistry(parent=None)[source]¶
Bases:
PySide6.QtCore.QObjectEvery job started through
make_thread(), while it runs.Deliberately tiny: a list plus a
changedsignal. 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.
- 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 useQThread.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.closeEventanswers 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_quitis 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.
- is_busy() bool[source]¶
Whether any run is still registered.
What the window asks before quitting.
- Returns:
True while a run is tracked.
- 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.
-1andNonemean “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_KEYSentry present is clamped to between 1 and the available workers, withNone,-1and 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
PauseGatefor the calling thread, orNone.
- spacr.qt.bridge.drain_thread(thread, worker=None, timeout_ms: int = 3000) bool[source]¶
Ask
threadto stop, wait for it, and never terminate it.QThread.terminate()ispthread_cancelunder 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 whyRunRegistry.cancel_all()documents “refuse to destroy the GUI rather than useQThread.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.
Noneis 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 –
Truewhen emission succeeds;Falsewhen PySide raisesRuntimeErrorbecause 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 callsthread.start().The worker remains Python-owned and is not connected to
deleteLater; its last strong reference releases it. The QThread itself keepsdeleteLaterbecause 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 offthread.finished— which has the caller’s affinity — rather thanworker.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
fnbyresolve_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 tothread.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_threadalso wiresthread.finished -> thread.deleteLater, and the same pump that delivers this call flushes that deferred delete first — soisRunning()raisesRuntimeError: Internal C++ object already deletedfrom 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 whatdeleteLaterjust did. A screen keyed on it silently retires nothing.
So: a wrapper whose C++ half is gone is proof of retirement (
deleteLateris only ever posted fromthread.finished), and afinishedsender 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 stampedpausable()— seePauseGate.- Parameters:
app_key – the app key, e.g.
'mask'or'measure'; keys outside the built-in chain fall back to theentry=an app registered, then to a plugin app of that key.
- spacr.qt.bridge.thread_has_stopped(thread) bool[source]¶
True when
threadis finished, never started, or already deleted.Nonecounts as stopped, and so does a QThread whose C++ half PySide6 has already taken away: asking it anything raisesRuntimeError, and an object that no longer exists is certainly not still running. That case is not exotic —make_threadwiresthread.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
QThreadduring 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.
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_updateare 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_scoresproduces 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