Files
windy-git/scripts/telemetry_emit.py
Grant Whitmer 00ec963f82
Some checks failed
check / gate (push) Has been cancelled
telemetry: fix ci.run completeness — cursor on (finish time, job id)
Telemetry Boss found jobs_finished=43 vs 8 ci.run rows. Root cause: the
high-water mark was the job id, but jobs FINISH out of id order, so every
long job that started before the mark and finished after it was silently
never emitted. Now a (finish time, id) cursor; finish = stopped, or
updated for skipped jobs with no stop time. Heartbeat finished/failed/
cancelled counts are derived from exactly the rows emitted, so
sum(jobs_finished) == count(ci.run) by construction (dry run on real
data: 97 == 97, failed 2 == 2, cancelled 13 == 13). posted_to_github now
set from the bridge's own rules. duration_ms = Gitea whole seconds x 1000.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:42:36 -04:00

257 lines
9.8 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
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
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()[:300]
if r.status >= 300:
raise urllib.error.HTTPError(
INGEST, r.status, body.decode(errors="replace"), None, None
)
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}, 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())