diff --git a/gateway/status.py b/gateway/status.py index 66a56517d2..d391e27863 100644 --- a/gateway/status.py +++ b/gateway/status.py @@ -1110,6 +1110,17 @@ def write_runtime_status( # continuous retry episode; None clears it on reconnect. platform_payload["retrying_since"] = retrying_since platform_payload["updated_at"] = _utc_now_iso() + # Writer identity: which PROCESS wrote this entry. The top-level + # pid/start_time are refreshed on every write, so they only identify + # the file's most recent writer — per-entry provenance is what lets + # a reader (the /api/status cross-profile aggregation) distinguish + # "written by the current live process" from "preserved from a prior + # process" with exact (pid, start_time) equality instead of clock + # heuristics. start_time is the same PID-reuse fingerprint the + # liveness checks use, so a recycled PID never masquerades as the + # original writer. + platform_payload["writer_pid"] = current_record["pid"] + platform_payload["writer_start_time"] = current_record["start_time"] payload["platforms"][platform] = platform_payload _write_json_file(path, payload) diff --git a/hermes_cli/web_server.py b/hermes_cli/web_server.py index 5f59ea8e20..91394d6416 100644 --- a/hermes_cli/web_server.py +++ b/hermes_cli/web_server.py @@ -3143,42 +3143,39 @@ def _profile_platform_ports(profile_home: Path, runtime: Optional[dict]) -> Dict return ports -# Small grace for comparing a platform entry's wall-clock ``updated_at`` -# against the gateway process's psutil-derived create time (both come from -# the host clock, but coarse rounding can skew them by a moment). -_PROFILE_PLATFORM_FRESHNESS_SLACK_S = 2.0 - - -def _profile_gateway_started_at( +def _profile_gateway_writer_identity( profile_home: Path, runtime: Optional[dict] -) -> Optional[float]: - """Wall-clock create time of the profile's LIVE gateway process, or None. +) -> Optional[tuple]: + """``(pid, start_time)`` identity of the profile's LIVE gateway, or None. - The record's own ``start_time`` field is a PID-reuse fingerprint, not a - timestamp (clock ticks since boot on Linux), so it cannot anchor a - freshness comparison. Instead reuse the validated-liveness helper — - which checks the recorded PID against the live process table, the - fingerprint, and the profile's home — and ask psutil for that process's - epoch create time. None when the record doesn't belong to a live - gateway (then nothing in it is current by definition). + Reuses the validated-liveness helper — recorded PID checked against the + live process table, the start-time PID-reuse fingerprint, and the + profile's home — then reads the live process's fingerprint via the same + ``_get_process_start_time`` that stamped it, so equality is exact (no + unit or clock-source mismatch). None when the record doesn't belong to + a live gateway; nothing in it is current by definition then. """ try: - from gateway.status import get_runtime_status_running_pid + from gateway.status import ( + _get_process_start_time, + get_runtime_status_running_pid, + ) pid = get_runtime_status_running_pid(runtime, expected_home=profile_home) if pid is None: return None - import psutil # type: ignore - - return float(psutil.Process(pid).create_time()) + start_time = _get_process_start_time(pid) + if start_time is None: + return None + return (pid, start_time) except Exception: return None -def _fresh_profile_platforms( - started_at: Optional[float], platforms: dict +def _owned_profile_platforms( + writer_identity: Optional[tuple], platforms: dict ) -> dict: - """Keep only platform entries written by the profile's CURRENT process. + """Keep only platform entries the profile's CURRENT process wrote. Gateway startup deliberately preserves plain platform entries in ``gateway_state.json`` across restarts (the dashboard keeps showing @@ -3186,34 +3183,32 @@ def _fresh_profile_platforms( endpoint compensates by filtering them against the current configuration. The cross-profile aggregation has no equivalent config context (a profile's platform set depends on tokens in that profile's - ``.env`` behind its secret scope), so it uses freshness instead: an - entry is aggregatable only when its ``updated_at`` is at/after the live - gateway process's create time — i.e. the current process wrote it. A - fatal entry left behind by a platform the operator has since - disabled/removed therefore stops degrading fleet health as soon as that - profile's gateway restarts (a config change requires that restart to - take effect anyway). Fail closed: entries without a parseable - ``updated_at``, or records with no live process, are excluded — - aggregation is a supplement, and a false "degraded forever" is the - worse failure mode. + ``.env`` behind its secret scope), so it demands strict process + ownership instead: ``write_runtime_status`` stamps every platform write + with the writer's ``(pid, start_time)`` identity, and an entry is + aggregatable only when that identity equals the profile's live gateway + process — exact match, no clock heuristics, so an entry written moments + before a fast restart can never masquerade as current. A fatal entry + left behind by a platform the operator has since disabled/removed thus + stops degrading fleet health as soon as that profile's gateway restarts + (a config change requires that restart to take effect anyway). Fail + closed: entries without a writer identity (legacy records) or records + with no live process are excluded — aggregation is a supplement, and a + false "degraded forever" is the worse failure mode. """ - if started_at is None: + if writer_identity is None: return {} - cutoff = float(started_at) - _PROFILE_PLATFORM_FRESHNESS_SLACK_S - fresh: Dict[str, dict] = {} + live_pid, live_start = writer_identity + owned: Dict[str, dict] = {} for key, value in platforms.items(): if not isinstance(value, dict): continue - updated_at = value.get("updated_at") - if not isinstance(updated_at, str): - continue - try: - written = datetime.fromisoformat(updated_at).timestamp() - except (ValueError, TypeError): - continue - if written >= cutoff: - fresh[key] = value - return fresh + if ( + value.get("writer_pid") == live_pid + and value.get("writer_start_time") == live_start + ): + owned[key] = value + return owned def _collect_profile_gateway_topology() -> Dict[str, Any]: @@ -3232,9 +3227,9 @@ def _collect_profile_gateway_topology() -> Dict[str, Any]: live gateway, ``"multiple"`` for independent per-profile gateways, ``"none"`` when nothing is running. * ``profile_platforms`` — ``{profile: platforms}`` runtime platform maps - for each LIVE gateway, freshness-filtered to entries written by that + for each LIVE gateway, ownership-filtered to entries stamped by that profile's current process (stale preserved entries for since-removed - platforms are excluded — see ``_fresh_profile_platforms``). Internal + platforms are excluded — see ``_owned_profile_platforms``). Internal aggregation input for ``/api/status`` (independent per-profile gateways write failures to their own ``gateway_state.json``, which the unparameterized endpoint would otherwise never see). Never exposed @@ -3272,16 +3267,17 @@ def _collect_profile_gateway_topology() -> Dict[str, Any]: multiplex = True plats = (runtime or {}).get("platforms") if isinstance(plats, dict) and plats: - # Freshness filter: gateway startup preserves plain platform + # Ownership filter: gateway startup preserves plain platform # entries across restarts, so the raw map can carry fatal state # for platforms the operator has since disabled/removed. Only - # entries written by the profile's current live process are - # aggregation candidates (see _fresh_profile_platforms). - fresh = _fresh_profile_platforms( - _profile_gateway_started_at(home, runtime), plats + # entries stamped with the profile's current live process's + # writer identity are aggregation candidates (see + # _owned_profile_platforms). + owned = _owned_profile_platforms( + _profile_gateway_writer_identity(home, runtime), plats ) - if fresh: - profile_platforms[name] = fresh + if owned: + profile_platforms[name] = owned entry: Dict[str, Any] = { "profile": name, "ports": _profile_platform_ports(home, runtime), @@ -3424,6 +3420,20 @@ def _status_platform_key_allowed( return configured is None or key in configured +# Per-entry writer-identity stamps (added by gateway.status.write_runtime_status +# for the aggregation ownership check) are process recon — the same class of +# detail as the auth-gated top-level ``gateway_pid`` — and must not project +# onto the public endpoint. +_PRIVATE_PLATFORM_ENTRY_KEYS = frozenset({"writer_pid", "writer_start_time"}) + + +def _public_platform_entry(value: Any) -> Any: + """Strip writer-identity stamps from a platform entry before projection.""" + if not isinstance(value, dict): + return value + return {k: v for k, v in value.items() if k not in _PRIVATE_PLATFORM_ENTRY_KEYS} + + def _merge_profile_gateway_platforms( gateway_platforms: dict, profile_platforms: dict ) -> dict: @@ -3454,7 +3464,7 @@ def _merge_profile_gateway_platforms( namespaced = f"{prof}:{key}" if not _is_profile_platform_status_key(namespaced): continue - merged.setdefault(namespaced, value) + merged.setdefault(namespaced, _public_platform_entry(value)) return merged @@ -3564,7 +3574,7 @@ async def get_status(profile: Optional[str] = None): # load must not fail open into projecting arbitrary keys from a # process-local JSON file onto this public endpoint. gateway_platforms = { - key: value + key: _public_platform_entry(value) for key, value in gateway_platforms.items() if _status_platform_key_allowed(key, configured_gateway_platforms) } diff --git a/tests/gateway/test_status.py b/tests/gateway/test_status.py index bfca704efc..b2e926ad02 100644 --- a/tests/gateway/test_status.py +++ b/tests/gateway/test_status.py @@ -211,15 +211,30 @@ class TestGatewayRuntimeStatus: clear_profile_platforms=True, ) - assert status.read_runtime_status()["platforms"] == { - "telegram": {"state": "connected"}, - "reviewer:slack": { - "state": "connected", - "updated_at": status.read_runtime_status()["platforms"][ - "reviewer:slack" - ]["updated_at"], - }, - } + platforms = status.read_runtime_status()["platforms"] + assert set(platforms) == {"telegram", "reviewer:slack"} + assert platforms["telegram"] == {"state": "connected"} + assert platforms["reviewer:slack"]["state"] == "connected" + + def test_platform_writes_are_stamped_with_writer_identity( + self, tmp_path, monkeypatch + ): + # The /api/status cross-profile aggregation must distinguish entries + # written by the CURRENT process from entries preserved across a + # restart (a wall-clock freshness window admits stale failures + # written moments before a fast restart). Every platform write is + # therefore stamped with the writer's (pid, start_time) identity — + # the same PID-reuse fingerprint the liveness checks use — so + # ownership is exact equality, not clock heuristics. + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + + status.write_runtime_status(platform="telegram", platform_state="connected") + + entry = status.read_runtime_status()["platforms"]["telegram"] + assert entry["writer_pid"] == os.getpid() + assert entry["writer_start_time"] == status._get_process_start_time( + os.getpid() + ) def test_clear_profile_platforms_repairs_malformed_platforms( self, tmp_path, monkeypatch diff --git a/tests/hermes_cli/test_web_server_gateway_topology.py b/tests/hermes_cli/test_web_server_gateway_topology.py index c14bcd50f8..4cfee0d76d 100644 --- a/tests/hermes_cli/test_web_server_gateway_topology.py +++ b/tests/hermes_cli/test_web_server_gateway_topology.py @@ -83,14 +83,15 @@ class TestCollectProfileGatewayTopology: # Independent per-profile gateways (gateway_mode == "multiple") each # write their own gateway_state.json; the collector surfaces every # LIVE profile's raw platform map for the /api/status merge (OOF-3). - # Entries must be fresh — written by the profile's current process. + # Entries must carry the current live process's writer identity. homes = [("default", tmp_path / "d"), ("coder", tmp_path / "c")] runtimes = { "default": { "platforms": { "telegram": { "state": "connected", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 100, + "writer_start_time": 111, } } }, @@ -99,7 +100,8 @@ class TestCollectProfileGatewayTopology: "discord": { "state": "fatal", "error_code": "duplicate_credential", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 200, + "writer_start_time": 222, } } }, @@ -107,12 +109,11 @@ class TestCollectProfileGatewayTopology: _patch_topology( monkeypatch, homes, running={"default", "coder"}, runtimes=runtimes ) - # Both gateways' live processes started before the entries were - # written (epoch for 2026-08-05T10:00:00Z). + identities = {tmp_path / "d": (100, 111), tmp_path / "c": (200, 222)} monkeypatch.setattr( web_server, - "_profile_gateway_started_at", - lambda home, runtime: 1785924000.0, + "_profile_gateway_writer_identity", + lambda home, runtime: identities.get(home), ) topo = _collect_profile_gateway_topology() assert topo["gateway_mode"] == "multiple" @@ -120,14 +121,16 @@ class TestCollectProfileGatewayTopology: "default": { "telegram": { "state": "connected", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 100, + "writer_start_time": 111, } }, "coder": { "discord": { "state": "fatal", "error_code": "duplicate_credential", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 200, + "writer_start_time": 222, } }, } @@ -135,28 +138,40 @@ class TestCollectProfileGatewayTopology: def test_stale_platform_entries_are_not_aggregated(self, tmp_path, monkeypatch): # Gateway startup preserves plain platform entries across restarts; # if the operator removed/disabled the platform and restarted, its - # old fatal entry must not keep degrading fleet health. Entries - # written before the live process started — or with no parseable - # updated_at — are excluded. + # old fatal entry must not keep degrading fleet health. Ownership + # is exact (pid, start_time) equality with the live process — an + # entry written by the PREVIOUS process moments before a fast + # restart, a legacy entry with no writer identity, and a recycled + # PID with a different start-time fingerprint are all excluded. homes = [("default", tmp_path / "d"), ("coder", tmp_path / "c")] runtimes = { "coder": { "platforms": { - # Written by the PRIOR process (before restart). + # Near-boundary: written by the PRIOR process (pid 199) + # immediately before a fast restart. A wall-clock + # freshness window would admit this; identity must not. "telegram": { "state": "fatal", "error_code": "duplicate_credential", - "updated_at": "2026-08-05T09:00:00+00:00", + "updated_at": "2026-08-05T10:00:00.900000+00:00", + "writer_pid": 199, + "writer_start_time": 110, }, - # No timestamp at all — fail closed. + # Legacy entry, no writer identity — fail closed. "discord": {"state": "fatal"}, - # Garbage timestamp — fail closed. - "slack": {"state": "fatal", "updated_at": "not-a-time"}, + # Same PID recycled, different start-time fingerprint — + # a different process, not the live writer. + "slack": { + "state": "fatal", + "writer_pid": 200, + "writer_start_time": 110, + }, # Written by the current process — kept. "signal": { "state": "fatal", "error_code": "duplicate_credential", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 200, + "writer_start_time": 222, }, } }, @@ -166,8 +181,8 @@ class TestCollectProfileGatewayTopology: ) monkeypatch.setattr( web_server, - "_profile_gateway_started_at", - lambda home, runtime: 1785924000.0, # 2026-08-05T10:00:00Z + "_profile_gateway_writer_identity", + lambda home, runtime: (200, 222), ) topo = _collect_profile_gateway_topology() assert set(topo["profile_platforms"].get("coder", {})) == {"signal"} @@ -181,14 +196,17 @@ class TestCollectProfileGatewayTopology: "platforms": { "telegram": { "state": "fatal", - "updated_at": "2026-08-05T10:00:05+00:00", + "writer_pid": 200, + "writer_start_time": 222, } } }, } _patch_topology(monkeypatch, homes, running={"coder"}, runtimes=runtimes) monkeypatch.setattr( - web_server, "_profile_gateway_started_at", lambda home, runtime: None + web_server, + "_profile_gateway_writer_identity", + lambda home, runtime: None, ) topo = _collect_profile_gateway_topology() assert topo["profile_platforms"] == {} @@ -384,7 +402,13 @@ class TestStatusEndpointTopology: "read_runtime_status", lambda path=None: { "gateway_state": "running", - "platforms": {"telegram": {"state": "connected"}}, + "platforms": { + "telegram": { + "state": "connected", + "writer_pid": 123, + "writer_start_time": 456, + } + }, }, ) monkeypatch.setattr( @@ -410,6 +434,8 @@ class TestStatusEndpointTopology: "telegram": { "state": "fatal", "error_code": "duplicate_credential", + "writer_pid": 789, + "writer_start_time": 1011, }, "bad key": {"state": "fatal"}, }, @@ -432,6 +458,12 @@ class TestStatusEndpointTopology: ) # Keys that don't survive the grammar are dropped, not projected. assert not any("bad key" in key for key in platforms) + # Writer-identity stamps (process recon, same class as the auth-gated + # gateway_pid) never project onto the endpoint — neither on the active + # profile's own entries nor on merged cross-profile entries. + for entry in platforms.values(): + assert "writer_pid" not in entry + assert "writer_start_time" not in entry # The platforms component rollup counts the merged failure. assert data["components"]["platforms"]["status"] == "degraded"