25 Commits

Author SHA1 Message Date
Kit OC5
a2daca61ef ci: mount windy-pro's read-only build inputs into dind; allow exactly that path for jobs
Non-secret inputs git-ignored in windy-pro (models, linux-x64 portable
bundle, enter-monitor build) that build-desktop needs. Mounted :ro into
dind; valid_volumes allows only /ci-inputs/windy-pro; refresh-ci-inputs.sh
copies them from the frozen release clone (read-only on the source).
Invariant I-5 narrowed, not dropped: exactly that one path, read-only in
dind, no other service mounts it, still no docker socket (proven to fail
on :rw). Orchestrator-approved (option a). Applied in an idle window.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:56:48 -04:00
Kit OC5
1bc55eaec9 compute guard: flag direct AI-provider use (Windy Mind is the only door), warn-only
All checks were successful
check / gate (push) Successful in 28s
canary / probe (push) Successful in 14s
Grant's rule (09-23): every model call goes through Windy Mind. The bridge
now posts windy-git/compute-guard on every PR head (lines the PR ADDS vs its
merge-base) and default-branch head (whole tree): provider hosts, provider
SDK imports/deps and raw provider key names. Warn-only: success + "⚠ WARN"
and a link to the first hit; COMPUTE_GUARD_MODE=block turns it red later.

Exceptions live in ci/compute-guard-allow.yml, each with a reason (Mind
itself, user-BYOK windy-agent / windy-code extension / windy-pro desktop +
MindPanel, windy-connect config writers). Tests, docs, comments, lockfiles,
vendored code and CI config are never scanned. Reads the sync's bare clones
(no docker exec); cached per (repo, sha, rules). Non-fatal; never a fake OK.

First cases = COMPUTE_BYPASS_AUDIT.md. Today on default branches: 38
findings in 3 repos (windy-chat audit #2, windy-pro account-server #3/#4,
windytalk reference/), 0 elsewhere.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:42:10 -04:00
Kit OC5
fdb0f5989e ci: run dind under Sysbox, not privileged (rollback override kept)
Some checks failed
canary / probe (push) Has been cancelled
check / gate (push) Has been cancelled
dind was privileged: true, so a job that escaped into dind was root on
Veron 1, which is Grant's workstation. Under sysbox-runc (sysbox-ce 0.7.1,
installed 09-23 with no docker restart) dind root is an unprivileged host
uid. Smoke-tested standalone: nested containers, internet, a services-style
postgres on a private network and a python image all pass unprivileged.
Fresh volume dind-storage-sysbox; the old dind-storage stays for
docker-compose.privileged.yml, the one-command rollback.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:23:55 -04:00
Kit OC5
9be2952ff8 rerun_ci: detect a running oneshot sync correctly (is-active lies for oneshot)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 14:22:59 -04:00
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
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
1a171eabd5 telemetry: CI emitter for admin.windyword.ai (inert until token)
Some checks failed
check / gate (push) Failing after 20s
canary / probe (push) Successful in 6s
ci.run (one row per finished job, exactly once via a high-water mark;
branch_kind default|pr|other so the dashboard can show "main is red") and
service.health (interval counts: finished/failed/cancelled, waiting,
running, runners online, oldest wait). Shapes declared with Windy
Telemetry 40; sends nothing until WINDYGIT_TELEMETRY_TOKEN exists.
State is only advanced after a 2xx, so a failed post retries.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:06:54 -04:00
28c31236b8 ci: janitor also clears jobs blocked forever on failed needs
When a needed job fails, Gitea leaves dependants BLOCKED (7) even after
the run finishes; eternitas build jobs sat there 8h. Mark them skipped
(what GitHub shows) once the run is done and 30 min have passed.
Found by the new telemetry dry run (oldest_waiting_s = 29160).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-23 11:06:33 -04:00
30 changed files with 2273 additions and 45 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

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

@@ -0,0 +1,184 @@
"""Compute guard: Windy Mind is the only door to AI compute (warn-only today)."""
from __future__ import annotations
import importlib.util
import subprocess
import sys
from pathlib import Path
import pytest
ROOT = Path(__file__).resolve().parents[2]
_spec = importlib.util.spec_from_file_location("compute_guard", ROOT / "scripts" / "compute_guard.py")
cg = importlib.util.module_from_spec(_spec)
sys.modules["compute_guard"] = cg
_spec.loader.exec_module(cg)
ALLOW = cg.load_allow(ROOT / "ci" / "compute-guard-allow.yml")
@pytest.mark.parametrize(
"path, text, kind",
[
# audit #1 (windy-search, closed) and #2 (windy-chat, live): the shapes they had
("service/app/anthropic_client.py", 'URL = "https://api.anthropic.com/v1/messages"', "provider host"),
("service/app/config.py", 'token = os.environ["ANTHROPIC_OAUTH_TOKEN"]', "provider key"),
("services/agent-roster/lib/llm.js", "const url = 'https://api.groq.com/openai/v1/chat/completions'", "provider host"),
("docker-compose.yml", " GROQ_API_KEY: ${GROQ_API_KEY}", "provider key"),
# audit #3/#4 (windy-pro account-server)
("account-server/src/routes/transcription.ts", "const r = await fetch('https://api.openai.com/v1/audio/transcriptions'", "provider host"),
("account-server/src/config.ts", "openaiKey: process.env.OPENAI_API_KEY,", "provider key"),
# SDKs and deps
("app/llm.py", "from anthropic import Anthropic", "provider SDK"),
("app/llm.py", "import openai", "provider SDK"),
("app/llm.py", "import google.generativeai as genai", "provider SDK"),
("src/ai.ts", 'import Anthropic from "@anthropic-ai/sdk";', "provider SDK"),
("src/ai.js", "const Groq = require('groq-sdk')", "provider SDK"),
("package.json", ' "openai": "^4.52.0",', "provider SDK dep"),
("requirements.txt", "anthropic>=0.40", "provider SDK dep"),
("pyproject.toml", ' "google-generativeai>=0.8",', "provider SDK dep"),
],
)
def test_audit_shapes_are_flagged(path, text, kind):
assert kind in [k for k, _ in cg.scan_line(path, text)]
@pytest.mark.parametrize(
"path, text",
[
("app/mind.py", 'MIND = "https://mind.windyword.ai/v1/chat/completions"'), # the door itself
("app/models.py", "openai_compatible = True # Mind speaks the OpenAI wire format"),
("app/x.py", "from app.openai_shim import x"), # a local module, not the SDK
("package.json", ' "openai-types-lite": "1.0.0",'), # a different package
("README.txt", "set OPENAI_API_KEY"), # scanned-by-rule, excluded by SKIP separately
],
)
def test_near_misses_are_not_flagged(path, text):
if cg.SKIP.search(path):
return
assert cg.scan_line(path, text) == []
@pytest.mark.parametrize(
"path",
["tests/test_llm.py", "api/tests/x.py", "src/ai.test.ts", "web/foo.spec.js", "docs/setup.md",
"README.md", "package-lock.json", "uv.lock", "node_modules/openai/index.js", ".github/workflows/ci.yml",
"conftest.py", "app/llm_test.py"],
)
def test_tests_docs_lockfiles_vendored_ci_are_never_scanned(path):
assert cg.SKIP.search(path)
def test_allow_list_needs_a_reason_per_entry(tmp_path):
bad = tmp_path / "a.yml"
bad.write_text("allow:\n - repo: x\n paths: ['*']\n")
with pytest.raises(ValueError):
cg.load_allow(bad)
@pytest.mark.parametrize(
"repo, path, ok",
[
("windy-mind", "app/providers/anthropic.py", True),
("windy-agent", "agent/providers.py", True),
("windy-code", "extensions/windy-ai/src/aiProvider.ts", True),
("windy-code", "web/server/llm.ts", False), # BYOK is the extension only
("windy-connect", "backend/src/writers/claude_code.py", True),
("windy-chat", "services/agent-roster/lib/llm.js", False), # audit #2: must be flagged
("windy-pro", "account-server/src/routes/translations.ts", False),
],
)
def test_allow_list_entries(repo, path, ok):
assert cg.allowed(repo, path, ALLOW) is ok
DIFF = """diff --git a/app/llm.py b/app/llm.py
--- a/app/llm.py
+++ b/app/llm.py
@@ -10,0 +11,2 @@
+import anthropic
+client = anthropic.Anthropic()
diff --git a/tests/test_llm.py b/tests/test_llm.py
--- /dev/null
+++ b/tests/test_llm.py
@@ -0,0 +1 @@
+import anthropic
@@ -40 +42 @@
-x = 1
+x = 2
"""
def test_only_added_non_test_lines_are_findings():
fs = cg.parse_added("windy-chat", DIFF, ALLOW)
assert [(f.path, f.line, f.kind) for f in fs] == [("app/llm.py", 11, "provider SDK")]
def _repo(tmp_path, files: dict[str, str]) -> tuple[Path, str]:
work = tmp_path / "w"
work.mkdir()
run = lambda *a: subprocess.run(["git", *a], cwd=work, check=True, capture_output=True) # noqa: E731
run("init", "-q", "-b", "main")
for p, text in files.items():
(work / p).parent.mkdir(parents=True, exist_ok=True)
(work / p).write_text(text)
run("add", "-A")
run("-c", "user.email=t@t", "-c", "user.name=t", "commit", "-qm", "x")
bare = tmp_path / "r.git"
subprocess.run(["git", "clone", "-q", "--bare", str(work), str(bare)], check=True)
sha = subprocess.run(["git", "--git-dir", str(bare), "rev-parse", "main"],
capture_output=True, text=True, check=True).stdout.strip()
return bare, sha
def test_tree_scan_on_a_real_git_repo(tmp_path):
bare, sha = _repo(tmp_path, {
"app/llm.py": "import os\nKEY = os.environ['OPENAI_API_KEY']\n",
"app/ok.py": "MIND = 'https://mind.windyword.ai'\n",
"tests/test_llm.py": "import anthropic\n",
"docs/x.md": "api.anthropic.com\n",
})
fs = cg.scan_tree("windy-chat", bare, sha, ALLOW)
assert [(f.path, f.line, f.kind) for f in fs] == [("app/llm.py", 2, "provider key")]
def test_warn_mode_never_turns_red(monkeypatch):
monkeypatch.setattr(cg, "MODE", "warn")
state, desc, f = cg.status_for([cg.Finding("a.py", 3, "provider host", "api.openai.com")], whole_tree=False)
assert state == "success" and desc.startswith("⚠ WARN (not blocking): 1 direct AI-provider use added")
assert "a.py:3" in desc and f.path == "a.py"
def test_block_mode_fails(monkeypatch):
monkeypatch.setattr(cg, "MODE", "block")
state, desc, _ = cg.status_for([cg.Finding("a.py", 3, "provider host", "x")], whole_tree=True)
assert state == "failure" and desc.startswith("BLOCKED")
def test_clean_is_ok():
assert cg.status_for([], whole_tree=True)[:2] == (
"success", "OK: no direct AI-provider use in tree (Windy Mind is the only door)")
@pytest.mark.parametrize(
"text",
[
" # The ANTHROPIC_OAUTH_TOKEN setting was removed on 2026-09-23 ON PURPOSE", # windy-search
"# ANTHROPIC_API_KEY=",
" // fallback used to call https://api.groq.com directly",
" * @see https://api.openai.com/v1/audio",
"<!-- api.anthropic.com -->",
],
)
def test_comments_are_not_calls(text):
assert cg.scan_line("service/app/config.py", text) == []
def test_code_with_a_trailing_comment_still_counts():
assert cg.scan_line("a.js", "fetch('https://api.openai.com/v1') // TODO move to Mind")
def test_windy_pro_desktop_is_byok_but_the_account_server_is_not():
assert cg.allowed("windy-pro", "src/client/desktop/main.js", ALLOW)
assert not cg.allowed("windy-pro", "account-server/src/routes/translations.ts", ALLOW)

View File

@@ -425,9 +425,28 @@ def test_i05_jobs_get_a_network_per_job_not_a_shared_bridge():
def test_i05_jobs_cannot_bind_mount_from_the_daemon_host():
cfg = (ROOT / "deploy" / "runner" / "config.yaml").read_text()
assert "valid_volumes: []" in cfg
assert 'docker_host: "-"' in cfg
"""Narrowed 2026-09-23 (orchestrator-approved): a job may bind-mount EXACTLY
one daemon path, windy-pro's non-secret build inputs, and only because dind
itself has that path READ-ONLY. Anything more (a second path, a writable
one, a glob) reopens the host to CI code. Still no docker socket for jobs."""
import re
import yaml
rd = ROOT / "deploy" / "runner"
cfg = yaml.safe_load((rd / "config.yaml").read_text())
allowed = cfg["container"]["valid_volumes"]
assert allowed in ([], ["/ci-inputs/windy-pro"]), f"I-5: jobs may mount nothing else: {allowed}"
assert cfg["container"]["docker_host"] == "-"
if allowed:
compose = yaml.safe_load((rd / "docker-compose.yml").read_text())
binds = [v for v in compose["services"]["dind"]["volumes"] if v.startswith("/")]
assert binds == ["/home/user1-gpu/ci-inputs/windy-pro:/ci-inputs/windy-pro:ro"], (
f"I-5: dind's only host bind must be the ci-inputs path, READ-ONLY: {binds}")
for name, svc in compose["services"].items():
if name != "dind":
for v in svc.get("volumes") or []:
assert not re.match(r"^/home/user1-gpu/ci-inputs", v), f"I-5: {name} mounts ci-inputs"
def test_i05_no_ci_container_can_reach_the_forge_network():

View File

@@ -8,7 +8,9 @@ no status at all.
from __future__ import annotations
import base64
import importlib.util
import sys
from pathlib import Path
import pytest
@@ -35,12 +37,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]}
@@ -67,10 +80,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
@@ -171,3 +185,200 @@ 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
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) == []
class _Guard:
def __init__(self, findings):
self.findings = findings
def check(self, repo, sha, default_branch, is_default_head):
return self.findings
@staticmethod
def status_for(findings, whole_tree):
if not findings:
return "success", "OK: clean", None
return "success", f"WARN {len(findings)}", findings[0]
class _F:
path, line = "app/llm.py", 7
def test_guard_posts_warn_with_a_link_to_the_first_finding(fake, monkeypatch):
f = fake()
monkeypatch.setitem(sys.modules, "compute_guard", _Guard([_F()]))
bridge.post_compute_guard("windy-chat", SHA, "main", False)
assert [(p["context"], p["state"], p["description"]) for p in f.posted] == [
("windy-git/compute-guard", "success", "WARN 1")]
assert f.posted[0]["target_url"].endswith(f"/src/commit/{SHA}/app/llm.py#L7")
def test_guard_same_status_is_not_reposted(fake, monkeypatch):
f = fake(statuses=[{"context": "windy-git/compute-guard", "state": "success", "description": "WARN 1"}])
monkeypatch.setitem(sys.modules, "compute_guard", _Guard([_F()]))
bridge.post_compute_guard("windy-chat", SHA, "main", False)
assert f.posted == []
def test_guard_that_cannot_run_posts_nothing(fake, monkeypatch):
f = fake()
monkeypatch.setitem(sys.modules, "compute_guard", _Guard(None))
bridge.post_compute_guard("windy-chat", SHA, "main", True)
assert f.posted == []

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

@@ -0,0 +1,40 @@
# Compute guard allow-list: code that MAY talk to an AI provider directly.
# Windy Mind is the ONLY door to AI compute (Grant, 2026-09-23). Every entry
# here is an exception to that rule and MUST say why. Paths are fnmatch globs
# relative to the repo root. Owner of this file: Windy Git lane (13); changes
# go through the orchestrator. Source of the first entries: COMPUTE_BYPASS_AUDIT.md.
allow:
- repo: windy-mind
paths: ["*"]
reason: "Windy Mind IS the door: provider clients belong here by definition."
- repo: windy-agent
paths: ["*"]
reason: >-
User BYOK: self-hosted agents call providers on the USER's own keys.
Mind stays opt-in there, or every self-hosted user's inference lands on
Grant's bill (no-cloud-cost-liability rule; audit #7).
- repo: windy-code
paths: ["extensions/windy-ai/*"]
reason: "User BYOK AI extension: the user's own provider keys; Mind is one opt-in provider (audit #8)."
- repo: windy-connect
paths: ["*writers/*"]
reason: "Writes client configs that NAME the user's own provider env vars; makes no provider calls (audit #11)."
- repo: windy-pro
paths: ["src/client/desktop/*"]
reason: >-
User BYOK desktop client: cloud STT/translate keys come from what the USER
enters (renderer localStorage -> electron-store; env var only for dev), and
the CSP line allows exactly those user-keyed hosts (audit #10). The
account-server is NOT covered: server-side calls go through Mind.
- repo: windy-pro
paths: ["src/client/web/src/pages/panels/MindPanel.jsx"]
reason: "Validates the USER's own OpenRouter key for BYOK (audit #10); spends no house money."
- repo: windy-git
paths: ["scripts/compute_guard.py", "ci/compute-guard-allow.yml"]
reason: "The guard's own pattern list and this file."

View File

@@ -54,6 +54,9 @@ container:
privileged: false
options:
workdir_parent: /workspace
valid_volumes: [] # a job cannot bind-mount anything from the daemon host
# A job may bind-mount exactly ONE daemon path: the read-only windy-pro build
# inputs (mounted :ro into dind itself). Per-runner, not per-repo (act_runner
# limit): any windyadmin repo could mount it; it is non-secret and read-only.
valid_volumes: ["/ci-inputs/windy-pro"]
docker_host: "-" # do NOT expose the runner's own docker socket to jobs
force_pull: false

View File

@@ -0,0 +1,15 @@
# ROLLBACK ONLY: the pre-Sysbox dind (privileged: true), kept one command away.
# Use it if CI breaks under Sysbox:
#
# cd /srv/windygit/src/deploy/runner
# sudo docker compose -f docker-compose.yml -f docker-compose.privileged.yml up -d dind
#
# (then restart the runners while idle). Compose merges `volumes` by container
# path, so this puts back the old `dind-storage` volume with its image cache.
# Going forward again: the same command without the second -f.
services:
dind:
runtime: runc
privileged: true
volumes:
- dind-storage:/var/lib/docker

View File

@@ -26,8 +26,12 @@
# * `dind` and every job container it spawns are UNTRUSTED. They are on a
# private network with no access to the forge, its database, or its .env.
#
# dind itself is privileged — that is the cost, and it is the reason a job
# escape lands in a disposable daemon rather than on Grant's workstation.
# dind is NOT privileged (2026-09-23): it runs under the Sysbox runtime
# (sysbox-ce on Veron, `runtime: sysbox-runc`), a user-namespaced system
# container whose root is an unprivileged host uid. A job that escapes its own
# container lands in dind as a nobody on the host, not as root on Grant's
# workstation. Before Sysbox, dind was `privileged: true`; that config is kept
# as docker-compose.privileged.yml (ROLLBACK ONLY, one command, see that file).
#
# ⚠️ Do NOT "simplify" this by mounting the host docker socket.
@@ -36,13 +40,20 @@ name: windy-git-runner
services:
dind:
image: docker.io/library/docker:27-dind
privileged: true
runtime: sysbox-runc # NOT privileged: see the I-5 note above
environment:
DOCKER_TLS_CERTDIR: "" # plain TCP on an isolated network, no host route
command: ["dockerd", "--host=tcp://0.0.0.0:2375", "--tls=false"]
networks: [jobs]
volumes:
- dind-storage:/var/lib/docker
# A fresh volume: Sysbox shifts ownership to its own uid range. The old
# `dind-storage` is kept untouched for the privileged rollback.
- dind-storage-sysbox:/var/lib/docker
# READ-ONLY, non-secret build inputs for windy-pro's desktop jobs (models,
# linux-x64 portable bundle, enter-monitor build), copied from the frozen
# release clone by deploy/runner/refresh-ci-inputs.sh. Jobs may mount ONLY
# this path (config.yaml valid_volumes). Orchestrator-approved 09-23.
- /home/user1-gpu/ci-inputs/windy-pro:/ci-inputs/windy-pro:ro
# G1.5 — bounded so a fork-bomb workflow cannot starve Grant's interactive
# session. Veron 1 is his workstation, not a dedicated build box.
cpus: 12.0 # 12 of 24 cores
@@ -159,6 +170,7 @@ networks:
volumes:
dind-storage:
dind-storage-sysbox:
runner-data:
runner-data-2:
runner-data-3:

View File

@@ -0,0 +1,21 @@
#!/usr/bin/env bash
# Refresh the READ-ONLY CI input cache for windy-pro desktop jobs from the FROZEN
# release clone on Veron. Reads ~/windy-pro-release only; never writes to it.
# Run when the release lane says engines / wheels / the portable bundle changed.
# Linux inputs only: 3 models + requirements-bundle.txt, bundled-portable/linux-x64,
# native/enter-monitor/build. Result is chmod a-w and mounted :ro into dind.
set -euo pipefail
SRC=/home/user1-gpu/windy-pro-release
DST=/home/user1-gpu/ci-inputs/windy-pro
models=$(ls "$SRC/extraResources/model" | grep -E '^windy-(nano|lite|core)-ct2$')
[ "$(wc -w <<<"$models")" = 3 ] || { echo "expected 3 models, got: $models"; exit 1; }
mkdir -p "$DST/extraResources/model" "$DST/bundled-portable" "$DST/native-enter-monitor-build"
chmod -R u+w "$DST"
R="ionice -c3 nice -n 19 rsync -a --delete"
for m in $models; do $R "$SRC/extraResources/model/$m/" "$DST/extraResources/model/$m/"; done
$R "$SRC/extraResources/requirements-bundle.txt" "$DST/extraResources/requirements-bundle.txt"
$R "$SRC/bundled-portable/linux-x64/" "$DST/bundled-portable/linux-x64/"
$R "$SRC/native/enter-monitor/build/" "$DST/native-enter-monitor-build/"
date -u +%FT%TZ > "$DST/.refreshed-from-windy-pro-release"
chmod -R a-w "$DST"
du -sh --apparent-size "$DST"

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

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

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

@@ -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,11 +80,16 @@ 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:
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)
@@ -88,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})",
@@ -107,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(
@@ -184,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
@@ -260,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

@@ -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,15 +15,39 @@ 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)) 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,
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;
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;
-- mark them skipped (4), which is what GitHub shows for the same situation.
UPDATE action_run_job j
SET status = 4, updated = extract(epoch from now())::bigint
FROM action_run r
WHERE r.id = j.run_id
AND j.status = 7
AND r.status IN (1, 2, 3)
AND to_timestamp(j.created) < now() - interval '30 minutes'
RETURNING j.run_id;
COMMIT;

270
scripts/compute_guard.py Normal file
View File

@@ -0,0 +1,270 @@
#!/usr/bin/env python3
"""Compute guard: Windy Mind is the ONLY door to AI compute (Grant, 2026-09-23).
Flags code that talks to an AI provider directly instead of through Windy Mind:
a provider API host, a provider SDK import or dependency, or a raw provider key
name. Direct calls skip Mind's metering, caps and live-model routing, and they
spend whichever key happens to be lying around (the audit found Grant's personal
Max OAuth token inside a platform container).
WARN-ONLY for now: the bridge posts `windy-git/compute-guard` as success with a
"⚠ WARN" description, so nothing turns red. `COMPUTE_GUARD_MODE=block` flips
findings to failure once the repos are clean (orchestrator's call).
- PR heads: only lines the PR ADDS (vs its merge-base with the default branch).
- Default-branch head: the whole tree (the baseline, and what `report` prints).
Exceptions live in ONE file, ci/compute-guard-allow.yml, each with a reason.
Tests, docs, lockfiles, vendored code and CI config are never scanned.
Reads the sync's bare GitHub clones on Veron (no docker exec: IO-stall lesson).
python3 scripts/compute_guard.py report [repo ...] # whole-tree findings on each default branch
"""
from __future__ import annotations
import fnmatch
import hashlib
import json
import os
import re
import subprocess
import sys
from dataclasses import dataclass
from pathlib import Path
import yaml
ROOT = Path(__file__).resolve().parents[1]
ALLOW_FILE = Path(os.environ.get("COMPUTE_GUARD_ALLOW", ROOT / "ci" / "compute-guard-allow.yml"))
WORK = Path(os.environ.get("SYNC_WORK", "/srv/windygit/sync"))
CACHE = Path(os.environ.get("COMPUTE_GUARD_CACHE", "/var/lib/windy-git/compute-guard-cache.json"))
MODE = os.environ.get("COMPUTE_GUARD_MODE", "warn") # warn | block
HOSTS = [
"api.anthropic.com", "api.openai.com", "api.groq.com",
"generativelanguage.googleapis.com", "api.mistral.ai", "api.perplexity.ai",
"openrouter.ai", "api.together.xyz", "api.together.ai", "api.cerebras.ai",
"api.sambanova.ai", "api.deepseek.com", "api.x.ai", "api.cohere.ai",
"api.cohere.com", "api.fireworks.ai", "api.replicate.com",
"api-inference.huggingface.co",
]
KEYS = [
"ANTHROPIC_API_KEY", "ANTHROPIC_OAUTH_TOKEN", "ANTHROPIC_AUTH_TOKEN",
"OPENAI_API_KEY", "GROQ_API_KEY", "GEMINI_API_KEY", "GOOGLE_GENERATIVE_AI_API_KEY",
"GOOGLE_AI_API_KEY", "MISTRAL_API_KEY", "PERPLEXITY_API_KEY", "PPLX_API_KEY",
"OPENROUTER_API_KEY", "TOGETHER_API_KEY", "CEREBRAS_API_KEY", "SAMBANOVA_API_KEY",
"DEEPSEEK_API_KEY", "XAI_API_KEY", "COHERE_API_KEY", "FIREWORKS_API_KEY",
"REPLICATE_API_TOKEN",
]
PY_SDKS = r"anthropic|openai|groq|mistralai|cohere|google\.generativeai|google\.genai|together|cerebras|litellm"
JS_SDKS = (r"@anthropic-ai/sdk|openai|groq-sdk|@google/generative-ai|@google/genai|@mistralai/mistralai"
r"|cohere-ai|together-ai|@ai-sdk/(?:anthropic|openai|groq|google|mistral)")
RULES: list[tuple[str, re.Pattern]] = [
("provider host", re.compile("|".join(re.escape(h) for h in HOSTS))),
("provider key", re.compile(r"\b(?:" + "|".join(KEYS) + r")\b")),
("provider SDK", re.compile(rf"^\s*(?:from|import)\s+(?:{PY_SDKS})(?:\s|\.|$|,)")),
("provider SDK", re.compile(rf"""(?:from\s+|require\(\s*|import\(\s*)['"](?:{JS_SDKS})(?:/[^'"]*)?['"]""")),
# dependency manifests: package.json keys, requirements / pyproject lines
("provider SDK dep", re.compile(rf'''^\s*"(?:{JS_SDKS})"\s*:''')),
("provider SDK dep", re.compile(rf'''^\s*["']?(?:{PY_SDKS.replace(chr(92) + ".", "-")})(?:\[[^\]]*\])?\s*(?:[<>=~!]=?|["',]|$)''')),
]
DEP_FILES = re.compile(r"(^|/)(package\.json|requirements[^/]*\.txt|pyproject\.toml|setup\.cfg|Pipfile)$")
# Never scanned: tests, docs, lockfiles, vendored/built code, CI config.
SKIP = re.compile(
r"(^|/)(tests?|__tests__|spec|docs?|node_modules|vendor|dist|build|\.github|\.gitea)/"
r"|(^|/)(test_[^/]*|[^/]*_test\.py|conftest\.py|[^/]*\.(test|spec)\.[cm]?[jt]sx?)$"
r"|\.(md|mdx|rst|txt|lock|snap|svg|png|jpg|pdf)$"
r"|(^|/)(package-lock\.json|pnpm-lock\.yaml|yarn\.lock|uv\.lock|poetry\.lock|Cargo\.lock)$"
)
@dataclass(frozen=True)
class Finding:
path: str
line: int
kind: str
match: str
def load_allow(path: Path = ALLOW_FILE) -> list[dict]:
data = yaml.safe_load(path.read_text()) or {}
entries = data.get("allow") or []
for e in entries: # a reason per entry is the whole point of the file
if not (e.get("repo") and e.get("paths") and str(e.get("reason", "")).strip()):
raise ValueError(f"allow entry needs repo, paths and a reason: {e}")
return entries
def allowed(repo: str, path: str, allow: list[dict]) -> bool:
for e in allow:
if e["repo"] == repo and any(fnmatch.fnmatch(path, g) for g in e["paths"]):
return True
return False
COMMENT = re.compile(r"^\s*(?:#|//|/\*|\*|<!--)")
def scan_line(path: str, text: str) -> list[tuple[str, str]]:
# A comment is not a call: "the ANTHROPIC_OAUTH_TOKEN setting was removed"
# (windy-search) must not count, nor a commented-out `# OPENAI_API_KEY=`.
if COMMENT.match(text):
return []
hits = []
for kind, rx in RULES:
if kind == "provider SDK dep" and not DEP_FILES.search(path):
continue
m = rx.search(text)
if m:
hits.append((kind, m.group(0).strip()[:60]))
return hits
def _git(bare: Path, *args: str) -> str:
return subprocess.run(
["git", "--git-dir", str(bare), *args],
capture_output=True, text=True, check=True, timeout=120,
).stdout
def scan_tree(repo: str, bare: Path, sha: str, allow: list[dict]) -> list[Finding]:
"""Every line in the tree at `sha` (default branch: the baseline)."""
# A cheap prefilter by git, then the real rules in Python.
pre = "|".join([re.escape(h) for h in HOSTS] + KEYS + ["anthropic", "openai", "groq", "mistral",
"generativeai", "genai", "cohere", "together", "cerebras", "litellm"])
try:
out = _git(bare, "grep", "-nIE", "-e", pre, sha, "--", ".")
except subprocess.CalledProcessError as e:
if e.returncode == 1: # no matches
return []
raise
found = []
for raw in out.splitlines():
# <sha>:<path>:<line>:<text>
try:
_, path, line, text = raw.split(":", 3)
except ValueError:
continue
if SKIP.search(path) or allowed(repo, path, allow):
continue
for kind, match in scan_line(path, text):
found.append(Finding(path, int(line), kind, match))
return found
def scan_added(repo: str, bare: Path, base_ref: str, sha: str, allow: list[dict]) -> list[Finding]:
"""Only the lines a PR adds, vs its merge-base with the default branch."""
mb = _git(bare, "merge-base", base_ref, sha).strip()
diff = _git(bare, "diff", "-U0", "--no-color", "--no-ext-diff", mb, sha)
return parse_added(repo, diff, allow)
HUNK = re.compile(r"^@@ -\d+(?:,\d+)? \+(\d+)(?:,\d+)? @@")
def parse_added(repo: str, diff: str, allow: list[dict]) -> list[Finding]:
found, path, line = [], None, 0
for raw in diff.splitlines():
if raw.startswith("+++ "):
p = raw[4:]
path = None if p == "/dev/null" else p[2:] if p.startswith("b/") else p
continue
m = HUNK.match(raw)
if m:
line = int(m.group(1))
continue
if path is None or raw.startswith("--- "):
continue
if raw.startswith("+"):
if not (SKIP.search(path) or allowed(repo, path, allow)):
for kind, match in scan_line(path, raw[1:]):
found.append(Finding(path, line, kind, match))
line += 1
return found
# ---- cache: a tree scan runs once per (repo, sha, rules+allow) --------------
def _fingerprint(allow: list[dict]) -> str:
return hashlib.sha256(
json.dumps([HOSTS, KEYS, PY_SDKS, JS_SDKS, SKIP.pattern, allow], sort_keys=True).encode()
).hexdigest()[:16]
def cached_scan(key: str, fn) -> list[Finding]:
try:
cache = json.loads(CACHE.read_text())
except (OSError, ValueError):
cache = {}
if key in cache:
return [Finding(**f) for f in cache[key]]
result = fn()
cache[key] = [f.__dict__ for f in result]
if len(cache) > 2000: # keep it small: newest entries win
cache = dict(list(cache.items())[-1000:])
try:
CACHE.parent.mkdir(parents=True, exist_ok=True)
tmp = CACHE.with_suffix(".tmp")
tmp.write_text(json.dumps(cache))
tmp.replace(CACHE)
except OSError:
pass
return result
def check(repo: str, sha: str, default_branch: str, is_default_head: bool) -> list[Finding] | None:
"""Findings for one commit, or None when the guard can't run (never a fake OK)."""
bare = WORK / f"{repo}.git"
if not bare.is_dir():
return None
allow = load_allow()
fp = _fingerprint(allow)
if is_default_head:
return cached_scan(f"tree:{repo}:{sha}:{fp}", lambda: scan_tree(repo, bare, sha, allow))
return cached_scan(
f"pr:{repo}:{sha}:{fp}",
lambda: scan_added(repo, bare, f"refs/heads/{default_branch}", sha, allow),
)
def status_for(findings: list[Finding], whole_tree: bool) -> tuple[str, str, Finding | None]:
"""(state, description, first finding) for the GitHub commit status."""
scope = "in tree" if whole_tree else "added"
if not findings:
what = "no direct AI-provider use in tree" if whole_tree else "no direct AI-provider use added"
return "success", f"OK: {what} (Windy Mind is the only door)", None
f = findings[0]
n = len(findings)
state = "failure" if MODE == "block" else "success"
lead = "BLOCKED" if MODE == "block" else "⚠ WARN (not blocking)"
desc = f"{lead}: {n} direct AI-provider use{'s' if n > 1 else ''} {scope}, e.g. {f.path}:{f.line} {f.match}"
return state, desc[:140], f
def report(repos: list[str]) -> int:
allow = load_allow()
total = 0
for repo in repos:
bare = WORK / f"{repo}.git"
if not bare.is_dir():
print(f"## {repo}: no sync clone, skipped")
continue
head = _git(bare, "symbolic-ref", "--short", "HEAD").strip()
sha = _git(bare, "rev-parse", head).strip()
fs = scan_tree(repo, bare, sha, allow)
total += len(fs)
print(f"## {repo} ({head} {sha[:7]}): {len(fs)} finding(s)")
for f in fs:
print(f" {f.path}:{f.line} [{f.kind}] {f.match}")
print(f"TOTAL {total}")
return 0
if __name__ == "__main__":
if len(sys.argv) >= 2 and sys.argv[1] == "report":
default = os.environ.get("BRIDGE_REPOS", "").split() or sorted(
p.name.removesuffix(".git") for p in WORK.glob("*.git"))
sys.exit(report(sys.argv[2:] or default))
sys.exit(__doc__)

View File

@@ -28,13 +28,18 @@ repo code, and no secret is handed to any repo.
from __future__ import annotations
import base64
import json
import os
import re
import subprocess
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 +90,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 +159,21 @@ def _call(base: str, token_header: str, method: str, path: str, body=None):
"User-Agent": "windy-git-pr-bridge/1",
},
)
# 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):
@@ -153,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.
@@ -173,7 +291,19 @@ 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:
# 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
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
@@ -181,6 +311,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:
@@ -198,6 +343,37 @@ def post_statuses(repo: str, sha: str) -> None:
print(f" {repo}@{sha[:7]} {ctx} = {state} -> {st}")
GUARD_CTX = "windy-git/compute-guard"
def post_compute_guard(repo: str, sha: str, default_branch: str, is_default_head: bool) -> None:
"""Windy Mind is the only door to AI compute: flag direct provider use (warn-only).
Non-fatal and never a fake OK: if the guard can't run, nothing is posted.
"""
try:
import compute_guard as cg # same directory; loaded lazily so the bridge never depends on it
findings = cg.check(repo, sha, default_branch, is_default_head)
except Exception as e: # noqa: BLE001 — the guard must never break CI signals
print(f" {repo}@{sha[:7]} compute-guard skipped ({type(e).__name__}: {str(e)[:80]})")
return
if findings is None:
return
state, desc, first = cg.status_for(findings, whole_tree=is_default_head)
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
for s in existing or []: # newest first: compare the latest guard status only
if s["context"] == GUARD_CTX:
if (s["state"], s.get("description")) == (state, desc):
return
break
url = (f"{PUBLIC}/{WG_OWNER}/{repo}/src/commit/{sha}/{first.path}#L{first.line}"
if first else f"{PUBLIC}/{WG_OWNER}/{repo}/src/commit/{sha}")
st, _ = github("POST", f"/repos/{GH_OWNER}/{repo}/statuses/{sha}",
{"state": state, "context": GUARD_CTX, "description": desc, "target_url": url})
print(f" {repo}@{sha[:7]} {GUARD_CTX} = {state} ({len(findings)} finding(s)) -> {st}")
def main() -> int:
if not (GITEA_TOKEN and GITHUB_TOKEN):
sys.exit("GITEA_ADMIN_TOKEN and GITHUB_TOKEN are required")
@@ -205,13 +381,17 @@ def main() -> int:
for repo in REPOS:
try:
shas = sync_prs(repo)
default_branch, default_head = "main", None
st, br = github("GET", f"/repos/{GH_OWNER}/{repo}")
if st == 200:
st, b = github("GET", f"/repos/{GH_OWNER}/{repo}/branches/{br['default_branch']}")
default_branch = br["default_branch"]
st, b = github("GET", f"/repos/{GH_OWNER}/{repo}/branches/{default_branch}")
if st == 200:
shas.append(b["commit"]["sha"])
default_head = b["commit"]["sha"]
shas.append(default_head)
for sha in dict.fromkeys(shas):
post_statuses(repo, sha)
post_compute_guard(repo, sha, default_branch, sha == default_head)
except Exception as e: # one repo's failure must not hide the others'
print(f" FAILED {repo}: {e}")
failed = 1

49
scripts/rerun_ci.sh Executable file
View File

@@ -0,0 +1,49 @@
#!/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}^")
# NOT `systemctl is-active`: the sync is Type=oneshot, which reads "activating"
# (exit 3) for its whole run, so is-active says "idle" mid-run.
busy() { case "$(systemctl show windygit-sync -p ActiveState --value)" in
activating|active|deactivating|reloading) return 0;; esac; return 1; }
while busy; 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

@@ -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.
@@ -90,5 +94,9 @@ if ! python3 "$(dirname "$0")/pr_status_bridge.py"; then
log "FAILED pr status bridge"; FAILED=1
fi
# CI telemetry -> admin.windyword.ai (shapes declared with Windy Telemetry 40).
# Sends nothing until WINDYGIT_TELEMETRY_TOKEN is set; never fails the sync.
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"

385
scripts/telemetry_emit.py Normal file
View File

@@ -0,0 +1,385 @@
#!/usr/bin/env python3
"""Emit Windy Git CI telemetry to admin.windyword.ai (Windy Telemetry 40's ledger).
Runs on Veron after every sync (root; reads the gitea DB via `docker exec`).
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 — 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
messages, no logs, no author names.
--dry-run print the batch instead of posting (and don't advance STATE)
"""
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 UTC, datetime
INGEST = os.environ.get("TELEMETRY_INGEST_URL", "https://admin.windyword.ai/v1/events")
TOKEN = os.environ.get("WINDYGIT_TELEMETRY_TOKEN", "")
STATE = os.environ.get("TELEMETRY_STATE", "/var/lib/windy-git/telemetry-state.json")
PLATFORM, SERVICE = "windy-git", "ci"
OUTCOME = {1: "success", 2: "failure", 3: "cancelled", 4: "skipped"}
EVENTS = {"push", "pull_request", "pull_request_sync", "schedule", "workflow_dispatch"}
RUNNERS_EXPECTED = 6
def sql(query: str) -> list[dict]:
"""Rows as dicts, via psql's json_agg — no driver needed on the host."""
wrapped = f"select coalesce(json_agg(t), '[]'::json) from ({query}) t;"
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=wrapped,
capture_output=True,
text=True,
check=True,
).stdout.strip()
return json.loads(out or "[]")
def load_state() -> dict:
try:
with open(STATE) as f:
return json.load(f)
except (OSError, ValueError):
return {}
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()
now = time.time()
since = float(state.get("last_ts", now - 300))
# 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, {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.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 ""
if ref.startswith("refs/pull/"):
kind = "pr"
elif ref == f"refs/heads/{j['default_branch']}":
kind = "default"
else:
kind = "other"
# 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,
"service": SERVICE,
"event_type": "ci.run",
"actor_type": "system",
"metadata": {
"repo": j["repo"],
"workflow": (j["workflow_id"] or "").removesuffix(".yml").removesuffix(".yaml"),
"job": j["job"],
"outcome": OUTCOME[j["status"]],
"event": j["event"] if j["event"] in EVENTS else "other",
"branch_kind": kind,
"run": j["run"],
"sha": j["sha"],
},
}
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 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
meta[k] = int(v)
events.append(
{
"ts": iso(now),
"platform": PLATFORM,
"service": SERVICE,
"event_type": "service.health",
"actor_type": "system",
"metadata": meta,
}
)
# 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:
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,
data=json.dumps({"events": events[i : i + 500]}).encode(),
method="POST",
headers={
"Authorization": f"Bearer {TOKEN}",
"Content-Type": "application/json",
"User-Agent": "windy-git-telemetry/1",
},
)
try:
with urllib.request.urlopen(req, timeout=30) as r:
body = r.read()
if r.status >= 300:
raise urllib.error.HTTPError(
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
except urllib.error.URLError as e:
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_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,
"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
if __name__ == "__main__":
sys.exit(main())