Compare commits
8 Commits
c83f808a60
...
push-veloc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
44e824c762 | ||
|
|
4e2db1e2e0 | ||
|
|
54c63b4d89 | ||
|
|
ac83c4721e | ||
|
|
2d4fadb090 | ||
|
|
e3b69fa759 | ||
|
|
4acf50d9ef | ||
|
|
b7a7e94df0 |
@@ -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."""
|
||||
|
||||
@@ -8,6 +8,7 @@ no status at all.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import importlib.util
|
||||
from pathlib import Path
|
||||
|
||||
@@ -35,12 +36,23 @@ def _run(i, wf, job, status, sha=SHA, n=1):
|
||||
|
||||
|
||||
class Fake:
|
||||
def __init__(self, runs=(), statuses=(), gh_prs=(), wg_prs=()):
|
||||
def __init__(self, runs=(), statuses=(), gh_prs=(), wg_prs=(), workflows=None):
|
||||
self.runs, self.statuses = list(runs), list(statuses)
|
||||
self.workflows = workflows or {} # {path: yaml text} at every commit
|
||||
self.gh_prs, self.wg_prs = list(gh_prs), list(wg_prs)
|
||||
self.posted, self.opened, self.closed = [], [], []
|
||||
|
||||
def gitea(self, method, path, body=None):
|
||||
if "/contents/" in path:
|
||||
want = path.split("/contents/", 1)[1].split("?", 1)[0]
|
||||
if want in self.workflows:
|
||||
return 200, {"content": base64.b64encode(self.workflows[want].encode()).decode()}
|
||||
files = [
|
||||
{"type": "file", "name": k.rsplit("/", 1)[1], "path": k}
|
||||
for k in self.workflows
|
||||
if k.rsplit("/", 1)[0] == want
|
||||
]
|
||||
return (200, files) if files else (404, None)
|
||||
if "/actions/tasks" in path:
|
||||
page = int(path.rsplit("page=", 1)[1])
|
||||
return 200, {"workflow_runs": self.runs[(page - 1) * 50 : page * 50]}
|
||||
@@ -209,3 +221,54 @@ def test_transport_blips_are_retried_but_http_errors_are_not(monkeypatch):
|
||||
monkeypatch.setattr(bridge.urllib.request, "urlopen", forbidden)
|
||||
assert bridge._call("http://x", "t", "GET", "/p") == (403, None)
|
||||
assert calls["n"] == 1
|
||||
|
||||
|
||||
GOOD = "on: push\njobs:\n test:\n runs-on: ubuntu-latest\n steps: []\n"
|
||||
BROKEN = "on: push\njobs:\n test:\n runs-on: x\n steps: [\n"
|
||||
|
||||
|
||||
def test_invalid_workflow_gets_an_error_status_even_with_no_runs(fake):
|
||||
# Gitea fires NO run for an invalid file: without this the PR shows nothing.
|
||||
f = fake(workflows={".github/workflows/ci.yml": BROKEN})
|
||||
bridge.post_statuses("windy-chat", SHA)
|
||||
assert [(p["context"], p["state"]) for p in f.posted] == [("windy-git/ci/workflow", "error")]
|
||||
assert "invalid YAML at line 5" in f.posted[0]["description"]
|
||||
assert f.posted[0]["target_url"].endswith(f"/src/commit/{SHA}/.github/workflows/ci.yml")
|
||||
|
||||
|
||||
def test_valid_workflows_post_nothing_extra(fake):
|
||||
f = fake(runs=[_run(1, "ci.yml", "test", "success")], workflows={".github/workflows/ci.yml": GOOD})
|
||||
bridge.post_statuses("windy-chat", SHA)
|
||||
assert [p["context"] for p in f.posted] == ["windy-git/ci/test"]
|
||||
|
||||
|
||||
def test_workflow_error_is_not_reposted(fake):
|
||||
f = fake(
|
||||
workflows={".github/workflows/ci.yml": BROKEN},
|
||||
statuses=[{"context": "windy-git/ci/workflow", "state": "error"}],
|
||||
)
|
||||
bridge.post_statuses("windy-chat", SHA)
|
||||
assert f.posted == []
|
||||
|
||||
|
||||
def test_gitea_dir_wins_over_github_dir(fake):
|
||||
# Gitea runs .gitea/workflows when it has files and ignores .github/workflows.
|
||||
f = fake(workflows={".gitea/workflows/ci.yml": GOOD, ".github/workflows/old.yml": BROKEN})
|
||||
bridge.post_statuses("windy-chat", SHA)
|
||||
assert f.posted == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"text, problem",
|
||||
[
|
||||
(GOOD, None),
|
||||
("on: push\njobs:\n a:\n uses: ./x.yml\n", None),
|
||||
(BROKEN, "invalid YAML at line 5"),
|
||||
("jobs:\n a:\n runs-on: x\n", "no `on:` trigger"),
|
||||
("on: push\n", "no `jobs:`"),
|
||||
("on: push\njobs:\n a:\n steps: []\n", "job `a` has no `runs-on:`"),
|
||||
("- a\n", "not a YAML mapping"),
|
||||
],
|
||||
)
|
||||
def test_workflow_problem(text, problem):
|
||||
assert bridge.workflow_problem(text) == problem
|
||||
|
||||
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
|
||||
|
||||
@@ -31,7 +31,7 @@ dependencies = [
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
dev = ["pytest>=8.3", "pytest-asyncio>=0.24", "ruff>=0.7", "mypy>=1.13"]
|
||||
dev = ["pytest>=8.3", "pyyaml>=6.0", "pytest-asyncio>=0.24", "ruff>=0.7", "mypy>=1.13"]
|
||||
|
||||
[tool.ruff]
|
||||
line-length = 100
|
||||
|
||||
@@ -28,6 +28,7 @@ repo code, and no secret is handed to any repo.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
@@ -36,6 +37,8 @@ import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
import yaml
|
||||
|
||||
GITEA = os.environ.get("BRIDGE_GITEA_URL", "http://localhost:3080").rstrip("/")
|
||||
PUBLIC = "https://app.windygit.com"
|
||||
GITEA_TOKEN = os.environ.get("GITEA_ADMIN_TOKEN", "")
|
||||
@@ -86,6 +89,62 @@ for _entry in os.environ.get(
|
||||
NON_BLOCKING[_repo.strip()] = {j.strip() for j in _jobs.split(",") if j.strip()}
|
||||
|
||||
|
||||
# Gitea reads the FIRST of these dirs that has workflow files at a commit (1.24).
|
||||
WORKFLOW_DIRS = (".gitea/workflows", ".github/workflows")
|
||||
|
||||
|
||||
def workflow_problem(text: str) -> str | None:
|
||||
"""Why Gitea would drop this workflow file, or None if it looks runnable.
|
||||
|
||||
Gitea skips an invalid workflow with one log line and fires no run at all,
|
||||
so on GitHub the PR just shows nothing, and people wait for CI that is never
|
||||
coming. These are the shapes we have actually hit, not a full schema.
|
||||
"""
|
||||
try:
|
||||
doc = yaml.safe_load(text)
|
||||
except yaml.YAMLError as e:
|
||||
mark = getattr(e, "problem_mark", None)
|
||||
return f"invalid YAML at line {mark.line + 1}" if mark else "invalid YAML"
|
||||
if not isinstance(doc, dict):
|
||||
return "not a YAML mapping"
|
||||
if "on" not in doc and True not in doc: # YAML 1.1 reads a bare `on` as True
|
||||
return "no `on:` trigger"
|
||||
jobs = doc.get("jobs")
|
||||
if not isinstance(jobs, dict) or not jobs:
|
||||
return "no `jobs:`"
|
||||
for name, job in jobs.items():
|
||||
if not isinstance(job, dict):
|
||||
return f"job `{name}` is not a mapping"
|
||||
if "runs-on" not in job and "uses" not in job:
|
||||
return f"job `{name}` has no `runs-on:`"
|
||||
return None
|
||||
|
||||
|
||||
def invalid_workflows(repo: str, sha: str) -> dict[str, tuple[str, str]]:
|
||||
"""{context: (path, problem)} for each workflow file at `sha` that won't run."""
|
||||
for d in WORKFLOW_DIRS:
|
||||
st, entries = gitea("GET", f"/repos/{WG_OWNER}/{repo}/contents/{d}?ref={sha}")
|
||||
if st == 404:
|
||||
continue
|
||||
if st != 200:
|
||||
raise RuntimeError(f"{repo}: Windy Git {d}@{sha[:7]} -> {st}")
|
||||
files = [e for e in entries or [] if e.get("type") == "file"
|
||||
and e["name"].endswith((".yml", ".yaml"))]
|
||||
if not files:
|
||||
continue
|
||||
bad = {}
|
||||
for e in files:
|
||||
st, f = gitea("GET", f"/repos/{WG_OWNER}/{repo}/contents/{e['path']}?ref={sha}")
|
||||
if st != 200:
|
||||
raise RuntimeError(f"{repo}: Windy Git {e['path']}@{sha[:7]} -> {st}")
|
||||
problem = workflow_problem(base64.b64decode(f["content"]).decode("utf-8", "replace"))
|
||||
if problem:
|
||||
stem = re.sub(r"\.ya?ml$", "", e["name"])
|
||||
bad[f"windy-git/{stem}/workflow"] = (e["path"], problem)
|
||||
return bad
|
||||
return {}
|
||||
|
||||
|
||||
def _call(base: str, token_header: str, method: str, path: str, body=None):
|
||||
req = urllib.request.Request(
|
||||
base + path,
|
||||
@@ -183,7 +242,8 @@ def post_statuses(repo: str, sha: str) -> None:
|
||||
ctx = f"windy-git/{r['workflow_id'].removesuffix('.yml')}/{r['name']}"
|
||||
if ctx not in latest or r["id"] > latest[ctx]["id"]:
|
||||
latest[ctx] = r
|
||||
if not latest:
|
||||
bad = invalid_workflows(repo, sha)
|
||||
if not (latest or bad):
|
||||
return
|
||||
|
||||
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
|
||||
@@ -191,6 +251,21 @@ def post_statuses(repo: str, sha: str) -> None:
|
||||
for s in existing or []: # newest first
|
||||
current.setdefault(s["context"], s["state"])
|
||||
|
||||
for ctx, (path, problem) in sorted(bad.items()):
|
||||
if current.get(ctx) == "error":
|
||||
continue
|
||||
st, _ = github(
|
||||
"POST",
|
||||
f"/repos/{GH_OWNER}/{repo}/statuses/{sha}",
|
||||
{
|
||||
"state": "error",
|
||||
"context": ctx,
|
||||
"description": f"Windy Git ignored this workflow, no CI ran: {problem}"[:140],
|
||||
"target_url": f"{PUBLIC}/{WG_OWNER}/{repo}/src/commit/{sha}/{path}",
|
||||
},
|
||||
)
|
||||
print(f" {repo}@{sha[:7]} {ctx} = error ({problem}) -> {st}")
|
||||
|
||||
for ctx, r in sorted(latest.items()):
|
||||
state = STATE.get(r["status"])
|
||||
if state is None or current.get(ctx) == state:
|
||||
|
||||
@@ -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