diff --git a/CHANGELOG-INTERNAL.md b/CHANGELOG-INTERNAL.md index 8dce8dd..6b02372 100644 --- a/CHANGELOG-INTERNAL.md +++ b/CHANGELOG-INTERNAL.md @@ -42,6 +42,39 @@ a `### BREAKING` section placed FIRST in that version, and each such line is pre - Help text and the flag tables in both READMEs now state that the flag is valid on its own. - Version bump deliberately NOT taken (convention 34): the owner closes it in `VERSION.txt`. +### Tests: the concurrency chaos suite now includes the SocketWatcher + +The suite's own charter is "many threads hammering one engine" and it names the failure it exists to +catch - threads tested in isolation. Since the SOCKET-layer work the session runs a FOURTH thread +and, since the connection-log fix, the capture thread reads the watcher's live map WITHOUT a lock - +and no test in that file ever passed the engine a `socket_source`, so neither the watcher's lifecycle +nor that lock-free read took part in any chaos. + +- `test_the_socket_watcher_survives_start_stop_cycles` - `CYCLES` start/stop rounds with a real + event stream running: fail-open (running implies a live capture thread AND a watcher), the watcher + cleared and its handle released on stop, and no `bean-socket-watcher` thread outliving its session. + Checked after a 0.3 s grace, because `SocketWatcher.stop()` only joins for 0.25 s. +- `test_the_capture_thread_reads_the_live_socket_map_under_churn` - the real new surface: the capture + thread resolving pids from the map while the watcher thread mutates it, the watchdog republishes it + wholesale, and settings and targeting churn underneath. Runs on `FastDivert` and on the SAME ports + the event source announces, or `pid_for` would always miss and a green run would prove nothing. +- **The crashlog watch is the test, not decoration.** `_pid_for` / `_process_for` swallow into + `crashlog.once` by design, so a read that started raising would leave every other assertion green. + Patching `once` also defeats its `_once_seen` dedupe, so repeated failures stay visible instead of + collapsing into one entry. +- Conclusiveness is a CONDITION the test waits for (`MIN_STAMPED` rows carrying the event stream's + pid), never a duration - the same reasoning as `MIN_BUILDS`/`MIN_ROWS`, and for the same reason: a + wall-clock budget lets machine speed decide whether the test proved anything. +- Both claims MUTATION-CONFIRMED: a `SocketWatcher.pid_for` that raises turns it red through the + crashlog watch, and an engine reverted to poller-only stamps 0 rows out of 2.8 million packets. + The first mutation also independently confirmed the failure-domain split below - it reported + `engine.ports` AND `engine.ports.pid`, not one collapsed entry. +- Stated in the docstring, so nobody assumes otherwise: this does NOT catch putting the lock back + into `pid_for`. A lock contends, it does not raise, and this test does not measure contention - + that property has its own guard in `test_socketwatch.py`. +- `_LiveSocketSource` is paced on purpose. Unpaced it saturates a core and starves the very threads + the test is about, which would make a green run meaningless. + ### Fixed: the connection log resolves the owner from the live socket map, not the poller `_pid_for` / `_process_for` asked `portmap.PortTable` - a snapshot refreshed a few times a second - diff --git a/tests/test_concurrency_chaos.py b/tests/test_concurrency_chaos.py index 096cfbe..35b13c0 100644 --- a/tests/test_concurrency_chaos.py +++ b/tests/test_concurrency_chaos.py @@ -17,6 +17,13 @@ * **no swallowed worker exception** - a worker that dies quietly is a worker whose job is not being done. +Since the SOCKET-layer work the session also runs a FOURTH thread, the +``SocketWatcher``, and the connection log reads its live ``port -> pid`` map from the +CAPTURE thread without taking a lock. That is a concurrency surface of its own, so +two tests at the bottom of this file give the engine a real event source: one drives +its lifecycle across start/stop cycles, the other reads the map from the capture +thread while the watcher mutates it and the watchdog republishes it. + Kept deliberately short (a few seconds); it is a smoke alarm, not a soak test. A note on the traffic these tests run on. ``SyntheticDivert`` sleeps once per @@ -35,6 +42,7 @@ from beantester.engine import BeanEngine from beantester.matchers import KIND_PROCESS, parse_matcher from beantester.settings import DEFAULT_SETTINGS, apply_settings +from beantester.socketwatch import CLOSE, CONNECT, SocketEvent from beantester.synthetic import SyntheticDivert, _SyntheticPacket, _SyntheticTCP from beantester.views import filter_sort_connections, traffic_totals from fakes import check @@ -48,6 +56,13 @@ MIN_BUILDS = 10 MIN_ROWS = 1000 +# Conclusiveness for the socket-map test below: enough connection rows stamped with +# the pid the event stream announces to prove the CAPTURE THREAD really read the live +# map. Like MIN_BUILDS/MIN_ROWS this is a condition the test waits for, never a speed +# assertion - a fixed duration would let machine speed decide whether it proved +# anything. +MIN_STAMPED = 50 + def _watch_worker_exceptions(): """Collect anything a thread raises (threads swallow exceptions by default).""" @@ -403,3 +418,176 @@ def test_stopping_joins_every_worker_thread(): check("no worker thread outlives stop()", not leaked, f"({leaked})") check("the thread handles are cleared", engine._t_cap is None and engine._t_inj is None and engine._t_wd is None) + + +# -- the SOCKET-layer watcher, which is the session's fourth thread ------------ # +class _LiveSocketSource: + """A socket-event source that keeps producing until it is closed. + + The fakes in ``tests/test_socketwatch*.py`` yield a fixed list and then park, + which is right for asserting a mapping but useless here: chaos needs the event + stream to still be MOVING while everything else moves. This one cycles through + its ports for ever, and closes a stale one every few events so the map both grows + and shrinks underneath whoever is reading it. + + The pacing matters. Without the sleep this saturates a core and starves the very + threads the test is about, which would make a green run meaningless. + """ + + def __init__(self, ports, pid=4242, delay=0.0005): + self._ports = list(ports) + self._pid = pid + self._delay = delay + self._closed = threading.Event() + + def __iter__(self): + i = 0 + while not self._closed.is_set(): + i += 1 + yield SocketEvent(CONNECT, self._pid, self._ports[i % len(self._ports)]) + if i % 4 == 0: + stale = self._ports[(i - 3) % len(self._ports)] + yield SocketEvent(CLOSE, self._pid, stale) + if self._delay: + time.sleep(self._delay) + + def close(self): + self._closed.set() + + +def test_the_socket_watcher_survives_start_stop_cycles(): + """The watcher is the session's FOURTH thread and it holds a WinDivert handle, so + it has to come and go exactly as capture, inject and watchdog do. Nothing in this + file used to start the engine with a socket source at all, so the watcher - and + the poller-vs-watcher table swap behind it - never took part in the chaos.""" + errors, restore = _watch_worker_exceptions() + try: + violations = [] + for cycle in range(CYCLES): + engine = BeanEngine() + engine.start("true", divert=SyntheticDivert(seed=cycle), + socket_source=_LiveSocketSource(range(4000, 4010))) + time.sleep(0.02) + + if engine.is_running(): + if engine._socketwatch is None: + violations.append(f"cycle {cycle}: running without a watcher") + capture = engine._t_cap + if capture is None or not capture.is_alive(): + violations.append(f"cycle {cycle}: running with no capture thread") + + engine.stop() + if engine._socketwatch is not None: + violations.append(f"cycle {cycle}: the watcher outlived stop()") + if engine._divert is not None: + violations.append(f"cycle {cycle}: the divert was not released") + + check("no fail-open or lifecycle violation across cycles", not violations, + f"({violations[:3]})") + check("no worker thread raised", not errors, f"({errors[:3]})") + # the THREAD is what actually leaks, and stop() only joins for 0.25 s + time.sleep(0.3) + leaked = [t.name for t in threading.enumerate() if "socket-watcher" in t.name] + check("no watcher thread outlived its session", not leaked, f"({leaked})") + finally: + restore() + + +def test_the_capture_thread_reads_the_live_socket_map_under_churn(): + """The connection log resolves the owning pid from the live SOCKET map, and it + does so ON THE CAPTURE THREAD, without a lock, while the watcher thread mutates + that map and the watchdog republishes it wholesale. That is a concurrency surface + no test covered: ``test_socketwatch.py`` hammers the map with a synthetic reader + in isolation, and the tests in this file never gave the engine a socket source. + + Why the crashlog watch is the whole test, not decoration: ``engine._pid_for`` and + ``_process_for`` SWALLOW their exceptions into ``crashlog.once`` by design (a + broken port table must not kill a session over a display name). So a read that + started raising under concurrency would leave every assertion below green and the + only trace would be a crash entry nobody looked at. Patching ``once`` also defeats + its ``_once_seen`` dedupe, so repeated failures are visible rather than collapsing + into one. + + What this WOULD catch (mutation-confirmed): a ``SocketWatcher.pid_for`` that + raises, and an engine that stops consulting the live map at all (rows then never + get stamped and the conclusiveness check fails). + + What it does NOT catch, so nobody assumes otherwise: putting the LOCK back into + ``pid_for``. A lock does not raise - it contends, and this test does not measure + contention. That property has its own guard, + ``test_socketwatch.py::test_pid_for_takes_no_lock_because_the_capture_thread_calls_it``. + """ + from beantester import crashlog + + errors, restore_hook = _watch_worker_exceptions() + swallowed = [] + real_once = crashlog.once + crashlog.once = lambda subsystem, exc: swallowed.append(f"{subsystem}: {exc!r}") + + pid = 4242 + engine = BeanEngine() + apply_settings(engine, DEFAULT_SETTINGS, lambda *_: None) + # the SAME ports FastDivert generates, or pid_for would always miss and a green + # run would prove nothing at all + engine.start("true", divert=FastDivert(), + socket_source=_LiveSocketSource(range(3000, 3500), pid=pid)) + stop = threading.Event() + problems = [] + + def churn(): + """Settings and targeting move under the capture thread, as in the GUI.""" + i = 0 + while not stop.is_set(): + try: + i += 1 + s = dict(DEFAULT_SETTINGS) + s.update(loss=i % 40, latency=i % 150, dup=i % 6, down=i % 400, + dst_port="80,443,!8080" if i % 2 else "") + apply_settings(engine, s, lambda *_: None) + if i % 3 == 0: + matcher = parse_matcher("python,!nonexistent_xyz", KIND_PROCESS) + engine.set_target(True, engine.target_for(matcher)) + else: + engine.set_target(False) + except Exception as exc: # pragma: no cover - the bug + problems.append(f"churn: {type(exc).__name__}: {exc}") + return + time.sleep(0.01) + + churner = threading.Thread(target=churn, name="applier", daemon=True) + churner.start() + + # Wait for the CONDITION (rows actually stamped from the map), not for a duration. + stamped = 0 + soft_deadline = time.monotonic() + STRESS_SECONDS + hard_deadline = time.monotonic() + 30.0 + try: + while time.monotonic() < hard_deadline: + stamped = sum(1 for c in engine.connections_snapshot(limit=None) + if c.get("pid") == pid) + if time.monotonic() >= soft_deadline and stamped > MIN_STAMPED: + break + if not engine.is_running(): + problems.append(f"the engine stopped by itself (fault={engine.fault})") + break + time.sleep(0.05) + finally: + stop.set() + churner.join(timeout=10) + seen = engine.stats_snapshot()["seen"] + engine.stop() + crashlog.once = real_once + restore_hook() + + # Diagnostics first: a swallowed read explains anything else that looks odd. + check("nothing was swallowed into the crash log", not swallowed, f"({swallowed[:3]})") + check("no thread raised", not errors, f"({errors[:3]})") + check("nothing went wrong on the driving side", not problems, f"({problems[:3]})") + # Then conclusiveness: prove the capture thread really read the LIVE map. + check("connection rows were stamped from the live socket map", + stamped > MIN_STAMPED, + f"({stamped} rows carry the event stream's pid; {seen} packets seen)") + check("the engine did not fault", engine.fault is None, f"({engine.fault})") + time.sleep(0.3) + leaked = [t.name for t in threading.enumerate() if "socket-watcher" in t.name] + check("no watcher thread outlived the test", not leaked, f"({leaked})")