6 Commits

Author SHA1 Message Date
Kit OC5
44e824c762 push velocity: correlate the SSO-id subquery on the grouped column
Dry-run against the real gitea DB: 'subquery uses ungrouped column u.id'.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:06:10 -04:00
Kit OC5
4e2db1e2e0 push velocity: isolate its query so a failure never costs ci.run rows
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 13:06:20 -04:00
Kit OC5
54c63b4d89 push velocity: key humans on windy_identity_id (SSO link); no id -> system + caller
Telemetry UPDATE 2 actor rule: agent/human rows without actor_id are
quarantined. Forge humans sign in only via Windy SSO, so Gitea's
external_login_user.external_id is their windy_identity_id.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 13:06:17 -04:00
Kit OC5
ac83c4721e 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 <noreply@anthropic.com>
2026-09-23 13:06:17 -04:00
Kit OC5
2d4fadb090 sync: bound the janitor and telemetry steps (docker exec hangs in an IO stall)
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
09-23 16:43Z the Veron data2 SMR stall left runc exec in D state; the
janitor's docker exec never returned, so the sync sat 'activating' and no
repo mirrored or got a status for any lane. Both steps are non-fatal;
now they time out (120 s / 180 s) and the run continues.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:54:34 -04:00
Kit OC5
e3b69fa759 telemetry: UPDATE 7 — read the ingest body; count quarantined + dropped on heartbeats
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
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
5 changed files with 282 additions and 8 deletions

View File

@@ -120,6 +120,9 @@ class Telemetry:
self.window_start = time.time() self.window_start = time.time()
self.requests = self.errors_5xx = self.errors_4xx = self.refusals_4xx = 0 self.requests = self.errors_5xx = self.errors_4xx = self.refusals_4xx = 0
self.latencies_ms: list[float] = [] 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) -------------------------- # ---- recording (never raises into a request) --------------------------
def record_request(self, status: int, duration_ms: float, *, refused: bool = False) -> None: 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: if len(self.buffer) > MAX_BUFFER:
self.dropped += len(self.buffer) - MAX_BUFFER
del self.buffer[: len(self.buffer) - MAX_BUFFER] del self.buffer[: len(self.buffer) - MAX_BUFFER]
def boot(self) -> None: def boot(self) -> None:
@@ -188,6 +192,8 @@ class Telemetry:
"errors_5xx": self.errors_5xx, "errors_5xx": self.errors_5xx,
"errors_4xx": self.errors_4xx, "errors_4xx": self.errors_4xx,
"refusals_4xx": self.refusals_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 if self.latencies_ms: # no traffic = no p95, not a fake 0
s = sorted(self.latencies_ms) s = sorted(self.latencies_ms)
@@ -199,7 +205,7 @@ class Telemetry:
self._reset_window() self._reset_window()
# ---- sending ------------------------------------------------------------ # ---- sending ------------------------------------------------------------
def _post(self, batch: list[dict]) -> int: def _post(self, batch: list[dict]) -> tuple[int, dict]:
req = urllib.request.Request( req = urllib.request.Request(
self.url, self.url,
data=json.dumps({"events": batch}).encode(), data=json.dumps({"events": batch}).encode(),
@@ -211,19 +217,32 @@ class Telemetry:
}, },
) )
with urllib.request.urlopen(req, timeout=20) as r: 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: async def flush(self) -> None:
if not self.enabled or not self.buffer: if not self.enabled or not self.buffer:
return return
batch = self.buffer[:500] batch = self.buffer[:500]
try: 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 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) log.warning("telemetry flush failed (%d rows kept): %s", len(self.buffer), exc)
return return
if 200 <= status < 300: if 200 <= status < 300:
del self.buffer[: len(batch)] 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: async def run(self) -> None:
"""The one in-process timer: flush every minute, heartbeat every hour.""" """The one in-process timer: flush every minute, heartbeat every hour."""

View File

@@ -0,0 +1,82 @@
"""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, wid=None):
return {"uid": uid, "login": login, "wid": wid, "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_sso_human_is_keyed_on_windy_identity_id():
ev, _ = te.push_velocity_events([row(login="u-5e1b9569abc", wid="5e1b9569-full-id", d24h=11)], NOW, {})
assert [(e["actor_type"], e["actor_id"], e["metadata"]["rule"]) for e in ev] == [
("human", "5e1b9569-full-id", "ref_deletes_24h")
]
assert "caller" not in ev[0]["metadata"]
def test_no_provable_id_is_system_plus_caller_never_an_invented_id():
# UPDATE 2 actor rule: agent/human rows without an actor_id are quarantined.
for login in ("u-nolink", "agent-weird"):
ev, _ = te.push_velocity_events([row(login=login, p24h=501)], NOW, {})
assert ev[0]["actor_type"] == "system" and "actor_id" not in ev[0]
assert ev[0]["metadata"]["caller"] == "unknown"
def test_passport_round_trip():
assert te.passport_from_login("agent-et26p1zgttp8") == "ET26-P1ZG-TTP8"
assert te.passport_from_login("u-abc") is None

View File

@@ -160,3 +160,43 @@ def test_synthetic_is_forwarded_downstream_only_for_synthetic_requests():
finally: finally:
tmod.SYNTHETIC.reset(token) tmod.SYNTHETIC.reset(token)
assert tmod.synthetic_headers() == {} 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

View File

@@ -82,7 +82,11 @@ done
# Jobs that name labels no runner has (ubuntu/macos/windows-latest) would wait # Jobs that name labels no runner has (ubuntu/macos/windows-latest) would wait
# forever and invisibly; cancel them after 30 min. Never fails the sync. # forever and invisibly; cancel them after 30 min. Never fails the sync.
bash "$(dirname "$0")/cancel_unrunnable.sh" || log "janitor failed (non-fatal)" # Both DB steps go through `docker exec`, which hangs outright while the host
# is in an IO stall (09-23: data2 SMR cliff wedged this sync for 10+ min and
# stopped mirroring + the bridge for every lane). They are optional; mirroring
# and the bridge are not. Bound them so a stuck exec costs one step, not the run.
timeout -k 10 120 bash "$(dirname "$0")/cancel_unrunnable.sh" || log "janitor failed or timed out (non-fatal)"
# Private repos can't run GitHub Actions; mirror their open PRs here so CI # Private repos can't run GitHub Actions; mirror their open PRs here so CI
# fires, and post the verdicts back to GitHub as commit statuses. # fires, and post the verdicts back to GitHub as commit statuses.
@@ -92,7 +96,7 @@ fi
# CI telemetry -> admin.windyword.ai (shapes declared with Windy Telemetry 40). # CI telemetry -> admin.windyword.ai (shapes declared with Windy Telemetry 40).
# Sends nothing until WINDYGIT_TELEMETRY_TOKEN is set; never fails the sync. # Sends nothing until WINDYGIT_TELEMETRY_TOKEN is set; never fails the sync.
python3 "$(dirname "$0")/telemetry_emit.py" || log "telemetry emit failed (non-fatal)" timeout -k 10 180 python3 "$(dirname "$0")/telemetry_emit.py" || log "telemetry emit failed or timed out (non-fatal)"
[[ "$FAILED" -ne 0 ]] && { log "COMPLETED WITH FAILURES"; exit 1; } [[ "$FAILED" -ne 0 ]] && { log "COMPLETED WITH FAILURES"; exit 1; }
log "all repos in step with GitHub" log "all repos in step with GitHub"

View File

@@ -20,6 +20,7 @@ from __future__ import annotations
import json import json
import os import os
import re
import subprocess import subprocess
import sys import sys
import time import time
@@ -69,6 +70,100 @@ def iso(epoch: float) -> str:
return datetime.fromtimestamp(epoch, UTC).isoformat().replace("+00:00", "Z") 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,
(select el.external_id from external_login_user el
where el.user_id = a.act_user_id order by el.external_id limit 1) as wid,
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 "<uid>:<rule>" -> 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",
"metadata": {
"rule": rule,
"window_s": window,
"count": n,
"threshold": limit,
"repos": int(r["repos"]),
"gitea_user_id": int(r["uid"]),
},
}
# Actor rule (telemetry UPDATE 2): agent/human rows MUST carry an
# actor_id. Humans sign in to the forge only via Windy SSO, so the
# external login id IS their windy_identity_id. No id we can prove
# -> actor_type system + metadata.caller, never an invented id (I-12).
actor_id = passport_from_login(login) if agent else (r.get("wid") or None)
if actor_id:
ev["actor_type"], ev["actor_id"] = ("agent" if agent else "human"), str(actor_id)
else:
ev["actor_type"] = "system"
ev["metadata"]["caller"] = "unknown"
events.append(ev)
return events, keep
def main() -> int: def main() -> int:
dry = "--dry-run" in sys.argv dry = "--dry-run" in sys.argv
state = load_state() state = load_state()
@@ -161,6 +256,12 @@ def main() -> int:
h["jobs_cancelled"] = sum(1 for j in jobs if j["status"] == 3) 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
# 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"): for k in ("repos_synced", "repos_sync_failed", "statuses_posted", "bridge_errors"):
v = os.environ.get(f"TELEMETRY_{k.upper()}") v = os.environ.get(f"TELEMETRY_{k.upper()}")
if v is not None and v.isdigit(): # absent = couldn't count; never invent 0 if v is not None and v.isdigit(): # absent = couldn't count; never invent 0
@@ -176,6 +277,20 @@ def main() -> int:
} }
) )
# Isolated: a failing push-velocity query must never cost the ci.run rows.
pv_alerted = state.get("pv_alerted", {})
try:
pv_rows = sql(PV_QUERY.format(h1=int(now) - 3600, h24=int(now) - 86400))
pv_events, pv_alerted = push_velocity_events(pv_rows, now, pv_alerted)
except (subprocess.CalledProcessError, ValueError, KeyError) as e:
print(f"[telemetry] push velocity check FAILED (non-fatal): {type(e).__name__}")
pv_events = []
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. # 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") spool = os.environ.get("JANITOR_SPOOL", "/var/lib/windy-git/janitor-cancelled.jsonl")
spooled = 0 spooled = 0
@@ -216,6 +331,7 @@ def main() -> int:
print("[telemetry] WINDYGIT_TELEMETRY_TOKEN unset — not sending (not a failure)") print("[telemetry] WINDYGIT_TELEMETRY_TOKEN unset — not sending (not a failure)")
return 0 return 0
quarantined = 0
for i in range(0, len(events), 500): for i in range(0, len(events), 500):
req = urllib.request.Request( req = urllib.request.Request(
INGEST, INGEST,
@@ -229,11 +345,20 @@ def main() -> int:
) )
try: try:
with urllib.request.urlopen(req, timeout=30) as r: with urllib.request.urlopen(req, timeout=30) as r:
body = r.read()[:300] body = r.read()
if r.status >= 300: if r.status >= 300:
raise urllib.error.HTTPError( 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: except urllib.error.HTTPError as e:
print(f"[telemetry] FAILED ingest HTTP {e.code}: {e.read()[:200]!r}") print(f"[telemetry] FAILED ingest HTTP {e.code}: {e.read()[:200]!r}")
return 1 # state NOT advanced: the same rows retry next run return 1 # state NOT advanced: the same rows retry next run
@@ -246,7 +371,11 @@ def main() -> int:
os.makedirs(os.path.dirname(STATE), exist_ok=True) 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) 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_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, "pv_alerted": pv_alerted},
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