Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
88 changes: 70 additions & 18 deletions beantester/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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".
Expand Down Expand Up @@ -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
Expand All @@ -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,
Expand Down Expand Up @@ -1210,14 +1249,19 @@ 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
# it. Short enough that a racing external STOP is never delayed noticeably,
# 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:
Expand All @@ -1244,22 +1288,27 @@ 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
# moment earlier, overwriting it here would throw away the only useful half
# 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.

Expand All @@ -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()

Expand Down Expand Up @@ -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():
Expand Down Expand Up @@ -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
Expand Down
5 changes: 4 additions & 1 deletion beantester/scenario_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading
Loading