fix(status): strict writer-identity ownership for aggregated platform entries (OOF-3)
The freshness window (updated_at >= live process create_time - 2s) had a P1 boundary hole: a stale failure written by the PREVIOUS process immediately before a fast restart landed inside the slack and was aggregated; if that platform was then removed, the new process never replaces the entry and NAS stays degraded indefinitely. Replace clock heuristics with persisted writer identity: - write_runtime_status now stamps every platform entry with the writing process's (writer_pid, writer_start_time) — the same PID-reuse fingerprint the liveness checks use, so a recycled PID never masquerades as the original writer. - The aggregation ownership filter requires exact equality between an entry's stamp and the profile's validated live gateway process (get_runtime_status_running_pid + _get_process_start_time). No slack, no timestamps. Legacy entries without a stamp fail closed. - Writer stamps are process recon (same class as the auth-gated gateway_pid) and are stripped from all /api/status projections, both active-profile and merged cross-profile entries. Near-boundary regression test: prior-process entry stamped 100ms before restart is excluded; recycled-pid-different-fingerprint excluded; legacy no-stamp excluded; current-process entry kept.
This commit is contained in:
@@ -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)
|
||||
|
||||
+67
-57
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user