diff --git a/api/app/telemetry.py b/api/app/telemetry.py index 6d7f2d5..9e8ef9a 100644 --- a/api/app/telemetry.py +++ b/api/app/telemetry.py @@ -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.""" diff --git a/api/tests/test_telemetry.py b/api/tests/test_telemetry.py index 2a85d58..cf540fb 100644 --- a/api/tests/test_telemetry.py +++ b/api/tests/test_telemetry.py @@ -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 diff --git a/scripts/telemetry_emit.py b/scripts/telemetry_emit.py index 512c5d3..c1af4b1 100644 --- a/scripts/telemetry_emit.py +++ b/scripts/telemetry_emit.py @@ -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