8 Commits

Author SHA1 Message Date
Kit OC5
44e824c762 push velocity: correlate the SSO-id subquery on the grouped column
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:10 -04:00
Kit OC5
4e2db1e2e0 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 13:06:20 -04:00
Kit OC5
54c63b4d89 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 13:06:17 -04:00
Kit OC5
ac83c4721e 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 13:06:17 -04:00
Kit OC5
2d4fadb090 sync: bound the janitor and telemetry steps (docker exec hangs in an IO stall)
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
09-23 16:43Z the Veron data2 SMR stall left runc exec in D state; the
janitor's docker exec never returned, so the sync sat 'activating' and no
repo mirrored or got a status for any lane. Both steps are non-fatal;
now they time out (120 s / 180 s) and the run continues.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:54:34 -04:00
Kit OC5
e3b69fa759 telemetry: UPDATE 7 — read the ingest body; count quarantined + dropped on heartbeats
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
The ledger answers 202 even when it quarantines rows. Both emitters now log
a warning with the reasons and report service.health.telemetry_quarantined
and telemetry_dropped (API: buffer overflow; sync: 0 by construction, since
a failed send keeps cursor + spool). HOLD until Telemetry Boss declares both
keys on windy-git's two service.health shapes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:32:28 -04:00
Kit OC5
4acf50d9ef bridge tests: fake serves workflow contents; cover invalid-workflow status
All checks were successful
check / gate (push) Successful in 23s
b7a7e94 made the bridge read workflow files, which the strict fake Gitea
refused (7 red). The fake now serves contents (404 when absent), and new
tests cover: error posted with no runs, valid files add nothing, no repost,
.gitea/workflows wins over .github/workflows, and each workflow_problem
shape. pyyaml declared in dev extras (the bridge imports it).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:31:19 -04:00
Kit OC5
b7a7e94df0 bridge: post an error status when Windy Git ignores an invalid workflow
Gitea drops an invalid workflow file with one log line and fires no run, so
the GitHub PR showed nothing and lanes waited for CI that never came
(windytalk #100). The bridge now reads each workflow file at the commit it
reports on and posts windy-git/<wf>/workflow = error with the reason.
Verified: 0 false positives on all 23 bridged repos' main; catches
windytalk #100's broken commits (invalid YAML at line 12), fix commit clean.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:29:12 -04:00
8 changed files with 423 additions and 11 deletions

View File

@@ -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."""

View File

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

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

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

View File

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

View File

@@ -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:

View File

@@ -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"

View File

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