Files
windy-git/scripts/telemetry_emit.py
Kit OC5 e3b69fa759
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
telemetry: UPDATE 7 — read the ingest body; count quarantined + dropped on heartbeats
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

277 lines
11 KiB
Python

#!/usr/bin/env python3
"""Emit Windy Git CI telemetry to admin.windyword.ai (Windy Telemetry 40's ledger).
Runs on Veron after every sync (root; reads the gitea DB via `docker exec`).
Shapes are declared with Telemetry 40 (2026-09-23) — do not add keys or enum
values without re-declaring: a declared family quarantines any row that
doesn't match.
ci.run one row per FINISHED job, exactly once — cursor on
(finish time, job id) in STATE; jobs finish out of id order
service.health one row per invocation: CI plane counts for the interval
Privacy: ids, names of repos/jobs, codes, counts, durations. No commit
messages, no logs, no author names.
--dry-run print the batch instead of posting (and don't advance STATE)
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
import time
import urllib.error
import urllib.request
from datetime import UTC, datetime
INGEST = os.environ.get("TELEMETRY_INGEST_URL", "https://admin.windyword.ai/v1/events")
TOKEN = os.environ.get("WINDYGIT_TELEMETRY_TOKEN", "")
STATE = os.environ.get("TELEMETRY_STATE", "/var/lib/windy-git/telemetry-state.json")
PLATFORM, SERVICE = "windy-git", "ci"
OUTCOME = {1: "success", 2: "failure", 3: "cancelled", 4: "skipped"}
EVENTS = {"push", "pull_request", "pull_request_sync", "schedule", "workflow_dispatch"}
RUNNERS_EXPECTED = 6
def sql(query: str) -> list[dict]:
"""Rows as dicts, via psql's json_agg — no driver needed on the host."""
wrapped = f"select coalesce(json_agg(t), '[]'::json) from ({query}) t;"
out = subprocess.run(
[
"docker",
"exec",
"-i",
"windy-git-db-1",
"sh",
"-c",
'psql -U "$POSTGRES_USER" -d gitea -At -v ON_ERROR_STOP=1',
],
input=wrapped,
capture_output=True,
text=True,
check=True,
).stdout.strip()
return json.loads(out or "[]")
def load_state() -> dict:
try:
with open(STATE) as f:
return json.load(f)
except (OSError, ValueError):
return {}
def iso(epoch: float) -> str:
return datetime.fromtimestamp(epoch, UTC).isoformat().replace("+00:00", "Z")
def main() -> int:
dry = "--dry-run" in sys.argv
state = load_state()
now = time.time()
since = float(state.get("last_ts", now - 300))
# Cursor = (finish time, job id), NOT job id alone: jobs finish out of id
# order, so an id high-water mark silently drops every long job that started
# before the mark and finished after it (Telemetry Boss caught this: 43
# finished vs 8 ci.run rows). Finish time = stopped, or updated for jobs
# Gitea/the janitor skipped without a stop time.
if "last_fin" in state:
last_fin, last_id = int(state["last_fin"]), int(state["last_id"])
else: # first run or pre-cursor state: start now, never replay history
last_fin, last_id = int(state.get("last_ts", now)), 0
cutoff = int(now) - 5 # leave the current second alone; late writers land next run
FIN = "coalesce(nullif(j.stopped, 0), j.updated)"
jobs = sql(f"""
select j.id, j.name as job, j.status, j.started, j.stopped, {FIN} as fin,
p.lower_name as repo, p.default_branch, r.workflow_id, r.event,
r.ref, r.index as run, left(r.commit_sha, 7) as sha
from action_run_job j
join action_run r on r.id = j.run_id
join repository p on p.id = r.repo_id
where j.status in (1, 2, 3, 4)
and ({FIN}, j.id) > ({last_fin}, {last_id})
and {FIN} <= {cutoff}
order by {FIN}, j.id
limit 2000""")
try: # posted_to_github: the bridge's own rules, from the same checkout
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import pr_status_bridge as bridge
except Exception: # noqa: BLE001
bridge = None
events = []
for j in jobs:
ref = j["ref"] or ""
if ref.startswith("refs/pull/"):
kind = "pr"
elif ref == f"refs/heads/{j['default_branch']}":
kind = "default"
else:
kind = "other"
# Gitea stores whole seconds; duration_ms is seconds*1000 (so 10000 = 10 s).
dur = (j["stopped"] - j["started"]) * 1000 if j["started"] and j["stopped"] else None
ev = {
"ts": iso(j["stopped"]),
"platform": PLATFORM,
"service": SERVICE,
"event_type": "ci.run",
"actor_type": "system",
"metadata": {
"repo": j["repo"],
"workflow": (j["workflow_id"] or "").removesuffix(".yml").removesuffix(".yaml"),
"job": j["job"],
"outcome": OUTCOME[j["status"]],
"event": j["event"] if j["event"] in EVENTS else "other",
"branch_kind": kind,
"run": j["run"],
"sha": j["sha"],
},
}
if dur is not None and dur >= 0:
ev["duration_ms"] = int(dur)
if bridge is not None:
wf = ev["metadata"]["workflow"]
ev["metadata"]["posted_to_github"] = bool(
j["repo"] in {r.lower() for r in bridge.REPOS}
and kind in ("default", "pr")
and j["status"] != 4
and not bridge.NO_DAEMON_JOB.search(j["job"])
and f"{wf}/{j['job']}" not in bridge.NON_BLOCKING.get(j["repo"], set())
)
events.append(ev)
# --- heartbeat: counts since the previous invocation --------------------
# Interval counts come from EXACTLY the rows emitted above, so
# sum(jobs_finished) over any window == count(ci.run) in it, by construction.
h = sql(f"""
select
(select count(*) from action_run_job where status in (5, 7)) as jobs_waiting,
(select count(*) from action_run_job where status = 6) as jobs_running,
(select count(*) from action_runner where deleted is null and last_online >= {int(now) - 120}) as runners_online,
(select coalesce(extract(epoch from now())::bigint - min(created), 0)
from action_run_job where status in (5, 7)) as oldest_waiting_s""")[0]
h["jobs_finished"] = len(jobs)
h["jobs_failed"] = sum(1 for j in jobs if j["status"] == 2)
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
meta[k] = int(v)
events.append(
{
"ts": iso(now),
"platform": PLATFORM,
"service": SERVICE,
"event_type": "service.health",
"actor_type": "system",
"metadata": meta,
}
)
# ci.job_cancelled: spooled by the janitor (cancel_unrunnable.sh), one JSON per job.
spool = os.environ.get("JANITOR_SPOOL", "/var/lib/windy-git/janitor-cancelled.jsonl")
spooled = 0
try:
with open(spool) as f:
for line in f:
try:
m = json.loads(line)
except ValueError:
continue
events.append(
{
"ts": iso(now),
"platform": PLATFORM,
"service": SERVICE,
"event_type": "ci.job_cancelled",
"actor_type": "system",
"metadata": {
k: m[k]
for k in ("repo", "workflow", "job", "reason", "runs_on", "waited_s")
},
}
)
spooled += 1
except OSError:
pass
if dry:
out = os.environ.get("TELEMETRY_DRY_OUT")
if out:
with open(out, "w") as f:
json.dump({"events": events}, f)
else:
print(json.dumps({"events": events}, indent=1)[:4000])
print(f"[telemetry] DRY RUN: {len(events)} events ({len(jobs)} ci.run)")
return 0
if not TOKEN:
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,
data=json.dumps({"events": events[i : i + 500]}).encode(),
method="POST",
headers={
"Authorization": f"Bearer {TOKEN}",
"Content-Type": "application/json",
"User-Agent": "windy-git-telemetry/1",
},
)
try:
with urllib.request.urlopen(req, timeout=30) as r:
body = r.read()
if r.status >= 300:
raise urllib.error.HTTPError(
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
except urllib.error.URLError as e:
print(f"[telemetry] FAILED ingest: {e.reason}")
return 1
if spooled:
open(spool, "w").close() # only after every batch was accepted
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,
"quarantined_unreported": quarantined},
f,
)
os.replace(STATE + ".tmp", STATE)
print(f"[telemetry] sent {len(events)} events ({len(jobs)} ci.run)")
return 0
if __name__ == "__main__":
sys.exit(main())