Source code for spacr.qt.screens.queue

"""
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)