diff --git a/.issueflows/01-current-issues/auto_status.md b/.issueflows/01-current-issues/auto_status.md index b18889d0..a7a021ef 100644 --- a/.issueflows/01-current-issues/auto_status.md +++ b/.issueflows/01-current-issues/auto_status.md @@ -14,3 +14,13 @@ stop_reason: > (public protocol, not a small change). #780 and #164 were not reached. Did not re-queue: another loop would hit the same safeguard. Did not open duplicate issues. Stage 3 was not started (epoch_gated). + +## Manual follow-up (2026-09-26) + +Stage 2 and 3 processed by hand on the user's request (no yolo merge): +- #779 → PR #1100 (merged earlier) +- #780 → PR #1101 (branch `780-load-since`, base master) +- #164 → PR #1102 (branch `164-cell-update`, base `780-load-since`) +- #781 → PR #1103 (branch `781-live-poll`, base `164-cell-update`) +- #782 → PR #1104 (branch `782-batch-live`, base `781-live-poll`) +Stacked; merge in order and retarget each PR to master as the base lands. diff --git a/.issueflows/01-current-issues/drive_status.md b/.issueflows/01-current-issues/drive_status.md index 15439da3..8a9f73b4 100644 --- a/.issueflows/01-current-issues/drive_status.md +++ b/.issueflows/01-current-issues/drive_status.md @@ -14,3 +14,5 @@ findings: last_outcome: pending auto_stop: stage 2 halted on #779 (yolo: no). #780 and #164 not reached. #781 and #782 not started. + +manual_follow_up (2026-09-26): all Epic L issues implemented by hand; PRs #1101 → #1102 → #1103 → #1104 stacked, awaiting review/merge. Remaining drive steps (final_review, cleanup, status) are the user's call after merges. diff --git a/.issueflows/03-solved-issues/issue782_original.md b/.issueflows/03-solved-issues/issue782_original.md new file mode 100644 index 00000000..12ba3c21 --- /dev/null +++ b/.issueflows/03-solved-issues/issue782_original.md @@ -0,0 +1,12 @@ +# Issue #782: L5 — batch live refresh + +- GitHub: https://github.com/jepegit/cellpy/issues/782 +- Epic: #783 (Epic L), Stage 3. Depends on: #164. + +## Original description + +Epic L of cellpy 2.2 (Stage 5). Design: live-incremental §6. Depends on L3 (#164). + +`b.update(live=True)` iterates the journal cells calling `c.update()`; +`b.poll(interval=, until=)` wraps the loop and re-runs the collectors/report +each tick. Rides entirely on the cell-level `update()` — no new core. diff --git a/.issueflows/03-solved-issues/issue782_plan.md b/.issueflows/03-solved-issues/issue782_plan.md new file mode 100644 index 00000000..3726f482 --- /dev/null +++ b/.issueflows/03-solved-issues/issue782_plan.md @@ -0,0 +1,17 @@ +# Plan: #782 batch live refresh + +Autonomous run under #783 (user: "process the issues"). + +1. `Batch.refresh(labels=None, raise_errors=False, **update_kwargs)` → + `{label: changed | Exception}` over loaded cells; clears summary cache on + change. +2. `Batch.update(live=True, **kw)` delegates to `refresh` and returns the + existing `BatchResult` (no reload). +3. `Batch.poll(interval, on_update, until, max_polls, timeout, + stop_when_complete, raise_errors, sleep, **update_kwargs)` reusing + `cellpy.utils.live.PollStatus`; on change rebuild `summaries`, recompute + `report()` → `last_report`, call `on_update(batch, outcome)`. +4. Tests `tests/test_batch_live.py` on a `from_cells` batch of two truncated + neware copies. +5. Docs: `docs/agents/index.md`, root `AGENTS.md`, design doc section, + HISTORY, registry. diff --git a/.issueflows/03-solved-issues/issue782_status.md b/.issueflows/03-solved-issues/issue782_status.md new file mode 100644 index 00000000..27f08308 --- /dev/null +++ b/.issueflows/03-solved-issues/issue782_status.md @@ -0,0 +1,20 @@ +# Status: #782 batch live refresh + +- [x] Done + +## Done + +- `Batch.refresh`, `Batch.update(live=True)`, `Batch.poll` in + `cellpy/batch/facade.py`. +- `tests/test_batch_live.py`: 6 essential tests. +- Docs: `docs/agents/index.md`, root `AGENTS.md`, `incremental-load-protocol.md` + (#782 section), HISTORY, test registry. + +## Notes + +- Branch `782-batch-live` stacked on `781-live-poll` (PR #1103). +- Persisting refreshed cells to `.cellpy` left to the caller. + +## Remaining + +- None. Epic L (#783) code complete pending PR merges #1101 → #1102 → #1103 → #1104. diff --git a/.issueflows/04-designs-and-guides/incremental-load-protocol.md b/.issueflows/04-designs-and-guides/incremental-load-protocol.md index 63332615..aa54133f 100644 --- a/.issueflows/04-designs-and-guides/incremental-load-protocol.md +++ b/.issueflows/04-designs-and-guides/incremental-load-protocol.md @@ -123,8 +123,24 @@ with `"interrupted"`. Run bookkeeping lives on `cell.poll_status` `batch.runner` already has `executor="threads"`. The `batch_core.py` `lstrip` bug is moot: `utils/batch_tools/` is gone since batch v3. +## Batch live refresh (#782) + +`cellpy/batch/facade.py`. `Batch.refresh(labels=None, raise_errors=False, +**update_kwargs) -> {label: bool | Exception}` calls `c.update()` on the +**loaded** cells of the lazy store only (a refresh never triggers a first +load) and clears the combined-summary cache when any cell changed. +`Batch.update(live=True, **kw)` is the spec'd spelling: it delegates to +`refresh` and returns the existing `BatchResult` unchanged (no reload, no +runner). `Batch.poll(...)` mirrors `live.poll` (same `PollStatus`, same +stop conditions; "complete" = every loaded cell has `source_complete`), and +on a changed tick rebuilds `summaries`, recomputes the QC `report()` into +`last_report`, and calls `on_update(batch, outcome)`. A failing cell ends +the poll with `stopped_by="error"` unless `raise_errors=True`. Persisting +refreshed cells to `.cellpy` is left to the caller (`c.save`); not folded +into `refresh`. + ## Link Design §3 in `cellpy-design-and-development/active/cellpy2-live-incremental-design.md`. Tests: `tests/test_load_since.py` (#780), `tests/test_cell_update.py` (#164), -`tests/test_live_poll.py` (#781). Consumer: batch live refresh (#782). +`tests/test_live_poll.py` (#781), `tests/test_batch_live.py` (#782). diff --git a/.issueflows/04-designs-and-guides/test-registry.md b/.issueflows/04-designs-and-guides/test-registry.md index 6969c639..42a37e20 100644 --- a/.issueflows/04-designs-and-guides/test-registry.md +++ b/.issueflows/04-designs-and-guides/test-registry.md @@ -217,6 +217,12 @@ current issue**. `/iflow-doctor` may audit the whole suite against this table. | tests/test_live_poll.py::test_poll_records_update_errors_unless_raise_errors | yes | yes | utils.live.poll | #781 | | | tests/test_live_poll.py::test_poll_keyboard_interrupt_returns_the_cell | yes | yes | utils.live.poll | #781 | | | tests/test_live_poll.py::test_processor_module_is_gone | yes | yes | utils.processor (deleted) | #781 | guards against resurrection | +| tests/test_batch_live.py::test_refresh_reports_per_cell_and_updates_summaries | yes | yes | Batch.refresh | #782 | one of two cells grows; summaries cache cleared | +| tests/test_batch_live.py::test_refresh_subset_and_error_capture | yes | yes | Batch.refresh(labels, raise_errors) | #782 | | +| tests/test_batch_live.py::test_update_live_does_not_reload | yes | yes | Batch.update(live=True) | #782 | returns same BatchResult | +| tests/test_batch_live.py::test_poll_refreshes_and_reruns_report | yes | yes | Batch.poll | #782 | fake clock; last_report rebuilt | +| tests/test_batch_live.py::test_poll_stops_on_until_and_complete | yes | yes | Batch.poll stop conditions | #782 | | +| tests/test_batch_live.py::test_poll_stops_on_cell_error | yes | yes | Batch.poll | #782 | | | tests/test_dbreader.py::test_missing_column_warns_once | yes | yes | readers.dbreader.Reader._pick_info | #1008 | warn-once per missing header | | tests/test_dbreader.py::test_nom_cap_specifics_column_reaches_pages | yes | yes | batch._dbengine._create_pages_dict | #1008 | db value → pages | | tests/test_dbreader.py::test_simple_db_engine_skip_file_search_excel_reader | yes | yes | batch._dbengine.simple_db_engine / find_files | #1017 | skip_file_search frames one row per cell | diff --git a/AGENTS.md b/AGENTS.md index 722bb490..54f06693 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -377,6 +377,8 @@ Quick facts: - Live/running test: `c.update()` re-reads only the new raw rows (incremental loaders) or reloads fully, returns `True` when frames changed; `False` if the raw file did not change on disk. Also works after `cellpy.get(".cellpy")`. + Follow on an interval: `cellpy.utils.live.poll(c, interval=60, on_update=cb)`; + batch: `b.refresh()` / `b.poll(interval=, on_update=)`. - Frames: `c.data.raw` / `.steps` / `.summary`; columns via `c.schema.*`. After a raw load, each cycle's raw capacity starts at 0. A forgotten tester reset that 1.x plotted as doubled capacity is rebased on load for every diff --git a/HISTORY.md b/HISTORY.md index 9bf9d665..8ffe46d0 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -25,6 +25,12 @@ (its thread-pool fan-out lives in `batch.runner`'s `executor="threads"`). (#781) +* Batch live refresh: `b.refresh()` (alias `b.update(live=True)`) calls + `c.update()` on every loaded cell and returns `{label: changed}`; + `b.poll(interval=, until=, max_polls=, on_update=)` repeats it, rebuilding + `b.summaries` and the QC report (`b.last_report`) on ticks that changed. + (#782) + ## [2.1.5.post6] - 2026-09-25 * `summary_collector(...).plot()` keeps a lone charge or discharge series diff --git a/cellpy/batch/facade.py b/cellpy/batch/facade.py index 570d2c51..91a4777b 100644 --- a/cellpy/batch/facade.py +++ b/cellpy/batch/facade.py @@ -234,10 +234,15 @@ def update( on_progress=None, executor: str = "serial", progress=None, + live: bool = False, **overrides ) -> BatchResult: """Load every cell, caching them in the store. + ``live=True`` does not reload: it calls ``CellpyCell.update()`` on + every loaded cell instead (see `refresh`) and returns the unchanged + ``result``; ``overrides`` are then forwarded to the cells' ``update``. + ``executor`` is ``"serial"`` (default), ``"threads"`` or ``"processes"``. ``"threads"`` mainly speeds up *reopening* cells from local ``.cellpy`` files; a first load of remote raw files does not overlap on the wire, @@ -251,6 +256,9 @@ def update( unknown (legacy) kwargs like ``testing`` are forwarded to the loader (``cellpy.get``) via ``loader_kwargs``. """ + if live: + self.refresh(**overrides) + return self._result policy = self.policy if overrides: known = {f.name for f in fields(LoadPolicy)} @@ -285,6 +293,113 @@ def recalc(self, **overrides) -> BatchResult: """ return self.update(recalc=True, **overrides) + # -- live refresh (#782) --------------------------------------------- + def refresh(self, labels: Sequence[str] | None = None, raise_errors: bool = False, **update_kwargs) -> dict: + """Pick up appended raw data on the loaded cells via ``CellpyCell.update()``. + + Only cells already in the store are touched (a lazy store never loads + a cell just to refresh it); use `update` / `load` for a first load. + The combined-summary cache is cleared when any cell changed, so + ``summaries`` / ``plot()`` / collectors see the new data. + + Args: + labels: subset of cell labels (default: every loaded cell). + raise_errors: re-raise a cell's ``update`` error instead of + recording it (as the exception) in the returned map. + **update_kwargs: forwarded to each ``CellpyCell.update`` (for + example ``model=`` when the loader needs it, or ``force=True``). + + Returns: + ``{label: changed}`` with ``True`` / ``False`` per cell, or the + exception object for a cell whose refresh failed. + """ + targets = list(labels) if labels is not None else [lbl for lbl in self._store if self._store.is_loaded(lbl)] + outcome: dict[str, Any] = {} + any_changed = False + for label in targets: + cell = self._store[label] + try: + changed = bool(cell.update(**update_kwargs)) + except Exception as exc: # noqa: BLE001 - reported per cell + if raise_errors: + raise + _log.error("refresh: %s failed (%s)", label, exc) + outcome[label] = exc + continue + outcome[label] = changed + any_changed = any_changed or changed + if any_changed: + self._summaries = None + return outcome + + def poll( + self, + interval: float = 60.0, + on_update=None, + until=None, + max_polls: int | None = None, + timeout: float | None = None, + stop_when_complete: bool = True, + raise_errors: bool = False, + sleep=None, + **update_kwargs, + ): + """Repeat `refresh` on an interval; rebuild summaries and report each tick. + + The batch counterpart of ``cellpy.utils.live.poll``. Every tick calls + `refresh`; when any cell changed, ``summaries`` are rebuilt, the QC + `report` is recomputed into ``last_report``, and ``on_update(batch, + outcome)`` runs (``outcome`` is the `refresh` map). + + Stop conditions, checked before each wait: every loaded cell reports + ``source_complete`` (``stop_when_complete``), ``until(batch)`` is true, + ``max_polls`` ticks, or ``timeout`` seconds. ``KeyboardInterrupt`` + stops cleanly. Bookkeeping in the returned ``PollStatus`` (also on + ``poll_status``). + """ + import time + + from cellpy.utils.live import PollStatus + + sleep = sleep or time.sleep + status = PollStatus() + self.poll_status = status + started = time.monotonic() + + def _all_complete() -> bool: + loaded = [lbl for lbl in self._store if self._store.is_loaded(lbl)] + return bool(loaded) and all(getattr(self._store[lbl], "source_complete", False) for lbl in loaded) + + try: + while status.stopped_by is None: + if stop_when_complete and _all_complete(): + status.stopped_by = "complete" + elif until is not None and until(self): + status.stopped_by = "until" + elif max_polls is not None and status.polls >= max_polls: + status.stopped_by = "max_polls" + elif timeout is not None and time.monotonic() - started >= timeout: + status.stopped_by = "timeout" + if status.stopped_by is not None: + break + sleep(interval) + status.polls += 1 + outcome = self.refresh(raise_errors=raise_errors, **update_kwargs) + failed = [exc for exc in outcome.values() if isinstance(exc, BaseException)] + if failed and not raise_errors: + status.stopped_by, status.error = "error", failed[0] + break + if any(v is True for v in outcome.values()): + status.updates += 1 + self.combine_summaries() + self.last_report = self.report() + if on_update is not None: + on_update(self, outcome) + except KeyboardInterrupt: + status.stopped_by = "interrupted" + _log.info("poll: finished %s", status) + return status + @property def summaries(self) -> pl.DataFrame: """Combined per-cycle summary frame across the batch (cached).""" diff --git a/cellpy/utils/live.py b/cellpy/utils/live.py index d56e9245..2eefbc94 100644 --- a/cellpy/utils/live.py +++ b/cellpy/utils/live.py @@ -1,160 +1,160 @@ -"""Follow a running test: poll a cell's raw source and refresh it (#781). - -The loop is a thin layer over ``CellpyCell.update()`` (#164): each tick asks the -cell to re-check its raw source and append what is new; ``on_update`` runs -after every tick that changed the frames. Nothing here talks to loaders or -cellpy-core directly. - -Examples: - ```python - from cellpy.utils import live - - def show(c): - print(c.data.summary.tail(1)) - - c = live.poll("running_test.csv", interval=60, on_update=show, max_polls=10, - instrument="neware_txt", model="UIO") - ``` - - Stop on your own condition instead of a tick budget: - - ```python - c = live.poll(c, interval=30, until=lambda cell: cell.get_cycle_numbers()[-1] >= 50) - ``` -""" - -from __future__ import annotations - -import logging -import time -from typing import Callable, Optional - -import cellpy -from cellpy.readers.cellreader import CellpyCell - -logging.captureWarnings(True) - - -class PollStatus: - """Outcome of one ``poll`` run (also exposed as ``cell.poll_status``). - - Attributes: - polls: ticks run (each tick is one ``update()`` call). - updates: ticks where the frames changed. - stopped_by: ``"complete"``, ``"until"``, ``"max_polls"``, - ``"timeout"``, ``"interrupted"``, or ``"error"``. - error: the exception when ``stopped_by == "error"``. - """ - - def __init__(self): - self.polls = 0 - self.updates = 0 - self.stopped_by: Optional[str] = None - self.error: Optional[BaseException] = None - - def __repr__(self): - return f"PollStatus(polls={self.polls}, updates={self.updates}, stopped_by={self.stopped_by!r})" - - -def poll( - cell_or_path, - interval: float = 30.0, - on_update: Optional[Callable] = None, - stop_when_complete: bool = True, - until: Optional[Callable] = None, - max_polls: Optional[int] = None, - timeout: Optional[float] = None, - raise_errors: bool = False, - sleep: Callable[[float], None] = time.sleep, - **get_kwargs, -): - """Refresh a cell from its raw source on an interval until a stop condition. - - Args: - cell_or_path: a ``CellpyCell`` or a path handed to ``cellpy.get`` - (with ``**get_kwargs``, e.g. ``instrument=``, ``model=``, - ``mass=``). When a path is given, the initial load counts as the - first update and ``on_update`` fires for it. - interval: seconds to wait between ticks. - on_update: ``on_update(cell)``; called after every tick that changed - the frames (and after the initial load from a path). - stop_when_complete: stop when the loader reported the test has ended - (``cell.source_complete``; only loaders that can tell set it). - until: ``until(cell) -> bool``; checked after every tick, stops when - true. - max_polls: stop after this many ticks (``None`` = no limit). - timeout: stop after this many seconds (``None`` = no limit). - raise_errors: re-raise exceptions from ``update``/``on_update`` - instead of stopping the loop with ``stopped_by="error"``. - sleep: the wait function (injectable for tests / event loops). - **get_kwargs: forwarded to ``cellpy.get`` when a path is given, and - (loader-related keys) to ``cell.update`` on every tick. - - Returns: - The cell, with the run recorded on ``cell.poll_status``. - - ``KeyboardInterrupt`` stops the loop cleanly (``stopped_by="interrupted"``). - With no stop condition at all the loop runs until interrupted. - """ - status = PollStatus() - update_kwargs = {k: v for k, v in get_kwargs.items() if k in _UPDATE_KWARGS} - - if isinstance(cell_or_path, CellpyCell): - cell = cell_or_path - else: - cell = cellpy.get(cell_or_path, **get_kwargs) - status.updates += 1 - _fire(on_update, cell, status, raise_errors) - - cell.poll_status = status - started = time.monotonic() - - try: - while status.stopped_by is None: - if _should_stop(cell, status, stop_when_complete, until, max_polls, timeout, started): - break - sleep(interval) - status.polls += 1 - try: - changed = cell.update(**update_kwargs) - except Exception as exc: # noqa: BLE001 - reported on the status - if raise_errors: - raise - logging.error(f"poll: update failed ({exc})") - status.stopped_by, status.error = "error", exc - break - if changed: - status.updates += 1 - _fire(on_update, cell, status, raise_errors) - except KeyboardInterrupt: - status.stopped_by = "interrupted" - logging.info(f"poll: finished {status}") - return cell - - -#: ``cellpy.get`` kwargs that ``update()`` must see too (loader recreation). -_UPDATE_KWARGS = ("model",) - - -def _fire(on_update, cell, status, raise_errors): - if on_update is None: - return - try: - on_update(cell) - except Exception as exc: # noqa: BLE001 - reported on the status - if raise_errors: - raise - logging.error(f"poll: on_update failed ({exc})") - status.stopped_by, status.error = "error", exc - - -def _should_stop(cell, status, stop_when_complete, until, max_polls, timeout, started) -> bool: - if stop_when_complete and getattr(cell, "source_complete", False): - status.stopped_by = "complete" - elif until is not None and until(cell): - status.stopped_by = "until" - elif max_polls is not None and status.polls >= max_polls: - status.stopped_by = "max_polls" - elif timeout is not None and time.monotonic() - started >= timeout: - status.stopped_by = "timeout" - return status.stopped_by is not None +"""Follow a running test: poll a cell's raw source and refresh it (#781). + +The loop is a thin layer over ``CellpyCell.update()`` (#164): each tick asks the +cell to re-check its raw source and append what is new; ``on_update`` runs +after every tick that changed the frames. Nothing here talks to loaders or +cellpy-core directly. + +Examples: + ```python + from cellpy.utils import live + + def show(c): + print(c.data.summary.tail(1)) + + c = live.poll("running_test.csv", interval=60, on_update=show, max_polls=10, + instrument="neware_txt", model="UIO") + ``` + + Stop on your own condition instead of a tick budget: + + ```python + c = live.poll(c, interval=30, until=lambda cell: cell.get_cycle_numbers()[-1] >= 50) + ``` +""" + +from __future__ import annotations + +import logging +import time +from typing import Callable, Optional + +import cellpy +from cellpy.readers.cellreader import CellpyCell + +logging.captureWarnings(True) + + +class PollStatus: + """Outcome of one ``poll`` run (also exposed as ``cell.poll_status``). + + Attributes: + polls: ticks run (each tick is one ``update()`` call). + updates: ticks where the frames changed. + stopped_by: ``"complete"``, ``"until"``, ``"max_polls"``, + ``"timeout"``, ``"interrupted"``, or ``"error"``. + error: the exception when ``stopped_by == "error"``. + """ + + def __init__(self): + self.polls = 0 + self.updates = 0 + self.stopped_by: Optional[str] = None + self.error: Optional[BaseException] = None + + def __repr__(self): + return f"PollStatus(polls={self.polls}, updates={self.updates}, stopped_by={self.stopped_by!r})" + + +def poll( + cell_or_path, + interval: float = 30.0, + on_update: Optional[Callable] = None, + stop_when_complete: bool = True, + until: Optional[Callable] = None, + max_polls: Optional[int] = None, + timeout: Optional[float] = None, + raise_errors: bool = False, + sleep: Callable[[float], None] = time.sleep, + **get_kwargs, +): + """Refresh a cell from its raw source on an interval until a stop condition. + + Args: + cell_or_path: a ``CellpyCell`` or a path handed to ``cellpy.get`` + (with ``**get_kwargs``, e.g. ``instrument=``, ``model=``, + ``mass=``). When a path is given, the initial load counts as the + first update and ``on_update`` fires for it. + interval: seconds to wait between ticks. + on_update: ``on_update(cell)``; called after every tick that changed + the frames (and after the initial load from a path). + stop_when_complete: stop when the loader reported the test has ended + (``cell.source_complete``; only loaders that can tell set it). + until: ``until(cell) -> bool``; checked after every tick, stops when + true. + max_polls: stop after this many ticks (``None`` = no limit). + timeout: stop after this many seconds (``None`` = no limit). + raise_errors: re-raise exceptions from ``update``/``on_update`` + instead of stopping the loop with ``stopped_by="error"``. + sleep: the wait function (injectable for tests / event loops). + **get_kwargs: forwarded to ``cellpy.get`` when a path is given, and + (loader-related keys) to ``cell.update`` on every tick. + + Returns: + The cell, with the run recorded on ``cell.poll_status``. + + ``KeyboardInterrupt`` stops the loop cleanly (``stopped_by="interrupted"``). + With no stop condition at all the loop runs until interrupted. + """ + status = PollStatus() + update_kwargs = {k: v for k, v in get_kwargs.items() if k in _UPDATE_KWARGS} + + if isinstance(cell_or_path, CellpyCell): + cell = cell_or_path + else: + cell = cellpy.get(cell_or_path, **get_kwargs) + status.updates += 1 + _fire(on_update, cell, status, raise_errors) + + cell.poll_status = status + started = time.monotonic() + + try: + while status.stopped_by is None: + if _should_stop(cell, status, stop_when_complete, until, max_polls, timeout, started): + break + sleep(interval) + status.polls += 1 + try: + changed = cell.update(**update_kwargs) + except Exception as exc: # noqa: BLE001 - reported on the status + if raise_errors: + raise + logging.error(f"poll: update failed ({exc})") + status.stopped_by, status.error = "error", exc + break + if changed: + status.updates += 1 + _fire(on_update, cell, status, raise_errors) + except KeyboardInterrupt: + status.stopped_by = "interrupted" + logging.info(f"poll: finished {status}") + return cell + + +#: ``cellpy.get`` kwargs that ``update()`` must see too (loader recreation). +_UPDATE_KWARGS = ("model",) + + +def _fire(on_update, cell, status, raise_errors): + if on_update is None: + return + try: + on_update(cell) + except Exception as exc: # noqa: BLE001 - reported on the status + if raise_errors: + raise + logging.error(f"poll: on_update failed ({exc})") + status.stopped_by, status.error = "error", exc + + +def _should_stop(cell, status, stop_when_complete, until, max_polls, timeout, started) -> bool: + if stop_when_complete and getattr(cell, "source_complete", False): + status.stopped_by = "complete" + elif until is not None and until(cell): + status.stopped_by = "until" + elif max_polls is not None and status.polls >= max_polls: + status.stopped_by = "max_polls" + elif timeout is not None and time.monotonic() - started >= timeout: + status.stopped_by = "timeout" + return status.stopped_by is not None diff --git a/docs/agents/index.md b/docs/agents/index.md index b56f59ba..a7e9fdf3 100644 --- a/docs/agents/index.md +++ b/docs/agents/index.md @@ -334,6 +334,10 @@ fig = summary_collector(b, family="fullcell_standard_gravimetric").plot() # database is NOT re-read (a UserWarning fires if you also passed db args). # Use allow_from_journal=False to rebuild the journal from the db. b = batch.load(name="my_experiment", project="my_project", allow_from_journal=False) +# Tests still running? Pick up appended raw rows without a reload: +changed = b.refresh() # {label: True/False} via c.update() per cell +# or keep following them; summaries and b.last_report are rebuilt on change +status = b.poll(interval=120, max_polls=30, on_update=lambda b, out: print(out)) ``` Dropping cells: diff --git a/tests/test_batch_live.py b/tests/test_batch_live.py new file mode 100644 index 00000000..7766b153 --- /dev/null +++ b/tests/test_batch_live.py @@ -0,0 +1,121 @@ +"""Batch live refresh (#782): ``Batch.refresh`` / ``update(live=True)`` / ``poll``.""" + +from __future__ import annotations + +import shutil + +import pytest + +import cellpy +from cellpy.batch.facade import from_cells +from tests.incremental_support import ( + NEWARE_KWARGS, + NEWARE_UIO, + assert_cell_frames_equal, + truncate_text_file, +) + +pytestmark = pytest.mark.essential + +HEAD_ROWS = 6000 + + +@pytest.fixture(scope="module") +def full_cell(): + return cellpy.get(NEWARE_UIO, testing=True, **NEWARE_KWARGS) + + +@pytest.fixture +def live_batch(tmp_path): + paths = {} + cells = {} + for label in ("cell_a", "cell_b"): + path = truncate_text_file(NEWARE_UIO, tmp_path / f"{label}.csv", HEAD_ROWS) + paths[label] = path + cells[label] = cellpy.get(path, testing=True, **NEWARE_KWARGS) + return from_cells(cells), paths + + +def _grow(path): + shutil.copyfile(NEWARE_UIO, path) + + +def test_refresh_reports_per_cell_and_updates_summaries(live_batch, full_cell): + b, paths = live_batch + n_before = len(b.summaries) + assert b.refresh() == {"cell_a": False, "cell_b": False} + _grow(paths["cell_b"]) + assert b.refresh() == {"cell_a": False, "cell_b": True} + assert len(b.summaries) > n_before # cache was cleared and rebuilt + assert_cell_frames_equal(b.cells["cell_b"], full_cell) + assert len(b.cells["cell_a"].data.raw) == HEAD_ROWS + + +def test_refresh_subset_and_error_capture(live_batch, monkeypatch): + b, paths = live_batch + + def boom(**kwargs): + raise RuntimeError("tester unplugged") + + monkeypatch.setattr(b.cells["cell_a"], "update", boom) + outcome = b.refresh() + assert isinstance(outcome["cell_a"], RuntimeError) + assert outcome["cell_b"] is False + with pytest.raises(RuntimeError): + b.refresh(labels=["cell_a"], raise_errors=True) + assert b.refresh(labels=["cell_b"]) == {"cell_b": False} + + +def test_update_live_does_not_reload(live_batch, full_cell): + b, paths = live_batch + result_before = b.result + _grow(paths["cell_a"]) + assert b.update(live=True) is result_before + assert_cell_frames_equal(b.cells["cell_a"], full_cell) + + +def test_poll_refreshes_and_reruns_report(live_batch, full_cell): + b, paths = live_batch + ticks = [] + seen = [] + + def clock(seconds): + ticks.append(seconds) + if len(ticks) == 1: + _grow(paths["cell_a"]) + elif len(ticks) == 3: + _grow(paths["cell_b"]) + + status = b.poll(interval=5, max_polls=4, on_update=lambda batch, outcome: seen.append(outcome), sleep=clock) + assert ticks == [5] * 4 + assert status.polls == 4 + assert status.updates == 2 + assert status.stopped_by == "max_polls" + assert seen == [{"cell_a": True, "cell_b": False}, {"cell_a": False, "cell_b": True}] + assert b.poll_status is status + assert set(b.last_report["cell"].to_list()) == {"cell_a", "cell_b"} + assert_cell_frames_equal(b.cells["cell_a"], full_cell) + assert_cell_frames_equal(b.cells["cell_b"], full_cell) + + +def test_poll_stops_on_until_and_complete(live_batch): + b, paths = live_batch + status = b.poll(interval=1, until=lambda batch: True, sleep=lambda s: None) + assert status.stopped_by == "until" and status.polls == 0 + for cell in b.cells.values(): + cell.source_complete = True + status = b.poll(interval=1, sleep=lambda s: None) + assert status.stopped_by == "complete" and status.polls == 0 + + +def test_poll_stops_on_cell_error(live_batch, monkeypatch): + b, paths = live_batch + + def boom(**kwargs): + raise RuntimeError("tester unplugged") + + monkeypatch.setattr(b.cells["cell_b"], "update", boom) + status = b.poll(interval=1, max_polls=3, sleep=lambda s: None) + assert status.stopped_by == "error" + assert isinstance(status.error, RuntimeError) + assert status.polls == 1