From af28f1b33fc3a0247bf9bbaa043311932f0b13ef Mon Sep 17 00:00:00 2001 From: DonislawDev Date: Sat, 25 Jul 2026 10:50:11 +0200 Subject: [PATCH] fix(engine): stamp the connection log from the live socket map The connection log asked the POLLER who owns a local port - a snapshot refreshed a few times a second - so a flow that opened AND finished inside one refresh interval left a row with no owner at all. Short-lived connections are what this tool gets pointed at (a browser opens hundreds a minute), so that was the common case, not an edge one. The engine now asks the SOCKET-event map first and falls back to the poller, and the CONNECT event lands ~0.1 ms before the SYN reaches the NETWORK layer, so the row is stamped from its very first packet. The read had to become LOCK-FREE first, and that is the substance of this change. pid_for used to take _lock, which the watcher thread holds on every socket event and reconcile holds across a whole snapshot merge - calling that from _log_conn would have let the CAPTURE THREAD queue behind maintenance, which is the stall convention 20 exists to prevent. So pid_for reads the reference once and does a C-level get on int keys (the idiom PortTable.pid_for already uses), and reconcile builds the new state to the side and publishes it by REASSIGNMENT, so its O(n) pass is atomic to a reader instead of being observed half-applied. - names stay cheap=True (cache or nothing), so a pid can reach a row before its name does; that is written down rather than glossed over - the shared lookup is _live_pid, deliberately WITHOUT a handler: the first attempt had _process_for delegate to _pid_for, which swallows, so a broken port table reported one failure instead of two and the name domain lost its voice. test_processes.py caught it; the two callers are two failure domains and each wraps the raising helper itself - _Broken in that test now models what the engine actually calls, and the insight in its comment (a fake missing the keyword raises TypeError, so the test passes while exercising the wrong failure) moved to the keyword that now matters - four new tests, each mutation-checked: poller-only turns the wiring test red, restoring the lock turns the no-lock guard red, mutating reconcile in place turns the identity guard red - the lock-free safety claim is RUN, not asserted: a reader hammering pid_for while another thread inserts, deletes and republishes the map never raised and only ever saw the real pid - test_hot_path.py (packet threads must never reach the OS) stays green: both lookups are dict reads Co-Authored-By: Claude Opus 5 --- CHANGELOG-INTERNAL.md | 48 ++++++++++++++++++++++ CHANGELOG.md | 10 +++++ beantester/engine.py | 64 ++++++++++++++++++++++------- beantester/socketwatch.py | 35 +++++++++++++--- tests/test_processes.py | 13 +++--- tests/test_socketwatch.py | 69 +++++++++++++++++++++++++++++++ tests/test_socketwatch_wiring.py | 70 +++++++++++++++++++++++++++++++- 7 files changed, 283 insertions(+), 26 deletions(-) diff --git a/CHANGELOG-INTERNAL.md b/CHANGELOG-INTERNAL.md index 5eb8658..8dce8dd 100644 --- a/CHANGELOG-INTERNAL.md +++ b/CHANGELOG-INTERNAL.md @@ -42,6 +42,54 @@ 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`. +### 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 - +so a flow that opened AND finished inside one refresh interval left a row with no owner at all. +Short-lived connections are what this tool gets pointed at, so that was the common case rather than +an edge one. The engine now asks the SOCKET-event map first (`SocketWatcher.pid_for`) and falls back +to the poller, so a row is stamped from its FIRST packet: the CONNECT event lands ~0.1 ms before the +SYN reaches the NETWORK layer (measured 2026-07-22). + +- **The read had to become LOCK-FREE first, and that is the substance of this change.** `pid_for` + used to take `_lock` - which the watcher thread holds on every socket event, and which + `reconcile` holds across a whole snapshot merge. Calling that from `_log_conn` would have let the + CAPTURE THREAD queue behind maintenance, which is exactly the stall convention 20 exists to + prevent. `pid_for` now reads the reference once and does a C-level `get` on int keys (the same + idiom as `PortTable.pid_for`), and `reconcile` builds the new state to the side and publishes it + by REASSIGNMENT, so its O(n) pass is atomic to a reader instead of being observed half-applied. +- Name resolution is unchanged and still `cheap=True` (cache or nothing), so a brand-new pid can + reach a row BEFORE its name does - names are warmed by the watchdog. That is written into the + docstring rather than glossed over: a PID with no name yet is still an answer, and `_log_conn` + keeps retrying the name while packets arrive. +- Deliberately NOT done: warming names from the watcher's map as well. It would close the remaining + name lag, but it adds per-pid OS calls to the watchdog and deserves its own measurement first. +- The shared lookup lives in `_live_pid`, which has NO handler of its own on purpose. The first + attempt had `_process_for` delegate to `_pid_for`, and because that one swallows and records under + `engine.ports.pid`, a broken port table reported ONE failure instead of two - the name domain + could no longer speak for itself. `test_processes.py::test_engine_records_a_broken_port_table_instead_of_going_quiet` + caught it. The two callers are two failure domains and each wraps the raising helper itself, which + is the same principle the watchdog already applies to refresh-vs-trim ("different jobs, different + failure domains"). +- `_Broken` in that test now models what the engine actually calls (`pid_for` + `name_of(cheap=)`) + instead of `process_for_port`, which the engine no longer touches. The insight in its comment - + that a fake missing the keyword raises TypeError and the test then passes while exercising the + wrong failure - still holds, so it moved to the keyword that now matters rather than being + deleted. +- New tests, each MUTATION-CHECKED rather than trusted: + `test_socketwatch_wiring.py::test_a_fresh_socket_stamps_the_connection_row_from_the_live_map` + (reverting the engine to poller-only turns it red), plus in `test_socketwatch.py` + `::test_pid_for_takes_no_lock_because_the_capture_thread_calls_it` (restoring the lock turns it + red) and `::test_reconcile_publishes_a_new_map_instead_of_mutating_in_place` (mutating in place + turns it red). +- `test_socketwatch.py::test_a_lock_free_reader_survives_writes_in_flight` RUNS the safety claim + instead of asserting it in prose: a reader hammering `pid_for` while another thread inserts, + deletes and republishes the whole map never raised and only ever saw the real pid. +- The engine test drives a GATED divert, so it asserts that the live map is consulted rather than + that the capture thread happened to lose a race with the watcher thread. +- `test_hot_path.py` (packet threads must never reach the OS) was run explicitly and stays green: + both lookups are dict reads. + ### Changed: SocketEvent carries only what the map is for (closes socket-event-fields) The SOCKET-layer event used to carry `proto`, `remote_ip`, `remote_port` and `outbound` "for the diff --git a/CHANGELOG.md b/CHANGELOG.md index a5f1c97..46a0d59 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,6 +36,16 @@ The format follows [Keep a Changelog](https://keepachangelog.com/); versions fol ### Fixed +- **Short-lived connections now show which program they belong to.** The Connections table works + out the owning program by asking Windows which application holds each socket, and that answer + used to come from a list refreshed a few times a second. A connection that opened and finished + in between two of those refreshes was never on the list, so its row stayed blank in the Process + and PID columns - and brief connections like that are exactly what you get when a browser, or an + app you are testing, opens hundreds of them a minute. The tool is now told who owns a connection + the moment it is created, so the row is filled in from its very first packet. The program's name + can still appear a moment after its process number, because names are looked up in the + background, but the row no longer stays empty. + - **A window you have moved now keeps its place when you switch language.** Changing the language rebuilds the whole main window, and a smaller window open at the time - Settings, for example - was torn down with it. It came back at the size and position it had the last time you closed it, diff --git a/beantester/engine.py b/beantester/engine.py index d3627f7..7443465 100644 --- a/beantester/engine.py +++ b/beantester/engine.py @@ -431,18 +431,25 @@ def _process_for(self, local_port): column was mostly "?" even when running as Administrator. """ try: - # allow_refresh=False is the whole point: process_for_port() otherwise - # calls refresh_if_stale(miss=True) when the port is unknown, which is - # four iphlpapi calls (and sometimes a psutil walk) ON THE CAPTURE - # THREAD - measured at ~16 a second against synthetic traffic. This is - # a SECOND path that did what targeting used to do; moving targeting off - # the hot path did nothing for it. The watchdog keeps the table fresh - # instead, exactly as it already does eviction and flow rotation. + # NOTHING here may reach the OS - this is the capture thread. The pid comes + # from _pid_for (live socket map first, poller second) and the NAME comes + # from the cache with cheap=True: never a refresh, never a psutil call. + # process_for_port() used to be called instead, and its allow_refresh=False + # was load-bearing for the same reason - left on, it would call + # refresh_if_stale(miss=True) for every unknown port, i.e. four iphlpapi + # calls (sometimes a psutil walk) in the packet path, measured at ~16 a + # second against synthetic traffic. The watchdog keeps the table fresh and + # warms the names instead, exactly as it already does eviction and rotation. # - # The cost is that a brand-new socket may read as "" for up to one - # refresh interval. _log_conn already retries while packets keep coming, - # so the row fills itself in rather than staying "?" for ever. - return self._ports.process_for_port(local_port, allow_refresh=False) + # A brand-new pid can therefore reach the row BEFORE its name does: the + # watcher supplies the pid instantly, while the name waits for the + # watchdog's warm_names. That is the honest split - a PID with no name yet + # is still an answer, and _log_conn keeps retrying the name while packets + # arrive, so the row fills itself in rather than staying "?" for ever. + pid = self._live_pid(local_port) + if pid is None: + return "" + return self._ports.name_of(pid, cheap=True) except Exception as _exc: # once(), not note(): this is the capture thread. A port table that # started failing turns every row's process into "?" - worth one @@ -453,15 +460,44 @@ def _process_for(self, local_port): def _pid_for(self, local_port): """PID owning ``local_port`` right now (None when unknown). - Same reasoning as ``_process_for``: resolved at capture time and stored, - because the socket is usually gone by the time the row is displayed. + Resolved at capture time and stored, because the socket is usually gone by the + time the row is displayed - and asked of the LIVE socket-event map first, the + poller second. That order is the point: the poller is a snapshot taken a few + times a second, so a flow that opens AND finishes inside one refresh interval + left a row with no owner at all, and short-lived connections are exactly what + this tool gets pointed at. The watcher is told the owner ~0.1 ms before the SYN + reaches the NETWORK layer (measured), so the row can be stamped from its first + packet. + + Neither path may touch the OS from here - this is the capture thread. The + watcher lookup is a lock-free dict read (see ``SocketWatcher.pid_for``) and the + poller's is a plain ``get`` on an already-cached map. Guarded by + ``tests/test_hot_path.py``. """ try: - return self._ports.pid_for(local_port) + return self._live_pid(local_port) except Exception as _exc: crashlog.once("engine.ports.pid", _exc) return None + def _live_pid(self, local_port): + """The owning pid: live socket map first, poller second. MAY RAISE. + + Deliberately without a handler of its own, because its two callers are two + different FAILURE DOMAINS and each has to be able to report for itself: a + broken pid lookup is ``engine.ports.pid``, a broken name lookup is + ``engine.ports``. Wrapping it here collapsed both into one record and + ``tests/test_processes.py::test_engine_records_a_broken_port_table_instead_of_going_quiet`` + caught exactly that - a port table that broke would have reported half of what + it does now. + """ + watcher = self._socketwatch # read ONCE: stop() clears it concurrently + if watcher is not None: + pid = watcher.pid_for(local_port) + if pid is not None: + return pid + return self._ports.pid_for(local_port) + def stats_snapshot(self): with self._slock: s = dict(self.st) diff --git a/beantester/socketwatch.py b/beantester/socketwatch.py index ef43894..813cae2 100644 --- a/beantester/socketwatch.py +++ b/beantester/socketwatch.py @@ -112,16 +112,23 @@ def reconcile(self, port_pid): longer lists are pruned only after being absent for TWO reconciles running, which spares a socket opened microseconds before the snapshot was taken (present via its event, not yet in that snapshot) from being evicted by it. + + The new state is built to the side and PUBLISHED BY REASSIGNMENT rather than + mutated in place, because ``pid_for`` reads this map WITHOUT a lock from the + capture thread: a reader has to see this pass either not at all or completely, + never with half the snapshot folded in and half the prunes applied. """ with self._lock: + merged = dict(self._ports) for port, pid in port_pid.items(): if port and pid and pid > 0: - self._ports[port] = pid - absent = set(self._ports) - set(port_pid) + merged[port] = pid + absent = set(merged) - set(port_pid) doomed = absent & self._suspect # absent twice running for port in doomed: - self._ports.pop(port, None) + merged.pop(port, None) self._suspect = absent - doomed # first-time absentees wait one pass + self._ports = merged # atomic swap for lock-free readers self._reconciles += 1 def snapshot(self): @@ -129,10 +136,28 @@ def snapshot(self): return dict(self._ports) def pid_for(self, port): + """Owning pid for a local port (``None`` when unknown). Takes NO LOCK. + + The CAPTURE THREAD reads this (``engine._pid_for``), and a lock here would be + precisely the thing this module must not do: the watcher thread holds ``_lock`` + on every socket event, and ``reconcile`` holds it across a whole snapshot + merge, so the packet path would queue behind maintenance. A stalled capture + thread means WinDivert is diverting into a queue nobody drains (convention 20). + + Lock-free is safe for the same reason it is in + :meth:`beantester.portmap.PortTable.pid_for`, with one addition. The reference + is read ONCE into a local, and a dict lookup on INT keys is C code that neither + releases the GIL nor calls back into Python, so it cannot interleave with + another thread's insert or delete: a reader sees the map either before or after + that write, never mid-resize. ``reconcile`` then publishes a NEW dict by + reassignment, so its O(n) pass is atomic to a reader instead of being observed + halfway through. Verified under load, not assumed - see + ``test_socketwatch.py::test_a_lock_free_reader_survives_writes_in_flight``. + """ if port is None: return None - with self._lock: - return self._ports.get(int(port)) + ports = self._ports # one atomic reference read, then a C-level get + return ports.get(int(port)) # -- name resolution: delegated, never duplicated -------------------------- # def refresh(self, now=None, force=False): diff --git a/tests/test_processes.py b/tests/test_processes.py index a2f21c3..e613cb6 100644 --- a/tests/test_processes.py +++ b/tests/test_processes.py @@ -336,14 +336,15 @@ def test_engine_records_a_broken_port_table_instead_of_going_quiet(monkeypatch): """The capture thread keeps going (a blank name beats a dead session), but the reason no longer disappears. ``once()``, not ``note()``: this is the hot path.""" class _Broken: - # the signature MATTERS: the engine reads with allow_refresh=False (it must - # never make the capture thread rebuild the table). A fake missing the - # keyword would raise TypeError instead, and the test would pass while - # exercising the wrong failure entirely. - def process_for_port(self, port, now=None, allow_refresh=True): + # The engine reads the PID (the live socket map first, this poller second) and + # then the NAME from the cache with cheap=True - it does not call + # process_for_port() any more. The signatures MATTER: a fake missing the + # ``cheap`` keyword would raise TypeError instead, and the test would pass + # while exercising the wrong failure entirely. + def pid_for(self, port): raise RuntimeError("boom") - def pid_for(self, port): + def name_of(self, pid, cheap=False): raise RuntimeError("boom") recorded = _spy_on_crashlog(monkeypatch) diff --git a/tests/test_socketwatch.py b/tests/test_socketwatch.py index 7857d91..63cc79d 100644 --- a/tests/test_socketwatch.py +++ b/tests/test_socketwatch.py @@ -229,3 +229,72 @@ def close(self): time.sleep(0.05) check("a mid-run failure IS recorded", len(recorded) == 1, f"({recorded})") w2.stop() + + +# -- the lock-free read the capture thread depends on ------------------------- # +def test_pid_for_takes_no_lock_because_the_capture_thread_calls_it(): + """A lock here would let the packet path queue behind the watcher thread or a + whole reconcile. Asserted mechanically rather than by reading the code: with a + lock that refuses to be taken, ``pid_for`` must still answer. + """ + w = _watcher() + w.apply(ev(CONNECT, 100, 5000)) # while the real lock still works + + class _Explodes: + def __enter__(self): + raise AssertionError("pid_for must not take the lock") + + def __exit__(self, *exc): + return False + + w._lock = _Explodes() + check("pid_for answers without taking the lock", w.pid_for(5000) == 100) + check("and still reports an unknown port as None", w.pid_for(9999) is None) + + +def test_reconcile_publishes_a_new_map_instead_of_mutating_in_place(): + """The identity swap is what makes the lock-free read safe across an O(n) pass: + mutating in place would let a reader observe half a reconcile.""" + w = _watcher() + w.apply(ev(CONNECT, 100, 5000)) + before = w._ports + w.reconcile({5000: 100, 80: 1}) + check("reconcile published a NEW dict", w._ports is not before) + check("with the merged content", w.snapshot() == {5000: 100, 80: 1}, + f"({w.snapshot()})") + + +def test_a_lock_free_reader_survives_writes_in_flight(): + """"Lock-free is safe here" is a SAFETY claim, so it is run in the conditions that + would break it instead of being asserted in a docstring: one thread reading a port + while another inserts, deletes and republishes the whole map underneath it. A torn + read would surface as an exception, or as a value that is neither the pid nor None. + """ + w = _watcher() + w.apply(ev(CONNECT, 100, 5000)) + errors, values, reads = [], set(), [0] + stop = threading.Event() + + def reader(): + try: + while not stop.is_set(): + values.add(w.pid_for(5000)) # a set, so this stays tiny + reads[0] += 1 + except Exception as exc: # a torn read would land here + errors.append(exc) + + t = threading.Thread(target=reader, daemon=True) + t.start() + try: + for i in range(300): + w.apply(ev(CONNECT, 100 + (i % 7), 6000 + i)) # inserts... + w.apply(ev(CLOSE, 100 + (i % 7), 6000 + i)) # ...and deletes + w.reconcile({5000: 100, 80: 1}) # whole-map republish + finally: + stop.set() + t.join(timeout=5) + + check("the lock-free reader never raised", not errors, f"({errors[:3]})") + check("it actually kept reading", reads[0] > 0) + check("and every value it saw was the real pid, never junk", values == {100}, + f"({values})") diff --git a/tests/test_socketwatch_wiring.py b/tests/test_socketwatch_wiring.py index 5fbb633..1f18918 100644 --- a/tests/test_socketwatch_wiring.py +++ b/tests/test_socketwatch_wiring.py @@ -14,7 +14,7 @@ from beantester.engine import BeanEngine from beantester.socketwatch import BIND, CONNECT, SocketEvent -from fakes import FakeDivert, check +from fakes import FakeDivert, FakePacket, check def _wait(pred, timeout=5.0): @@ -129,6 +129,74 @@ def boom(): eng.stop() +class _GatedDivert: + """Hands over its packet only once the test opens the gate. + + Without the gate the assertion would be a race against the watcher thread: the + packet could reach the capture loop before the CONNECT event was applied, and the + test would be measuring thread scheduling instead of whether the connection log + consults the live map at all. + """ + + def __init__(self, packets): + self.inbox = list(packets) + self.gate = threading.Event() + self.sent = [] + self.closed = False + + def open(self): + pass + + def recv(self): + self.gate.wait() + if self.inbox: + return self.inbox.pop(0) + while not self.closed: + time.sleep(0.003) + raise OSError("closed") + + def send(self, packet): + self.sent.append(packet) + + def close(self): + self.closed = True + self.gate.set() # release a parked recv() so stop() can join + + +def _row_for(engine, local_port): + for row in engine.connections_snapshot(limit=None): + if row["local_port"] == local_port: + return row + return None + + +def test_a_fresh_socket_stamps_the_connection_row_from_the_live_map(): + """Why chunk 2 reaches the connection log at all. + + The poller knows NOTHING about this port - that is the short-lived-connection + case, where a flow opens and finishes inside one refresh interval and the row used + to end up with no owner at all. The watcher was told the owner by its CONNECT + event, so the row gets both a pid and a process name. + """ + ports = _FakePorts({}) # the poller never sees this socket + eng = BeanEngine() + eng._ports = ports + divert = _GatedDivert([FakePacket(port=5000)]) + eng.start("true", divert=divert, socket_source=_FakeSocketSource([ev(CONNECT, 100, 5000)])) + try: + check("the poller cannot answer for this port", ports.pid_for(5000) is None) + check("but the watcher was told the owner", + _wait(lambda: eng._socketwatch.pid_for(5000) == 100)) + divert.gate.set() # only now let the packet through + check("the row was stamped from the live map", + _wait(lambda: (_row_for(eng, 5000) or {}).get("pid") == 100)) + row = _row_for(eng, 5000) + check("...and named through the delegated name cache", + row["proc"] == "chrome.exe", f"({row})") + finally: + eng.stop() + + def test_stopping_the_engine_leaves_no_watcher_thread_behind(): before = {t.name for t in threading.enumerate()} eng = BeanEngine()