16 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
c83f808a60 bridge: retry transport blips (TLS timeout/reset), never HTTP errors
All checks were successful
check / gate (push) Successful in 20s
canary / probe (push) Successful in 9s
A single GitHub TLS handshake timeout failed the whole sync, flipped its
windy-job heartbeat to ok:false and would page for nothing. Up to 3
attempts with backoff for URLError/timeout/reset; HTTP errors return
immediately as before. Test covers both.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:15:22 -04:00
eb27e63db3 telemetry: adopt the end-to-end synthetic convention (UPDATE 4)
All checks were successful
check / gate (push) Successful in 24s
canary / probe (push) Successful in 9s
Replaces the keyed marker from 1c3b5b0 with the ecosystem convention:
any X-Windy-Synthetic value marks the request synthetic; the flag lives in
a per-request contextvar, labels this request's rows, and is FORWARDED on
downstream calls (Eternitas trust lookup, Gitea API). The canary sends
"1". Rows are still recorded; the label separates, never suppresses.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 12:06:51 -04:00
1c3b5b0638 telemetry: synthetic:true on canary refusals (keyed, not a bare flag)
All checks were successful
check / gate (push) Successful in 25s
canary / probe (push) Successful in 7s
The canary deliberately sends forged tokens every 10 min; those refusal
rows read as attacks. It now sends X-Windy-Synthetic carrying a shared
secret (Gitea repo secret CANARY_SYNTHETIC_KEY = WINDYGIT_SYNTHETIC_KEY in
Veron .env); the API marks the row synthetic only on a constant-time
match, so an attacker cannot label their own refusals synthetic to hide.
synthetic is declared on forge.auth.failed (Telemetry Boss, UPDATE 3).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:54:56 -04:00
90643fe48e telemetry step 2: API boot/health + forge.auth.failed (declared)
All checks were successful
check / gate (push) Successful in 25s
canary / probe (push) Successful in 6s
Membrane first: I-2 and MEMBRANE.v1 now list the windy-admin ledger
(POST /v1/events). api/app/telemetry.py: service.boot once per start
(commit_sha omitted when unknown, I-12), an hourly in-process
service.health with the shared keys (requests, errors_5xx/4xx,
refusals_4xx, p95_ms only when there was traffic), and one
forge.auth.failed row per refused request: declared 13-code enum,
http_status, caller class, route TEMPLATE (never the concrete path),
actor_type system with no actor_id (all-lanes rule). No token = nothing
sent or buffered; flush failures keep rows (bounded) and never raise.
Token from root-only /etc/windygit/telemetry.env (optional env_file).
8 behavioural tests.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:47:29 -04:00
00ec963f82 telemetry: fix ci.run completeness — cursor on (finish time, job id)
Some checks failed
check / gate (push) Has been cancelled
Telemetry Boss found jobs_finished=43 vs 8 ci.run rows. Root cause: the
high-water mark was the job id, but jobs FINISH out of id order, so every
long job that started before the mark and finished after it was silently
never emitted. Now a (finish time, id) cursor; finish = stopped, or
updated for skipped jobs with no stop time. Heartbeat finished/failed/
cancelled counts are derived from exactly the rows emitted, so
sum(jobs_finished) == count(ci.run) by construction (dry run on real
data: 97 == 97, failed 2 == 2, cancelled 13 == 13). posted_to_github now
set from the bridge's own rules. duration_ms = Gitea whole seconds x 1000.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:42:36 -04:00
b5e4eaf57a telemetry: ci.job_cancelled from the janitor; interval_s on heartbeat
All checks were successful
check / gate (push) Successful in 22s
canary / probe (push) Successful in 6s
The janitor now returns one JSON line per job it cancels (repo, workflow,
job, reason, runs_on, waited_s) into a spool; the emitter ships them as
ci.job_cancelled (declared with Telemetry Boss) and truncates the spool
only after a 2xx. Run status recompute folded into the same statement.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:27:53 -04:00
baaa542bae docs: audit disposition 09-23 — R2 god token replaced by bucket-scoped token
All checks were successful
check / gate (push) Successful in 43s
canary / probe (push) Successful in 7s
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:15:35 -04:00
0634a6cb1b style: ruff fix in telemetry_emit
All checks were successful
check / gate (push) Successful in 29s
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:08:21 -04:00
19 changed files with 1078 additions and 54 deletions

View File

@@ -81,7 +81,7 @@ Numbered because code cites them. Changing one requires an ADR that names it.
1. **I-1 · Gitea is a component, never a merged tree.** Our code lives in our services and calls Gitea's REST API. Any patch to Gitea source lives in `patches/` as a numbered, rebasable diff with a one-line justification, and `make check` fails if `patches/` grows past **3** files without an ADR.
2. **I-2 · The membrane is ENUMERATED.**
**Calls out:** `windy-cloud` kernel `GET /api/v1/storage/objects` + `HEAD` (read user objects to version them) · `windy-cloud` `POST /api/v1/storage/quota/check` · `eternitas` `GET /api/v1/trust/{passport}` (band + allowed_actions) · `eternitas` `GET /api/v1/registry/{passport}/integrity` · `account-server` OIDC discovery + JWKS · `windy-cloud-sites` `POST /api/v1/sites/{id}/versions` (publish docs from a repo).
**Calls out:** `windy-cloud` kernel `GET /api/v1/storage/objects` + `HEAD` (read user objects to version them) · `windy-cloud` `POST /api/v1/storage/quota/check` · `eternitas` `GET /api/v1/trust/{passport}` (band + allowed_actions) · `eternitas` `GET /api/v1/registry/{passport}/integrity` · `account-server` OIDC discovery + JWKS · `windy-cloud-sites` `POST /api/v1/sites/{id}/versions` (publish docs from a repo). · windy-admin ledger `POST https://admin.windyword.ai/v1/events` (field telemetry, 2026-09-23: `ci.run`, `ci.job_cancelled`, `service.boot`, `service.health`, `forge.auth.failed` — shapes declared with the ledger owner first; no content, no passports, no emails)
**Calls in:** `POST /internal/repo-from-folder` (Cloud portal: git-enable a folder) · `POST /internal/mirror-status` (ops).
**Events out:** `repo.created`, `repo.pushed`, `release.published`, `model.published`, `ci.completed`.
**Events in:** `passport.revoked` (fail-closed), `storage.quota.exceeded`, `identity.created`.

View File

@@ -28,6 +28,7 @@ from api.app.config import Settings
from api.app.ept import EptInvalid, looks_like_ept, verify_ept
from api.app.errors import RepairPointer, passport_unresolvable
from api.app.hub_jwt import HubTokenInvalid, verify_hub_token
from api.app.telemetry import synthetic_headers
log = logging.getLogger(__name__)
@@ -125,7 +126,7 @@ async def resolve_passport(settings: Settings, passport: str) -> tuple[str, tupl
)
url = f"{settings.eternitas_base_url}/api/v1/trust/{passport}"
headers = {"X-API-Key": settings.eternitas_platform_api_key}
headers = {"X-API-Key": settings.eternitas_platform_api_key, **synthetic_headers()}
last_status = 0
for attempt in range(3):
async with httpx.AsyncClient(timeout=httpx.Timeout(8.0, connect=3.0)) as client:

View File

@@ -66,6 +66,12 @@ class Settings(BaseSettings):
# Flip to True once the hub emits aud on every access token.
hub_require_aud: bool = False
# ---- field telemetry (admin.windyword.ai ledger) ----------------------
# Unset token = nothing sent, nothing buffered. The token lives in the
# root-only /etc/windygit/telemetry.env on Veron, never in the repo.
windygit_telemetry_token: str = ""
telemetry_ingest_url: str = "https://admin.windyword.ai/v1/events"
# Internal callers (the Cloud portal calling /internal/*). A first-class
# caller class, not a bypass: unset means service calls are REFUSED.
service_token: str = ""

View File

@@ -6,11 +6,13 @@ component and is reached only over its REST API (D-2 / I-1).
from __future__ import annotations
import asyncio
import logging
import socket
import time
from contextlib import asynccontextmanager
from fastapi import FastAPI
from fastapi import FastAPI, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import JSONResponse
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
@@ -25,6 +27,7 @@ from api.app.providers.registry import (
R2Provider,
)
from api.app.routes import health, repos, webhooks
from api.app.telemetry import SYNTHETIC, Telemetry, caller_class, is_synthetic
logging.basicConfig(
level=logging.INFO,
@@ -44,9 +47,7 @@ def _refuse_kit_zero(settings) -> None:
if not settings.is_production:
return
try:
local_ips = {
info[4][0] for info in socket.getaddrinfo(socket.gethostname(), None)
}
local_ips = {info[4][0] for info in socket.getaddrinfo(socket.gethostname(), None)}
except socket.gaierror:
return
if settings.kit0_host in local_ips:
@@ -98,8 +99,23 @@ async def lifespan(app: FastAPI):
# systemd Restart=always, plus the runbook's `systemctl status`.
]
telemetry = Telemetry(
settings.telemetry_ingest_url,
settings.windygit_telemetry_token,
environment=settings.environment,
commit_sha=info.commit_sha,
version=info.version,
)
app.state.telemetry = telemetry
telemetry.boot()
await telemetry.flush()
task = asyncio.create_task(telemetry.run()) if telemetry.enabled else None
yield
if task is not None:
task.cancel()
await telemetry.flush()
if engine is not None:
await engine.dispose()
@@ -119,8 +135,44 @@ app.include_router(repos.router)
app.include_router(webhooks.router)
@app.middleware("http")
async def _count_requests(request: Request, call_next):
"""Heartbeat counts (requests, 4xx/5xx, refusals, p95). Never raises."""
start = time.perf_counter()
marker = SYNTHETIC.set(is_synthetic(request.headers))
try:
response = await call_next(request)
finally:
SYNTHETIC.reset(marker)
tel = getattr(request.app.state, "telemetry", None)
if tel is not None:
tel.record_request(
response.status_code,
(time.perf_counter() - start) * 1000,
refused=getattr(request.state, "refused", False),
)
return response
@app.exception_handler(RepairPointer)
async def _repair_pointer_handler(_, exc: RepairPointer) -> JSONResponse:
async def _repair_pointer_handler(request: Request, exc: RepairPointer) -> JSONResponse:
tel = getattr(request.app.state, "telemetry", None)
detail = exc.detail if isinstance(exc.detail, dict) else {}
code = detail.get("code")
if tel is not None and code in tel.auth_codes:
# A refusal is a failure row (field-visibility rule 1). The caller is
# unauthenticated by definition, so: system actor, no actor_id, and
# the route TEMPLATE, never the concrete path.
request.state.refused = True
route = request.scope.get("route")
tel.auth_failed(
code=code,
http_status=exc.status_code,
caller=caller_class(request.headers),
route=getattr(route, "path", None),
upstream_status=getattr(exc, "upstream_status", None),
synthetic=is_synthetic(request.headers),
)
return JSONResponse(status_code=exc.status_code, content=exc.detail)

View File

@@ -20,6 +20,7 @@ import httpx
from api.app.config import Settings
from api.app.errors import RepairPointer, provider_unconfigured
from api.app.telemetry import synthetic_headers
_TIMEOUT = httpx.Timeout(20.0, connect=5.0)
@@ -36,6 +37,7 @@ class GiteaClient:
return {
"Authorization": f"token {self._s.gitea_admin_token}",
"Content-Type": "application/json",
**synthetic_headers(), # end-to-end synthetic convention (Telemetry UPDATE 4)
}
async def _request(self, method: str, path: str, **kw: Any) -> httpx.Response:

255
api/app/telemetry.py Normal file
View File

@@ -0,0 +1,255 @@
"""Field telemetry to the admin ledger (admin.windyword.ai) — step 2, 2026-09-23.
Shapes are DECLARED with Telemetry Boss (the ledger owner); the server
quarantines any row that doesn't match, so never add a key or a code here
without re-declaring it first:
service.boot once per process start {commit_sha, version, environment}
service.health hourly, in-process interval_s, uptime_s, requests,
errors_5xx, errors_4xx,
refusals_4xx, p95_ms
forge.auth.failed every refused request {code, http_status, caller,
route?, upstream_status?}
Refusals come first: a refused caller is the most expensive silent failure
("the button did nothing"). An UNAUTHENTICATED caller has no trustworthy id, so
per the all-lanes actor rule the row is actor_type "system", no actor_id, and
the caller class goes in metadata.caller.
Privacy: codes, statuses, route TEMPLATES, counts, durations. Never a passport
number, an email, a token fragment or a concrete path with names in it.
No token → nothing is sent and nothing is buffered. A failed flush keeps the
rows (bounded) and retries on the next tick; it never raises into a request.
"""
from __future__ import annotations
import asyncio
import contextvars
import json
import logging
import time
import urllib.request
from datetime import UTC, datetime
log = logging.getLogger("windy-git.telemetry")
PLATFORM, SERVICE = "windy-git", "api"
FLUSH_EVERY_S = 60
HEALTH_EVERY_S = 3600
MAX_BUFFER = 5000
MAX_LATENCY_SAMPLES = 20000
# The declared forge.auth.failed code enum (Telemetry Boss, 2026-09-23). A code
# outside this set is NOT a refusal row — it counts in errors_* instead.
AUTH_CODES = frozenset(
{
"not_signed_in",
"token_invalid",
"token_unrecognised",
"ept_invalid",
"passport_revoked",
"passport_unresolvable",
"agent_read_only",
"agent_rate_limited",
"quota_exceeded",
"trust_unavailable",
"throttle_unavailable",
"service_token_invalid",
"service_auth_unconfigured",
}
)
def _iso(epoch: float) -> str:
return datetime.fromtimestamp(epoch, UTC).isoformat().replace("+00:00", "Z")
# Ecosystem convention (Telemetry UPDATE 4): synthetic traffic travels END TO
# END. Originators (canaries, probes, journeys) send `X-Windy-Synthetic: 1`;
# every service marks all of that request's rows synthetic:true AND forwards the
# header on every downstream call. Absent = real. Never strip it, never set it
# on real traffic. The label separates rows — it never suppresses them.
SYNTHETIC: contextvars.ContextVar[bool] = contextvars.ContextVar("windy_synthetic", default=False)
def is_synthetic(headers) -> bool:
return bool((headers.get("x-windy-synthetic") or "").strip())
def synthetic_headers() -> dict:
"""Merge into every downstream request made while serving this one."""
return {"X-Windy-Synthetic": "1"} if SYNTHETIC.get() else {}
def caller_class(headers) -> str:
"""Declared values: anonymous_human | anonymous_agent | unknown."""
from api.app.ept import looks_like_ept
if headers.get("x-service-token"):
return "unknown"
auth = headers.get("authorization") or ""
if not auth.lower().startswith("bearer "):
return "unknown"
return "anonymous_agent" if looks_like_ept(auth.split(" ", 1)[1].strip()) else "anonymous_human"
class Telemetry:
def __init__(
self,
url: str,
token: str,
*,
environment: str = "",
commit_sha: str | None = None,
version: str = "",
) -> None:
self.url, self.token = url, token
self.environment, self.commit_sha, self.version = environment, commit_sha, version
self.started = time.time()
self.buffer: list[dict] = []
self.auth_codes = AUTH_CODES
self._reset_window()
@property
def enabled(self) -> bool:
return bool(self.token)
def _reset_window(self) -> None:
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:
self.requests += 1
if status >= 500:
self.errors_5xx += 1
elif refused:
self.refusals_4xx += 1
elif status >= 400:
self.errors_4xx += 1
if len(self.latencies_ms) < MAX_LATENCY_SAMPLES:
self.latencies_ms.append(duration_ms)
def _event(self, event_type: str, metadata: dict, *, ts: float | None = None) -> None:
if not self.enabled:
return
self.buffer.append(
{
"ts": _iso(ts or time.time()),
"platform": PLATFORM,
"service": SERVICE,
"event_type": event_type,
"actor_type": "system",
"metadata": metadata,
}
)
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:
meta = {"version": self.version, "environment": self.environment}
if self.commit_sha: # unknown is absent, never invented (I-12)
meta["commit_sha"] = self.commit_sha
self._event("service.boot", meta, ts=self.started)
def auth_failed(
self,
*,
code: str,
http_status: int,
caller: str,
route: str | None = None,
upstream_status: int | None = None,
synthetic: bool = False,
) -> None:
if code not in AUTH_CODES:
return
meta: dict = {
"code": code,
"http_status": int(http_status),
"caller": caller,
"synthetic": bool(synthetic),
}
if route:
meta["route"] = route
if upstream_status is not None:
meta["upstream_status"] = int(upstream_status)
self._event("forge.auth.failed", meta)
def health_row(self) -> dict:
now = time.time()
meta = {
"interval_s": int(now - self.window_start),
"uptime_s": int(now - self.started),
"requests": self.requests,
"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)
meta["p95_ms"] = int(s[min(len(s) - 1, int(0.95 * (len(s) - 1) + 0.5))])
return meta
def health(self) -> None:
self._event("service.health", self.health_row())
self._reset_window()
# ---- sending ------------------------------------------------------------
def _post(self, batch: list[dict]) -> tuple[int, dict]:
req = urllib.request.Request(
self.url,
data=json.dumps({"events": batch}).encode(),
method="POST",
headers={
"Authorization": f"Bearer {self.token}",
"Content-Type": "application/json",
"User-Agent": "windy-git-api-telemetry/1",
},
)
with urllib.request.urlopen(req, timeout=20) as r:
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, 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."""
last_health = time.monotonic()
while True:
await asyncio.sleep(FLUSH_EVERY_S)
if time.monotonic() - last_health >= HEALTH_EVERY_S:
self.health()
last_health = time.monotonic()
await self.flush()

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]}
@@ -171,3 +183,92 @@ def test_default_non_blocking_is_grants_ruling():
"ci/test-installer",
"ci/reality-check",
}
def test_transport_blips_are_retried_but_http_errors_are_not(monkeypatch):
import urllib.error
calls = {"n": 0}
class _R:
status = 200
def read(self):
return b"{}"
def __enter__(self):
return self
def __exit__(self, *a):
return False
def flaky(req, timeout):
calls["n"] += 1
if calls["n"] < 3:
raise urllib.error.URLError("_ssl.c:983: The handshake operation timed out")
return _R()
monkeypatch.setattr(bridge.urllib.request, "urlopen", flaky)
monkeypatch.setattr(bridge.time, "sleep", lambda s: None)
assert bridge._call("http://x", "t", "GET", "/p") == (200, {})
assert calls["n"] == 3
def forbidden(req, timeout):
calls["n"] += 1
raise urllib.error.HTTPError("http://x/p", 403, "no", {}, None)
calls["n"] = 0
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

202
api/tests/test_telemetry.py Normal file
View File

@@ -0,0 +1,202 @@
"""Field telemetry from the API (step 2): behavioural, no network.
A refusal must become exactly one forge.auth.failed row in the DECLARED shape
(the ledger quarantines anything else); non-refusal errors must not; the
heartbeat must count what happened and never invent a p95 for no traffic.
"""
from __future__ import annotations
import httpx
import pytest
from fastapi import Depends, FastAPI
from api.app import telemetry as tmod
from api.app.errors import RepairPointer
DECLARED_AUTH_KEYS = {"code", "http_status", "caller", "route", "upstream_status", "synthetic"}
def _app(tel: tmod.Telemetry) -> FastAPI:
"""The real middleware + handler, re-registered on a bare app (no DB)."""
from api.app import main
app = FastAPI()
app.state.telemetry = tel
app.middleware("http")(main._count_requests)
app.exception_handler(RepairPointer)(main._repair_pointer_handler)
def refuse(code: str, status: int):
def dep():
raise RepairPointer(
status_code=status,
code=code,
speak="no",
machine_cause="test",
remediation_tool=None,
)
return dep
@app.get("/api/v1/repos/{repo}/grants", dependencies=[Depends(refuse("passport_revoked", 403))])
async def grants(repo: str):
return {}
@app.get("/api/v1/nope", dependencies=[Depends(refuse("repo_not_found", 404))])
async def nope():
return {}
@app.get("/ok")
async def ok():
return {"ok": True}
return app
async def _get(app, path, headers=None):
async with httpx.AsyncClient(transport=httpx.ASGITransport(app=app), base_url="http://t") as c:
return await c.get(path, headers=headers or {})
def _tel():
return tmod.Telemetry(
"http://ledger.invalid/v1/events",
"tok",
environment="test",
commit_sha="abc1234",
version="0.1.0",
)
@pytest.mark.asyncio
async def test_refusal_emits_one_declared_row_with_route_template_not_path():
tel = _tel()
r = await _get(
_app(tel),
"/api/v1/repos/grandmas-secret-project/grants",
{"Authorization": "Bearer eyJhbGciOiJFUzI1NiIsInR5cCI6IkVQVCJ9.e30.x"},
)
assert r.status_code == 403
rows = [e for e in tel.buffer if e["event_type"] == "forge.auth.failed"]
assert len(rows) == 1
row = rows[0]
assert row["actor_type"] == "system" and "actor_id" not in row
assert set(row["metadata"]) <= DECLARED_AUTH_KEYS
assert row["metadata"]["code"] == "passport_revoked"
assert row["metadata"]["http_status"] == 403
assert row["metadata"]["caller"] == "anonymous_agent"
assert row["metadata"]["route"] == "/api/v1/repos/{repo}/grants"
assert "grandmas-secret-project" not in str(row)
@pytest.mark.asyncio
async def test_non_auth_errors_are_counted_but_not_refusal_rows():
tel = _tel()
await _get(_app(tel), "/api/v1/nope")
assert not [e for e in tel.buffer if e["event_type"] == "forge.auth.failed"]
assert tel.errors_4xx == 1 and tel.refusals_4xx == 0
@pytest.mark.asyncio
async def test_heartbeat_counts_requests_refusals_and_p95():
tel = _tel()
app = _app(tel)
for _ in range(3):
await _get(app, "/ok")
await _get(app, "/api/v1/repos/x/grants")
meta = tel.health_row()
assert meta["requests"] == 4 and meta["refusals_4xx"] == 1 and meta["errors_5xx"] == 0
assert isinstance(meta["p95_ms"], int)
assert {"interval_s", "uptime_s"} <= set(meta)
def test_no_traffic_means_no_p95_not_a_fake_zero():
assert "p95_ms" not in _tel().health_row()
def test_no_token_sends_and_buffers_nothing():
tel = tmod.Telemetry("http://ledger.invalid", "")
tel.boot()
tel.auth_failed(code="token_invalid", http_status=401, caller="unknown")
assert tel.buffer == []
def test_unknown_code_is_never_sent_as_a_refusal():
tel = _tel()
tel.auth_failed(code="made_up_code", http_status=401, caller="unknown")
assert tel.buffer == []
def test_boot_omits_an_unknown_commit_rather_than_inventing_one():
tel = tmod.Telemetry("http://x", "tok", commit_sha=None, version="0.1.0")
tel.boot()
assert "commit_sha" not in tel.buffer[0]["metadata"]
def test_caller_classes_are_the_declared_three():
assert tmod.caller_class({}) == "unknown"
assert tmod.caller_class({"x-service-token": "s"}) == "unknown"
assert (
tmod.caller_class({"authorization": "Bearer eyJhbGciOiJSUzI1NiJ9.e30.x"})
== "anonymous_human"
)
@pytest.mark.asyncio
async def test_synthetic_header_marks_the_row_and_absent_means_real():
async def refusal(headers):
tel = _tel()
await _get(_app(tel), "/api/v1/repos/x/grants", headers)
return [e for e in tel.buffer if e["event_type"] == "forge.auth.failed"][0]["metadata"]["synthetic"]
assert await refusal({"X-Windy-Synthetic": "1"}) is True
assert await refusal({}) is False
def test_synthetic_is_forwarded_downstream_only_for_synthetic_requests():
token = tmod.SYNTHETIC.set(True)
try:
assert tmod.synthetic_headers() == {"X-Windy-Synthetic": "1"}
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

@@ -16,7 +16,11 @@ services:
# I-12: baked at build time. A runtime COMMIT_SHA override is ignored.
COMMIT_SHA: ${COMMIT_SHA_BUILD:-}
BUILT_AT: ${BUILT_AT:-}
env_file: [.env]
env_file:
- .env
# WINDYGIT_TELEMETRY_TOKEN (root-only on Veron). Optional: no file = no telemetry.
- path: /etc/windygit/telemetry.env
required: false
environment:
DATABASE_URL: postgresql+asyncpg://windygit:${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD}@db:5432/windygit
GITEA_BASE_URL: http://gitea:3000

View File

@@ -93,3 +93,24 @@ mattered most, *is an agent really that agent*, shipped inverted and untested.
The remedy is not more process; it is **behavioral tests and canary probes for
the security-critical paths**, so verification persists instead of living in a
transcript.
## Disposition update — 2026-09-23 (lane 13)
**Privileged dind beside the tokens — materially reduced, not closed.**
- The account-wide R2 token is **gone from Veron**. `.env` now carries a token scoped
to Workers R2 Bucket Item Read/Write on `windy-git-lfs` + `windy-git-backups` only,
minted by API (verified: works on both buckets, refused on any other). A CI escape
now reaches Windy Git's own two buckets, not every bucket and zone in the account.
- Runners take jobs **only from windyadmin-owned repos** (`action_runner.owner_id`), and
forge self-registration is off, so no stranger's workflow can run here.
- A host egress filter (`deploy/runner/egress.sh`) stops job containers reaching Veron,
the LAN, WireGuard or Tailscale.
- Still open: dind runs `--privileged` (next: Sysbox); the host still holds a GitHub
token and the Gitea admin token.
**Revocation / webhook secret** — `ETERNITAS_WEBHOOK_SECRET` recovered from the
Eternitas platform row and set; signed deliveries verify.
**Tests are string asserts** — the security paths now have behavioural suites
(`test_hub_jwt.py`, `test_webhooks_behavior.py`, `test_pr_status_bridge.py`);
a mutation check showed the old grep invariant passing a broken HMAC prefix strip.

View File

@@ -15,6 +15,7 @@ Mirrored into `windy-cloud` and `eternitas` on change.
| eternitas | `GET /api/v1/trust/{passport}` | band + allowed_actions |
| eternitas | `GET /api/v1/registry/{passport}/integrity` | ⚠️ note the path — `windy-registry` calls `/api/v1/passports/{p}/status`, which 404s, which is why the integrity index has never been populated |
| account-server | OIDC discovery + JWKS | human identity (G3.1) |
| windy-admin ledger | `POST /v1/events` (admin.windyword.ai) | field telemetry: `ci.run`, `ci.job_cancelled`, `service.boot`, `service.health`, `forge.auth.failed`. Shapes are declared with the ledger owner BEFORE shipping (the server quarantines undeclared keys). Codes, counts, route templates only |
| windy-cloud-sites | `POST /api/v1/sites/{id}/versions` | publish docs from a repo |
## Calls IN

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

@@ -73,6 +73,9 @@ class Check:
def _probe(c: Check) -> Result:
data = json.dumps(c.body).encode() if c.body else None
headers = {"User-Agent": "windy-git-canary/1.0", **c.headers}
# Our own probes are synthetic traffic (ecosystem convention, Telemetry
# UPDATE 4): every service they touch labels the resulting rows.
headers["X-Windy-Synthetic"] = "1"
if data:
headers["Content-Type"] = "application/json"
req = urllib.request.Request(c.url, data=data, method=c.method, headers=headers)

View File

@@ -1,6 +1,10 @@
#!/usr/bin/env bash
# Cancel jobs no runner can ever take (see cancel_unrunnable.sql). Run on Veron as root.
set -euo pipefail
n=$(docker exec -i windy-git-db-1 sh -c 'psql -U "$POSTGRES_USER" -d gitea -At -v ON_ERROR_STOP=1' \
< "$(dirname "$0")/cancel_unrunnable.sql" | grep -cE '^[0-9]+$' || true)
echo "[janitor] cancelled unrunnable jobs in ${n} run(s)"
SPOOL="${JANITOR_SPOOL:-/var/lib/windy-git/janitor-cancelled.jsonl}"
mkdir -p "$(dirname "$SPOOL")"
out=$(docker exec -i windy-git-db-1 sh -c 'psql -U "$POSTGRES_USER" -d gitea -At -v ON_ERROR_STOP=1' \
< "$(dirname "$0")/cancel_unrunnable.sql")
printf '%s\n' "$out" | grep '^{' >> "$SPOOL" || true
n=$(printf '%s\n' "$out" | grep -c '^{' || true)
echo "[janitor] cancelled ${n} unrunnable job(s)"

View File

@@ -15,17 +15,29 @@ WITH dead AS (
AND to_timestamp(j.created) < now() - interval '30 minutes'
AND EXISTS (SELECT 1 FROM jsonb_array_elements_text(j.runs_on::jsonb) l
WHERE l NOT IN ('veron-1', 'linux-x64', 'self-hosted', 'linux', 'x64'))
RETURNING j.run_id
RETURNING j.id, j.run_id, j.name, j.runs_on, j.created
), runs AS (
UPDATE action_run r
SET status = CASE
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status = 2) THEN 2
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status IN (5, 6, 7)
AND x.id NOT IN (SELECT id FROM dead)) THEN r.status
ELSE 3 END,
stopped = CASE WHEN r.stopped = 0 THEN extract(epoch from now())::bigint ELSE r.stopped END
WHERE r.id IN (SELECT DISTINCT run_id FROM dead)
RETURNING r.id
)
UPDATE action_run r
SET status = CASE
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status = 2) THEN 2
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status IN (5, 6, 7)) THEN r.status
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status = 3) THEN 3
ELSE 1 END,
stopped = CASE WHEN r.stopped = 0 THEN extract(epoch from now())::bigint ELSE r.stopped END
WHERE r.id IN (SELECT DISTINCT run_id FROM dead)
RETURNING r.id;
-- One JSON line per cancelled job: the telemetry emitter ships these as
-- ci.job_cancelled (declared with Telemetry Boss, 2026-09-23).
SELECT json_build_object(
'repo', p.lower_name,
'workflow', regexp_replace(r.workflow_id, '\.ya?ml$', ''),
'job', d.name,
'reason', 'unrunnable_label',
'runs_on', (SELECT string_agg(l, ',') FROM jsonb_array_elements_text(d.runs_on::jsonb) l),
'waited_s', (extract(epoch from now())::bigint - d.created))::text
FROM dead d JOIN action_run r ON r.id = d.run_id JOIN repository p ON p.id = r.repo_id
WHERE (SELECT count(*) FROM runs) >= 0;
-- Jobs BLOCKED on `needs:` inside a run that has already finished (a needed job
-- failed): Gitea leaves them status 7 forever. They were never going to run;

View File

@@ -28,13 +28,17 @@ repo code, and no secret is handed to any repo.
from __future__ import annotations
import base64
import json
import os
import re
import sys
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", "")
@@ -85,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,
@@ -98,12 +158,21 @@ def _call(base: str, token_header: str, method: str, path: str, body=None):
"User-Agent": "windy-git-pr-bridge/1",
},
)
try:
with urllib.request.urlopen(req, timeout=60) as r:
raw = r.read()
return r.status, (json.loads(raw) if raw else None)
except urllib.error.HTTPError as e:
return e.code, None
# Transport errors (TLS handshake timeout, reset) are retried: one GitHub
# blip used to fail the whole sync, flip its heartbeat to ok:false and page
# someone for nothing. HTTP errors are answers, not blips — never retried.
for attempt in range(3):
try:
with urllib.request.urlopen(req, timeout=60) as r:
raw = r.read()
return r.status, (json.loads(raw) if raw else None)
except urllib.error.HTTPError as e:
return e.code, None
except (urllib.error.URLError, TimeoutError, ConnectionError):
if attempt == 2:
raise
time.sleep(2 * (attempt + 1))
raise AssertionError("unreachable")
def gitea(method, path, body=None):
@@ -173,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")
@@ -181,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

@@ -6,8 +6,8 @@ Shapes are declared with Telemetry 40 (2026-09-23) — do not add keys or enum
values without re-declaring: a declared family quarantines any row that
doesn't match.
ci.run one row per FINISHED job, exactly once (high-water mark on
action_run_job.id in STATE)
ci.run one row per FINISHED job, exactly once — cursor on
(finish time, job id) in STATE; jobs finish out of id order
service.health one row per invocation: CI plane counts for the interval
Privacy: ids, names of repos/jobs, codes, counts, durations. No commit
@@ -20,12 +20,13 @@ from __future__ import annotations
import json
import os
import re
import subprocess
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime, timezone
from datetime import UTC, datetime
INGEST = os.environ.get("TELEMETRY_INGEST_URL", "https://admin.windyword.ai/v1/events")
TOKEN = os.environ.get("WINDYGIT_TELEMETRY_TOKEN", "")
@@ -66,32 +67,139 @@ def load_state() -> dict:
def iso(epoch: float) -> str:
return datetime.fromtimestamp(epoch, timezone.utc).isoformat().replace("+00:00", "Z")
return datetime.fromtimestamp(epoch, UTC).isoformat().replace("+00:00", "Z")
# ---- push velocity: DETECT + ALERT ONLY (G3.4, 2026-09-23) -----------------
# `git push` goes straight to Gitea and never touches our API, so throttle.py
# cannot see it (NOT_ENFORCED_HERE). Gitea's own `action` table does record
# every push, so we read it here, emit `forge.push_velocity` when an account
# crosses a threshold, and let Telemetry Boss's detector page. Nothing here sits
# in the push path; nothing is ever refused (orchestrator, 09-23).
#
# Thresholds = the STANDARD-band bases from config.py (500 pushes/day; the
# force-push base of 10/day is used for ref deletes, the closest thing we can
# see). EI band multipliers are NOT applied: a platinum agent over 500/day is
# still flagged, for a human to look at, not blocked. Gitea records no
# "forced" flag, so force pushes cannot be told apart from pushes: named, not
# guessed.
PV_RULES = ( # (rule, row key, window_s, threshold)
("pushes_1h", "p1h", 3600, 60),
("pushes_24h", "p24h", 86400, 500),
("ref_deletes_24h", "d24h", 86400, 10),
)
# The GitHub -> Windy Git sync pushes as windyadmin every 5 min, by design.
PV_EXEMPT = {"windyadmin"}
# Gitea op_type: 5 commit push, 9 tag push, 16 tag delete, 17 branch delete.
# One action row per WATCHER is written for each push; user_id = act_user_id
# keeps exactly the actor's own copy.
PV_QUERY = """
select a.act_user_id as uid, u.lower_name as login,
(select el.external_id from external_login_user el
where el.user_id = a.act_user_id order by el.external_id limit 1) as wid,
count(*) filter (where a.op_type in (5, 9) and a.created_unix > {h1}) as p1h,
count(*) filter (where a.op_type in (5, 9)) as p24h,
count(*) filter (where a.op_type in (16, 17)) as d24h,
count(distinct a.repo_id) as repos
from action a join "user" u on u.id = a.act_user_id
where a.created_unix > {h24} and a.user_id = a.act_user_id
and a.op_type in (5, 9, 16, 17)
group by 1, 2"""
def passport_from_login(login: str) -> str | None:
"""agent-et26abcd1234 -> ET26-ABCD-1234 (repos.py _owner_login, reversed)."""
m = re.fullmatch(r"agent-([a-z0-9]{4})([a-z0-9]{4})([a-z0-9]{4})", login)
return "-".join(g.upper() for g in m.groups()) if m else None
def push_velocity_events(rows: list[dict], now: float, alerted: dict) -> tuple[list[dict], dict]:
"""(events, alerted') — one row per account per rule per window while over.
`alerted` maps "<uid>:<rule>" -> epoch of the last row. An account still over
the line is re-reported once per window, not every 5 minutes; one that drops
back under is forgotten, so a later burst reports again.
"""
events, keep = [], {}
for r in rows:
login = str(r["login"])
if login in PV_EXEMPT:
continue
agent = login.startswith("agent-")
for rule, key, window, limit in PV_RULES:
n = int(r[key])
if n <= limit:
continue
k = f"{r['uid']}:{rule}"
last = alerted.get(k)
if last is not None and now - float(last) < window:
keep[k] = last
continue
keep[k] = now
ev = {
"ts": iso(now),
"platform": PLATFORM,
"service": "forge",
"event_type": "forge.push_velocity",
"metadata": {
"rule": rule,
"window_s": window,
"count": n,
"threshold": limit,
"repos": int(r["repos"]),
"gitea_user_id": int(r["uid"]),
},
}
# Actor rule (telemetry UPDATE 2): agent/human rows MUST carry an
# actor_id. Humans sign in to the forge only via Windy SSO, so the
# external login id IS their windy_identity_id. No id we can prove
# -> actor_type system + metadata.caller, never an invented id (I-12).
actor_id = passport_from_login(login) if agent else (r.get("wid") or None)
if actor_id:
ev["actor_type"], ev["actor_id"] = ("agent" if agent else "human"), str(actor_id)
else:
ev["actor_type"] = "system"
ev["metadata"]["caller"] = "unknown"
events.append(ev)
return events, keep
def main() -> int:
dry = "--dry-run" in sys.argv
state = load_state()
now = time.time()
last_job = int(state.get("last_job_id", 0))
since = float(state.get("last_ts", now - 300))
if not last_job:
# First run: start at the current high-water mark rather than replaying
# a month of history into the ledger as if it happened now.
last_job = int(sql("select coalesce(max(id),0) as m from action_run_job")[0]["m"])
# Cursor = (finish time, job id), NOT job id alone: jobs finish out of id
# order, so an id high-water mark silently drops every long job that started
# before the mark and finished after it (Telemetry Boss caught this: 43
# finished vs 8 ci.run rows). Finish time = stopped, or updated for jobs
# Gitea/the janitor skipped without a stop time.
if "last_fin" in state:
last_fin, last_id = int(state["last_fin"]), int(state["last_id"])
else: # first run or pre-cursor state: start now, never replay history
last_fin, last_id = int(state.get("last_ts", now)), 0
cutoff = int(now) - 5 # leave the current second alone; late writers land next run
FIN = "coalesce(nullif(j.stopped, 0), j.updated)"
jobs = sql(f"""
select j.id, j.name as job, j.status, j.started, j.stopped,
select j.id, j.name as job, j.status, j.started, j.stopped, {FIN} as fin,
p.lower_name as repo, p.default_branch, r.workflow_id, r.event,
r.ref, r.index as run, left(r.commit_sha, 7) as sha
from action_run_job j
join action_run r on r.id = j.run_id
join repository p on p.id = r.repo_id
where j.id > {last_job} and j.status in (1, 2, 3, 4) and j.stopped > 0
order by j.id
where j.status in (1, 2, 3, 4)
and ({FIN}, j.id) > ({last_fin}, {last_id})
and {FIN} <= {cutoff}
order by {FIN}, j.id
limit 2000""")
try: # posted_to_github: the bridge's own rules, from the same checkout
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import pr_status_bridge as bridge
except Exception: # noqa: BLE001
bridge = None
events = []
for j in jobs:
ref = j["ref"] or ""
@@ -101,7 +209,8 @@ def main() -> int:
kind = "default"
else:
kind = "other"
dur = (j["stopped"] - j["started"]) * 1000 if j["started"] else None
# Gitea stores whole seconds; duration_ms is seconds*1000 (so 10000 = 10 s).
dur = (j["stopped"] - j["started"]) * 1000 if j["started"] and j["stopped"] else None
ev = {
"ts": iso(j["stopped"]),
"platform": PLATFORM,
@@ -121,20 +230,38 @@ def main() -> int:
}
if dur is not None and dur >= 0:
ev["duration_ms"] = int(dur)
if bridge is not None:
wf = ev["metadata"]["workflow"]
ev["metadata"]["posted_to_github"] = bool(
j["repo"] in {r.lower() for r in bridge.REPOS}
and kind in ("default", "pr")
and j["status"] != 4
and not bridge.NO_DAEMON_JOB.search(j["job"])
and f"{wf}/{j['job']}" not in bridge.NON_BLOCKING.get(j["repo"], set())
)
events.append(ev)
# --- heartbeat: counts since the previous invocation --------------------
# Interval counts come from EXACTLY the rows emitted above, so
# sum(jobs_finished) over any window == count(ci.run) in it, by construction.
h = sql(f"""
select
(select count(*) from action_run_job where stopped >= {int(since)} and status in (1,2,3,4)) as jobs_finished,
(select count(*) from action_run_job where stopped >= {int(since)} and status = 2) as jobs_failed,
(select count(*) from action_run_job where stopped >= {int(since)} and status = 3) as jobs_cancelled,
(select count(*) from action_run_job where status in (5, 7)) as jobs_waiting,
(select count(*) from action_run_job where status = 6) as jobs_running,
(select count(*) from action_runner where deleted is null and last_online >= {int(now) - 120}) as runners_online,
(select coalesce(extract(epoch from now())::bigint - min(created), 0)
from action_run_job where status in (5, 7)) as oldest_waiting_s""")[0]
h["jobs_finished"] = len(jobs)
h["jobs_failed"] = sum(1 for j in jobs if j["status"] == 2)
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
@@ -150,14 +277,61 @@ 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
try:
with open(spool) as f:
for line in f:
try:
m = json.loads(line)
except ValueError:
continue
events.append(
{
"ts": iso(now),
"platform": PLATFORM,
"service": SERVICE,
"event_type": "ci.job_cancelled",
"actor_type": "system",
"metadata": {
k: m[k]
for k in ("repo", "workflow", "job", "reason", "runs_on", "waited_s")
},
}
)
spooled += 1
except OSError:
pass
if dry:
print(json.dumps({"events": events}, indent=1)[:4000])
out = os.environ.get("TELEMETRY_DRY_OUT")
if out:
with open(out, "w") as f:
json.dump({"events": events}, f)
else:
print(json.dumps({"events": events}, indent=1)[:4000])
print(f"[telemetry] DRY RUN: {len(events)} events ({len(jobs)} ci.run)")
return 0
if not TOKEN:
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,
@@ -171,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
@@ -183,10 +366,16 @@ def main() -> int:
print(f"[telemetry] FAILED ingest: {e.reason}")
return 1
if spooled:
open(spool, "w").close() # only after every batch was accepted
os.makedirs(os.path.dirname(STATE), exist_ok=True)
new_last = max([j["id"] for j in jobs], default=last_job)
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_job_id": new_last, "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