5 Commits

Author SHA1 Message Date
Kit OC5
60708dd8db push velocity: correlate the SSO-id subquery on the grouped column
All checks were successful
check / gate (push) Successful in 21s
canary / probe (push) Successful in 9s
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:27 -04:00
Kit OC5
aedc772d29 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 14:06:27 -04:00
Kit OC5
acb8ec3a16 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 14:06:27 -04:00
Kit OC5
0051c72037 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 14:06:27 -04:00
Kit OC5
4c76b1b12a canary: end the hub session the login probe opens (journey cleanup rule)
All checks were successful
check / gate (push) Successful in 32s
canary / probe (push) Successful in 12s
identity.login created a live hub session every 10 min and never ended it.
It now logs out with the token it got: retried on 5xx / no response
(8 x 15 s), 401/404/410 = already over, any other 4xx fails fast, and a
cleanup it can't finish is reported as identity.logout DOWN "CLEANUP
FAILED" (alerts + red run). The hub's /auth/logout revokes every refresh
token of the account (verified live), so the next run's logout heals a
leftover; no ledger needed. Proven end to end: login 200, logout 200,
10/10 checks.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 13:39:55 -04:00
4 changed files with 356 additions and 4 deletions

View File

@@ -0,0 +1,92 @@
"""The canary's login probe must end the session it opens (journey cleanup rule)."""
from __future__ import annotations
import importlib.util
import io
import sys
import urllib.error
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
_spec = importlib.util.spec_from_file_location("canary", ROOT / "scripts" / "canary.py")
canary = importlib.util.module_from_spec(_spec)
sys.modules["canary"] = canary # dataclasses resolve their module by name
_spec.loader.exec_module(canary)
class _Resp:
def __init__(self, status=200, body=b"{}"):
self.status, self._body = status, body
def read(self):
return self._body
def __enter__(self):
return self
def __exit__(self, *a):
return False
def _err(code):
return urllib.error.HTTPError(canary.LOGOUT_URL, code, "x", {}, io.BytesIO(b""))
def _script(monkeypatch, outcomes):
calls = []
def fake(req, timeout=None):
calls.append((req.get_method(), req.full_url, req.get_header("Authorization")))
o = outcomes.pop(0)
if isinstance(o, Exception):
raise o
return o
monkeypatch.setattr(canary.urllib.request, "urlopen", fake)
return calls
def test_logout_ends_the_session(monkeypatch):
calls = _script(monkeypatch, [_Resp(200)])
r = canary.logout("tok", sleep=lambda s: None)
assert r.status == "ok"
assert calls == [("POST", canary.LOGOUT_URL, "Bearer tok")]
def test_5xx_and_no_response_are_retried_then_succeed(monkeypatch):
calls = _script(monkeypatch, [_err(502), OSError("reset"), _Resp(200)])
assert canary.logout("tok", sleep=lambda s: None).status == "ok"
assert len(calls) == 3
def test_already_over_counts_as_done(monkeypatch):
_script(monkeypatch, [_err(401)])
assert canary.logout("tok", sleep=lambda s: None).status == "ok"
def test_other_4xx_fails_fast_and_honestly(monkeypatch):
calls = _script(monkeypatch, [_err(400)])
r = canary.logout("tok", sleep=lambda s: None)
assert r.status == "down" and r.detail.startswith("CLEANUP FAILED") and len(calls) == 1
def test_retries_are_bounded_and_reported(monkeypatch):
calls = _script(monkeypatch, [_err(503)] * 8)
r = canary.logout("tok", attempts=8, sleep=lambda s: None)
assert r.status == "down" and "CLEANUP FAILED after 8 tries" in r.detail and len(calls) == 8
def test_login_probe_logs_out_with_the_token_it_got(monkeypatch):
calls = _script(monkeypatch, [_Resp(200, b'{"token": "abc"}'), _Resp(200)])
c = canary.Check("identity.login", "https://account.windyword.ai/api/v1/auth/login", "x",
method="POST", body={"email": "e", "password": "p"},
after=canary._logout_after_login)
r = canary._probe(c)
assert r.status == "ok" and [f.status for f in r.followups] == ["ok"]
assert calls[1] == ("POST", canary.LOGOUT_URL, "Bearer abc")
def test_login_without_token_is_a_cleanup_failure_not_a_pass():
[f] = canary._logout_after_login(b"{}")
assert f.status == "down" and "CLEANUP FAILED" in f.detail

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

@@ -33,6 +33,7 @@ import sys
import time import time
import urllib.error import urllib.error
import urllib.request import urllib.request
from collections.abc import Callable
from dataclasses import dataclass, field from dataclasses import dataclass, field
STATE_PATH = os.environ.get("CANARY_STATE", "canary-state.json") STATE_PATH = os.environ.get("CANARY_STATE", "canary-state.json")
@@ -46,6 +47,16 @@ ALERT_FROM = os.environ.get("CANARY_ALERT_FROM", "office@thewindstorm.uk")
LOGIN_WARN_SECONDS = float(os.environ.get("CANARY_LOGIN_WARN_S", "35")) LOGIN_WARN_SECONDS = float(os.environ.get("CANARY_LOGIN_WARN_S", "35"))
TIMEOUT = float(os.environ.get("CANARY_TIMEOUT_S", "60")) TIMEOUT = float(os.environ.get("CANARY_TIMEOUT_S", "60"))
# Journey cleanup rule (orchestrator, 2026-09-23). The login probe creates a hub
# session (access + refresh token) every run, so it must end it. The hub's
# /auth/logout revokes the token AND every refresh token of the account
# (verified live: access 401, refresh 401 after it). So the next successful
# logout also heals anything a failed run left behind; no ledger needed.
LOGOUT_URL = "https://account.windyword.ai/api/v1/auth/logout"
LOGOUT_ATTEMPTS = 8 # retried on 5xx / no response only
LOGOUT_GAP_S = 15.0
LOGOUT_GONE = (401, 404, 410) # the session is already over = done
@dataclass @dataclass
class Result: class Result:
@@ -54,6 +65,7 @@ class Result:
detail: str detail: str
seconds: float = 0.0 seconds: float = 0.0
user_visible: str = "" user_visible: str = ""
followups: list[Result] = field(default_factory=list)
@dataclass @dataclass
@@ -68,6 +80,8 @@ class Check:
# When True this check INVERTS: a 2xx is a critical failure (a security # When True this check INVERTS: a 2xx is a critical failure (a security
# control opened) and a 401/403/503 is the healthy, expected outcome. # control opened) and a 401/403/503 is the healthy, expected outcome.
must_refuse: bool = False must_refuse: bool = False
# Runs on a 2xx with the response body; returns follow-up results (cleanup).
after: Callable[[bytes], list[Result]] | None = None
def _probe(c: Check) -> Result: def _probe(c: Check) -> Result:
@@ -91,13 +105,18 @@ def _probe(c: Check) -> Result:
elapsed, c.what_it_proves) elapsed, c.what_it_proves)
if r.status >= 400: if r.status >= 400:
return Result(c.name, "down", f"HTTP {r.status}", elapsed, c.what_it_proves) return Result(c.name, "down", f"HTTP {r.status}", elapsed, c.what_it_proves)
raw = r.read()
warn = c.warn_seconds warn = c.warn_seconds
if warn and elapsed > warn: if warn and elapsed > warn:
return Result( res = Result(
c.name, "slow", f"HTTP {r.status} in {elapsed:.1f}s (warn >{warn:.0f}s)", c.name, "slow", f"HTTP {r.status} in {elapsed:.1f}s (warn >{warn:.0f}s)",
elapsed, c.what_it_proves, elapsed, c.what_it_proves,
) )
return Result(c.name, "ok", f"HTTP {r.status} in {elapsed:.1f}s", elapsed, c.what_it_proves) else:
res = Result(c.name, "ok", f"HTTP {r.status} in {elapsed:.1f}s", elapsed, c.what_it_proves)
if c.after:
res.followups = c.after(raw)
return res
except urllib.error.HTTPError as e: except urllib.error.HTTPError as e:
if c.must_refuse and e.code in (401, 403, 503): if c.must_refuse and e.code in (401, 403, 503):
return Result(c.name, "ok", f"correctly refused (HTTP {e.code})", return Result(c.name, "ok", f"correctly refused (HTTP {e.code})",
@@ -110,6 +129,52 @@ def _probe(c: Check) -> Result:
) )
def logout(token: str, *, attempts: int = LOGOUT_ATTEMPTS, gap: float = LOGOUT_GAP_S,
sleep: Callable[[float], None] = time.sleep) -> Result:
"""End the session the login probe opened. Honest: never ok unless proven."""
what = "the canary leaves no live session behind (journey cleanup rule)"
headers = {
"User-Agent": "windy-git-canary/1.0",
"X-Windy-Synthetic": "1",
"Authorization": f"Bearer {token}",
}
start = time.monotonic()
last = "no attempt"
for i in range(attempts):
if i:
sleep(gap)
req = urllib.request.Request(LOGOUT_URL, data=b"", method="POST", headers=headers)
try:
with urllib.request.urlopen(req, timeout=TIMEOUT) as r:
return Result("identity.logout", "ok", f"session ended (HTTP {r.status})",
time.monotonic() - start, what)
except urllib.error.HTTPError as e:
if e.code in LOGOUT_GONE:
return Result("identity.logout", "ok", f"session already over (HTTP {e.code})",
time.monotonic() - start, what)
if e.code < 500: # a 4xx won't change on retry: fail fast
return Result("identity.logout", "down", f"CLEANUP FAILED: HTTP {e.code}",
time.monotonic() - start, what)
last = f"HTTP {e.code}"
except Exception as e: # noqa: BLE001 — no response / timeout: retry
last = f"{type(e).__name__}"
return Result("identity.logout", "down",
f"CLEANUP FAILED after {attempts} tries: {last} (next run's logout heals it)",
time.monotonic() - start, what)
def _logout_after_login(raw: bytes) -> list[Result]:
try:
token = (json.loads(raw or b"{}") or {}).get("token")
except ValueError:
token = None
if not token:
return [Result("identity.logout", "down",
"CLEANUP FAILED: login returned no token to log out with",
0.0, "the canary leaves no live session behind (journey cleanup rule)")]
return [logout(token)]
def build_checks() -> list[Check]: def build_checks() -> list[Check]:
checks = [ checks = [
Check( Check(
@@ -187,6 +252,7 @@ def build_checks() -> list[Check]:
method="POST", method="POST",
body={"email": email, "password": pw}, body={"email": email, "password": pw},
warn_seconds=LOGIN_WARN_SECONDS, warn_seconds=LOGIN_WARN_SECONDS,
after=_logout_after_login,
) )
) )
return checks return checks
@@ -263,7 +329,10 @@ def main() -> int:
args = ap.parse_args() args = ap.parse_args()
previous = load_state() previous = load_state()
results = [_probe(c) for c in build_checks()] results = []
for c in build_checks():
r = _probe(c)
results += [r, *r.followups]
print(f"windy canary — {time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())}\n") print(f"windy canary — {time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())}\n")
for r in results: for r in results:

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()
@@ -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)