"""
Queue screen — dashboard for the plate queue.
Layout:
┌───────────────────────────────────────────────────────────────┐
│ [Add plate] [Import CSV…] [Clear finished] [Run] [Stop] │
├───────────────────────────────────────────────────────────────┤
│ ID │ App │ Label │ Status │ Elapsed │ │
│ ... │ mask │ /data/plateA │ success │ 42.1 s │ [Remove] │
│ ... │ mask │ /data/plateB │ running │ 08.5 s │ │
│ ... │ mask │ /data/plateC │ queued │ │ [Remove] │
└───────────────────────────────────────────────────────────────┘
Each row reflects a :class:`spacr.qt.plate_queue.QueueItem`. The
screen owns a :class:`spacr.qt.plate_queue.PlateQueue` and a single
worker QThread that walks items sequentially.
"""
from __future__ import annotations
import logging
import time
from typing import Optional
from PySide6.QtCore import Qt, QThread, QTimer, Signal
from PySide6.QtWidgets import (
QAbstractItemView, QFileDialog, QHBoxLayout, QHeaderView, QLabel,
QMessageBox, QPushButton, QTableWidget, QTableWidgetItem,
QVBoxLayout, QWidget,
)
from spacr.cancellation import (
CancellationToken, PipelineCancelled, installed_token,
)
from ..plate_queue import (
PlateQueue, QueueItem, Status, import_plates_from_csv,
)
from ..widgets.collapsible_splitter import FoldSection
from ..widgets.sortable_table import install_sorting, table_item
LOG = logging.getLogger("spacr.qt.queue_screen")
class _QueueRunner(QThread):
"""Walks the queue sequentially in a background thread.
Emits :attr:`item_state_changed` on every status transition so
the table can refresh a single row without a full rebuild.
:param queue: the plate queue to run. Held, not copied, and read as the
run proceeds -- so items may be added to it after the runner starts.
:param parent: parent object; ownership only. It does NOT stop the
thread: see :meth:`abort`, which is what teardown must call.
"""
item_state_changed = Signal(str)
queue_finished = Signal()
def __init__(self, queue: PlateQueue, parent=None):
"""Hold the queue and start not stopped."""
super().__init__(parent)
self._queue = queue
self._stop = False
#: Cooperative Stop for the item that is running *right now*.
#: :meth:`stop` deliberately does not touch it — that control means
#: "no more items after this one". :meth:`abort` does, and is what
#: teardown needs: without it, closing the screen left an
#: hours-long pipeline running inside a QThread parented to a
#: widget Qt was about to destroy, which is a core dump, not an
#: exception.
self._token = CancellationToken()
def stop(self) -> None:
"""Ask the runner to exit after the current item completes."""
self._stop = True
def abort(self, reason: str = "queue screen closed") -> None:
"""Stop after the current item's next declared safe boundary.
Shipped pipeline entry points poll :func:`spacr.cancellation.
checkpoint` at field/trial/job boundaries, so this stops the run
without ever interrupting a half-written mask or a partial insert.
"""
self._stop = True
self._token.cancel(reason)
def run(self) -> None:
"""Run each queued plate in turn until the queue empties or is stopped.
EVERY EMIT GOES THROUGH ``emit_safely``. A queue run outlives the screen
that started it -- closing the window mid-run leaves this thread
emitting at a destroyed C++ object, which raises out of a ``QThread.run``
override, and an exception out of a virtual override ABORTS the process
rather than failing the run.
The database updates stay unguarded on purpose: a queue item that
finished must be recorded as finished whether or not anyone is watching,
and sqlite does not care that the window closed.
"""
from ..bridge import emit_safely, resolve_pipeline_entry
while not self._stop:
item = self._queue.next_queued()
if item is None:
break
self._queue.update(item.id, status=Status.RUNNING,
start_ts=time.time())
emit_safely(self.item_state_changed, item.id)
try:
fn = resolve_pipeline_entry(item.app_key)
if fn is None:
raise RuntimeError(
f"no pipeline for app_key={item.app_key!r}")
with installed_token(self._token):
self._token.checkpoint()
fn(item.settings)
except PipelineCancelled as e:
self._queue.update(item.id, status=Status.QUEUED,
end_ts=None, error="")
LOG.info("queue item %s cancelled: %s", item.id, str(e))
emit_safely(self.item_state_changed, item.id)
break
except Exception as e:
LOG.warning("queue item %s failed: %s", item.id, e,
exc_info=True)
self._queue.update(item.id, status=Status.FAILED,
end_ts=time.time(), error=str(e))
emit_safely(self.item_state_changed, item.id)
continue
self._queue.update(item.id, status=Status.SUCCESS,
end_ts=time.time())
emit_safely(self.item_state_changed, item.id)
emit_safely(self.queue_finished)
_COLUMNS = ("ID", "App", "Label", "Status", "Elapsed", "")
[docs]
class QueueScreen(QWidget):
"""Main widget rendering the plate queue.
:param queue: the queue to render. ``None`` loads spaCR's own, which is
the ordinary case; a test passes one built on a temporary path so it
does not disturb the user's real queue.
:param parent: parent widget.
"""
queue_size_changed = Signal(int)
def __init__(self, queue: Optional[PlateQueue] = None, parent=None):
"""Build the queue screen and start its elapsed-time tick.
:param queue: the queue to show; ``None`` builds an empty one.
:param parent: parent widget, or ``None``.
"""
super().__init__(parent)
self._queue = queue if queue is not None else PlateQueue()
self._runner: Optional[_QueueRunner] = None
self._build_ui()
if queue is None:
self._claim_the_queue()
from ..dnd import install_dropzone
from ..dnd_handlers import get_handler
install_dropzone(self, get_handler("queue"), self)
self._refresh_table()
self._tick = QTimer(self)
self._tick.setInterval(1000)
self._tick.timeout.connect(self._refresh_elapsed_only)
self._tick.start()
def _build_ui(self):
"""Lay out the toolbar and the plate table.
The table folds under a "Queue" heading (item 471), which then sits
at the bottom of the screen.
``Add current plate`` is deliberately left unwired here: the settings it
adds belong to whichever app screen is active, so ``MainWindow`` connects
it -- see ``wire_add_current``.
"""
outer = QVBoxLayout(self)
outer.setContentsMargins(24, 24, 24, 24)
outer.setSpacing(12)
header = QLabel("Plate Queue")
header.setObjectName("DisplayHeading")
outer.addWidget(header)
subtitle = QLabel(
"Chain multiple plates through the same pipeline. "
"The runner processes one plate at a time and picks up "
"where it left off if you close the app.")
subtitle.setObjectName("Muted")
outer.addWidget(subtitle)
bar = QHBoxLayout()
self._btn_add = QPushButton("Add current plate", self)
self._btn_import = QPushButton("Import CSV…", self)
self._btn_clear = QPushButton("Clear finished", self)
self._btn_run = QPushButton("Run queue", self)
self._btn_stop = QPushButton("Stop", self)
self._btn_stop.setEnabled(False)
for b in (self._btn_add, self._btn_import, self._btn_clear,
self._btn_run, self._btn_stop):
bar.addWidget(b)
bar.addStretch(1)
outer.addLayout(bar)
self._btn_import.clicked.connect(self._on_import)
self._btn_clear.clicked.connect(self._on_clear_finished)
self._btn_run.clicked.connect(self.start_runner)
self._btn_stop.clicked.connect(self.stop_runner)
self._table = QTableWidget(self)
install_sorting(self._table)
self._table.setColumnCount(len(_COLUMNS))
self._table.setHorizontalHeaderLabels(_COLUMNS)
self._table.horizontalHeader().setStretchLastSection(True)
self._table.horizontalHeader().setSectionResizeMode(
2, QHeaderView.Stretch)
self._table.verticalHeader().setVisible(False)
self._table.setEditTriggers(QAbstractItemView.NoEditTriggers)
self._table.setSelectionBehavior(QAbstractItemView.SelectRows)
self._table_section = FoldSection(self._table, "Queue",
persist_key="queue/Queue")
outer.addWidget(self._table_section, 1)
def _claim_the_queue(self) -> None:
"""Lock the shared queue file, or open it read-only if another
spaCR window already has it.
Read-only shows the queue but never writes it or runs it: the
buttons that would are disabled and the queue stops saving.
"""
from ..crash_recovery import _claim_project
from ..i18n import tr
answer = _claim_project(self, getattr(self._queue, "_path", None))
if answer in ("locked", "continue"):
return
self._queue.read_only = True
for button in (self._btn_add, self._btn_import, self._btn_clear,
self._btn_run):
button.setEnabled(False)
button.setToolTip(tr(
"Read-only: another spaCR window has this queue open."))
[docs]
def wire_add_current(self, callback):
"""Route the "Add current plate" button through ``callback``.
MainWindow supplies a callback that snapshots the currently
active AppScreen's settings dict and returns
``(app_key, settings_dict)``. This screen calls
:meth:`add_item` with that pair.
:param callback: a no-argument callable returning ``(app_key,
settings_dict)``. It is called on every click; an exception it
raises, or a settings dict without ``src``, is shown in a message
box instead of adding an item.
"""
def _on_click():
"""Add the current screen's settings to the queue."""
try:
app_key, settings = callback()
except Exception as e:
QMessageBox.warning(self, "Queue",
f"Couldn't add current plate: {e}")
return
if not settings.get("src"):
QMessageBox.information(self, "Queue",
"The active app has no `src` set — nothing to enqueue.")
return
self.add_item(app_key, settings)
self._btn_add.clicked.connect(_on_click)
[docs]
def add_item(self, app_key: str, settings: dict) -> QueueItem:
"""Build a queue item from settings and add it.
:param app_key: which module the item runs.
:param settings: the settings that run uses.
:returns: the queued item.
"""
item = QueueItem.build(app_key, settings)
self._queue.add(item)
self._refresh_table()
self.queue_size_changed.emit(len(self._queue))
return item
[docs]
def queue(self) -> PlateQueue:
"""The plate queue this screen shows.
:returns: the queue.
"""
return self._queue
[docs]
def start_runner(self):
"""Start working through the queue, unless it is already running.
IDEMPOTENT ON PURPOSE. The button can be pressed twice, and a second
runner over one queue would run every item twice.
"""
if self._runner is not None and self._runner.isRunning():
return
if self._queue.next_queued() is None:
QMessageBox.information(self, "Queue",
"Nothing to run — every item is already finished.")
return
self._runner = _QueueRunner(self._queue, self)
self._runner.item_state_changed.connect(self._on_item_changed)
self._runner.queue_finished.connect(self._on_runner_done)
self._btn_run.setEnabled(False)
self._btn_stop.setEnabled(True)
self._runner.start()
[docs]
def stop_runner(self):
"""Ask the runner to stop after the item it is on.
AFTER, not during: a half-written plate is worse than a queue that
takes another minute to come to rest.
"""
if self._runner is not None and self._runner.isRunning():
self._runner.stop()
def _on_runner_done(self):
"""Swap Run back for Stop and redraw the table once the runner stops."""
self._btn_run.setEnabled(True)
self._btn_stop.setEnabled(False)
self._refresh_table()
[docs]
def closeEvent(self, event):
"""Stop the queue runner before Qt destroys this screen.
``_QueueRunner`` is parented to this widget and executes whole
pipelines, so its thread can be hours long. Destroying the screen
with it running deletes a live QThread, which Qt answers with
``qFatal("QThread: Destroyed while thread is still running")``.
It also does not go through :func:`spacr.qt.bridge.make_thread`, so
the process-wide run registry — and therefore
``MainWindow.closeEvent``'s drain — has never been able to see it.
:param event: the close event; it is passed on to the base class once
the queue runner has been aborted and drained.
"""
from ..bridge import drain_thread
self._tick.stop()
runner, self._runner = self._runner, None
if runner is not None:
try:
runner.abort()
except (AttributeError, RuntimeError):
pass
if not drain_thread(runner, timeout_ms=5000):
runner.setParent(None)
super().closeEvent(event)
def _on_item_changed(self, item_id: str):
"""Redraw the table after one item changed.
:param item_id: which item changed; the whole table is re-read either
way, so it is not used to narrow the redraw.
"""
self._refresh_table()
def _on_import(self):
"""Import plates from a CSV and add them to the queue.
A file that cannot be read is reported in a dialog rather than raised:
picking the wrong CSV is a normal mistake, not a crash.
"""
path, _ = QFileDialog.getOpenFileName(
self, "Import plates from CSV", "", "CSV files (*.csv)")
if not path:
return
try:
items = import_plates_from_csv(path, base_settings={},
app_key="mask")
except Exception as e:
QMessageBox.warning(self, "Queue import",
f"Couldn't import {path}:\n{e}")
return
for it in items:
self._queue.add(it)
self._refresh_table()
self.queue_size_changed.emit(len(self._queue))
QMessageBox.information(self, "Queue import",
f"Added {len(items)} plate(s) from {path}.")
def _on_clear_finished(self):
"""Drop every finished plate from the queue."""
n = self._queue.clear_finished()
self._refresh_table()
self.queue_size_changed.emit(len(self._queue))
if n:
self._table.selectRow(-1)
def _refresh_table(self):
"""Rebuild the plate table, one row per queued item.
Each row carries its own Remove button, disabled while that plate is
running.
"""
items = self._queue.items()
self._table.setRowCount(len(items))
for row, item in enumerate(items):
self._table.setItem(row, 0, table_item(item.id))
self._table.setItem(row, 1, table_item(item.app_key))
self._table.setItem(row, 2, table_item(item.label))
status_item = table_item(item.status.value)
self._set_status_color(status_item, item.status)
self._table.setItem(row, 3, status_item)
elapsed = item.elapsed_s
self._table.setItem(row, 4, table_item(
"" if elapsed is None else f"{elapsed:.1f} s"))
btn = QPushButton("Remove", self)
btn.setEnabled(item.status != Status.RUNNING)
btn.clicked.connect(
lambda _c=False, iid=item.id: self._on_remove(iid))
self._table.setCellWidget(row, 5, btn)
def _refresh_elapsed_only(self):
"""Tick the elapsed column of the running plates, and only that column.
Driven by a one-second timer. Rebuilding the whole table every second
would churn it and lose the user's selection with it.
"""
items = self._queue.items()
for row, item in enumerate(items):
if item.status != Status.RUNNING or row >= self._table.rowCount():
continue
e = item.elapsed_s
if e is not None:
self._table.setItem(row, 4, table_item(f"{e:.1f} s"))
def _on_remove(self, item_id: str):
"""Remove one plate from the queue.
A running plate is refused with a dialog: removing the row would leave
the runner working on an item the queue no longer knows about.
:param item_id: the plate to remove.
"""
item = self._queue.find(item_id)
if item is not None and item.status == Status.RUNNING:
QMessageBox.warning(self, "Queue",
"Can't remove a plate while it's running — stop the queue first.")
return
self._queue.remove(item_id)
self._refresh_table()
self.queue_size_changed.emit(len(self._queue))
@staticmethod
def _set_status_color(item: QTableWidgetItem, status: Status):
"""Colour a status cell by its status.
:param item: the cell to colour.
:param status: the item's status; an unrecognised one falls back to
black rather than leaving the previous colour in place.
"""
colors = {
Status.QUEUED: Qt.darkGray,
Status.RUNNING: Qt.blue,
Status.SUCCESS: Qt.darkGreen,
Status.FAILED: Qt.darkRed,
Status.SKIPPED: Qt.gray,
}
c = colors.get(status, Qt.black)
item.setForeground(c)