Compare commits
6 Commits
4acf50d9ef
...
push-veloc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
44e824c762 | ||
|
|
4e2db1e2e0 | ||
|
|
54c63b4d89 | ||
|
|
ac83c4721e | ||
|
|
2d4fadb090 | ||
|
|
e3b69fa759 |
@@ -120,6 +120,9 @@ class Telemetry:
|
||||
self.window_start = time.time()
|
||||
self.requests = self.errors_5xx = self.errors_4xx = self.refusals_4xx = 0
|
||||
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) --------------------------
|
||||
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:
|
||||
self.dropped += len(self.buffer) - MAX_BUFFER
|
||||
del self.buffer[: len(self.buffer) - MAX_BUFFER]
|
||||
|
||||
def boot(self) -> None:
|
||||
@@ -188,6 +192,8 @@ class Telemetry:
|
||||
"errors_5xx": self.errors_5xx,
|
||||
"errors_4xx": self.errors_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
|
||||
s = sorted(self.latencies_ms)
|
||||
@@ -199,7 +205,7 @@ class Telemetry:
|
||||
self._reset_window()
|
||||
|
||||
# ---- sending ------------------------------------------------------------
|
||||
def _post(self, batch: list[dict]) -> int:
|
||||
def _post(self, batch: list[dict]) -> tuple[int, dict]:
|
||||
req = urllib.request.Request(
|
||||
self.url,
|
||||
data=json.dumps({"events": batch}).encode(),
|
||||
@@ -211,19 +217,32 @@ class Telemetry:
|
||||
},
|
||||
)
|
||||
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:
|
||||
if not self.enabled or not self.buffer:
|
||||
return
|
||||
batch = self.buffer[:500]
|
||||
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
|
||||
log.warning("telemetry flush failed (%d rows kept): %s", len(self.buffer), exc)
|
||||
return
|
||||
if 200 <= status < 300:
|
||||
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:
|
||||
"""The one in-process timer: flush every minute, heartbeat every hour."""
|
||||
|
||||
82
api/tests/test_push_velocity.py
Normal file
82
api/tests/test_push_velocity.py
Normal 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
|
||||
@@ -160,3 +160,43 @@ def test_synthetic_is_forwarded_downstream_only_for_synthetic_requests():
|
||||
finally:
|
||||
tmod.SYNTHETIC.reset(token)
|
||||
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
|
||||
|
||||
@@ -82,7 +82,11 @@ done
|
||||
|
||||
# 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.
|
||||
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
|
||||
# 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).
|
||||
# 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; }
|
||||
log "all repos in step with GitHub"
|
||||
|
||||
@@ -20,6 +20,7 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
@@ -69,6 +70,100 @@ 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,
|
||||
(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:
|
||||
dry = "--dry-run" in sys.argv
|
||||
state = load_state()
|
||||
@@ -161,6 +256,12 @@ def main() -> int:
|
||||
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
|
||||
@@ -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.
|
||||
spool = os.environ.get("JANITOR_SPOOL", "/var/lib/windy-git/janitor-cancelled.jsonl")
|
||||
spooled = 0
|
||||
@@ -216,6 +331,7 @@ def main() -> int:
|
||||
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,
|
||||
@@ -229,11 +345,20 @@ def main() -> int:
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as r:
|
||||
body = r.read()[:300]
|
||||
body = r.read()
|
||||
if r.status >= 300:
|
||||
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:
|
||||
print(f"[telemetry] FAILED ingest HTTP {e.code}: {e.read()[:200]!r}")
|
||||
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)
|
||||
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)
|
||||
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)
|
||||
print(f"[telemetry] sent {len(events)} events ({len(jobs)} ci.run)")
|
||||
return 0
|
||||
|
||||
Reference in New Issue
Block a user