Compare commits
1 Commits
4acf50d9ef
...
telemetry-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e3b69fa759 |
@@ -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."""
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user