Compare commits
4 Commits
8690c4f8d3
...
push-veloc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
44e824c762 | ||
|
|
4e2db1e2e0 | ||
|
|
54c63b4d89 | ||
|
|
ac83c4721e |
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
|
||||||
@@ -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()
|
||||||
@@ -182,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
|
||||||
@@ -264,7 +373,7 @@ def main() -> int:
|
|||||||
with open(STATE + ".tmp", "w") as f:
|
with open(STATE + ".tmp", "w") as f:
|
||||||
json.dump(
|
json.dump(
|
||||||
{"last_fin": new_fin, "last_id": new_id, "last_ts": now,
|
{"last_fin": new_fin, "last_id": new_id, "last_ts": now,
|
||||||
"quarantined_unreported": quarantined},
|
"quarantined_unreported": quarantined, "pv_alerted": pv_alerted},
|
||||||
f,
|
f,
|
||||||
)
|
)
|
||||||
os.replace(STATE + ".tmp", STATE)
|
os.replace(STATE + ".tmp", STATE)
|
||||||
|
|||||||
Reference in New Issue
Block a user