Compare commits
6 Commits
ci-sysbox
...
60708dd8db
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
60708dd8db | ||
|
|
aedc772d29 | ||
|
|
acb8ec3a16 | ||
|
|
0051c72037 | ||
|
|
4c76b1b12a | ||
|
|
2d4fadb090 |
92
api/tests/test_canary_cleanup.py
Normal file
92
api/tests/test_canary_cleanup.py
Normal 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
|
||||
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
|
||||
@@ -1,15 +0,0 @@
|
||||
# ROLLBACK ONLY: the pre-Sysbox dind (privileged: true), kept one command away.
|
||||
# Use it if CI breaks under Sysbox:
|
||||
#
|
||||
# cd /srv/windygit/src/deploy/runner
|
||||
# sudo docker compose -f docker-compose.yml -f docker-compose.privileged.yml up -d dind
|
||||
#
|
||||
# (then restart the runners while idle). Compose merges `volumes` by container
|
||||
# path, so this puts back the old `dind-storage` volume with its image cache.
|
||||
# Going forward again: the same command without the second -f.
|
||||
services:
|
||||
dind:
|
||||
runtime: runc
|
||||
privileged: true
|
||||
volumes:
|
||||
- dind-storage:/var/lib/docker
|
||||
@@ -26,12 +26,8 @@
|
||||
# * `dind` and every job container it spawns are UNTRUSTED. They are on a
|
||||
# private network with no access to the forge, its database, or its .env.
|
||||
#
|
||||
# dind is NOT privileged (2026-09-23): it runs under the Sysbox runtime
|
||||
# (sysbox-ce on Veron, `runtime: sysbox-runc`), a user-namespaced system
|
||||
# container whose root is an unprivileged host uid. A job that escapes its own
|
||||
# container lands in dind as a nobody on the host, not as root on Grant's
|
||||
# workstation. Before Sysbox, dind was `privileged: true`; that config is kept
|
||||
# as docker-compose.privileged.yml (ROLLBACK ONLY, one command, see that file).
|
||||
# dind itself is privileged — that is the cost, and it is the reason a job
|
||||
# escape lands in a disposable daemon rather than on Grant's workstation.
|
||||
#
|
||||
# ⚠️ Do NOT "simplify" this by mounting the host docker socket.
|
||||
|
||||
@@ -40,15 +36,13 @@ name: windy-git-runner
|
||||
services:
|
||||
dind:
|
||||
image: docker.io/library/docker:27-dind
|
||||
runtime: sysbox-runc # NOT privileged: see the I-5 note above
|
||||
privileged: true
|
||||
environment:
|
||||
DOCKER_TLS_CERTDIR: "" # plain TCP on an isolated network, no host route
|
||||
command: ["dockerd", "--host=tcp://0.0.0.0:2375", "--tls=false"]
|
||||
networks: [jobs]
|
||||
volumes:
|
||||
# A fresh volume: Sysbox shifts ownership to its own uid range. The old
|
||||
# `dind-storage` is kept untouched for the privileged rollback.
|
||||
- dind-storage-sysbox:/var/lib/docker
|
||||
- dind-storage:/var/lib/docker
|
||||
# G1.5 — bounded so a fork-bomb workflow cannot starve Grant's interactive
|
||||
# session. Veron 1 is his workstation, not a dedicated build box.
|
||||
cpus: 12.0 # 12 of 24 cores
|
||||
@@ -165,7 +159,6 @@ networks:
|
||||
|
||||
volumes:
|
||||
dind-storage:
|
||||
dind-storage-sysbox:
|
||||
runner-data:
|
||||
runner-data-2:
|
||||
runner-data-3:
|
||||
|
||||
@@ -33,6 +33,7 @@ import sys
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
from collections.abc import Callable
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
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"))
|
||||
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
|
||||
class Result:
|
||||
@@ -54,6 +65,7 @@ class Result:
|
||||
detail: str
|
||||
seconds: float = 0.0
|
||||
user_visible: str = ""
|
||||
followups: list[Result] = field(default_factory=list)
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -68,6 +80,8 @@ class Check:
|
||||
# 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.
|
||||
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:
|
||||
@@ -91,13 +105,18 @@ def _probe(c: Check) -> Result:
|
||||
elapsed, c.what_it_proves)
|
||||
if r.status >= 400:
|
||||
return Result(c.name, "down", f"HTTP {r.status}", elapsed, c.what_it_proves)
|
||||
raw = r.read()
|
||||
warn = c.warn_seconds
|
||||
if warn and elapsed > warn:
|
||||
return Result(
|
||||
res = Result(
|
||||
c.name, "slow", f"HTTP {r.status} in {elapsed:.1f}s (warn >{warn:.0f}s)",
|
||||
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:
|
||||
if c.must_refuse and e.code in (401, 403, 503):
|
||||
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]:
|
||||
checks = [
|
||||
Check(
|
||||
@@ -187,6 +252,7 @@ def build_checks() -> list[Check]:
|
||||
method="POST",
|
||||
body={"email": email, "password": pw},
|
||||
warn_seconds=LOGIN_WARN_SECONDS,
|
||||
after=_logout_after_login,
|
||||
)
|
||||
)
|
||||
return checks
|
||||
@@ -263,7 +329,10 @@ def main() -> int:
|
||||
args = ap.parse_args()
|
||||
|
||||
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")
|
||||
for r in results:
|
||||
|
||||
@@ -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()
|
||||
@@ -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.
|
||||
spool = os.environ.get("JANITOR_SPOOL", "/var/lib/windy-git/janitor-cancelled.jsonl")
|
||||
spooled = 0
|
||||
@@ -264,7 +373,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)
|
||||
|
||||
Reference in New Issue
Block a user