diff --git a/CHANGELOG.md b/CHANGELOG.md index b6d3c6c..aa60d45 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,13 @@ The format follows [Keep a Changelog](https://keepachangelog.com/); versions fol ### Fixed +- **A command-line run stops on time even while its console window is paused.** + Selecting text in the console window pauses the program's output until the selection + ends. A run that reached its `--duration`, or hit a failure, used to wait for that + output before stopping, so traffic stayed impaired for as long as the console was + paused. The run now stops first and prints the reason afterwards. When a scenario + fails, its "Scenario stopped" line now comes after "Stop.". + - **Editing "Target process" during a session changes nothing until you click "Apply changes".** The field used to reach the running session on its own. Clearing it to type a new name, or stopping halfway through an expression such as `re:^fire(`, switched diff --git a/beantester/engine.py b/beantester/engine.py index d52de28..b946690 100644 --- a/beantester/engine.py +++ b/beantester/engine.py @@ -21,6 +21,12 @@ ``atexit`` hook guarantees the handle is released even on an abrupt shutdown. The same watchdog enforces the session deadline (``duration``). +The log is the caller's code, so it can block (a console paused by a text +selection, a pipe nobody reads) or raise. Nothing that reaches it may stand +between noticing a stop condition and closing the divert: every stop path +closes first and SAYS why afterwards (``_stop_locked``'s ``say``), and +``log`` never lets the sink's exception out. + What this actually sustains (the one place that number lives) ------------------------------------------------------------- "150 000 packets a second" appears in several cost arguments around this @@ -131,7 +137,7 @@ def _stop_live_engines(): class BeanEngine: def __init__(self, log_fn=lambda *_: None): - self.log = log_fn + self._log_fn = log_fn self.core = BeanCore() self._divert = None self._running = False @@ -212,6 +218,25 @@ def set_seed(self, seed): def effective_seed(self): return self._effective_seed + def log(self, msg): + """Hand one line to the log the owner passed in - which may raise. + + The sink is foreign code and it is called from the worker threads and + from every stop. MEASURED 2026-09-28: a sink that raised on the deadline + line killed the watchdog, and the session then never stopped at all, with + nothing left watching it. Whatever the sink raises is recorded and ends + here. A sink that BLOCKS cannot be fixed from this side, which is why the + stop paths close the divert before they say anything (``_stop_locked``). + """ + try: + self._log_fn(msg) + except Exception as _exc: + crashlog.note(_exc, "engine.log") + + def _say(self, lines): + for line in lines: + self.log(line) + def log_event(self, kind, text): now = time.monotonic() elapsed = (now - self._start_mono) if self._start_mono else 0.0 @@ -975,8 +1000,10 @@ def _start_locked(self, filt, divert, duration, socket_source=None, narrow=False # stop() closes the divert, stops/joins whatever DID start, and clears # _running. Then re-raise so the caller (GUI _finish_start, CLI) reports the # failure instead of believing the session is live. Convention 20. - self.log(T("log.engine_fault", e=str(exc))) - self.stop(reason="fault") + # The fault is said BY the stop, after the divert is closed: said first, a + # log that blocks held the divert open for as long as it blocked. Called + # directly because start() already holds _stop_lock. + self._stop_locked("fault", say=(T("log.engine_fault", e=str(exc)),)) raise # Two extra tries, ~0.45 s in total, and ONLY for "the device does not exist". @@ -1092,7 +1119,7 @@ def stop(self, reason="user"): with self._stop_lock: self._stop_locked(reason) - def _worker_stop(self, reason): + def _worker_stop(self, reason, say=()): """Stop initiated BY one of the engine's own worker threads. It must not BLOCK on ``_stop_lock``. A concurrent external ``stop()`` holds @@ -1105,18 +1132,30 @@ def _worker_stop(self, reason): and do the stop ourselves (the uncontended deadline / fault case); if another stop already holds it, return AT ONCE so its join of this thread completes and this thread dies. + + ``say`` is said either way: by ``_stop_locked`` once the divert is closed, or + here when bowing out, where the stop holding the lock does the closing. """ if not self._stop_lock.acquire(blocking=False): + self._say(say) return try: - self._stop_locked(reason) + self._stop_locked(reason, say) finally: self._stop_lock.release() - def _stop_locked(self, reason): + def _stop_locked(self, reason, say=()): """The stop body. The caller MUST hold ``_stop_lock`` - ``stop()`` blocks to - take it, ``_worker_stop`` takes it without blocking.""" + take it, ``_worker_stop`` takes it without blocking. + + ``say`` - the lines explaining WHY (a fault, the deadline). They are said + here, at the very end of the teardown just before "Stop.", and never by the + caller before the call: the log is the caller's code and can block, and + MEASURED 2026-09-28 on every stop path, a log said first held the divert + open for exactly as long as the log blocked. Said even when there is + nothing left to stop, so a caller that lost the race still reports.""" if not self._running: + self._say(say) return self._running = False # Cleared EARLY (it used to be set near the end): once the deadline is gone, @@ -1210,6 +1249,11 @@ def _stop_locked(self, reason): # balance - they were dropped BY the shutdown, not lost in transit. self._bump("drop_shutdown", discarded) _LIVE_ENGINES.discard(self) + # Only now, with nothing of the session left: see the docstring. Said at the + # divert close, a held log kept the rest of the teardown waiting behind it - + # the system-wide timer request, the switch interval, the queued packets and + # the atexit registration all stayed until the log moved (measured). + self._say(say) self.log(T("log.stop")) # How long the capture thread waits on _stop_lock before re-checking WHO holds @@ -1217,7 +1261,7 @@ def _stop_locked(self, reason): # long enough not to spin. See _fault_stop_blocking. FAULT_LOCK_POLL_S = 0.05 - def _fail_stop(self, error, blocking=True): + def _fail_stop(self, error, blocking=True, lead=()): """A worker died: stop the session so the network is never left impaired. ``blocking`` picks HOW the stop is taken, and the two callers genuinely differ: @@ -1244,8 +1288,13 @@ def _fail_stop(self, error, blocking=True): 2 s join timeout. Never a deadlock - but "cannot" was too strong, and a sentence like that is what stops the next session from looking. NOT reproduced: found by reading, and the window is a few instructions wide. + + ``lead`` - lines the caller has to say ahead of the fault line (the recv + error). Handed over instead of said, and the fault line with them: the stop + says them once the divert is closed (``_stop_locked``). """ if not self._running: + self._say(lead) return # First fault wins. The watchdog's "worker thread died unexpectedly" is a # SYMPTOM of the real error - if the capture thread recorded the cause a @@ -1253,13 +1302,13 @@ def _fail_stop(self, error, blocking=True): # of the report. Both are still LOGGED: two failures are two events. if not self.fault: self.fault = str(error) - self.log(T("log.engine_fault", e=str(error))) + say = (*lead, T("log.engine_fault", e=str(error))) if blocking: - self._fault_stop_blocking() + self._fault_stop_blocking(say) else: - self._worker_stop(reason="fault") + self._worker_stop("fault", say) - def _fault_stop_blocking(self): + def _fault_stop_blocking(self, say=()): """Take ``_stop_lock`` for a capture-thread fault: wait for a start, never for another stop. @@ -1269,13 +1318,14 @@ def _fault_stop_blocking(self): the case the blocking path exists for), and no longer running means a stop already owns the teardown, is closing the divert, and is joining THIS thread with a 2.0 s timeout. Waiting on that one buys nothing and costs the user a - two-second STOP. + two-second STOP. ``say`` as in ``_worker_stop``. """ while not self._stop_lock.acquire(timeout=self.FAULT_LOCK_POLL_S): if not self._running: + self._say(say) return try: - self._stop_locked("fault") + self._stop_locked("fault", say) finally: self._stop_lock.release() @@ -1338,13 +1388,15 @@ def _watchdog_loop(self): except Exception as _exc: crashlog.note(_exc, "engine") if deadline_reached(self._deadline, time.monotonic()): - self.log(T("log.duration_reached", v=f"{self._duration:g}")) # _worker_stop, not stop(): this runs on the watchdog thread, and a # user pressing STOP at the same instant holds _stop_lock while joining # this very thread. Blocking on the lock here would hang STOP for its # 2 s join timeout (measured 2.09 s); the user's stop already closes # the divert, so we can just bow out. - self._worker_stop(reason="duration") + # The line goes WITH the stop, not before it: a console paused by a + # text selection kept the session going past its --duration. + self._worker_stop( + "duration", (T("log.duration_reached", v=f"{self._duration:g}"),)) return for t in (self._t_cap, self._t_inj): if t is not None and not t.is_alive(): @@ -1420,8 +1472,8 @@ def _capture_loop(self): # The divert is still open but nothing drains it any more: # WinDivert would keep queueing (and then dropping) the user's # packets. Fail OPEN - stop the session and release the driver. - self.log(f"{T('log.recv_error')}: {e}") - self._fail_stop(e) + # The line is handed to the stop, which says it after closing. + self._fail_stop(e, lead=(f"{T('log.recv_error')}: {e}",)) break now = time.monotonic() # BEAT FIRST, then clear the flag, and the order is the whole safety of diff --git a/beantester/scenario_runner.py b/beantester/scenario_runner.py index 599589a..67c69cb 100644 --- a/beantester/scenario_runner.py +++ b/beantester/scenario_runner.py @@ -116,8 +116,11 @@ def _loop(self, scenario, base, log): except Exception as exc: # noqa: BLE001 - the net itself try: crashlog.record(exc, "scenario_runner") - log(T("log.scenario_failed", e=f"{type(exc).__name__}: {exc}")) + # Stop FIRST, then say why: `log` is the caller's and can block (a + # paused console), and said first it kept the session impairing + # traffic for as long as it blocked (measured 2026-09-28). self.engine.worker_failed(exc) + log(T("log.scenario_failed", e=f"{type(exc).__name__}: {exc}")) except Exception as _exc: # A safety net that can itself fall through is not one. crashlog.note(_exc, "scenario_runner") diff --git a/tests/test_failsafe.py b/tests/test_failsafe.py index b731cef..c010b1b 100644 --- a/tests/test_failsafe.py +++ b/tests/test_failsafe.py @@ -9,11 +9,14 @@ * a session stops itself at its ``duration`` deadline, * a dead worker thread makes the engine stop (= release the divert) and say so, * the GUI survives a broken tick, never calls Tcl from a worker thread, and - always releases the engine when the window closes. + always releases the engine when the window closes, + * no stop waits for the log, which is the caller's code and can block or raise. """ import threading import time +import pytest + from beantester.engine import _LIVE_ENGINES, BeanEngine, deadline_reached from beantester.i18n import T from fakes import FakePacket, check @@ -1317,6 +1320,184 @@ def test_the_stall_check_answers_no_for_every_state_that_is_not_one(): eng._capture_has_stalled() is False) +# --- the log is the caller's code: it can block, and it can raise ------------- # +# +# On the CLI the log is a plain write to stderr. A console whose text is being +# selected with the mouse holds that write until the selection ends, and so does a +# pipe nobody reads any more. Every stop path used to SAY why before it stopped, so +# the divert stayed open for exactly as long as the log was held - measured +# 2026-09-28 on all six paths below, and on a session past its --duration. + + +class _FrozenLog: + """A log that stops returning at the first line containing ``needle``.""" + + def __init__(self): + self.needle = None + self.lines = [] + self.frozen = threading.Event() + self.thaw = threading.Event() + self.raised = None # the error a start() on another thread got + + def __call__(self, line): + line = str(line) + self.lines.append(line) + if self.needle and self.needle in line and not self.thaw.is_set(): + self.frozen.set() + self.thaw.wait(10) + + +def _frozen_at_the_deadline(eng, log, monkeypatch): + log.needle = T("log.duration_reached", v="0.2") + divert = QuietDivert() + eng.start("test", divert=divert, duration=0.2) + return divert + + +def _frozen_at_a_recv_error(eng, log, monkeypatch): + log.needle = "driver went away" + divert = ExplodingDivert(packets=0) + eng.start("test", divert=divert) + return divert + + +def _frozen_at_a_dead_worker(eng, log, monkeypatch): + log.needle = "died unexpectedly" + monkeypatch.setattr(eng, "_capture_loop", lambda: None) # ends at once + divert = QuietDivert() + eng.start("test", divert=divert) + return divert + + +def _frozen_at_a_stall(eng, log, monkeypatch): + log.needle = "stopped making progress" + eng.CAPTURE_STALL_S = 0.3 + divert = _StallingDivert() + eng.start("test", divert=divert) + return divert + + +def _frozen_at_a_foreign_worker(eng, log, monkeypatch): + log.needle = "the timeline broke" + divert = QuietDivert() + eng.start("test", divert=divert) + threading.Thread(target=eng.worker_failed, + args=(RuntimeError("the timeline broke"),), daemon=True).start() + return divert + + +def _frozen_at_a_failed_start(eng, log, monkeypatch): + log.needle = "would not start" + + def refuse(): + raise RuntimeError("the resolver would not start") + + monkeypatch.setattr(eng._resolver, "start", refuse) + divert = QuietDivert() + + def start(): + # start() blocks in the log it says on its way out, so it runs here + try: + eng.start("test", divert=divert) + except RuntimeError as exc: + log.raised = str(exc) + + threading.Thread(target=start, daemon=True).start() + return divert + + +@pytest.mark.parametrize("path", [ + _frozen_at_the_deadline, _frozen_at_a_recv_error, _frozen_at_a_dead_worker, + _frozen_at_a_stall, _frozen_at_a_foreign_worker, _frozen_at_a_failed_start, +], ids=lambda path: path.__name__.removeprefix("_frozen_at_")) +def test_a_stop_closes_the_divert_while_the_log_is_still_blocked(path, monkeypatch): + """Close first, SAY why afterwards - and still in the order a tester reads. + + The divert has to close while the log is still held: that is the whole fix. + The rest of the teardown too, so a held log keeps nothing of the session. + Then, once the log moves again, the reason has to come before "Stop.", the + order the lines had before (convention: a log that reads backwards cannot tell + a tester what happened when). + """ + log = _FrozenLog() + eng = BeanEngine(log_fn=log) + divert = path(eng, log, monkeypatch) + try: + check("the stop path reached the log and the log is held", + log.frozen.wait(5), f"({log.lines})") + check("the divert is closed while the log is still held (network restored)", + _wait_until(lambda: divert.closed, 3.0), + f"(running={eng.is_running()}, lines={log.lines})") + check("and the session is no longer running", eng.is_running() is False) + # The rest of the teardown does not wait for the log either: the + # system-wide timer request, the switch interval and the atexit entry. + check("and nothing of the session is still held while the log is", + eng not in set(_LIVE_ENGINES) and not eng._fine_timers + and not eng._fast_switch, + f"(tracked={eng in set(_LIVE_ENGINES)}, timers={eng._fine_timers}, " + f"switch={eng._fast_switch})") + finally: + log.thaw.set() + if hasattr(divert, "release"): + divert.release() + eng.stop() # waits for the stop that was held, so the log is complete + + reason = next((i for i, line in enumerate(log.lines) if log.needle in line), None) + stop = [i for i, line in enumerate(log.lines) if line == T("log.stop")] + check("the reason is said once the log moves again", reason is not None, + f"({log.lines})") + check("and before the one 'Stop.' line", len(stop) == 1 and reason < stop[0], + f"({log.lines})") + if path is _frozen_at_a_failed_start: + check("the caller still gets the start's own error", + _wait_until(lambda: log.raised is not None) + and log.raised == "the resolver would not start", f"({log.raised!r})") + + +def test_a_log_that_raises_cannot_cancel_a_stop(monkeypatch): + """A log that RAISES took the stop down with it. + + MEASURED 2026-09-28: a log raising on the deadline line killed the watchdog - + the session never stopped, and nothing was left to notice anything else. At + START it was worse: the announcement raised into the failure handler, whose own + fault line raised again before it could stop, so the divert stayed open and the + caller got the log's error instead of the real one. + """ + def broken(line): + raise RuntimeError("the log is gone") + + eng = BeanEngine(log_fn=broken) + divert = QuietDivert() + eng.start("test", divert=divert, duration=0.2) + check("deadline: a session with a broken log still starts", eng.is_running() is True) + + def stopped_completely(): + kinds = [(e[2], e[3]) for e in eng.events_snapshot()] + return (not eng.is_running() and divert.closed + and ("STOP", "events.duration_reached") in kinds) + + check("deadline: and still stops at its deadline, teardown and all", + _wait_until(stopped_completely, 3.0), + f"(running={eng.is_running()}, closed={divert.closed})") + + def refuse(): + raise RuntimeError("the resolver would not start") + + eng = BeanEngine(log_fn=broken) + monkeypatch.setattr(eng._resolver, "start", refuse) + divert = QuietDivert() + raised = None + try: + eng.start("test", divert=divert) + except RuntimeError as exc: + raised = str(exc) + check("failed start: the caller gets the start's error, not the log's", + raised == "the resolver would not start", f"({raised!r})") + check("failed start: the divert is closed", divert.closed is True) + check("failed start: nothing is left running or tracked", + eng.is_running() is False and eng not in set(_LIVE_ENGINES)) + + # --- START and STOP must survive the worker ending badly ---------------------- # diff --git a/tests/test_mutation_registry.py b/tests/test_mutation_registry.py index bec3f46..a5b50d5 100644 --- a/tests/test_mutation_registry.py +++ b/tests/test_mutation_registry.py @@ -123,6 +123,17 @@ "new": " pass", "test": "test_a_timeline_that_breaks_takes_the_session_down_with_it", }, + { + # Back to saying it first: a held console keeps the session impairing + # traffic for as long as it is held. + "label": "scenario: a broken timeline says so before the engine is told", + "file": "beantester/scenario_runner.py", + "old": " self.engine.worker_failed(exc)\n" + " log(T(\"log.scenario_failed\", e=f\"{type(exc).__name__}: {exc}\"))", + "new": " log(T(\"log.scenario_failed\", e=f\"{type(exc).__name__}: {exc}\"))\n" + " self.engine.worker_failed(exc)", + "test": "test_a_broken_timeline_stops_the_session_before_it_says_so", + }, { # The shipped stop(): a flag and a return. The thread is still between two # steps and applies one more set of settings after the caller moved on. @@ -444,6 +455,87 @@ "new": " if False:", "test": "test_a_capture_thread_that_is_alive_but_no_longer_moving_fails_open", }, + { + # The shipped order before 2026-09-28: say it, then stop. A console held by + # a text selection kept the session running past its --duration. + "label": "engine: the deadline is said before the session is stopped", + "file": "beantester/engine.py", + "old": " self._worker_stop(\n" + " \"duration\", (T(\"log.duration_reached\", " + "v=f\"{self._duration:g}\"),))", + "new": " self.log(T(\"log.duration_reached\", v=f\"{self._duration:g}\"))\n" + " self._worker_stop(\"duration\")", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # The fault line said by the worker before it asks for the stop: every + # watchdog fault and the foreign-worker door hold the divert open again. + "label": "engine: a fault is said before the session is stopped", + "file": "beantester/engine.py", + "old": " say = (*lead, T(\"log.engine_fault\", e=str(error)))", + "new": " self._say((*lead, T(\"log.engine_fault\", e=str(error))))\n" + " say = ()", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # The stop itself says why before it closes - the callers hand the lines + # over correctly and it still waits for the log. + "label": "engine: a stop says why before it closes the divert", + "file": "beantester/engine.py", + "old": " if self._divert is not None:\n" + " try:\n" + " self._divert.close()", + "new": " self._say(say)\n" + " say = ()\n" + " if self._divert is not None:\n" + " try:\n" + " self._divert.close()", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # The first version of this fix: said at the divert close, before the rest + # of the teardown - a held log kept the timer request, the switch interval + # and the atexit entry until it moved. + "label": "engine: a stop says why before the rest of its teardown", + "file": "beantester/engine.py", + "old": " self.log_event(\"STOP\", self.EVENT_BY_REASON.get(reason, " + "\"events.stopped\"))", + "new": " self._say(say)\n" + " say = ()\n" + " self.log_event(\"STOP\", self.EVENT_BY_REASON.get(reason, " + "\"events.stopped\"))", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # START's failure handler said its fault first, then stopped. + "label": "engine: a failed start says its fault before stopping", + "file": "beantester/engine.py", + "old": " self._stop_locked(\"fault\", say=(T(\"log.engine_fault\", e=str(exc)),))", + "new": " self.log(T(\"log.engine_fault\", e=str(exc)))\n" + " self._stop_locked(\"fault\")", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # The capture thread said the recv error itself before asking for the stop. + "label": "engine: a recv error is said before the stop", + "file": "beantester/engine.py", + "old": " self._fail_stop(e, lead=(f\"{T('log.recv_error')}: {e}\",))", + "new": " self.log(f\"{T('log.recv_error')}: {e}\")\n" + " self._fail_stop(e)", + "test": "test_a_stop_closes_the_divert_while_the_log_is_still_blocked", + }, + { + # The sink's exception reaches the watchdog again: it dies on the deadline + # line and the session never stops. + "label": "engine: a log that raises escapes into the caller again", + "file": "beantester/engine.py", + "old": " try:\n" + " self._log_fn(msg)\n" + " except Exception as _exc:\n" + " crashlog.note(_exc, \"engine.log\")", + "new": " self._log_fn(msg)", + "test": "test_a_log_that_raises_cannot_cancel_a_stop", + }, { # The other direction, and the more expensive one to get wrong: without the # phase, a thread parked in recv() on a link with no traffic looks exactly diff --git a/tests/test_scenario_runner.py b/tests/test_scenario_runner.py index 02693b0..2f3bc29 100644 --- a/tests/test_scenario_runner.py +++ b/tests/test_scenario_runner.py @@ -334,3 +334,31 @@ def explode(*_a, **_kw): any("TypeError" in str(line) for line in logged), f"({logged!r})") check("a broken timeline is still not a FINISHED one", not runner.finished, "'finished' means the timeline ran out - see its own docstring") + + +def test_a_broken_timeline_stops_the_session_before_it_says_so(monkeypatch): + """The engine is told FIRST; the line for the user comes after. + + ``log`` is the caller's and can block: on the CLI it is a write to a console + that a text selection holds until it ends. Said first, it kept the session + impairing traffic for as long as the console was held (measured 2026-09-28). + Recorded here as "how many failures had the engine been told about when the + line was said". + """ + def explode(*_a, **_kw): + raise TypeError("a value the engine cannot use") + + monkeypatch.setattr(scenario_runner, "apply_settings", explode) + monkeypatch.setattr(scenario_runner, "settings_summary", lambda s, lang: "summary") + + engine = FakeEngine() + said = [] + runner = ScenarioRunner(engine) + runner.start(FakeScenario(loop=False, duration=5.0), base_settings={}, + log=lambda line: said.append((len(engine.failures), str(line)))) + _join(runner) + + told_at = [told for told, line in said if "TypeError" in line] + check("the failure is said", told_at, f"({said!r})") + check("and only once the engine has been told to stop", told_at == [1], + f"({said!r})")