diff --git a/CLAUDE.md b/CLAUDE.md index f055d6c..7d32792 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -101,6 +101,7 @@ shape, so read the one closest to what you are writing: | runs the ROI manager's segment, evaluate and analyse as one | `npc_analysis.chain.yaml` (a shipped chain: an evaluator step runs over every ROI before the next; `docs/batch.md`, "ROIs: segment, evaluate, analyse") | | opens an interactive window of the GUI's | `bead_calibration.py` (`Plugin.window`: the button opens what the GUI registered in `session.window_openers`) | | runs an external program in its own Python | `uipsf.py` (`smappy.uipsf`: a worker script under that program's interpreter, progress read from its output, the result converted into smappy's own file) | +| is run by MicroClaw during an acquisition | `gui/live_session.py` (the GUI on a fit another program controls through `Context.stop` and `writer_finished`; the package and its protocol runner are `ries-lab/smappy-microclaw`) | A chain of plugins that runs as one, and running one over many files (`smappy-batch`, the batch window): `docs/batch.md`. @@ -184,6 +185,9 @@ written this way; it is also what makes the tests readable. never assume a session exists. * `ctx.report(text)` for progress, `ctx.emit(event, payload)` to hand partial results on. Both are no-ops when nobody is listening. +* `ctx.stop` and `ctx.writer_finished` -- `threading.Event`s or None, set from + another thread: end the run early keeping what it has, and "the live source + is complete". The Fit plugin honours both; a long plugin may check `stop`. * The **grouped** table (one row per blink) is a table of its own, with its own filter, at `session.layers[i].state.sets["grouped"]`; it exists only once the user has switched that layer to grouped, and `state.grouped_stale` says it no diff --git a/NOTES.md b/NOTES.md index ca88e24..28798aa 100644 --- a/NOTES.md +++ b/NOTES.md @@ -728,9 +728,30 @@ mid-write raises instead of returning half an image -- and the file list is re-globbed on each poll, so the `_1.ome.tif` that appears when the current file fills up is picked up without reopening anything. +**Unless somebody says when it ended.** A program that drives the microscope +(MicroClaw) knows when the writer has finished, and then a pause -- a refocus, +a laser change, a time-lapse interval -- must not end the fit: `finished`, an +event handed to the watch, replaces the timeout altogether, and once it is set +one more pass reads everything and the stream ends. A longer timeout would +only move the pause that ends the fit too early, which is why MicroClaw's +protocol (its `design/84`) asks for the event and not for a number. The +Fit plugin takes it from `Context.writer_finished`, and `Context.stop` ends a +run early with its file closed and nothing finished (no drift correction of +half an acquisition; the caller allows ten seconds). +`smappy.gui.live_session` is the GUI opened on such a run: the same session +and windows, the run started without its preflight question, and the fit's +file added to `Session.protected`, because the caller hands it on with a +digest while the window stays open and a Ctrl+S would change what that digest +vouches for. + **The engine is flushed on a timer as well as when its ROI buffer fills.** The buffer holds 15000 ROIs; a sparse sample would take minutes to fill it, and -nothing would appear in the meantime. +nothing would appear in the meantime. The timer is checked once per block, so +during a pause the watch sends empty blocks (`idle_blocks`) for it to be checked +on -- without them, what was buffered before a pause stayed off the screen +until the pause was over. The GUI's live fit flushes every 5 s +(`LIVE_FLUSH_SECONDS`); `smappy.live.LiveFit` still has its own loop and does +not send idle blocks yet. Checked against the offline path on 300 frames of the astigmatic dataset, replayed frame by frame into a growing two-file series: same 28724 diff --git a/pyproject.toml b/pyproject.toml index 03c7e76..19d37ae 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -12,7 +12,7 @@ build-backend = "setuptools.build_meta" # package (a Smappee energy-monitor wrapper, last released in 2018). The import # name is unaffected: `pip install smappy-smlm`, then `import smappy`. name = "smappy-smlm" -version = "0.2.0" +version = "0.3.0" description = "Single-molecule localization fitting pipeline (Python port of the SMAP fast-simple workflow)" readme = "README.md" license = "BSD-3-Clause" diff --git a/src/smappy/gui/app.py b/src/smappy/gui/app.py index 1879289..e0c553a 100644 --- a/src/smappy/gui/app.py +++ b/src/smappy/gui/app.py @@ -817,7 +817,7 @@ def open_image(self) -> None: def save(self) -> None: from ..io.formats import writer_for - if self.session.path is None: + if self.session.path is None or self.session.is_protected(self.session.path): return self.save_as() try: writer_for(self.session.path) @@ -830,11 +830,18 @@ def save(self) -> None: def save_as(self) -> None: start = str(self.session.path) if self.session.path else folders.start() + if self.session.is_protected(start): + here = Path(start) # offer a new name beside the kept file + start = str(here.with_name(f"{here.stem}_edited{here.suffix}")) path, _ = QFileDialog.getSaveFileName(self, "Save localizations", start, "HDF5 (*.h5 *.hdf5)") if path: folders.remember(path) - self.session.save(path) + try: + self.session.save(path) + except ValueError as error: + QMessageBox.warning(self, "Not saved", str(error)) + return self._on_session("locs") diff --git a/src/smappy/gui/live_session.py b/src/smappy/gui/live_session.py new file mode 100644 index 0000000..0558072 --- /dev/null +++ b/src/smappy/gui/live_session.py @@ -0,0 +1,252 @@ +"""The GUI, opened on a fit that is already running, for a program that drives it. + +MicroClaw runs the microscope and starts an analysis while the acquisition is +being written; the SMAPpy package it installs (`ries-lab/smappy-microclaw`) +calls this to open the ordinary GUI on that dataset. It is the same session, +the same two windows and the same Fit plugin a user would start by hand -- the +point is that the person at the microscope sees the reconstruction build up +and can work with it in the program they already know -- with three things a +hand-started fit does not need: + +* **The run is started here, with no question asked.** Nobody may be at the + keyboard, so the plugin's preflight is skipped; the caller checks what it + would have asked before calling. +* **It is controlled from another thread.** `stop` ends the fit early and + keeps the file; `writer_finished` says the acquisition is complete, after + which the watch reads what is left and the fit ends -- not after an idle + timeout, because a pause in an acquisition is not its end + (`smappy.io.watch`). Both are `threading.Event`s, so a reader thread can set + them without touching Qt. +* **The fit's file is frozen once it is done.** The caller hands it on with a + digest, and the window stays open, so the session refuses to save over it + (`Session.protected`); what the user does afterwards goes to a new file. + +`on_finished` is called exactly once, on the GUI thread, when the fit has +ended and its file is closed -- also when the windows were closed while it ran, +which is a cancellation. `run` returns only once the windows are closed. +""" +from __future__ import annotations + +import sys +import threading +import time +import traceback +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any, Callable, Dict, List, Optional + +from .. import plugins +from ..plugins import param_specs, settings_from, settings_values + + +@dataclass +class LiveOutcome: + """How a fit started by `LiveSession` ended.""" + + state: str # "succeeded", "cancelled" or "failed" + complete: bool # every frame the writer wrote was read + path: Optional[Path] = None # the fit's file, closed; None: nothing saved + stats: Dict[str, Any] = field(default_factory=dict) + error: str = "" # one line, for "failed" + traceback: str = "" + + +def fit_settings(plugin_path: str, values: Optional[Dict[str, Any]] = None): + """The settings of the fitter at ``plugin_path``: its defaults under ``values``. + + ``values`` is a flat dotted map (``{"fit.roisize": 13}``). Unlike + `plugins.settings_from`, which forgives a saved name the plugin no longer + has, this refuses one: the values come from a call somebody just made, and + a misspelt camera offset that silently kept its default would be a wrong + fit with nothing to say so. + """ + cls = plugins.get(plugin_path) + values = dict(values or {}) + known = set(_dotted_names(cls.Settings)) + unknown = sorted(set(values) - known) + if unknown: + raise ValueError(f"{plugin_path} has no setting {', '.join(unknown)}; " + f"it has {', '.join(sorted(known))}") + settings = settings_from(cls.Settings, values) + got = settings_values(settings) + refused = sorted(k for k, v in values.items() if got.get(k) != v) + if refused: + # settings_from falls back to the defaults on a value the class + # refuses; here that is an error, not a reset + raise ValueError(f"{plugin_path} refused {', '.join(refused)}") + return settings + + +def _dotted_names(settings_cls, prefix: str = "") -> List[str]: + out = [] + for name, spec in param_specs(settings_cls).items(): + if spec.children is not None: + out += _dotted_names(spec.type, f"{prefix}{name}.") + else: + out.append(prefix + name) + return out + + +class LiveSession: + """The smappy GUI with the fitter at ``plugin_path`` running ``settings``.""" + + def __init__(self, plugin_path: str, settings, *, + on_finished: Optional[Callable[[LiveOutcome], None]] = None, + on_progress: Optional[Callable[[str], None]] = None, + on_closed: Optional[Callable[[], None]] = None, + stop: Optional[threading.Event] = None, + writer_finished: Optional[threading.Event] = None): + self.plugin_path = plugin_path + self.settings = settings + self.on_finished = on_finished + self.on_progress = on_progress + self.on_closed = on_closed + # the caller's own events when it has them -- a protocol reader that + # sets them as the messages arrive -- else ours, set through these + self.stop = stop or threading.Event() + self.writer_finished = writer_finished or threading.Event() + self.outcome: Optional[LiveOutcome] = None + self.session = None + self.control = None + self.render = None + self.panel = None + self._window = None + + # ------------------------------------------------------------- running + def open(self) -> None: + """Build the windows and start the fit; the event loop is the caller's. + + `run` is this plus the loop. Split so a test can drive the loop + itself. + """ + from PySide6.QtWidgets import QApplication + + from ..session import Session + from .app import ControlWindow, RenderWindow + from .collector import collect_on_gui_thread + from .widgets import apply_style + + app = QApplication.instance() or QApplication([sys.argv[0]]) + apply_style(app) + collect_on_gui_thread(app) + self.session = Session() + out = self.settings.output.resolve(self.settings.source.path) + if out is not None: + self.session.protected.add(Path(out)) + self.render = RenderWindow(self.session) + self.control = ControlWindow(self.session, self.render) + screen = app.primaryScreen().availableGeometry() + self.control.move(screen.left(), screen.top()) + self.render.move(screen.left() + self.control.width() + 20, screen.top()) + self.render.show() + self.control.show() + + self.panel = self._panel() + self.panel.form.set(self.settings) + self.panel.ended.connect(self._ended) + self.panel.progressed.connect(self._progressed) + if not self.panel.start_run(ask=False, stop=self.stop, + writer_finished=self.writer_finished): + self._finish(LiveOutcome("failed", False, + error=self.panel.status.text() or + "the fit did not start")) + + def run(self) -> int: + """Open, and run the event loop until the windows are closed.""" + from PySide6.QtWidgets import QApplication + + try: + self.open() + except Exception as error: + # a failure to build the windows is the fit's failure: the caller + # is owed an outcome whatever happens + self._finish(LiveOutcome("failed", False, error=_one_line(error), + traceback=traceback.format_exc())) + if self.on_closed is not None: + self.on_closed() + return 1 + code = QApplication.instance().exec() + self.wait() + if self.on_closed is not None: + self.on_closed() + return code + + def wait(self) -> None: + """After the windows have closed: end a fit that is still running. + + Closing the windows before the fit has finished is a cancellation -- + the user has decided not to look -- and the fit closes its file and + reports as a stopped run does. Its last signals are queued for the + GUI thread, so they are delivered here by hand. + """ + from PySide6.QtWidgets import QApplication + + if self.outcome is None: + self.stop.set() + app = QApplication.instance() + while self.outcome is None: + app.processEvents() + time.sleep(0.02) + thread = getattr(self.panel, "_thread", None) + if thread is not None: + thread.wait() + + # ------------------------------------------------------------ internals + def _panel(self): + """The fitter's panel in the Localize tab, or a window of its own. + + The tab is where a user would look for it; a workspace that has + unpinned it still gets the fit, in the window the Plugins menu opens. + """ + for index in range(self.control.tabs.count()): + tab = self.control.tabs.widget(index) + slots = getattr(tab, "slots", None) + if not slots: + continue + for instance in tab.tab.instances: + if instance.plugin != self.plugin_path: + continue + panel = slots[instance.id].ensure() + if panel is not None: + self.control.tabs.setCurrentIndex(index) + tab.open_section(instance.id) + return panel + from .plugin_panel import PluginWindow + self._window = PluginWindow(plugins.get(self.plugin_path), self.session, + self.control) + self._window.show() + return self._window.panel + + def _progressed(self, text: str) -> None: + if self.on_progress is not None and self.outcome is None: + self.on_progress(text) + + def _ended(self, kind: str, payload) -> None: + if self.outcome is not None: + return # a later run the user started by hand + if kind == "failed": + lines = str(payload).strip().splitlines() + # stopped before it had a file -- while the dataset was still + # being waited for, say -- is still a stop, not a broken fit + state = "cancelled" if self.stop.is_set() else "failed" + self._finish(LiveOutcome(state, False, + error=lines[-1] if lines else "failed", + traceback=str(payload))) + return + data = payload.data or {} + stopped = bool(data.get("stopped")) + live = self.settings.source.live + complete = not stopped and (self.writer_finished.is_set() or not live) + path = data.get("path") + self._finish(LiveOutcome("cancelled" if stopped else "succeeded", complete, + path=Path(path) if path else None, + stats=dict(data.get("stats") or {}))) + + def _finish(self, outcome: LiveOutcome) -> None: + self.outcome = outcome + if self.on_finished is not None: + self.on_finished(outcome) + + +def _one_line(error: BaseException) -> str: + return f"{type(error).__name__}: {error}".splitlines()[0] diff --git a/src/smappy/gui/plugin_panel.py b/src/smappy/gui/plugin_panel.py index a02e24a..ddfdfd4 100644 --- a/src/smappy/gui/plugin_panel.py +++ b/src/smappy/gui/plugin_panel.py @@ -62,6 +62,10 @@ class PluginPanel(QWidget): # is what keeps `_on_progress` and `_on_stream` off that thread. progressed = Signal(str) streamed = Signal(str, object) + # ("done", Result) or ("failed", traceback) once a run or preview is over + # and the panel has dealt with it: what a caller that started the run + # itself (`gui.live_session`) waits for + ended = Signal(str, object) def __init__(self, plugin_cls: Type[Plugin], session: Session, parent=None): super().__init__(parent) @@ -318,16 +322,26 @@ def open_window(self) -> None: self.status.setText("") opener(settings) - def _start(self, job: str, message: str, **kwargs) -> None: + def start_run(self, ask: bool = True, **context) -> bool: + """Run with the form's settings; whether the run started. + + ``context`` goes into the run's `Context` -- `stop` and + `writer_finished` for a run something else controls -- and ``ask=False`` + skips the preflight question, for a run nobody is there to answer. + """ + return self._start("run", "running...", context=context, ask=ask) + + def _start(self, job: str, message: str, context=None, ask: bool = True, + **kwargs) -> bool: try: settings = self.form.value() except ValueError as e: self.status.setText(f"bad value: {e}") - return - if job == "run": + return False + if job == "run" and ask: go, chosen = self._preflight(settings) if not go: - return + return False if chosen is not settings: # a preflight choice may hand back other settings than the # form's; showing them is what keeps the run from being a @@ -338,7 +352,8 @@ def _start(self, job: str, message: str, **kwargs) -> None: # selection now, so the worker cannot race a live fit rebinding them context = self.session.context(progress=self.progressed.emit, stream=self.streamed.emit, - grouping=self.plugin.grouping) + grouping=self.plugin.grouping, + **(context or {})) self._job = job for button in self._buttons(): button.setEnabled(False) @@ -356,6 +371,7 @@ def _start(self, job: str, message: str, **kwargs) -> None: for sig in (self._worker.done, self._worker.failed): sig.connect(self._thread.quit) self._thread.start() + return True def _preflight(self, settings): """Ask the plugin whether this run wants agreeing to first. @@ -514,6 +530,7 @@ def _on_done(self, result: Result) -> None: "show the plugin's result figure (the drift curves, say)") if self._job == "preview" and result.figures(): self.plot(raise_window=not live) + self.ended.emit("done", result) def _on_failed(self, text: str) -> None: live, self._live_running = self._live_running, False @@ -526,9 +543,11 @@ def _on_failed(self, text: str) -> None: # few localizations in it, which is not a failure to report: the # next position will ask again self.status.setText("live: " + text.strip().splitlines()[-1][:60]) + self.ended.emit("failed", text) return self.output.appendPlainText(text.strip().splitlines()[-1]) self.status.setText("failed") + self.ended.emit("failed", text) def show_text(self) -> None: """What the run reported, in a window that can hold it.""" diff --git a/src/smappy/io/ndtiff.py b/src/smappy/io/ndtiff.py index 33417f8..1ebe449 100644 --- a/src/smappy/io/ndtiff.py +++ b/src/smappy/io/ndtiff.py @@ -242,14 +242,17 @@ def _block(self, first: int, last: int) -> np.ndarray: return np.stack([self.frame(i) for i in range(first, last)]) def watch(self, chunk: int = 100, settings=None, start: int = 0, - stop: Optional[int] = None, stop_event=None, on_wait=None + stop: Optional[int] = None, stop_event=None, on_wait=None, + finished=None, idle_blocks: bool = False ) -> Iterator[Tuple[int, np.ndarray]]: """Yield blocks as the acquisition writes them. The index gains a record only once an image is complete, so following it needs none of the care a growing TIFF does -- re-read it, and read what is new. It stops when nothing has arrived for ``settings.timeout``, - which is how an acquisition ends. + which is how an acquisition ends -- or, given a ``finished`` event, + only once that is set and what was written by then has been read + (`smappy.io.watch`, which also says what ``idle_blocks`` is for). """ from .watch import WatchSettings, _sleep, _stopped @@ -257,6 +260,9 @@ def watch(self, chunk: int = 100, settings=None, start: int = 0, index = start last_new = time.monotonic() while not _stopped(stop_event): + # before the reload, so the last pass sees every frame written + # before the writer said it was done + ending = _stopped(finished) self.reload() available = len(self.index) if stop is None else min(stop, len(self.index)) while index < available and not _stopped(stop_event): @@ -264,13 +270,15 @@ def watch(self, chunk: int = 100, settings=None, start: int = 0, yield index, self._block(index, last) index = last last_new = time.monotonic() - if stop is not None and index >= stop: + if ending or (stop is not None and index >= stop): break waited = time.monotonic() - last_new - if waited >= settings.timeout: + if finished is None and waited >= settings.timeout: break if on_wait is not None: on_wait(waited) + if idle_blocks: + yield index, np.empty((0, *self.shape), self.dtype) _sleep(settings.poll, stop_event) # -------------------------------------------------------------- refreshing @@ -316,9 +324,15 @@ def open_ndtiff(path) -> NDTiffSource: return source +# where ndstorage keeps the full-resolution images of a dataset that has (or +# could have) a pyramid; MicroClaw's acquisitions are written this way, and +# the folder a user or a caller names is the one above it +FULL_RESOLUTION = "Full resolution" + + def _dataset_folder(path) -> Optional[Path]: path = Path(path) - for folder in (path, path.parent): + for folder in (path, path / FULL_RESOLUTION, path.parent): if (folder / INDEX_NAME).is_file(): return folder return None diff --git a/src/smappy/io/watch.py b/src/smappy/io/watch.py index ed3878e..fdb6d4a 100644 --- a/src/smappy/io/watch.py +++ b/src/smappy/io/watch.py @@ -7,6 +7,15 @@ acquisition ends, since Micro-Manager does not announce it and the frame count in the summary metadata is what was *planned*, not what was written. +Unless somebody does announce it. A program that drives the microscope +(MicroClaw) knows when the writer has finished, and a pause in an acquisition +-- a refocus, a laser being changed, a long time-lapse interval -- is then not +an end: the ``finished`` event replaces the idle timeout altogether. While it +is unset the watch waits however long it takes; once it is set, one more pass +reads everything there is, the held-back last page included, and the stream +ends. Raising the timeout instead would only move the pause that ends a fit +too early. + Three things make this safe against a writer working on the same file: * **The last page of the file being written is held back** until another page @@ -48,8 +57,8 @@ class WatchSettings: def watch_stack(path, chunk: int = 100, settings: Optional[WatchSettings] = None, start: int = 0, stop: Optional[int] = None, - stop_event=None, on_wait=None, tags=None - ) -> Iterator[Tuple[int, np.ndarray]]: + stop_event=None, on_wait=None, tags=None, finished=None, + idle_blocks: bool = False) -> Iterator[Tuple[int, np.ndarray]]: """Yield ``(first_frame, block)`` from an acquisition as it is written. Blocks are ``chunk`` frames, except that whatever has accumulated is handed @@ -61,13 +70,21 @@ def watch_stack(path, chunk: int = 100, settings: Optional[WatchSettings] = None called on each idle poll with the time since the last new frame, for a progress line. ``tags`` is a `smappy.frametags.FrameTags` that is handed each frame's metadata as it is read. + + ``finished``, a :class:`threading.Event`, is set once the writer is known + to be done; given one, idle time never ends the stream (see the module + docstring). ``idle_blocks`` yields an empty block on every idle poll, so + that whoever consumes the stream gets to act while nothing arrives -- a + fit holding thousands of ROIs in its buffer fits them during a pause, + rather than at the next frame, which may be minutes away. """ settings = settings or WatchSettings() ndtiff = _wait_for_ndtiff(path, settings, stop_event) if ndtiff is not None: ndtiff.tags = tags yield from ndtiff.watch(chunk=chunk, settings=settings, start=start, - stop=stop, stop_event=stop_event, on_wait=on_wait) + stop=stop, stop_event=stop_event, on_wait=on_wait, + finished=finished, idle_blocks=idle_blocks) return first = _wait_for_file(path, settings, stop_event) @@ -84,8 +101,13 @@ def watch_stack(path, chunk: int = 100, settings: Optional[WatchSettings] = None # the acquisition has clearly ended nothing is going to follow it, so the # timeout triggers one more pass that does read it final = False + empty = None # an idle block, once the frame shape is known while not _stopped(stop_event): + # read before the pages are: a frame written just before the writer + # said it was done is then still in this pass + if _stopped(finished): + final = True files = _series_files(first) arrived = 0 @@ -96,6 +118,8 @@ def watch_stack(path, chunk: int = 100, settings: Optional[WatchSettings] = None images, page_i, exhausted = _read_new_pages(files[file_i], page_i, hold_last=hold, metadata=planes) + if images and empty is None: + empty = np.empty((0, *images[0].shape), images[0].dtype) for k, image in enumerate(images): if (stop is not None and index >= stop) or _stopped(stop_event): done = True @@ -127,11 +151,13 @@ def watch_stack(path, chunk: int = 100, settings: Optional[WatchSettings] = None break waited = time.monotonic() - last_new - if waited >= settings.timeout: + if finished is None and waited >= settings.timeout: final = True continue if on_wait is not None: on_wait(waited) + if idle_blocks and empty is not None: + yield buffer_start, empty _sleep(settings.poll, stop_event) @@ -149,6 +175,8 @@ def open_growing_stack(path, settings: Optional[WatchSettings] = None, return ndtiff first = _wait_for_file(path, settings, stop_event) + if _stopped(stop_event): + raise InterruptedError(f"stopped while waiting for {path} to appear") if first is None: raise TimeoutError(f"no readable TIFF appeared at {path} within " f"{settings.appear_timeout} s") @@ -197,7 +225,8 @@ def _looks_like_ndtiff(directory: Path) -> bool: NDTiffStorage writes its stack file before the index, so between the two there is a moment that would otherwise be read as a Micro-Manager series. """ - return directory.is_dir() and any(directory.glob("*NDTiffStack*.tif")) + return directory.is_dir() and (any(directory.glob("*NDTiffStack*.tif")) or + any(directory.glob("Full resolution/*NDTiffStack*.tif"))) def _read_new_pages(path: Path, page_i: int, hold_last: bool, diff --git a/src/smappy/pipeline.py b/src/smappy/pipeline.py index bed39e5..43c14ad 100644 --- a/src/smappy/pipeline.py +++ b/src/smappy/pipeline.py @@ -262,12 +262,19 @@ def fit_stack(frames: Iterable[Tuple[int, np.ndarray]], camera: CameraMetadata, def drive(engine, frames: Iterable[Tuple[int, np.ndarray]], sink: Optional[Callable[[Localizations], None]] = None, progress: Optional[Callable[[object], None]] = None, - read_ahead: int = 2) -> Localizations: + read_ahead: int = 2, + flush_seconds: Optional[float] = None) -> Localizations: """Push blocks through anything with ``push`` and ``flush``. Split out of `fit_stack` so that the dual-channel engine (`smappy.dualfit.DualChannelEngine`), which buffers pairs of ROIs rather than single ones, is driven by the same loop. + + ``flush_seconds`` fits what is buffered at least that often, for a live + source: the buffer holds thousands of ROIs, and a sparse acquisition would + otherwise show nothing for minutes (`smappy.live.LiveFit` does the same). + It is checked once per block, so a source that pauses should send empty + blocks while it waits (`watch_stack(idle_blocks=True)`). """ collected = Localizations() @@ -279,8 +286,15 @@ def emit(locs: Optional[Localizations]) -> None: else: collected.extend(locs) + last_flush = time.monotonic() for first_frame, block in prefetch(frames, read_ahead): - emit(engine.push(block, first_frame)) + # an empty block is a live source saying nothing has arrived + # (`watch_stack(idle_blocks=True)`): only the clock below has work + if len(block): + emit(engine.push(block, first_frame)) + if flush_seconds is not None and time.monotonic() - last_flush >= flush_seconds: + emit(engine.flush()) + last_flush = time.monotonic() if progress is not None: progress(engine) emit(engine.flush()) diff --git a/src/smappy/plugins/__init__.py b/src/smappy/plugins/__init__.py index 5505990..9e696e9 100644 --- a/src/smappy/plugins/__init__.py +++ b/src/smappy/plugins/__init__.py @@ -259,7 +259,8 @@ def __init__(self, session=None, locs: Optional[Localizations] = None, stream: Optional[Callable[[str, Any], None]] = None, site=None, site_table: Optional[Sequence[Dict[str, Any]]] = None, rois=None, grouping: Optional[str] = None, - evaluations: Sequence[Tuple[str, str]] = ()): + evaluations: Sequence[Tuple[str, str]] = (), + stop=None, writer_finished=None): self.session = session self.layer = layer # "grouped", "ungrouped", or None for each layer's own: what a chain @@ -276,6 +277,14 @@ def __init__(self, session=None, locs: Optional[Localizations] = None, # chain's results unless told otherwise (`roi.choose_evaluation`) self.evaluations = list(evaluations) self._rois = rois # a project without a session + # Two `threading.Event`s for a run that lasts as long as an + # acquisition, set from another thread; None when nobody can. + # `stop` asks it to end early and keep what it has -- a fit then + # closes its file and returns the table so far. `writer_finished` + # says a live source is complete, which ends the watch on it after + # one more read instead of after an idle timeout (`smappy.io.watch`). + self.stop = stop + self.writer_finished = writer_finished self._progress = progress self._stream = stream if locs is None: @@ -347,7 +356,8 @@ def for_site(self, site, locs: Optional[Localizations] = None, selection=self.selection if selection is None else selection, layer=self.layer, progress=self._progress, stream=self._stream, site=site, site_table=self.site_table, - rois=self._rois, grouping=self.grouping) + rois=self._rois, grouping=self.grouping, + stop=self.stop, writer_finished=self.writer_finished) @dataclass diff --git a/src/smappy/plugins/fit.py b/src/smappy/plugins/fit.py index d055f2b..25288d2 100644 --- a/src/smappy/plugins/fit.py +++ b/src/smappy/plugins/fit.py @@ -690,7 +690,7 @@ def run(self, ctx: Context, settings) -> Result: src = settings.source if not src.path: raise ValueError("choose a source file") - source = _open(src, watch=src.live) + source = _open(src, watch=src.live, stop_event=ctx.stop) camera = settings.camera.resolve(source) fit = settings.fit finder = self.finder(settings, camera) @@ -702,14 +702,18 @@ def run(self, ctx: Context, settings) -> Result: tags = FrameTags() if src.live: from ..io.watch import WatchSettings, watch_stack + # given `ctx.writer_finished`, the writer's word ends the watch + # and *stop after* is never consulted (`smappy.io.watch`) blocks = watch_stack(src.path, chunk=src.chunk, start=src.start, stop=src.stop, settings=WatchSettings(timeout=src.live_timeout), - tags=tags) + tags=tags, stop_event=ctx.stop, + finished=ctx.writer_finished, idle_blocks=True) total = None # a growing stack has no end to count to else: stop = min(src.stop, source.n_frames) if src.stop else None source.tags = tags - blocks = source.frames(chunk=src.chunk, start=src.start, stop=stop) + blocks = _until(source.frames(chunk=src.chunk, start=src.start, stop=stop), + ctx.stop) total = max((stop if stop is not None else source.n_frames) - src.start, 0) # the frames kept with the table: counted as they go past, so the # stack is read once, for the fit @@ -759,7 +763,8 @@ def report(engine) -> None: try: from ..pipeline import drive engine = self.engine(settings, camera, finder, model) - drive(engine, blocks, sink=sink, progress=report, read_ahead=2) + drive(engine, blocks, sink=sink, progress=report, read_ahead=2, + flush_seconds=LIVE_FLUSH_SECONDS if src.live else None) record["frame_tags"] = tags.table() if writer is not None: writer.set_metadata({"stats": dict(engine.stats), @@ -775,11 +780,17 @@ def report(engine) -> None: collected.metadata[key] = record[key] seconds = time.perf_counter() - started rate = stats["frames"] / seconds if seconds > 0 else 0.0 + stopped = ctx.stop is not None and ctx.stop.is_set() text = (f"{stats['localizations']} localizations from {stats['frames']} frames" f" in {seconds:.1f} s ({rate:,.1f} frames/s)" + + (", stopped early" if stopped else "") + (f", saved to {out}" if out else "")) - finished = self.finish(ctx, settings, collected.compact()) + # A stopped run keeps the raw fit and finishes nothing: whoever asked + # wants the file now (MicroClaw allows ten seconds), and a drift + # correction of part of an acquisition is not one anybody asked for. + finished = (Finished(collected.compact()) if stopped + else self.finish(ctx, settings, collected.compact())) # the table in hand says what it was fitted from, as its file does: a # comparison with a simulation's truth finds the recipe by it finished.locs.metadata.setdefault("source", str(src.path)) @@ -816,15 +827,29 @@ def report(engine) -> None: if shown: text = "\n".join([text, summary(record["frame_tags"], shown)]) return Result(locs=finished.locs, text=text, plots=plots, - data={"stats": stats, "path": out, + data={"stats": stats, "path": out, "stopped": stopped, "images": [raw] if raw is not None else []}, settings=settings) -def _open(src: SourceSettings, watch: bool): +# how often a live fit fits what it has buffered, whether or not the buffer is +# full: the interval at which a sparse acquisition appears on screen +LIVE_FLUSH_SECONDS = 5.0 + + +def _until(blocks, stop_event): + """``blocks`` until ``stop_event`` is set; all of them without one.""" + for item in blocks: + if stop_event is not None and stop_event.is_set(): + return + yield item + + +def _open(src: SourceSettings, watch: bool, stop_event=None): if watch: from ..io.watch import WatchSettings, open_growing_stack - return open_growing_stack(src.path, WatchSettings(timeout=src.live_timeout)) + return open_growing_stack(src.path, WatchSettings(timeout=src.live_timeout), + stop_event=stop_event) from ..io.tiff import open_stack return open_stack(src.path) diff --git a/src/smappy/session.py b/src/smappy/session.py index 6927f9f..dc35bb8 100644 --- a/src/smappy/session.py +++ b/src/smappy/session.py @@ -315,6 +315,11 @@ class Session: def __init__(self, locs: Optional[Localizations] = None, path=None): self.locs = locs if locs is not None else Localizations({}, {}) self.path: Optional[Path] = Path(path) if path else None + # Files `save` refuses to write over. A fit run for MicroClaw hands + # its file on as the job's result, with a digest taken there and then, + # and the window stays open on it: a Ctrl+S afterwards would change + # what the record vouches for. The session's work goes to a new file. + self.protected: set = set() self.layers: List[Layer] = [Layer(self.locs)] self.roi: Optional[Region] = None # the drawn 2D ROI (regions.py) self._rois = None # the ROI manager's project @@ -578,6 +583,13 @@ def full_view(self, margin_fraction: float = 0.01): def first_locs_layer(self) -> int: return next((i for i, l in enumerate(self.layers) if not l.is_image), 0) + def is_protected(self, path) -> bool: + """Whether ``path`` is one of the files `save` will not write over.""" + if not path: + return False + target = Path(path).resolve() + return any(target == Path(p).resolve() for p in self.protected) + def save(self, path=None, gui_state: Optional[bool] = None) -> Path: """Write the table, its history, its ROIs and -- optionally -- the GUI. @@ -588,6 +600,9 @@ def save(self, path=None, gui_state: Optional[bool] = None) -> Path: from .io.formats import writer_for from .io.hdf5 import save_gui_state, save_localizations, save_results path = Path(path or self.path) + if self.is_protected(path): + raise ValueError(f"{path.name} is kept as the fit wrote it: " + f"save under another name") # before anything is written: the path is the open file's by default, # and a SMAP _sml.mat, a csv or a MINFLUX export is not a file this # writes -- it used to be overwritten with HDF5 under its own name diff --git a/tests/test_live_session.py b/tests/test_live_session.py new file mode 100644 index 0000000..d023afe --- /dev/null +++ b/tests/test_live_session.py @@ -0,0 +1,212 @@ +"""The GUI opened on a running fit, as a program driving the microscope opens it. + +The dataset grows here the way NDTiffStorage writes one -- pixels appended to +the stack file, then a record appended to the index -- so the fit is reading +while the "microscope" writes, and the pause in the middle is longer than the +fit's idle timeout: only the writer's word may end it. +""" +import json +import struct +import time + +import numpy as np +import pytest + +from smappy.io.ndtiff import INDEX_NAME, SUMMARY_MARKER, SUMMARY_OFFSET + +PLUGIN = "Localize/Gaussian 2D" + + +class GrowingNDTiff: + """An NDTiff dataset written a frame at a time, by appending.""" + + NAME = "Stack_NDTiffStack.tif" + + def __init__(self, folder, shape): + self.folder = folder + folder.mkdir(parents=True) + summary = json.dumps({"PixelSize_um": 0.1, "Height": shape[0], + "Width": shape[1]}).encode() + (folder / self.NAME).write_bytes( + b"\0" * SUMMARY_OFFSET + struct.pack(" 0 + + # the file is handed on as it is; the session's work goes elsewhere + assert live.session.is_protected(out) + with pytest.raises(ValueError, match="another name"): + live.session.save(out) + live.session.save(tmp_path / "out" / "edited.hdf5") + finally: + close(live) + live.wait() + + +@pytest.mark.slow +def test_stopping_keeps_what_was_fitted_and_says_the_input_is_incomplete(tmp_path, frames): + pytest.importorskip("PySide6") + from smappy.gui.live_session import LiveSession + from smappy.io.hdf5 import load_localizations + + data = GrowingNDTiff(tmp_path / "ds", frames.shape[1:]) + data.write(frames[:20]) + out = tmp_path / "locs.hdf5" + ended = [] + live = LiveSession(PLUGIN, settings_for(data.folder, out), on_finished=ended.append) + live.open() + try: + assert pump(lambda: len(live.session.locs) > 0) + live.stop.set() + assert pump(lambda: ended) + outcome, = ended + assert (outcome.state, outcome.complete) == ("cancelled", False) + assert len(load_localizations(out)) == len(live.session.locs) + finally: + close(live) + live.wait() + + +@pytest.mark.slow +def test_closing_the_windows_while_it_fits_is_a_cancellation(tmp_path, frames): + pytest.importorskip("PySide6") + from smappy.gui.live_session import LiveSession + + data = GrowingNDTiff(tmp_path / "ds", frames.shape[1:]) + data.write(frames[:10]) + ended = [] + live = LiveSession(PLUGIN, settings_for(data.folder, tmp_path / "locs.hdf5"), + on_finished=ended.append) + live.open() + assert pump(lambda: len(live.session.locs) > 0) # fitting, not starting + close(live) + live.wait() # what `run` does once the loop ends + outcome, = ended + assert outcome.state == "cancelled" and not outcome.complete + assert (tmp_path / "locs.hdf5").exists() + + +@pytest.mark.slow +def test_a_fit_that_cannot_start_is_reported_once_as_failed(tmp_path): + pytest.importorskip("PySide6") + from smappy.gui.live_session import LiveSession + + ended = [] + settings = settings_for(tmp_path / "nothing-here", tmp_path / "locs.hdf5", + **{"source.live": False}) + live = LiveSession(PLUGIN, settings, on_finished=ended.append) + live.open() + try: + assert pump(lambda: ended) + outcome, = ended + assert outcome.state == "failed" and outcome.error + finally: + close(live) + live.wait() + assert len(ended) == 1 diff --git a/tests/test_ndtiff.py b/tests/test_ndtiff.py index 036f565..a891abd 100644 --- a/tests/test_ndtiff.py +++ b/tests/test_ndtiff.py @@ -206,3 +206,59 @@ def test_a_micro_manager_tiff_is_not_mistaken_for_ndtiff(tmp_path): source = open_growing_stack(directory, WatchSettings(poll=0.01, timeout=0.2, appear_timeout=1.0)) assert type(source).__name__ == "ImageSource" + + +def test_a_pause_does_not_end_a_watch_that_waits_for_the_writer(tmp_path): + """With ``finished`` given, idle time is not the end of the acquisition. + + The writer pauses for several timeouts -- a refocus, a time-lapse + interval -- and then writes more; only setting the event ends the watch, + and the frames written just before it are still read. + """ + from smappy.io.watch import WatchSettings + + frames = stack(n=6) + folder = write_ndtiff(tmp_path / "ds", frames[:2]) + source = open_ndtiff(folder) + finished = threading.Event() + settings = WatchSettings(poll=0.02, timeout=0.1) + + def microscope(): + time.sleep(0.5) # five timeouts with nothing new + write_ndtiff(folder, frames[:4]) + time.sleep(0.3) + write_ndtiff(folder, frames) + finished.set() # straight after the last frame + + writer = threading.Thread(target=microscope) + writer.start() + read = np.concatenate([block for _, block in + source.watch(chunk=2, settings=settings, + finished=finished)]) + writer.join() + assert np.array_equal(read, frames) + + +def test_a_watch_whose_writer_has_finished_reads_the_rest_and_ends(tmp_path): + from smappy.io.watch import WatchSettings + + frames = stack(n=5) + source = open_ndtiff(write_ndtiff(tmp_path / "ds", frames)) + finished = threading.Event() + finished.set() + t0 = time.monotonic() + read = np.concatenate([block for _, block in + source.watch(chunk=2, finished=finished, + settings=WatchSettings(poll=0.01, timeout=30))]) + assert np.array_equal(read, frames) + assert time.monotonic() - t0 < 5 # not the 30 s timeout + + +def test_a_dataset_kept_under_full_resolution_is_found_from_the_folder_above(tmp_path): + """ndstorage writes a dataset that may have a pyramid into + ``Full resolution/``; MicroClaw names the folder above it.""" + frames = stack() + write_ndtiff(tmp_path / "acq" / "Full resolution", frames) + assert is_ndtiff(tmp_path / "acq") + source = open_ndtiff(tmp_path / "acq") + assert np.array_equal(source.frame(3), frames[3]) diff --git a/tests/test_watch.py b/tests/test_watch.py index 96ef161..bdcbe25 100644 --- a/tests/test_watch.py +++ b/tests/test_watch.py @@ -186,3 +186,26 @@ def test_a_growing_stack_can_be_opened_for_its_metadata(tmp_path): assert source.shape == (16, 24) assert source.dtype == np.uint16 assert source.n_frames >= 1 # only what was there when it opened + + +def test_a_pause_longer_than_the_timeout_does_not_end_a_watch_given_finished(tmp_path): + """The TIFF path honours the writer's word as the NDTiff path does. + + The last page is held back while frames are coming; setting ``finished`` + makes the final pass read it too. + """ + scope = Microscope(tmp_path) + scope.write(3) + finished = threading.Event() + + def microscope(): + time.sleep(4 * FAST.timeout) # four timeouts of silence + scope.write(4) + finished.set() + + thread = threading.Thread(target=microscope, daemon=True) + thread.start() + seen = _collect(watch_stack(tmp_path, chunk=100, settings=FAST, + finished=finished)) + thread.join() + assert [n for _, n in seen] == list(range(7))