7 Commits

Author SHA1 Message Date
Kit OC5
fe3bce39ff bridge: post pending for queued jobs a runner can take
All checks were successful
check / gate (push) Successful in 28s
Gitea 1.24 lists only picked-up jobs, so a queued PR showed NOTHING on
GitHub and lanes asked whether their push was lost (Windy Mind #131,
Windy Cloud today). The bridge now reads waiting jobs from the gitea DB
and posts pending where nothing newer was picked up; a queued re-run
supersedes the stale failure it replaces.

Only status 5 jobs whose runs-on labels a live runner has: blocked jobs
often end skipped and label-unrunnable jobs are cancelled unpicked, and
neither ever reaches /actions/tasks, so their pending would never resolve.
Lookup is bounded (30 s) and non-fatal: the IO-stall lesson.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:20:32 -04:00
Kit OC5
6b0b60eb4a scripts: rerun_ci.sh — re-fire a PR's CI without the web button
Gitea 1.24 has no rerun API and the web button needs Grant's SSO identity.
Guarded branch rewind that the next sync undoes; restores the branch itself
on timeout. Used today for eternitas #166 and windy-mind #131 after the
Veron IO stall.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:17:10 -04:00
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
8 changed files with 543 additions and 5 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

@@ -79,10 +79,11 @@ class Fake:
@pytest.fixture
def fake(monkeypatch):
def make(**kw):
def make(queued=(), **kw):
f = Fake(**kw)
monkeypatch.setattr(bridge, "gitea", f.gitea)
monkeypatch.setattr(bridge, "github", f.github)
monkeypatch.setattr(bridge, "queued_jobs", lambda repo, sha: list(queued))
return f
return make
@@ -272,3 +273,70 @@ def test_gitea_dir_wins_over_github_dir(fake):
)
def test_workflow_problem(text, problem):
assert bridge.workflow_problem(text) == problem
def _q(n, wf, job):
return {"run_number": n, "workflow_id": wf, "name": job}
def test_queued_job_shows_pending_instead_of_nothing(fake):
f = fake(queued=[_q(5, "ci.yml", "test")])
bridge.post_statuses("windy-chat", SHA)
assert [(p["context"], p["state"]) for p in f.posted] == [("windy-git/ci/test", "pending")]
assert f.posted[0]["target_url"].endswith("/actions/runs/5")
def test_queued_rerun_supersedes_the_stale_failure(fake):
f = fake(runs=[_run(1, "ci.yml", "test", "failure", n=4)], queued=[_q(7, "ci.yml", "test")])
bridge.post_statuses("windy-chat", SHA)
assert [(p["context"], p["state"]) for p in f.posted] == [("windy-git/ci/test", "pending")]
def test_older_queued_job_never_overrides_a_newer_verdict(fake):
f = fake(runs=[_run(1, "ci.yml", "test", "success", n=9)], queued=[_q(3, "ci.yml", "test")])
bridge.post_statuses("windy-chat", SHA)
assert [(p["context"], p["state"]) for p in f.posted] == [("windy-git/ci/test", "success")]
def test_queued_docker_and_non_blocking_jobs_stay_unposted(fake):
f = fake(queued=[_q(2, "ci.yml", "docker-build"), _q(2, "ci.yml", "build-desktop")])
bridge.post_statuses("windy-pro", SHA)
assert f.posted == []
def test_queued_lookup_refuses_unsafe_input():
assert bridge.queued_jobs("x'; drop table t;--", SHA) == []
assert bridge.queued_jobs("windy-chat", "not-a-sha") == []
def _db(monkeypatch, jobs, labels):
import json as _json
import subprocess as _sp
payload = _json.dumps({"jobs": jobs, "labels": [_json.dumps(x) for x in labels]})
monkeypatch.setattr(
bridge.subprocess, "run",
lambda *a, **k: _sp.CompletedProcess(a, 0, stdout=payload, stderr=""),
)
RUNNER = ["veron-1", "linux-x64", "self-hosted", "linux", "x64"]
def test_only_jobs_a_runner_can_take_are_pending(monkeypatch):
# macos-latest is cancelled unpicked by the janitor: pending would never resolve.
_db(monkeypatch, [
{"run_number": 3, "workflow_id": "ci.yml", "name": "test", "runs_on": '["self-hosted","linux","x64"]'},
{"run_number": 3, "workflow_id": "ci.yml", "name": "mac", "runs_on": '["macos-latest"]'},
], [RUNNER])
assert [j["name"] for j in bridge.queued_jobs("windy-chat", SHA)] == ["test"]
def test_lookup_failure_is_non_fatal(monkeypatch):
import subprocess as _sp
def boom(*a, **k):
raise _sp.TimeoutExpired("docker", 30)
monkeypatch.setattr(bridge.subprocess, "run", boom)
assert bridge.queued_jobs("windy-chat", SHA) == []

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

@@ -152,3 +152,16 @@ disqualifying the moment a stranger depends on it. **The trigger is not a date
it is the first external push.** Move the control plane to a dedicated VPS (not
Kit 0), keep Veron 1 as the runner. It is an rsync, a Postgres dump and three
DNS record edits.
## Re-run a PR's CI (Gitea 1.24 has no rerun API)
```bash
ssh wg-veron
cd /srv/windygit/src && bash scripts/rerun_ci.sh <repo> <branch> <github-head-sha-prefix>
```
Moves the Windy Git branch back one commit; the next sync force-pushes the
GitHub head again and Gitea re-fires every workflow for that event on the same
commit. Guarded: refuses unless the branch is at the given sha, waits for a
sync that starts AFTER the rewind, restores the branch itself on timeout.
Don't use the web "Re-run" button: it needs a hub-SSO session as windyadmin,
which is Grant's identity.

View File

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

View File

@@ -32,6 +32,7 @@ import base64
import json
import os
import re
import subprocess
import sys
import time
import urllib.error
@@ -222,6 +223,54 @@ def sync_prs(repo: str) -> list[str]:
return heads
SAFE_NAME = re.compile(r"^[A-Za-z0-9._-]+$")
SAFE_SHA = re.compile(r"^[0-9a-f]{40}$")
def queued_jobs(repo: str, sha: str) -> list[dict]:
"""Jobs at `sha` that are waiting for a runner and that a runner CAN take.
Gitea 1.24's API lists only PICKED-UP jobs (/actions/tasks), so a queued PR
showed nothing on GitHub and people asked whether the push was lost. The
truth is in the gitea DB. Bounded + non-fatal: during the 09-23 IO stall
`docker exec` hung for an hour and must never wedge the bridge again.
Only status 5 (waiting) with labels some live runner has. A `pending` we
post must end in a verdict we will also see, or it sits yellow forever:
blocked jobs (7) often end SKIPPED, and jobs for labels no runner has
(macos-/windows-/ubuntu-latest) are cancelled by the janitor unpicked;
neither ever appears in /actions/tasks.
"""
if not (SAFE_NAME.match(repo) and SAFE_NAME.match(WG_OWNER) and SAFE_SHA.match(sha)):
return []
query = (
"select json_build_object("
" 'jobs', (select coalesce(json_agg(t), '[]'::json) from ("
" select ar.index as run_number, ar.workflow_id, j.name, j.runs_on"
" from action_run_job j join action_run ar on ar.id = j.run_id"
" join repository r on r.id = j.repo_id join \"user\" o on o.id = r.owner_id"
f" where o.lower_name = '{WG_OWNER.lower()}' and r.lower_name = '{repo.lower()}'"
f" and ar.commit_sha = '{sha}' and j.status = 5) t),"
" 'labels', (select coalesce(json_agg(agent_labels), '[]'::json)"
" from action_runner where coalesce(deleted, 0) = 0));"
)
try:
out = subprocess.run(
["docker", "exec", "-i", "windy-git-db-1", "sh", "-c",
'psql -U "$POSTGRES_USER" -d gitea -At -v ON_ERROR_STOP=1'],
input=query, capture_output=True, text=True, check=True, timeout=30,
).stdout.strip()
got = json.loads(out or "{}")
runners = [set(json.loads(x or "[]")) for x in got.get("labels") or []]
return [
j for j in got.get("jobs") or []
if any(set(json.loads(j.get("runs_on") or "[]")) <= r for r in runners)
]
except (subprocess.SubprocessError, OSError, ValueError) as e:
print(f" {repo}: queued-job lookup skipped ({type(e).__name__})")
return []
def post_statuses(repo: str, sha: str) -> None:
# Gitea caps a page at 50 (MAX_RESPONSE_ITEMS) whatever `limit` says, and a
# daily scheduled workflow can push a quiet main's runs off page 1.
@@ -242,6 +291,17 @@ 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
# Queued jobs: `pending` where nothing newer has been picked up. A re-run
# queued behind an old failure must read pending, not the stale red.
for q in queued_jobs(repo, sha):
if NO_DAEMON_JOB.search(q["name"]):
continue
wf = q["workflow_id"].removesuffix(".yml")
if f"{wf}/{q['name']}" in NON_BLOCKING.get(repo, ()):
continue
ctx = f"windy-git/{wf}/{q['name']}"
if ctx not in latest or q["run_number"] > latest[ctx]["run_number"]:
latest[ctx] = {"id": 0, "status": "waiting", "run_number": q["run_number"]}
bad = invalid_workflows(repo, sha)
if not (latest or bad):
return

45
scripts/rerun_ci.sh Executable file
View File

@@ -0,0 +1,45 @@
#!/usr/bin/env bash
# Re-run a PR's (or branch's) CI on Windy Git. Runs ON Veron 1.
#
# bash scripts/rerun_ci.sh <repo> <branch> <sha-prefix>
#
# Gitea 1.24 has NO rerun API; the web button needs a hub-SSO session as
# windyadmin, which is Grant's identity, so we don't use it. Instead: move the
# Windy Git branch back one commit, let the next sync force-push the GitHub head
# again, and Gitea fires an ordinary push / pull_request_sync event on the SAME
# commit. Every workflow on that event re-runs, not only the failed one.
#
# Safety: refuses unless the branch is exactly at <sha-prefix> (GitHub's head),
# never rewinds while a sync is running (a run already past this repo would
# not push it back), waits for a sync that STARTS after the rewind, and if the
# branch is not verifiably back at <sha-prefix> by the deadline, restores it
# itself, so Windy Git is never left behind GitHub.
set -euo pipefail
repo="${1:?repo}"; branch="${2:?branch}"; want="${3:?sha prefix}"
G="sudo docker exec -u git windy-git-gitea-1 git -C /data/git/repositories/windyadmin/${repo}.git"
head=$($G rev-parse "refs/heads/${branch}")
[[ "$head" == "$want"* ]] || { echo "refusing: ${branch} is at ${head:0:7}, not ${want}"; exit 1; }
parent=$($G rev-parse "${head}^")
while systemctl is-active -q windygit-sync; do sleep 5; done
$G update-ref "refs/heads/${branch}" "$parent" "$head"
mark=$(awk '{print int($1*1000000)}' /proc/uptime)
echo "rewound ${repo}:${branch} ${head:0:7} -> ${parent:0:7}"
deadline=$(( $(date +%s) + 900 ))
until [ "$(systemctl show windygit-sync -p ExecMainStartTimestampMonotonic --value)" -gt "$mark" ] \
&& [ "$($G rev-parse "refs/heads/${branch}")" = "$head" ]; do
if [ "$(date +%s)" -ge "$deadline" ]; then
$G update-ref "refs/heads/${branch}" "$head" "$($G rev-parse "refs/heads/${branch}")" || true
echo "TIMEOUT: restored ${branch} to ${head:0:7} by hand; NO new run fired"; exit 1
fi
sleep 10
done
echo "restored by sync: ${branch} = ${head:0:7}"
sleep 5
~/bin/wg-q <<SQL
select ar.index, ar.workflow_id, ar.event, ar.status, to_char(to_timestamp(ar.created),'HH24:MI:SS')
from action_run ar join repository r on r.id = ar.repo_id
where r.name = '${repo}' and ar.commit_sha = '${head}' order by ar.id desc limit 6;
SQL

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