"""Run spaCR pipelines in Qt workers and relay their progress to the GUI.
:class:`PipelineWorker` executes a pipeline call in a ``QThread`` and routes
captured standard output and errors through ``line_ready``. The process-wide
:data:`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
:class:`spacr.cancellation.CancellationToken`, which long-running pipelines
poll at safe boundaries. :class:`PauseGate` and :func:`checkpoint` provide the
pause protocol, but pipeline entries that do not call ``checkpoint`` report
``supports_pause=False`` and must expose Pause as unavailable.
"""
from __future__ import annotations
import atexit
import gc
import io
import logging
import os
import re
import sys
import threading
import time
import traceback
import functools
from typing import Any, Callable, Dict, List, Optional
from PySide6.QtCore import QObject, QThread, Qt, Signal
from spacr.qt.gil_priority import responsive_gui
from spacr.cancellation import (
CancellationToken,
PipelineCancelled,
checkpoint as cancellation_checkpoint,
installed_token,
)
from ..figures.style import figure_style, theme_target
LOG = logging.getLogger(__name__)
class _StreamRedirector(io.TextIOBase):
"""A file-like object that emits every write to a queue for the UI.
Pipeline libraries (cellpose especially) print progress WITHOUT a
trailing newline while they set up — model download, warmup, etc.
A pure "emit only on \\n" redirector holds those bytes hostage in
the buffer, and to the user it looks like the app hung after
"Starting mask…". Two mitigations here:
1. **Chunk cap** — if the buffer grows past ``_MAX_BUF_CHARS``
we emit it regardless of newline, then keep buffering. This
makes long dependency-import chatter visible instead of silent.
2. **Idle flush** — the caller can call :meth:`idle_flush` from a
QTimer to emit whatever partial line has been sitting quiet for
a while, so short-but-newline-less progress lines still surface.
"""
_MAX_BUF_CHARS = 1024
def __init__(self, on_write: Callable[[str], None]):
"""Buffer writes into lines and hand each one to ``on_write``.
:param on_write: called with each completed line. CALLED ON WHATEVER
THREAD WROTE -- the worker for a print, the pump thread for an
idle flush -- so anything touching Qt widgets from here has to
get itself onto the GUI thread.
"""
super().__init__()
self._buf = ""
self._on_write = on_write
self._lock = threading.Lock()
def write(self, s: str) -> int:
"""Buffer written text and emit it a line at a time.
Emitting happens OUTSIDE the lock, so a slow slot cannot block the
thread that is writing. A buffer that grows past the cap without a
newline is flushed anyway -- a progress bar that never emits one would
otherwise be held until the run ended.
:param s: the text written; coerced with :func:`str`, because a
non-string reaching a redirected stream is a caller's bug and losing
the output would hide it.
:returns: how many characters were accepted, as a stream must.
"""
if not isinstance(s, str):
s = str(s)
with self._lock:
self._buf += s
emits = []
while "\n" in self._buf:
line, self._buf = self._buf.split("\n", 1)
emits.append(line + "\n")
if len(self._buf) >= self._MAX_BUF_CHARS:
emits.append(self._buf)
self._buf = ""
for chunk in emits:
self._safe_emit(chunk)
return len(s)
def flush(self) -> None:
"""Emit whatever is buffered, newline or not."""
with self._lock:
pending, self._buf = self._buf, ""
if pending:
self._safe_emit(pending)
def idle_flush(self) -> None:
"""Emit any pending partial line — safe to call from the pump."""
with self._lock:
pending, self._buf = self._buf, ""
if pending:
self._safe_emit(pending)
def _safe_emit(self, s: str) -> None:
"""Hand one line to the callback, swallowing anything it raises.
A print must never fail because something downstream of the console did.
"""
try:
self._on_write(s)
except Exception:
pass
class _ThreadStreamRouter(io.TextIOBase):
"""Route writes from each worker thread to that worker's console.
Replacing ``sys.stdout`` independently in overlapping workers is unsafe:
the second worker saves the first worker's redirector, and whichever one
finishes first restores a stream belonging to the other run. This single
process-wide proxy keeps the public stream stable while selecting the
destination by thread identity.
"""
def __init__(self, original):
"""Wrap the real stream and route writes by thread identity.
:param original: the stream this proxy replaces, kept so a write
from a thread with no registered console still reaches the
terminal. It is the FALLBACK, not a tee: a write that finds a
target goes there instead, not as well.
"""
super().__init__()
self.original = original
self._targets: Dict[int, List[_StreamRedirector]] = {}
self._lock = threading.RLock()
def register(self, target: _StreamRedirector) -> None:
"""Route this thread's writes to a redirector.
Registered as a STACK per thread, so a nested run restores the outer
one's target when it unregisters rather than clearing the route.
:param target: the redirector to send this thread's output to.
"""
ident = threading.get_ident()
with self._lock:
self._targets.setdefault(ident, []).append(target)
def unregister(self, target: _StreamRedirector) -> None:
"""Stop routing this thread's writes to a redirector.
:param target: the redirector to remove; one already gone is ignored,
since teardown order is not this object's to guarantee.
"""
ident = threading.get_ident()
with self._lock:
stack = self._targets.get(ident, [])
if target in stack:
stack.remove(target)
if not stack:
self._targets.pop(ident, None)
def has_targets(self) -> bool:
"""Report whether any thread is currently routed.
:returns: ``True`` while at least one redirector is registered.
"""
with self._lock:
return bool(self._targets)
def _target(self):
"""The stream for the calling thread, or the original.
Reads the TOP of that thread's stack, so nested redirections unwind in
the order they were made rather than the last one winning for good.
"""
with self._lock:
stack = self._targets.get(threading.get_ident(), [])
return stack[-1] if stack else self.original
def write(self, value: str) -> int:
"""Write to whichever redirector this thread is routed to.
:param value: the text.
:returns: how many characters were accepted.
"""
return self._target().write(value)
def flush(self) -> None:
"""Flush this thread's redirector, tolerating one that has gone away."""
try:
self._target().flush()
except Exception:
pass
@property
def encoding(self):
"""The wrapped stream's encoding, or ``None``.
Delegated rather than declared: code that inspects a stream's encoding
is asking about the real one underneath.
"""
return getattr(self.original, "encoding", None)
def isatty(self) -> bool:
"""Whether the wrapped stream is a terminal.
:returns: the real stream's answer, and ``False`` for a stream that
does not implement it -- a router is never itself a terminal.
"""
return bool(getattr(self.original, "isatty", lambda: False)())
_STREAM_ROUTER_LOCK = threading.RLock()
_STDOUT_ROUTER: Optional[_ThreadStreamRouter] = None
_STDERR_ROUTER: Optional[_ThreadStreamRouter] = None
def _register_worker_streams(
target: _StreamRedirector,
) -> tuple[_ThreadStreamRouter, _ThreadStreamRouter]:
"""Install/reuse the process routers and register the calling thread."""
global _STDOUT_ROUTER, _STDERR_ROUTER
with _STREAM_ROUTER_LOCK:
if not isinstance(sys.stdout, _ThreadStreamRouter):
_STDOUT_ROUTER = _ThreadStreamRouter(sys.stdout)
sys.stdout = _STDOUT_ROUTER
else:
_STDOUT_ROUTER = sys.stdout
if not isinstance(sys.stderr, _ThreadStreamRouter):
_STDERR_ROUTER = _ThreadStreamRouter(sys.stderr)
sys.stderr = _STDERR_ROUTER
else:
_STDERR_ROUTER = sys.stderr
_STDOUT_ROUTER.register(target)
_STDERR_ROUTER.register(target)
return _STDOUT_ROUTER, _STDERR_ROUTER
def _unregister_worker_streams(
target: _StreamRedirector,
stdout_router: _ThreadStreamRouter,
stderr_router: _ThreadStreamRouter,
) -> None:
"""Remove the calling worker and restore original streams when idle."""
global _STDOUT_ROUTER, _STDERR_ROUTER
with _STREAM_ROUTER_LOCK:
stdout_router.unregister(target)
stderr_router.unregister(target)
if not stdout_router.has_targets() and sys.stdout is stdout_router:
sys.stdout = stdout_router.original
_STDOUT_ROUTER = None
if not stderr_router.has_targets() and sys.stderr is stderr_router:
sys.stderr = stderr_router.original
_STDERR_ROUTER = None
_MPL_SHOW_LOCK = threading.RLock()
_MPL_SHOW_TARGETS: Dict[int, List[Callable[..., Any]]] = {}
_MPL_ORIGINAL_SHOW: Optional[Callable[..., Any]] = None
_MPL_MODULE = None
def _matplotlib_show_router(*args, **kwargs):
"""Route ``plt.show()`` to the current thread's figure capture.
Calls without a registered capture are discarded on worker threads so
Matplotlib cannot start a Qt event loop outside the GUI thread. Calls on
the main thread fall back to the original ``show`` implementation.
"""
with _MPL_SHOW_LOCK:
stack = _MPL_SHOW_TARGETS.get(threading.get_ident(), [])
target = stack[-1] if stack else None
original = _MPL_ORIGINAL_SHOW
if target is not None:
return target(*args, **kwargs)
if threading.current_thread() is not threading.main_thread():
LOG.debug("plt.show() on %s with no capture: dropped rather than "
"entering a Qt event loop off the GUI thread",
threading.current_thread().name)
return None
return original(*args, **kwargs) if original is not None else None
def _register_matplotlib_show(plt, target: Callable[..., Any]) -> None:
"""Route ``plt.show`` by worker thread without cross-run restoration."""
global _MPL_ORIGINAL_SHOW, _MPL_MODULE
with _MPL_SHOW_LOCK:
if not _MPL_SHOW_TARGETS:
_MPL_ORIGINAL_SHOW = plt.show
_MPL_MODULE = plt
plt.show = _matplotlib_show_router
_MPL_SHOW_TARGETS.setdefault(threading.get_ident(), []).append(target)
def _unregister_matplotlib_show(target: Callable[..., Any]) -> None:
"""Remove one thread's ``show`` target, restoring the original when the last goes.
The module-level patch is undone only when no thread is routed any
more, and only if it is still this router that is installed -- something
else having patched ``show`` in the meantime is not ours to revert.
:param target: the show callable to unregister.
"""
global _MPL_ORIGINAL_SHOW, _MPL_MODULE
with _MPL_SHOW_LOCK:
ident = threading.get_ident()
stack = _MPL_SHOW_TARGETS.get(ident, [])
if target in stack:
stack.remove(target)
if not stack:
_MPL_SHOW_TARGETS.pop(ident, None)
if not _MPL_SHOW_TARGETS and _MPL_MODULE is not None:
if _MPL_MODULE.show is _matplotlib_show_router:
_MPL_MODULE.show = _MPL_ORIGINAL_SHOW
_MPL_ORIGINAL_SHOW = None
_MPL_MODULE = None
#: Attribute a pipeline entry point sets on itself to declare that it
#: polls :func:`checkpoint` often enough for a pause to mean something.
#: See :func:`pausable`.
PAUSABLE_ATTR = "__spacr_pausable__"
#: Attribute carrying the app key an entry point belongs to, stamped by
#: :func:`resolve_pipeline_entry` so :func:`make_thread` can name the job
#: without every caller having to pass it.
APP_KEY_ATTR = "__spacr_app_key__"
[docs]
class PauseGate:
"""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 ``INSERT``\\ s
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 :func:`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 :func:`pausable`. Stop is a separate contract:
:mod:`spacr.cancellation` is polled by shipped workflows at durable
boundaries and never claims to suspend an in-progress write.
Thread-safety: :meth:`pause` / :meth:`resume` are called from the GUI
thread; :meth:`wait_if_paused` blocks the worker thread. That is the
whole contract — a :class:`threading.Event` does the rest.
"""
def __init__(self) -> None:
"""Create the gate open, with nothing paused."""
self._running = threading.Event()
self._running.set()
self._paused_since: Optional[float] = None
[docs]
def pause(self) -> None:
"""Ask the worker to stop at its next checkpoint."""
if self._running.is_set():
self._paused_since = time.time()
self._running.clear()
[docs]
def resume(self) -> None:
"""Release a paused worker."""
self._paused_since = None
self._running.set()
[docs]
def is_paused(self) -> bool:
"""True once :meth:`pause` has been called and not yet released."""
return not self._running.is_set()
[docs]
def paused_for(self) -> float:
"""Seconds spent waiting, or ``0.0`` when not paused."""
since = self._paused_since
return 0.0 if since is None else max(0.0, time.time() - since)
[docs]
def wait_if_paused(self, timeout: Optional[float] = None) -> bool:
"""Block while paused. Returns True once running again.
Called by pipeline code through :func:`checkpoint`; safe to call
when not paused, where it returns immediately.
"""
return self._running.wait(timeout)
#: Per-thread gate, installed by :meth:`PipelineWorker.run` for the
#: duration of the call. Thread-local rather than global so two
#: overlapping workers cannot pause each other.
_LOCAL = threading.local()
[docs]
def current_gate() -> Optional[PauseGate]:
"""The :class:`PauseGate` for the calling thread, or ``None``."""
return getattr(_LOCAL, "gate", None)
[docs]
def checkpoint() -> None:
"""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
:class:`PipelineWorker`), so pipeline code can call it
unconditionally and stay importable from a plain script.
"""
cancellation_checkpoint()
gate = current_gate()
if gate is not None:
gate.wait_if_paused()
[docs]
def pausable(fn: Callable) -> Callable:
"""Mark an entry point as honouring :func:`checkpoint`.
Setting this on a function that does *not* actually call
:func:`checkpoint` is how you ship a Pause button that lies, so the
marker is deliberately explicit rather than inferred.
:param fn: the entry-point callable to mark; it is returned, and an object
that cannot take attributes is returned unmarked.
"""
try:
setattr(fn, PAUSABLE_ATTR, True)
except (AttributeError, TypeError):
pass
return fn
#: ``spacr.utils.print_progress`` writes exactly this shape, and it is
#: the only progress signal the pipelines emit. Parsing it off stdout is
#: read-only observation — it gives the Home screen a real "field 41 of
#: 96" without anything having to be threaded through the pipeline.
_PROGRESS_RE = re.compile(r"\bProgress:\s*(\d+)\s*/\s*(\d+)")
_MASK_GPU_RE = re.compile(
r"\[mask GPUs\] (\w+) Progress: (\d+)/(\d+) archives, failed (\d+)"
r"((?: \| GPU \S+ \w+ \d+/\d+)*)")
_MASK_GPU_WORKER_RE = re.compile(r"GPU (\S+) (\w+) (\d+)/(\d+)")
_WATCH_FOLDER_RE = re.compile(
r"watch_folder: (\d+) analysed, (\d+) waiting, (\d+) failed")
def _watch_folder_progress(text: str) -> Optional[dict]:
"""Read the newest folder-watch count line in ``text``, if any.
The folder watch of Make Masks prints the line whenever a count changes.
:param text: worker output.
:returns: ``{'done', 'waiting', 'failed'}`` counts, or None.
"""
matches = list(_WATCH_FOLDER_RE.finditer(text or ""))
if not matches:
return None
done, waiting, failed = (int(value) for value in matches[-1].groups())
return {"done": done, "waiting": waiting, "failed": failed}
def _mask_gpu_progress(text: str) -> Optional[dict]:
"""Read the newest parallel mask progress line in ``text``, if any.
``spacr._mask_workers._progress_line`` writes the line; this returns the
object role, overall ``done``/``total``/``failed`` archive counts and a
``workers`` list of ``(device, state, done, total)`` tuples.
"""
matches = list(_MASK_GPU_RE.finditer(text or ""))
if not matches:
return None
match = matches[-1]
workers = [(device, state, int(done), int(total))
for device, state, done, total
in _MASK_GPU_WORKER_RE.findall(match.group(5))]
return {"object_type": match.group(1), "done": int(match.group(2)),
"total": int(match.group(3)), "failed": int(match.group(4)),
"workers": workers}
WORKER_SETTING_KEYS = (
"n_jobs",
"n_workers",
"ultrack_n_workers",
"infection_xgb_n_jobs",
)
[docs]
def worker_capacity(total: Optional[int] = None) -> int:
"""Logical CPU capacity used by the cooperative run allocator."""
value = os.cpu_count() if total is None else total
try:
return max(1, int(value or 1))
except (TypeError, ValueError):
return 1
[docs]
def available_worker_count(total: Optional[int] = None) -> int:
"""Workers a newly-started run may use, never fewer than one."""
capacity = worker_capacity(total)
reserved = sum(
max(0, int(getattr(handle, "worker_count", 1)) - 1)
for handle in registry().active()
)
return max(1, capacity - reserved)
[docs]
def apply_worker_budget(
settings: Dict[str, Any],
total: Optional[int] = None,
) -> int:
"""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.
:param 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.
"""
available = available_worker_count(total)
allocations: List[int] = []
for key in WORKER_SETTING_KEYS:
if key not in settings:
continue
raw = settings.get(key)
try:
requested = int(raw) if raw is not None else available
except (TypeError, ValueError):
continue
if requested <= 0:
requested = available
resolved = max(1, min(requested, available))
settings[key] = resolved
allocations.append(resolved)
return max(allocations, default=1)
[docs]
class RunHandle(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 :data:`_PROGRESS_RE`) rather than reported by the
pipeline, because the pipelines have no reporting channel.
:param app_key: which module is running. Falls back to ``"job"`` so a
handle always has a name to show.
:param worker: the :class:`PipelineWorker` doing the work. Its
``worker_count`` is read once here rather than on every update.
:param thread: the thread the worker was moved to. Held so the handle
can wait on it, not so it can be restarted.
:param parent: parent object.
"""
changed = Signal()
def __init__(self, app_key: str, worker: "PipelineWorker",
thread: QThread, parent=None):
"""Wrap one running job for the run registry.
``blocks_shutdown`` and ``user_visible`` are read from the worker at
construction rather than on demand: :meth:`retire` drops the worker
reference, and the answers are still needed after that.
:param app_key: what the job is called.
:param worker: the pipeline worker doing the work.
:param thread: the thread it runs on.
:param parent: parent object, or ``None``.
"""
super().__init__(parent)
self.app_key = app_key or "job"
self.worker = worker
self.thread = thread
self.worker_count = max(1, int(getattr(worker, "worker_count", 1)))
self.started_at = time.time()
#: ``(done, total)`` from the last progress line, or ``None``.
self.progress: Optional[tuple] = None
#: Last non-blank line the job printed — what it is doing now.
self.last_line = ""
#: Whether a still-running instance of this job is a reason to
#: refuse to close the application. Cached at construction because
#: :meth:`retire` drops ``self.worker``. See
#: :attr:`PipelineWorker.blocks_shutdown`.
self.blocks_shutdown = bool(
getattr(worker, "blocks_shutdown", True))
#: False for housekeeping the user did not start and cannot act on.
#: Home's run banner and anything else that reports "a run is in
#: progress" must skip these. The usage poller submits every two
#: seconds, so without this Home flashes "<module> usage - running"
#: on and off for as long as a module screen is open.
self.user_visible = bool(getattr(worker, "user_visible", True))
worker.line_ready.connect(self._on_line)
@property
[docs]
def supports_pause(self) -> bool:
"""Whether pausing this job would actually pause it."""
return self.worker.supports_pause
@property
[docs]
def gate(self) -> PauseGate:
"""The pause gate this run checks between steps.
:returns: the gate.
"""
return self.worker.gate
[docs]
def elapsed(self) -> float:
"""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.
"""
return max(0.0, time.time() - self.started_at)
[docs]
def is_running(self) -> bool:
"""Whether this handle still owns a live QThread."""
thread = self.thread
if thread is None:
return False
try:
return bool(thread.isRunning())
except RuntimeError:
return False
[docs]
def request_cancel(self, reason: str = "cancelled by the user") -> None:
"""Cooperatively cancel this job and ask its thread to retire."""
worker = self.worker
thread = self.thread
if worker is not None:
worker.request_cancel(reason)
if thread is not None:
try:
thread.requestInterruption()
except RuntimeError:
pass
[docs]
def fraction(self) -> Optional[float]:
"""Completed fraction in ``0..1``, or ``None`` when unknown."""
if not self.progress:
return None
done, total = self.progress
return None if total <= 0 else max(0.0, min(1.0, done / total))
[docs]
def retire(self) -> None:
"""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.
"""
try:
self.worker.line_ready.disconnect(self._on_line)
except (RuntimeError, TypeError):
pass
self.worker = None
self.thread = None
registry().unregister(self)
def _on_line(self, chunk: str) -> None:
"""Record the newest output line, and any progress it carries.
:param chunk: whatever the worker just printed; blank output is ignored,
and the line is truncated so one enormous line cannot become the
status bar.
"""
text = chunk.strip()
if not text:
return
self.last_line = text.splitlines()[-1][:200]
matches = list(_PROGRESS_RE.finditer(text))
if matches:
match = matches[-1]
self.progress = (int(match.group(1)), int(match.group(2)))
self.changed.emit()
[docs]
class RunRegistry(QObject):
"""Every job started through :func:`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.
:param parent: parent widget.
"""
changed = Signal()
def __init__(self, parent=None):
"""Create the empty run registry.
:param parent: parent object, or ``None``.
"""
super().__init__(parent)
self._handles: List[RunHandle] = []
[docs]
def register(self, handle: RunHandle) -> RunHandle:
"""Take ownership of a run and start reporting it.
:param handle: the run to track.
:returns: the same handle, for chaining.
"""
self._handles.append(handle)
handle.changed.connect(self.changed)
self.changed.emit()
return handle
[docs]
def unregister(self, handle: RunHandle) -> None:
"""Stop tracking a run and hand ownership back to Python.
:param handle: the run to drop.
"""
if handle in self._handles:
self._handles.remove(handle)
handle.setParent(None)
self.changed.emit()
[docs]
def active(self) -> List[RunHandle]:
"""Handles for the jobs running right now, oldest first."""
return list(self._handles)
[docs]
def is_busy(self) -> bool:
"""Whether any run is still registered.
What the window asks before quitting.
:returns: True while a run is tracked.
"""
return bool(self._handles)
[docs]
def cancel_all(
self,
timeout_ms: int = 5000,
reason: str = "application shutdown",
) -> List[RunHandle]:
"""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
:attr:`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.
:param timeout_ms: total wait budget across all active threads.
:param reason: cancellation reason recorded by each worker.
:returns: handles that block shutdown and did not stop in the budget.
"""
handles = self.active()
for handle in handles:
if getattr(handle, "stop_on_quit", True):
handle.request_cancel(reason)
deadline = time.monotonic() + max(0, int(timeout_ms)) / 1000.0
for handle in handles:
thread = handle.thread
if thread is None or not handle.is_running():
continue
remaining_ms = max(
0, int((deadline - time.monotonic()) * 1000))
if remaining_ms <= 0:
break
try:
thread.wait(remaining_ms)
except RuntimeError:
pass
return [
handle for handle in handles
if handle.is_running() and handle.blocks_shutdown
]
[docs]
def clear(self) -> None:
"""Drop every handle. For tests — never call this on a live app."""
if self._handles:
self._handles.clear()
self.changed.emit()
_REGISTRY: Optional[RunRegistry] = None
[docs]
def registry() -> RunRegistry:
"""The process-wide :class:`RunRegistry` (created on first use)."""
global _REGISTRY
if _REGISTRY is None:
_REGISTRY = RunRegistry()
return _REGISTRY
class _ExternalJobHandle(QObject):
"""A job spaCR runs outside a :class:`PipelineWorker`, in the run registry.
A child process (a model install, a queued job that is its own
process) has no worker thread, yet the user started it and may want to
watch or stop it. This handle carries what the Jobs window, Home and
the quit check read from a :class:`RunHandle`: a name, progress, the
last output line, elapsed time and a Cancel that calls ``cancel``.
:param app_key: what the job is called.
:param cancel: called with the reason when the user cancels, or ``None``
when the job cannot be stopped.
:param blocks_shutdown: whether a still-running job is a reason to
refuse to close the window.
:param user_visible: whether the job is shown to the user at all.
:param stop_on_quit: whether closing spaCR cancels the job.
:param parent: parent object, or ``None``.
"""
changed = Signal()
def __init__(self, app_key: str, cancel: Optional[Callable[[str], Any]] = None,
*, blocks_shutdown: bool = False, user_visible: bool = True,
stop_on_quit: bool = True, parent=None):
"""Describe one running external job.
:param app_key: what the job is called.
:param cancel: the stop request, or ``None``.
:param blocks_shutdown: whether it vetoes closing the window.
:param user_visible: whether it is listed.
:param stop_on_quit: whether closing spaCR cancels it.
:param parent: parent object, or ``None``.
"""
super().__init__(parent)
self.app_key = app_key or "job"
self.worker = None
self.thread = None
self.worker_count = 1
self.started_at = time.time()
self.progress: Optional[tuple] = None
self.last_line = ""
self.blocks_shutdown = bool(blocks_shutdown)
self.user_visible = bool(user_visible)
self.supports_pause = False
self.stop_on_quit = bool(stop_on_quit)
self.gate = PauseGate()
self._cancel = cancel
self._live = True
def elapsed(self) -> float:
"""Seconds since the job started, never negative."""
return max(0.0, time.time() - self.started_at)
def is_running(self) -> bool:
"""Whether the job has not yet been retired."""
return self._live
def fraction(self) -> Optional[float]:
"""Completed fraction in ``0..1``, or ``None`` when unknown."""
if not self.progress:
return None
done, total = self.progress
return None if total <= 0 else max(0.0, min(1.0, done / total))
def request_cancel(self, reason: str = "cancelled by the user") -> None:
"""Ask the job to stop; nothing happens when it cannot be stopped.
:param reason: why, passed on to the stop request.
"""
if self._cancel is not None and self._live:
self._cancel(reason)
def report(self, done: Optional[int] = None, total: Optional[int] = None,
line: str = "") -> None:
"""Record progress and the newest output line, and tell the listeners.
:param done: steps finished, with ``total``; ``None`` leaves it.
:param total: steps in all.
:param line: what the job is doing now; blank leaves it.
"""
if done is not None and total:
self.progress = (int(done), int(total))
text = str(line or "").strip()
if text:
self.last_line = text.splitlines()[-1][:200]
self.changed.emit()
def retire(self) -> None:
"""The job has ended: drop out of the registry."""
if not self._live:
return
self._live = False
registry().unregister(self)
def _track_external_job(app_key: str,
cancel: Optional[Callable[[str], Any]] = None,
**options: Any) -> _ExternalJobHandle:
"""Register a job that runs outside :func:`make_thread` until it retires.
The caller reports progress through :meth:`_ExternalJobHandle.report`
and must call :meth:`_ExternalJobHandle.retire` when the job ends.
:param app_key: what the job is called in the Jobs window.
:param cancel: called with the reason on Cancel, or ``None``.
:param options: ``blocks_shutdown``, ``user_visible`` and
``stop_on_quit`` for the handle.
:returns: the registered handle.
"""
reg = registry()
handle = _ExternalJobHandle(app_key, cancel, parent=reg, **options)
reg.register(handle)
return handle
def _registered_handle(worker: Any) -> Optional["RunHandle"]:
"""The registry handle of a :func:`make_thread` worker, or ``None``.
:param worker: the worker :func:`make_thread` returned.
:returns: its handle while the job is registered.
"""
if worker is None:
return None
for handle in registry().active():
if getattr(handle, "worker", None) is worker:
return handle
return None
#: ``(thread, worker)`` pairs that outlived the widget which owned them.
#: Parking a stubborn thread here is not a leak-with-a-nice-name: it is the
#: only way to satisfy Qt's rule that a QThread must not be destroyed while
#: it is running, once the owner has decided to go away anyway. Entries are
#: released by :func:`prune_parked_threads`, which every :func:`drain_thread`
#: call runs first.
_PARKED_THREADS: List[tuple] = []
_PARKED_LOCK = threading.Lock()
#: How long the exit hook below waits for the parked threads, in milliseconds.
PARKED_EXIT_WAIT_MS = 20_000
_PARKED_EXIT_HOOK_INSTALLED = False
[docs]
def wait_for_parked_threads(timeout_ms: int = PARKED_EXIT_WAIT_MS) -> int:
"""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.
:param timeout_ms: Total waiting time, in milliseconds, shared by all
parked threads.
:returns: Number of threads still running when the time limit expires.
"""
deadline = time.monotonic() + max(0, int(timeout_ms)) / 1000.0
with _PARKED_LOCK:
pairs = list(_PARKED_THREADS)
for thread, _worker in pairs:
remaining = deadline - time.monotonic()
if remaining <= 0:
break
try:
thread.wait(int(remaining * 1000))
except RuntimeError:
continue
return prune_parked_threads()
def _drain_parked_threads_at_exit() -> None:
"""``atexit`` hook: give parked threads their last chance to finish."""
still_running = wait_for_parked_threads()
if still_running:
LOG.error(
"%d worker thread(s) are still running as the process exits. "
"Their QThread wrappers are about to be destroyed, which Qt "
"treats as fatal. Whatever they are inside did not answer a "
"stop request.", still_running)
def _install_parked_exit_hook() -> None:
"""Register the exit hook once, and only if something has been parked."""
global _PARKED_EXIT_HOOK_INSTALLED
if _PARKED_EXIT_HOOK_INSTALLED:
return
_PARKED_EXIT_HOOK_INSTALLED = True
atexit.register(_drain_parked_threads_at_exit)
[docs]
def prune_parked_threads() -> int:
"""Release parked ``(thread, worker)`` pairs whose thread has exited.
:returns: how many pairs are still parked afterwards.
"""
with _PARKED_LOCK:
alive = []
for pair in _PARKED_THREADS:
try:
if pair[0].isRunning():
alive.append(pair)
except RuntimeError:
pass
_PARKED_THREADS[:] = alive
return len(_PARKED_THREADS)
[docs]
def parked_thread_count() -> int:
"""How many stubborn threads are currently parked. For diagnostics."""
with _PARKED_LOCK:
return len(_PARKED_THREADS)
[docs]
def thread_has_stopped(thread) -> bool:
"""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.
:param thread: a QThread, or None.
"""
if thread is None:
return True
try:
return not thread.isRunning()
except RuntimeError:
return True
[docs]
def prune_job_pairs(pairs, finished=None) -> List[tuple]:
"""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.
:param pairs: the screen's ``(thread, worker)`` ownership list.
:param finished: the QThread that just finished, when known.
:returns: the pairs still worth holding a reference to.
"""
kept: List[tuple] = []
for pair in pairs:
thread = pair[0] if pair else None
if finished is not None and thread is finished:
continue
if not thread_has_stopped(thread):
kept.append(pair)
return kept
[docs]
def emit_safely(signal, *args) -> bool:
"""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.
"""
try:
signal.emit(*args)
return True
except RuntimeError:
return False
[docs]
def drain_thread(thread, worker=None, timeout_ms: int = 3000) -> bool:
"""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
:meth:`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 :data:`_PARKED_THREADS`)
so that nothing drops the last reference to a running QThread — which
is the abort this whole module is arranged around.
:param thread: the QThread to drain. ``None`` is accepted and is a no-op.
:param worker: object to keep alive alongside a parked thread.
:param timeout_ms: how long to wait for the thread to exit.
:returns: True when the thread has stopped (or was already gone).
"""
prune_parked_threads()
if thread is None:
return True
try:
if not thread.isRunning():
return True
except RuntimeError:
return True
try:
thread.quit()
stopped = bool(thread.wait(max(0, int(timeout_ms))))
except RuntimeError:
return True
if stopped:
return True
with _PARKED_LOCK:
_PARKED_THREADS.append((thread, worker))
_install_parked_exit_hook()
LOG.warning(
"A worker thread did not stop within %d ms; it is parked rather "
"than terminated so the process is not left with a corrupt heap.",
timeout_ms,
)
return False
class _SkipFigureCapture(Exception):
"""Select the no-Matplotlib path for read-only background jobs."""
class _OutputDrain(QObject):
"""Sends a worker's waiting output from the thread that built the worker.
Lives in that thread (the GUI thread for every screen), so the chunk
is emitted where its receivers live. A worker built where no event
loop runs still loses nothing: the run drains its output before it
reports that it finished.
:param worker: the :class:`PipelineWorker` whose output this sends.
"""
def __init__(self, worker: "PipelineWorker"):
"""Connect to ``worker``'s pending-output signal."""
super().__init__()
import weakref
self._worker = weakref.ref(worker)
worker._output_pending.connect(self._schedule, Qt.QueuedConnection)
def _schedule(self) -> None:
"""Drain the worker one interval from now."""
from PySide6.QtCore import QTimer
QTimer.singleShot(
int(PipelineWorker._OUTPUT_INTERVAL_S * 1000), self._drain)
def _drain(self) -> None:
"""Emit what the worker has waiting, if it still exists."""
worker = self._worker()
if worker is not None:
worker._drain_output()
[docs]
class PipelineWorker(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.
:param fn: the pipeline function to run. It is called on the worker
thread, so anything it constructs belongs to that thread.
:param settings: the settings dict handed to ``fn``.
:param worker_count: how many workers the run is allowed, passed through
to the pipeline rather than used here.
:param app_key: which module is running, for the journal and for routing
output back to the right screen.
:param journal: whether to record the run in the journal.
:param 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.
"""
line_ready = Signal(str)
finished = Signal(bool)
error = Signal(str)
figure_ready = Signal(object, str)
#: Whatever the pipeline function RETURNED. For the regression that is
#: {'results': coef_df, 'res_folder': ..., 'model': ...} -- the table the
#: run just fitted, in memory.
#:
#: It used to be thrown away, so the only way to see a finished run was to
#: guess where it had written and re-read the CSV. Guessing a path is how
#: a screen ends up showing last month's results, or none at all.
result_ready = Signal(object)
_output_pending = Signal()
_OUTPUT_INTERVAL_S = 0.05
def __init__(
self,
fn: Callable[..., Any],
settings: Dict[str, Any],
worker_count: int = 1,
app_key: str = "",
journal: bool = True,
capture_figures: bool = True,
):
"""Prepare to run ``fn(settings)`` in a worker thread.
:param fn: pipeline entry point (see :func:`resolve_pipeline_entry`).
:param settings: keyword-style dict passed as the sole argument.
:param worker_count: worker allocation reserved in the run registry.
:param app_key: optional explicit module name for run-history records.
:param journal: create a reproducibility manifest. Set false only for
read-only background UI maintenance such as refreshing history.
:param 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.
"""
super().__init__()
self._fn = fn
self._settings = settings
self._app_key_override = str(app_key or "")
self._journal_enabled = bool(journal)
self._capture_figures = bool(capture_figures)
self.worker_count = max(1, int(worker_count))
self.cancel_token = CancellationToken()
self.was_cancelled = False
#: Latch this worker waits on when the pipeline calls
#: :func:`checkpoint`. Always present; only *effective* when
#: :attr:`supports_pause` is True.
self.gate = PauseGate()
self._out_lock = threading.RLock()
self._out_parts: List[str] = []
self._out_last = 0.0
self._out_scheduled = False
self._out_drain = _OutputDrain(self)
def _push_output(self, text: str) -> None:
"""Hand worker output to ``line_ready``, joining bursts.
Emits at once when nothing is waiting and the last emission is
older than 50 ms; otherwise the text waits and the drain on the
thread that built this worker sends it within one interval, so a
run printing thousands of lines a second costs the GUI thread twenty
appends a second rather than one per line. Every other emission of
this worker drains first, so the order of the output is kept. The
lock is held only to move text, never across an emission, so a
process forked mid-run (a DataLoader worker) cannot inherit it held.
:param text: one chunk of the run's stdout or stderr.
"""
with self._out_lock:
now = time.monotonic()
immediate = (not self._out_parts
and now - self._out_last >= self._OUTPUT_INTERVAL_S)
if immediate:
self._out_last = now
else:
self._out_parts.append(text)
if self._out_scheduled:
return
self._out_scheduled = True
if immediate:
self.line_ready.emit(text)
else:
self._output_pending.emit()
def _drain_output(self) -> None:
"""Emit whatever output is waiting, as one chunk. Any thread."""
with self._out_lock:
self._out_scheduled = False
if not self._out_parts:
return
text = "".join(self._out_parts)
self._out_parts.clear()
self._out_last = time.monotonic()
self.line_ready.emit(text)
def _say(self, text: str) -> None:
"""Emit ``text`` after any output still waiting."""
self._drain_output()
with self._out_lock:
self._out_last = time.monotonic()
self.line_ready.emit(text)
def _say_error(self, tb: str) -> None:
"""Emit a traceback after any output still waiting."""
self._drain_output()
self.error.emit(tb)
[docs]
def request_cancel(self, reason: str = "cancelled by the user") -> bool:
"""Request a stop at the pipeline's next declared safe boundary."""
first = self.cancel_token.cancel(reason)
self.gate.resume()
return first
@property
[docs]
def app_key(self) -> str:
"""App key this job belongs to, or ``""`` for an ad-hoc job."""
return self._app_key_override or str(
getattr(self._fn, APP_KEY_ATTR, "") or ""
)
@property
[docs]
def supports_pause(self) -> bool:
"""True only when the entry point declares :func:`pausable`.
No shipped pipeline does, so this is False everywhere today —
see :class:`PauseGate` for why that is the honest answer rather
than a missing feature.
"""
return bool(getattr(self._fn, PAUSABLE_ATTR, False))
@property
[docs]
def blocks_shutdown(self) -> bool:
"""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 :mod:`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
:meth:`__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.
"""
return self._journal_enabled
[docs]
def run(self) -> None:
"""Invoked by QThread.started; runs the pipeline function to completion."""
if self.cancel_token.cancelled:
self.was_cancelled = True
self._say(
f"Cancelled before start: {self.cancel_token.reason}\n")
self.gate.resume()
self.finished.emit(False)
return
_LOCAL.gate = self.gate
journal_holder = [None]
def _forward_output(text: str) -> None:
"""Emit worker text and retain warning lines in its manifest."""
self._push_output(text)
run = journal_holder[0]
if run is not None and re.search(
r"\b(?:warning|warn)\b", text, flags=re.IGNORECASE,
):
run.record_warning(text)
redirect = _StreamRedirector(_forward_output)
stdout_router, stderr_router = _register_worker_streams(redirect)
capture_show = None
try:
if not self._capture_figures:
raise _SkipFigureCapture
import matplotlib
if matplotlib.get_backend().lower() != "agg":
matplotlib.use("Agg", force=True)
import matplotlib.pyplot as plt
from matplotlib._pylab_helpers import Gcf
worker = self
emitted_ids = set()
fig_counter = [0]
def _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.
"""
for number in list(plt.get_fignums()):
manager = Gcf.get_fig_manager(number)
if manager is not None:
yield manager.canvas.figure
preexisting_figures = tuple(_registered_figures())
preexisting_ids = {id(fig) for fig in preexisting_figures}
def _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.
"""
return (id(fig) in emitted_ids
and getattr(fig, "_spacr_emitted", False))
def _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.
"""
emitted_ids.add(id(fig))
try:
fig._spacr_emitted = True
except Exception: # noqa: BLE001
pass
def _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.
"""
for fig in _registered_figures():
if id(fig) in preexisting_ids:
continue
with figure_style(theme_target()):
already_emitted = _already_emitted(fig)
if already_emitted and not getattr(
fig, "_spacr_live_update", False):
continue
_mark_emitted(fig)
png_path = ""
try:
import tempfile
from .widgets.figure_queue import render_figure_to_png
fig_counter[0] += 1
tmp = os.path.join(
tempfile.gettempdir(),
f"spacr_fig_{os.getpid()}_{fig_counter[0]}.png")
if render_figure_to_png(fig, tmp):
png_path = tmp
except Exception:
png_path = ""
worker.figure_ready.emit(fig, png_path)
return None
capture_show = _capture_show
_register_matplotlib_show(plt, capture_show)
def _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.
"""
if _already_emitted(fig) and not getattr(
fig, "_spacr_live_update", False):
return
_mark_emitted(fig)
png_path = ""
try:
import tempfile
from .widgets.figure_queue import render_figure_to_png
fig_counter[0] += 1
tmp = os.path.join(
tempfile.gettempdir(),
f"spacr_fig_{os.getpid()}_{fig_counter[0]}.png")
if render_figure_to_png(fig, tmp):
png_path = tmp
except Exception:
png_path = ""
worker.figure_ready.emit(fig, png_path)
try:
from spacr.figure_sink import set_sink
set_sink(_publish_figure)
except Exception:
LOG.debug("could not install the figure sink", exc_info=True)
except _SkipFigureCapture:
plt = None
except Exception:
plt = None
journal_context = None
journal_run = None
try:
from spacr.run_journal import open_run
if self._journal_enabled:
if self._settings.get("hash_inputs", False):
self._say(
"Recording reproducibility input hashes…\n"
)
journal_context = open_run(
self.app_key or getattr(self._fn, "__name__", "job"),
self._settings,
)
journal_run = journal_context.__enter__()
journal_holder[0] = journal_run
self._say(
f"Reproducibility manifest: {journal_run.dir}\n"
)
lock = getattr(journal_run, "_analysis_lock", None)
if lock:
self._say(f"{lock.get('summary')}\n")
except Exception as exc:
journal_context = None
journal_run = None
message = (
"WARNING: could not open reproducibility manifest: "
f"{type(exc).__name__}: {exc}\n"
)
LOG.exception("Could not open run journal")
self._say(message)
ok = False
try:
with installed_token(self.cancel_token), responsive_gui():
self.cancel_token.checkpoint()
payload = self._fn(self._settings)
if journal_run is not None:
from spacr.run_journal import _raise_if_incomplete
_raise_if_incomplete(journal_run)
ok = True
if payload is not None:
try:
self.result_ready.emit(payload)
except Exception: # pragma: no cover - never fail a run
LOG.debug("could not deliver the pipeline result",
exc_info=True)
except PipelineCancelled as exc:
self.was_cancelled = True
message = f"Cancelled safely: {exc}\n"
LOG.info("Pipeline %s cancelled: %s", self.app_key or self._fn,
str(exc))
if journal_run is not None:
journal_run.set_status("cancelled")
journal_run.record_warning(message.strip())
self._say(message)
except SystemExit as exc:
ok = exc.code in (None, 0)
if not ok:
tb = traceback.format_exc()
LOG.error("Pipeline exited with status %r", exc.code)
if journal_run is not None:
journal_run.set_status("failed")
journal_run.error_traceback = tb
self._say_error(tb)
except Exception:
tb = traceback.format_exc()
LOG.exception("Pipeline worker failed")
if journal_run is not None:
journal_run.set_status("failed")
journal_run.error_traceback = tb
self._say_error(tb)
except BaseException:
tb = traceback.format_exc()
LOG.exception("Pipeline worker aborted")
if journal_run is not None:
journal_run.set_status("failed")
journal_run.error_traceback = tb
self._say_error(tb)
finally:
if journal_run is not None and ok:
try:
from spacr.run_journal import _raise_if_incomplete
_raise_if_incomplete(journal_run)
except Exception:
ok = False
journal_run.set_status("failed")
journal_run.error_traceback = traceback.format_exc()
self._say_error(journal_run.error_traceback)
else:
journal_run.set_status("success")
if journal_context is not None:
try:
journal_context.__exit__(None, None, None)
except Exception as exc:
LOG.exception("Could not close run journal")
self._say(
"WARNING: could not finalize reproducibility "
f"manifest: {type(exc).__name__}: {exc}\n"
)
try:
redirect.flush()
except Exception:
pass
_unregister_worker_streams(
redirect, stdout_router, stderr_router
)
try:
from spacr.figure_sink import clear_sink
clear_sink()
except Exception:
LOG.debug("could not clear the figure sink", exc_info=True)
if capture_show is not None and plt is not None:
try:
_unregister_matplotlib_show(capture_show)
except Exception:
pass
self.gate.resume()
journal_holder[0] = None
try:
delattr(_LOCAL, "gate")
except AttributeError:
pass
self._drain_output()
self.finished.emit(ok)
def _tag(app_key: str, fn: Optional[Callable]) -> Optional[Callable]:
"""Stamp ``fn`` with the app key it belongs to and return it.
:func:`make_thread` reads this back so the run registry can name the
job. Stamping beats adding an ``app_key`` argument to
``make_thread`` because every screen in the package already calls
``make_thread`` and none of them would have been updated.
The mapping is 1:1 — no two app keys resolve to the same callable —
so writing the attribute onto a shared module-level function is
unambiguous. Silently skipped for callables that reject attributes
(builtins, ``functools.partial``), which simply leaves the job
unnamed.
"""
if fn is None:
return None
try:
setattr(fn, APP_KEY_ATTR, app_key)
except (AttributeError, TypeError):
pass
return fn
def _say_what_is_wrong_with_the_settings(app_key, fn):
"""Wrap a pipeline entry so it reports settings problems before it runs.
``spacr.validate.validate_settings`` knows that ``n_job`` is not a
setting and that ``n_jobs`` is. The CLI and the batch runner both ask it.
The GUI did not: this bridge hands the settings dict straight to the
pipeline, so a settings CSV loaded into a screen ran with its typos
intact and the key did nothing. That is not hypothetical -- the
real `crop_measure_settings.csv` in use here asks for `n_job`, and a
measure run started thirty workers while the file said four.
It REPORTS and runs anyway rather than refusing. The panel itself cannot
produce an unknown key, so the only source is a file the user chose to
load, and a screen that silently declines to start would be a worse
failure than one that says what it ignored. The CLI still refuses.
Validation never breaks a run: anything it raises is swallowed, because
a broken checker must not stop work the checker was only advising on.
"""
if fn is None:
return None
@functools.wraps(fn)
def run(settings=None, *args, **kwargs):
"""Run the entry point, turning a settings problem into a message."""
if isinstance(settings, dict):
try:
from spacr.validate import (ERROR, coerce_expected_types,
validate_settings)
settings = coerce_expected_types(settings, app_key)
found = list(validate_settings(dict(settings), app_key))
except Exception: # noqa: BLE001
found = []
for problem in found:
mark = "ERROR" if problem.severity == ERROR else "WARNING"
where = f" [{problem.setting}]" if problem.setting else ""
print(f"[settings] {mark}{where}: {problem.message}")
if problem.fix:
print(f"[settings] {problem.fix}")
from spacr.resource_log import _ram_guard_scope
with _ram_guard_scope(settings):
if settings is None:
return fn(*args, **kwargs)
return fn(settings, *args, **kwargs)
return run
[docs]
def resolve_pipeline_entry(app_key: str) -> Callable[[Dict[str, Any]], Any] | None:
"""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
:func:`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 (:func:`_tag`) so the
run registry can say *which module* is running. Note that none of
these are stamped :func:`pausable` — see :class:`PauseGate`.
:param 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.
"""
from .verbose_logger import log_call
def _ret(fn):
"""Tag the entry point with its app key and its settings check."""
return _tag(app_key, _say_what_is_wrong_with_the_settings(app_key, fn))
try:
if app_key == "mask":
from spacr.core import preprocess_generate_masks
return _ret(log_call(preprocess_generate_masks))
if app_key == "timelapse":
from spacr.core import preprocess_generate_masks_timelapse
return _ret(log_call(preprocess_generate_masks_timelapse))
if app_key == "motility":
from spacr.timelapse import automated_motility_assay
return _ret(log_call(automated_motility_assay))
if app_key == "measure":
from spacr.measure import measure_crop
return _ret(log_call(measure_crop))
if app_key == "external_masks":
from spacr.external_masks import prepare_external_masks
return _ret(log_call(prepare_external_masks))
if app_key == "illumination":
from spacr.illumination import prepare_illumination_correction
return _ret(log_call(prepare_illumination_correction))
if app_key == "classify_merged":
from spacr.classify import classify
return _ret(log_call(classify))
if app_key == "classify":
from spacr.deep_spacr import deep_spacr
return _ret(log_call(deep_spacr))
if app_key == "umap":
from spacr.core import generate_image_umap
return _ret(log_call(generate_image_umap))
if app_key == "train_cellpose":
from spacr.submodules import train_cellpose
return _ret(log_call(train_cellpose))
if app_key == "cellpose_masks":
from spacr.spacr_cellpose import identify_masks_finetune
return _ret(log_call(identify_masks_finetune))
if app_key == "map_barcodes":
from spacr.sequencing import generate_barecode_mapping
return _ret(log_call(generate_barecode_mapping))
if app_key == "ml_analyze":
from spacr.ml import generate_ml_scores
return _ret(log_call(generate_ml_scores))
if app_key == "regression":
from spacr.ml import perform_regression
return _ret(log_call(perform_regression))
if app_key == "recruitment":
from spacr.submodules import analyze_recruitment
return _ret(log_call(analyze_recruitment))
if app_key == "activation":
from spacr.deep_spacr import generate_activation_map
return _ret(log_call(generate_activation_map))
if app_key == "foreign":
from spacr.foreign import import_project
return _ret(log_call(import_project))
if app_key == "align":
from spacr.align import align_folder
return _ret(log_call(align_folder))
if app_key == "ops":
from spacr.ops_engine import run_ops
return _ret(log_call(run_ops))
if app_key == "convert":
from spacr.convert import convert_folder
return _ret(log_call(convert_folder))
if app_key == "invasion":
from spacr.submodules import analyze_invasion
return _ret(log_call(analyze_invasion))
if app_key == "replication":
from spacr.submodules import analyze_replication
return _ret(log_call(analyze_replication))
if app_key == "analyze_plaques":
from spacr.submodules import analyze_plaques
return _ret(log_call(analyze_plaques))
if app_key == "barcode_qc":
from spacr.sequencing_qc import barcode_qc
return _ret(log_call(barcode_qc))
if app_key == "explain_cv":
from spacr.surrogate import run_explain_cv
return _ret(log_call(run_explain_cv))
if app_key == "anndata_export":
from spacr.anndata_export import run_anndata_export
return _ret(log_call(run_anndata_export))
from .app import registered_entry
registered = registered_entry(app_key)
if registered is not None:
return _ret(log_call(registered))
from spacr.plugins import get_app, load_object
plugin_app = get_app(app_key)
if plugin_app is not None:
entry = load_object(plugin_app.entrypoint)
if not callable(entry):
raise TypeError(
f"Plugin entry point {plugin_app.entrypoint!r} is not callable"
)
return _ret(log_call(entry))
except Exception:
LOG.exception("Could not resolve pipeline entry for %s", app_key)
return None
return None
#: Stack for a pipeline worker thread, in bytes.
#:
#: A pthread on macOS gets 512 KB by default and Qt does not raise it, so a
#: QThread starts with ~1/16th of the main thread's 8 MB. The pipeline is the
#: same code either way, and parts of it want a *lot* of stack: a classify run
#: died with SIGBUS, "Thread stack size exceeded due to excessive recursion",
#: in ``___chkstk_darwin`` under OpenBLAS's ``dgetrf_parallel``, reached from
#: ``np.linalg.inv`` at ``spacr/ml.py`` (the Mahalanobis inverse covariance).
#: The crash report put that thread's stack at 544 KB. Nothing was recursing —
#: ``chkstk`` is the probe that discovers the guard page, and the report's
#: "excessive recursion" wording is a guess macOS makes about any stack
#: overflow.
#:
#: This is address space, not memory: pages commit as they are touched, so a
#: generous number costs nothing until it is used. 64 MB is the main thread's
#: 8 MB with room for LAPACK's blocked kernels on a wide feature matrix.
#: ``SPACR_WORKER_STACK_MB`` overrides it for anyone who needs to.
WORKER_STACK_BYTES = 64 * 1024 * 1024
def _widen_worker_stack(thread: "QThread") -> None:
"""Give a worker thread a stack the pipeline can actually run in.
Must be called before ``start()`` — Qt ignores ``setStackSize`` on a
running thread. Wrapped: if the platform refuses the size, the thread
keeps the default and the run proceeds exactly as it did before, which is
the behaviour this replaces (INVARIANTS §10).
"""
megabytes = os.environ.get("SPACR_WORKER_STACK_MB", "").strip()
size = WORKER_STACK_BYTES
if megabytes:
try:
requested = int(megabytes)
except (TypeError, ValueError):
requested = 0
if requested > 0:
size = requested * 1024 * 1024
try:
thread.setStackSize(size)
except Exception:
pass
[docs]
def 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["QThread", PipelineWorker]:
"""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 :func:`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.
:param fn: callable invoked as ``fn(settings)`` on the worker thread.
:param settings: the single argument handed to ``fn``.
:param app_key: overrides the key stamped on ``fn`` by
:func:`resolve_pipeline_entry`. Only needed for ad-hoc jobs that
want to show up on Home under a name of their own.
:param journal: create a reproducibility record. Disable only for
read-only UI housekeeping that is not an analysis run.
:param 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.
"""
if capture_figures and "matplotlib.pyplot" not in sys.modules:
gc_was_enabled = gc.isenabled()
try:
if gc_was_enabled:
gc.disable()
import matplotlib.pyplot # noqa: F401
except Exception: # noqa: BLE001
LOG.debug("matplotlib.pyplot could not be pre-imported",
exc_info=True)
finally:
if gc_was_enabled:
gc.enable()
thread = QThread()
_widen_worker_stack(thread)
allocation = apply_worker_budget(settings)
worker = PipelineWorker(
fn, settings, worker_count=allocation, app_key=app_key,
journal=journal, capture_figures=capture_figures,
)
worker.user_visible = bool(user_visible)
worker.moveToThread(thread)
thread.started.connect(worker.run)
worker.finished.connect(thread.quit, Qt.DirectConnection)
thread.finished.connect(thread.deleteLater)
reg = registry()
handle = RunHandle(app_key or worker.app_key, worker, thread, parent=reg)
reg.register(handle)
thread.finished.connect(handle.retire)
return thread, worker