From be0dddaac6c2b6438f8f1fb78153466fb0ba4c5d Mon Sep 17 00:00:00 2001 From: Kit OC5 Date: Wed, 23 Sep 2026 12:51:37 -0400 Subject: [PATCH] telemetry: detect push velocity from Gitea's action table (alert only) git push never touches our API, so throttle.py can't see it. Gitea's action table records every push; the 5-min sync-side emitter now reads it and emits forge.push_velocity when an account crosses 60 pushes/1h, 500 pushes/24h (standard-band base) or 10 ref deletes/24h. One row per account per rule per window while over; windyadmin (the sync) exempt. Nothing sits in the push path and nothing is refused. HOLD until Telemetry Boss declares the shape. Co-Authored-By: Claude Opus 5.5 --- api/tests/test_push_velocity.py | 77 ++++++++++++++++++++++++++ scripts/telemetry_emit.py | 97 ++++++++++++++++++++++++++++++++- 2 files changed, 173 insertions(+), 1 deletion(-) create mode 100644 api/tests/test_push_velocity.py diff --git a/api/tests/test_push_velocity.py b/api/tests/test_push_velocity.py new file mode 100644 index 0000000..eb56c17 --- /dev/null +++ b/api/tests/test_push_velocity.py @@ -0,0 +1,77 @@ +"""Push-velocity detection (scripts/telemetry_emit.py): detect + alert only. + +Driven through the real function with rows shaped like the Gitea query's. +""" + +from __future__ import annotations + +import importlib.util +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT / "scripts")) +_spec = importlib.util.spec_from_file_location("telemetry_emit", ROOT / "scripts" / "telemetry_emit.py") +te = importlib.util.module_from_spec(_spec) +_spec.loader.exec_module(te) + +NOW = 1_800_000_000.0 + + +def row(login="agent-et26abcd1234", uid=7, p1h=0, p24h=0, d24h=0, repos=1): + return {"uid": uid, "login": login, "p1h": p1h, "p24h": p24h, "d24h": d24h, "repos": repos} + + +def test_under_every_threshold_emits_nothing(): + ev, keep = te.push_velocity_events([row(p1h=60, p24h=500, d24h=10)], NOW, {}) + assert ev == [] and keep == {} + + +def test_burst_emits_one_declared_row_with_the_passport(): + ev, keep = te.push_velocity_events([row(p1h=61, p24h=61, repos=3)], NOW, {}) + assert len(ev) == 1 + e = ev[0] + assert e["event_type"] == "forge.push_velocity" and e["service"] == "forge" + assert e["actor_type"] == "agent" and e["actor_id"] == "ET26-ABCD-1234" + assert e["metadata"] == { + "rule": "pushes_1h", "window_s": 3600, "count": 61, "threshold": 60, + "repos": 3, "gitea_user_id": 7, + } + assert keep == {"7:pushes_1h": NOW} + + +def test_still_over_is_reported_once_per_window_not_every_run(): + _, keep = te.push_velocity_events([row(p1h=90)], NOW, {}) + ev, keep = te.push_velocity_events([row(p1h=95)], NOW + 300, keep) + assert ev == [] and keep == {"7:pushes_1h": NOW} + ev, _ = te.push_velocity_events([row(p1h=95)], NOW + 3601, keep) + assert len(ev) == 1 + + +def test_dropping_back_under_rearms(): + _, keep = te.push_velocity_events([row(p1h=90)], NOW, {}) + _, keep = te.push_velocity_events([row(p1h=5)], NOW + 300, keep) + assert keep == {} + ev, _ = te.push_velocity_events([row(p1h=90)], NOW + 600, keep) + assert len(ev) == 1 + + +def test_the_sync_account_is_exempt(): + ev, _ = te.push_velocity_events([row(login="windyadmin", uid=1, p1h=9999, p24h=9999)], NOW, {}) + assert ev == [] + + +def test_human_rows_carry_no_invented_actor_id(): + ev, _ = te.push_velocity_events([row(login="u-5e1b9569abc", d24h=11)], NOW, {}) + assert [(e["actor_type"], e["metadata"]["rule"]) for e in ev] == [("human", "ref_deletes_24h")] + assert "actor_id" not in ev[0] + + +def test_unparseable_agent_login_keeps_agent_type_without_actor_id(): + ev, _ = te.push_velocity_events([row(login="agent-weird", p24h=501)], NOW, {}) + assert ev[0]["actor_type"] == "agent" and "actor_id" not in ev[0] + + +def test_passport_round_trip(): + assert te.passport_from_login("agent-et26p1zgttp8") == "ET26-P1ZG-TTP8" + assert te.passport_from_login("u-abc") is None diff --git a/scripts/telemetry_emit.py b/scripts/telemetry_emit.py index c1af4b1..c37d153 100644 --- a/scripts/telemetry_emit.py +++ b/scripts/telemetry_emit.py @@ -20,6 +20,7 @@ from __future__ import annotations import json import os +import re import subprocess import sys import time @@ -69,6 +70,92 @@ def iso(epoch: float) -> str: return datetime.fromtimestamp(epoch, UTC).isoformat().replace("+00:00", "Z") +# ---- push velocity: DETECT + ALERT ONLY (G3.4, 2026-09-23) ----------------- +# `git push` goes straight to Gitea and never touches our API, so throttle.py +# cannot see it (NOT_ENFORCED_HERE). Gitea's own `action` table does record +# every push, so we read it here, emit `forge.push_velocity` when an account +# crosses a threshold, and let Telemetry Boss's detector page. Nothing here sits +# in the push path; nothing is ever refused (orchestrator, 09-23). +# +# Thresholds = the STANDARD-band bases from config.py (500 pushes/day; the +# force-push base of 10/day is used for ref deletes, the closest thing we can +# see). EI band multipliers are NOT applied: a platinum agent over 500/day is +# still flagged, for a human to look at, not blocked. Gitea records no +# "forced" flag, so force pushes cannot be told apart from pushes: named, not +# guessed. +PV_RULES = ( # (rule, row key, window_s, threshold) + ("pushes_1h", "p1h", 3600, 60), + ("pushes_24h", "p24h", 86400, 500), + ("ref_deletes_24h", "d24h", 86400, 10), +) +# The GitHub -> Windy Git sync pushes as windyadmin every 5 min, by design. +PV_EXEMPT = {"windyadmin"} +# Gitea op_type: 5 commit push, 9 tag push, 16 tag delete, 17 branch delete. +# One action row per WATCHER is written for each push; user_id = act_user_id +# keeps exactly the actor's own copy. +PV_QUERY = """ + select a.act_user_id as uid, u.lower_name as login, + count(*) filter (where a.op_type in (5, 9) and a.created_unix > {h1}) as p1h, + count(*) filter (where a.op_type in (5, 9)) as p24h, + count(*) filter (where a.op_type in (16, 17)) as d24h, + count(distinct a.repo_id) as repos + from action a join "user" u on u.id = a.act_user_id + where a.created_unix > {h24} and a.user_id = a.act_user_id + and a.op_type in (5, 9, 16, 17) + group by 1, 2""" + + +def passport_from_login(login: str) -> str | None: + """agent-et26abcd1234 -> ET26-ABCD-1234 (repos.py _owner_login, reversed).""" + m = re.fullmatch(r"agent-([a-z0-9]{4})([a-z0-9]{4})([a-z0-9]{4})", login) + return "-".join(g.upper() for g in m.groups()) if m else None + + +def push_velocity_events(rows: list[dict], now: float, alerted: dict) -> tuple[list[dict], dict]: + """(events, alerted') — one row per account per rule per window while over. + + `alerted` maps ":" -> epoch of the last row. An account still over + the line is re-reported once per window, not every 5 minutes; one that drops + back under is forgotten, so a later burst reports again. + """ + events, keep = [], {} + for r in rows: + login = str(r["login"]) + if login in PV_EXEMPT: + continue + agent = login.startswith("agent-") + for rule, key, window, limit in PV_RULES: + n = int(r[key]) + if n <= limit: + continue + k = f"{r['uid']}:{rule}" + last = alerted.get(k) + if last is not None and now - float(last) < window: + keep[k] = last + continue + keep[k] = now + ev = { + "ts": iso(now), + "platform": PLATFORM, + "service": "forge", + "event_type": "forge.push_velocity", + "actor_type": "agent" if agent else "human", + "metadata": { + "rule": rule, + "window_s": window, + "count": n, + "threshold": limit, + "repos": int(r["repos"]), + "gitea_user_id": int(r["uid"]), + }, + } + passport = passport_from_login(login) if agent else None + if passport: # unknown is absent, never invented (I-12) + ev["actor_id"] = passport + events.append(ev) + return events, keep + + def main() -> int: dry = "--dry-run" in sys.argv state = load_state() @@ -182,6 +269,14 @@ def main() -> int: } ) + pv_rows = sql(PV_QUERY.format(h1=int(now) - 3600, h24=int(now) - 86400)) + pv_events, pv_alerted = push_velocity_events(pv_rows, now, state.get("pv_alerted", {})) + for e in pv_events: + m = e["metadata"] + print(f"[telemetry] WARNING push velocity: gitea user {m['gitea_user_id']} " + f"{m['rule']} = {m['count']} > {m['threshold']}") + events += pv_events + # 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 @@ -264,7 +359,7 @@ def main() -> int: with open(STATE + ".tmp", "w") as f: json.dump( {"last_fin": new_fin, "last_id": new_id, "last_ts": now, - "quarantined_unreported": quarantined}, + "quarantined_unreported": quarantined, "pv_alerted": pv_alerted}, f, ) os.replace(STATE + ".tmp", STATE)