1 Commits

Author SHA1 Message Date
Kit OC5
e3b69fa759 telemetry: UPDATE 7 — read the ingest body; count quarantined + dropped on heartbeats
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
The ledger answers 202 even when it quarantines rows. Both emitters now log
a warning with the reasons and report service.health.telemetry_quarantined
and telemetry_dropped (API: buffer overflow; sync: 0 by construction, since
a failed send keeps cursor + spool). HOLD until Telemetry Boss declares both
keys on windy-git's two service.health shapes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:32:28 -04:00
3 changed files with 85 additions and 6 deletions

View File

@@ -120,6 +120,9 @@ class Telemetry:
self.window_start = time.time()
self.requests = self.errors_5xx = self.errors_4xx = self.refusals_4xx = 0
self.latencies_ms: list[float] = []
# UPDATE 7: rows the ledger quarantined (it still answers 202) and rows
# this process lost (buffer overflow). Non-zero = a bug in this emitter.
self.quarantined = self.dropped = 0
# ---- recording (never raises into a request) --------------------------
def record_request(self, status: int, duration_ms: float, *, refused: bool = False) -> None:
@@ -147,6 +150,7 @@ class Telemetry:
}
)
if len(self.buffer) > MAX_BUFFER:
self.dropped += len(self.buffer) - MAX_BUFFER
del self.buffer[: len(self.buffer) - MAX_BUFFER]
def boot(self) -> None:
@@ -188,6 +192,8 @@ class Telemetry:
"errors_5xx": self.errors_5xx,
"errors_4xx": self.errors_4xx,
"refusals_4xx": self.refusals_4xx,
"telemetry_quarantined": self.quarantined,
"telemetry_dropped": self.dropped,
}
if self.latencies_ms: # no traffic = no p95, not a fake 0
s = sorted(self.latencies_ms)
@@ -199,7 +205,7 @@ class Telemetry:
self._reset_window()
# ---- sending ------------------------------------------------------------
def _post(self, batch: list[dict]) -> int:
def _post(self, batch: list[dict]) -> tuple[int, dict]:
req = urllib.request.Request(
self.url,
data=json.dumps({"events": batch}).encode(),
@@ -211,19 +217,32 @@ class Telemetry:
},
)
with urllib.request.urlopen(req, timeout=20) as r:
return r.status
try:
body = json.loads(r.read() or b"{}")
except ValueError:
body = {}
return r.status, body if isinstance(body, dict) else {}
async def flush(self) -> None:
if not self.enabled or not self.buffer:
return
batch = self.buffer[:500]
try:
status = await asyncio.to_thread(self._post, batch)
status, body = await asyncio.to_thread(self._post, batch)
except Exception as exc: # noqa: BLE001 - telemetry must never take the API down
log.warning("telemetry flush failed (%d rows kept): %s", len(self.buffer), exc)
return
if 200 <= status < 300:
del self.buffer[: len(batch)]
self.note_quarantine(body)
def note_quarantine(self, body: dict) -> None:
# 202 does NOT mean every row landed: refused rows are dead-lettered.
q = body.get("quarantined")
if isinstance(q, int) and q > 0:
self.quarantined += q
log.warning("telemetry: %d row(s) QUARANTINED by the ledger: %s", q,
"; ".join(map(str, body.get("rejections") or [])) or "no reason given")
async def run(self) -> None:
"""The one in-process timer: flush every minute, heartbeat every hour."""

View File

@@ -160,3 +160,43 @@ def test_synthetic_is_forwarded_downstream_only_for_synthetic_requests():
finally:
tmod.SYNTHETIC.reset(token)
assert tmod.synthetic_headers() == {}
# ---- UPDATE 7: the ledger answers 202 even when it quarantines rows ----------
@pytest.mark.asyncio
async def test_quarantined_rows_are_warned_and_counted_on_the_next_heartbeat(monkeypatch, caplog):
tel = _tel()
tel.boot()
monkeypatch.setattr(
tel, "_post", lambda b: (202, {"accepted": 0, "quarantined": 1, "rejections": ["undeclared key"]})
)
with caplog.at_level("WARNING", logger="windy-git.telemetry"):
await tel.flush()
assert tel.buffer == [] # sent; the ledger dead-lettered it, retrying won't help
assert "QUARANTINED" in caplog.text and "undeclared key" in caplog.text
tel.health()
assert tel.buffer[-1]["metadata"]["telemetry_quarantined"] == 1
assert tel.health_row()["telemetry_quarantined"] == 0 # reset per heartbeat window
@pytest.mark.asyncio
async def test_clean_send_reports_zero_and_logs_nothing(monkeypatch, caplog):
tel = _tel()
tel.boot()
monkeypatch.setattr(tel, "_post", lambda b: (202, {"accepted": 1, "quarantined": 0, "rejections": []}))
with caplog.at_level("WARNING", logger="windy-git.telemetry"):
await tel.flush()
assert caplog.text == ""
row = tel.health_row()
assert row["telemetry_quarantined"] == 0 and row["telemetry_dropped"] == 0
def test_buffer_overflow_is_counted_as_dropped(monkeypatch):
monkeypatch.setattr(tmod, "MAX_BUFFER", 3)
tel = _tel()
for _ in range(5):
tel.boot()
assert len(tel.buffer) == 3
assert tel.health_row()["telemetry_dropped"] == 2

View File

@@ -161,6 +161,12 @@ def main() -> int:
h["jobs_cancelled"] = sum(1 for j in jobs if j["status"] == 3)
meta = {k: int(v) for k, v in h.items()}
meta["interval_s"] = int(now - since) # ecosystem-standard key
# UPDATE 7. Quarantines seen on earlier sends (the ledger answers 202 anyway)
# are carried in the state file until a heartbeat reports them. Dropped is 0
# by construction: a failed send keeps the cursor and the spool, so every row
# is re-sent next run (a partial failure can duplicate, never lose).
meta["telemetry_quarantined"] = int(state.get("quarantined_unreported", 0))
meta["telemetry_dropped"] = 0
for k in ("repos_synced", "repos_sync_failed", "statuses_posted", "bridge_errors"):
v = os.environ.get(f"TELEMETRY_{k.upper()}")
if v is not None and v.isdigit(): # absent = couldn't count; never invent 0
@@ -216,6 +222,7 @@ def main() -> int:
print("[telemetry] WINDYGIT_TELEMETRY_TOKEN unset — not sending (not a failure)")
return 0
quarantined = 0
for i in range(0, len(events), 500):
req = urllib.request.Request(
INGEST,
@@ -229,11 +236,20 @@ def main() -> int:
)
try:
with urllib.request.urlopen(req, timeout=30) as r:
body = r.read()[:300]
body = r.read()
if r.status >= 300:
raise urllib.error.HTTPError(
INGEST, r.status, body.decode(errors="replace"), None, None
INGEST, r.status, body[:300].decode(errors="replace"), None, None
)
try:
resp = json.loads(body or b"{}")
except ValueError:
resp = {}
q = resp.get("quarantined") if isinstance(resp, dict) else None
if isinstance(q, int) and q > 0:
quarantined += q
reasons = "; ".join(map(str, resp.get("rejections") or [])) or "no reason given"
print(f"[telemetry] WARNING {q} row(s) QUARANTINED by the ledger: {reasons}")
except urllib.error.HTTPError as e:
print(f"[telemetry] FAILED ingest HTTP {e.code}: {e.read()[:200]!r}")
return 1 # state NOT advanced: the same rows retry next run
@@ -246,7 +262,11 @@ def main() -> int:
os.makedirs(os.path.dirname(STATE), exist_ok=True)
new_fin, new_id = (jobs[-1]["fin"], jobs[-1]["id"]) if jobs else (last_fin, last_id)
with open(STATE + ".tmp", "w") as f:
json.dump({"last_fin": new_fin, "last_id": new_id, "last_ts": now}, f)
json.dump(
{"last_fin": new_fin, "last_id": new_id, "last_ts": now,
"quarantined_unreported": quarantined},
f,
)
os.replace(STATE + ".tmp", STATE)
print(f"[telemetry] sent {len(events)} events ({len(jobs)} ci.run)")
return 0