Compare commits
26 Commits
5b16114b98
...
telemetry-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e3b69fa759 | ||
|
|
4acf50d9ef | ||
|
|
b7a7e94df0 | ||
| c83f808a60 | |||
| eb27e63db3 | |||
| 1c3b5b0638 | |||
| 90643fe48e | |||
| 00ec963f82 | |||
| b5e4eaf57a | |||
| baaa542bae | |||
| 0634a6cb1b | |||
| 1a171eabd5 | |||
| 28c31236b8 | |||
| cd5967031b | |||
| f246417095 | |||
| 50c1464043 | |||
| 7a63f90da3 | |||
| b8f97f0731 | |||
| 95c33c8004 | |||
| d5181f1c6d | |||
| e6530d3171 | |||
| 419443573a | |||
| 4a34b35441 | |||
| dfe5543eda | |||
| 40cb455d0d | |||
| 8c404eb410 |
14
AGENTS.md
14
AGENTS.md
@@ -2,11 +2,17 @@
|
|||||||
|
|
||||||
Read this before touching anything. Then read `DNA_STRAND_MASTER_PLAN.md`, which is the source of truth.
|
Read this before touching anything. Then read `DNA_STRAND_MASTER_PLAN.md`, which is the source of truth.
|
||||||
|
|
||||||
## Current state
|
## Current state (2026-09-23)
|
||||||
|
|
||||||
**GENESIS.** No code. No `make dev` yet — building it is codon **G0.8**.
|
**LIVE on Veron 1** — `app.windygit.com` (Gitea 1.24.6, Windy SSO only),
|
||||||
|
`api.windygit.com` (our plane: humans via hub JWKS, agents via Eternitas EPT),
|
||||||
|
`models.windygit.com`. Strands G0–G5, G7, G11 done; see the plan for the rest.
|
||||||
|
|
||||||
The next work is Strand **G0** (cell substrate), then **G1** (Veron 1 host + Cloudflare Tunnel), then **G2** (Gitea, stock and branded), then **G3** (identity), then **G4** (storage). G0–G4 are sequential. G5–G12 are concurrent once G4 lands.
|
It is also **the permanent CI for the private platform repos** (GitHub Actions
|
||||||
|
cannot run on them): `scripts/sync_from_github.sh` + `scripts/pr_status_bridge.py`,
|
||||||
|
onboarding in `docs/CUTOVER.md`, operations in `docs/RUNBOOK-VERON.md`.
|
||||||
|
|
||||||
|
Standing dev checkout: **OC5 `~/windy-git`**. Deploy copy: Veron `/srv/windygit/src`.
|
||||||
|
|
||||||
## The rules that will get you reverted if you break them
|
## The rules that will get you reverted if you break them
|
||||||
|
|
||||||
@@ -28,7 +34,7 @@ The next work is Strand **G0** (cell substrate), then **G1** (Veron 1 host + Clo
|
|||||||
- Errors are 4-field repair pointers: `{code, speak, machine_cause, remediation_tool}`. No exceptions, including validation errors.
|
- Errors are 4-field repair pointers: `{code, speak, machine_cause, remediation_tool}`. No exceptions, including validation errors.
|
||||||
- Every tool response carries `state_proof` + `next_actions`.
|
- Every tool response carries `state_proof` + `next_actions`.
|
||||||
- Telemetry `actor_type` comes from the enum `{human, agent, system}`. **`'service'` is not legal** — it 422s and silently drops the whole batch. A sibling service is losing telemetry to exactly this today.
|
- Telemetry `actor_type` comes from the enum `{human, agent, system}`. **`'service'` is not legal** — it 422s and silently drops the whole batch. A sibling service is losing telemetry to exactly this today.
|
||||||
- Runner labels are explicit and pinned. **`ubuntu-latest` is banned** — all four `windy-registry` workflows use it and every run fails.
|
- Runner labels are explicit and pinned: `[self-hosted, linux, x64]` or `veron-1`. **`ubuntu-latest` is banned** — no runner here has it, so the job queues forever.
|
||||||
|
|
||||||
## Membrane
|
## Membrane
|
||||||
|
|
||||||
|
|||||||
@@ -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.
|
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.**
|
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).
|
**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 out:** `repo.created`, `repo.pushed`, `release.published`, `model.published`, `ci.completed`.
|
||||||
**Events in:** `passport.revoked` (fail-closed), `storage.quota.exceeded`, `identity.created`.
|
**Events in:** `passport.revoked` (fail-closed), `storage.quota.exceeded`, `identity.created`.
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ from api.app.config import Settings
|
|||||||
from api.app.ept import EptInvalid, looks_like_ept, verify_ept
|
from api.app.ept import EptInvalid, looks_like_ept, verify_ept
|
||||||
from api.app.errors import RepairPointer, passport_unresolvable
|
from api.app.errors import RepairPointer, passport_unresolvable
|
||||||
from api.app.hub_jwt import HubTokenInvalid, verify_hub_token
|
from api.app.hub_jwt import HubTokenInvalid, verify_hub_token
|
||||||
|
from api.app.telemetry import synthetic_headers
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
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}"
|
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
|
last_status = 0
|
||||||
for attempt in range(3):
|
for attempt in range(3):
|
||||||
async with httpx.AsyncClient(timeout=httpx.Timeout(8.0, connect=3.0)) as client:
|
async with httpx.AsyncClient(timeout=httpx.Timeout(8.0, connect=3.0)) as client:
|
||||||
|
|||||||
@@ -66,6 +66,12 @@ class Settings(BaseSettings):
|
|||||||
# Flip to True once the hub emits aud on every access token.
|
# Flip to True once the hub emits aud on every access token.
|
||||||
hub_require_aud: bool = False
|
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
|
# Internal callers (the Cloud portal calling /internal/*). A first-class
|
||||||
# caller class, not a bypass: unset means service calls are REFUSED.
|
# caller class, not a bypass: unset means service calls are REFUSED.
|
||||||
service_token: str = ""
|
service_token: str = ""
|
||||||
|
|||||||
@@ -6,11 +6,13 @@ component and is reached only over its REST API (D-2 / I-1).
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
import socket
|
import socket
|
||||||
|
import time
|
||||||
from contextlib import asynccontextmanager
|
from contextlib import asynccontextmanager
|
||||||
|
|
||||||
from fastapi import FastAPI
|
from fastapi import FastAPI, Request
|
||||||
from fastapi.exceptions import RequestValidationError
|
from fastapi.exceptions import RequestValidationError
|
||||||
from fastapi.responses import JSONResponse
|
from fastapi.responses import JSONResponse
|
||||||
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
|
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
|
||||||
@@ -25,6 +27,7 @@ from api.app.providers.registry import (
|
|||||||
R2Provider,
|
R2Provider,
|
||||||
)
|
)
|
||||||
from api.app.routes import health, repos, webhooks
|
from api.app.routes import health, repos, webhooks
|
||||||
|
from api.app.telemetry import SYNTHETIC, Telemetry, caller_class, is_synthetic
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
@@ -44,9 +47,7 @@ def _refuse_kit_zero(settings) -> None:
|
|||||||
if not settings.is_production:
|
if not settings.is_production:
|
||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
local_ips = {
|
local_ips = {info[4][0] for info in socket.getaddrinfo(socket.gethostname(), None)}
|
||||||
info[4][0] for info in socket.getaddrinfo(socket.gethostname(), None)
|
|
||||||
}
|
|
||||||
except socket.gaierror:
|
except socket.gaierror:
|
||||||
return
|
return
|
||||||
if settings.kit0_host in local_ips:
|
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`.
|
# 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
|
yield
|
||||||
|
|
||||||
|
if task is not None:
|
||||||
|
task.cancel()
|
||||||
|
await telemetry.flush()
|
||||||
if engine is not None:
|
if engine is not None:
|
||||||
await engine.dispose()
|
await engine.dispose()
|
||||||
|
|
||||||
@@ -119,8 +135,44 @@ app.include_router(repos.router)
|
|||||||
app.include_router(webhooks.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)
|
@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)
|
return JSONResponse(status_code=exc.status_code, content=exc.detail)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ import httpx
|
|||||||
|
|
||||||
from api.app.config import Settings
|
from api.app.config import Settings
|
||||||
from api.app.errors import RepairPointer, provider_unconfigured
|
from api.app.errors import RepairPointer, provider_unconfigured
|
||||||
|
from api.app.telemetry import synthetic_headers
|
||||||
|
|
||||||
_TIMEOUT = httpx.Timeout(20.0, connect=5.0)
|
_TIMEOUT = httpx.Timeout(20.0, connect=5.0)
|
||||||
|
|
||||||
@@ -36,6 +37,7 @@ class GiteaClient:
|
|||||||
return {
|
return {
|
||||||
"Authorization": f"token {self._s.gitea_admin_token}",
|
"Authorization": f"token {self._s.gitea_admin_token}",
|
||||||
"Content-Type": "application/json",
|
"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:
|
async def _request(self, method: str, path: str, **kw: Any) -> httpx.Response:
|
||||||
|
|||||||
255
api/app/telemetry.py
Normal file
255
api/app/telemetry.py
Normal 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()
|
||||||
@@ -734,3 +734,21 @@ def test_g23_brand_css_filename_is_versioned():
|
|||||||
assert m, "brand CSS must carry a version in its FILENAME"
|
assert m, "brand CSS must carry a version in its FILENAME"
|
||||||
assert (ROOT / "deploy" / "branding" / "public" / "assets" / "css"
|
assert (ROOT / "deploy" / "branding" / "public" / "assets" / "css"
|
||||||
/ f"theme-windy.v{m.group(1)}.css").exists()
|
/ f"theme-windy.v{m.group(1)}.css").exists()
|
||||||
|
|
||||||
|
|
||||||
|
def test_backup_never_bundles_credential_repos_to_r2():
|
||||||
|
"""kit-army-config (the lockbox) and the *-soul / anima repos carry
|
||||||
|
credentials; the R2 bundles are plaintext. Behavioural: run the script's
|
||||||
|
own exclusion function against the names."""
|
||||||
|
import subprocess
|
||||||
|
|
||||||
|
script = (ROOT / "scripts" / "backup.sh").read_text()
|
||||||
|
fn = script[script.index('EXCLUDE="'):script.index("cleanup()")]
|
||||||
|
# A file named like a pattern in cwd must not break the match (glob expansion).
|
||||||
|
probe = "cd \"$(mktemp -d)\" && touch x-soul && " + fn + (
|
||||||
|
'for n in kit-army-config anima windy-0-soul kit-0c5-soul herm-0-soul '
|
||||||
|
'soulsafe windy-chat eternitas; do excluded "$n" && echo "X $n" || echo "- $n"; done'
|
||||||
|
)
|
||||||
|
out = subprocess.run(["bash", "-c", probe], capture_output=True, text=True, check=True).stdout
|
||||||
|
skipped = {ln[2:] for ln in out.splitlines() if ln.startswith("X ")}
|
||||||
|
assert skipped == {"kit-army-config", "anima", "windy-0-soul", "kit-0c5-soul", "herm-0-soul"}
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ no status at all.
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import base64
|
||||||
import importlib.util
|
import importlib.util
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -35,12 +36,23 @@ def _run(i, wf, job, status, sha=SHA, n=1):
|
|||||||
|
|
||||||
|
|
||||||
class Fake:
|
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.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.gh_prs, self.wg_prs = list(gh_prs), list(wg_prs)
|
||||||
self.posted, self.opened, self.closed = [], [], []
|
self.posted, self.opened, self.closed = [], [], []
|
||||||
|
|
||||||
def gitea(self, method, path, body=None):
|
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:
|
if "/actions/tasks" in path:
|
||||||
page = int(path.rsplit("page=", 1)[1])
|
page = int(path.rsplit("page=", 1)[1])
|
||||||
return 200, {"workflow_runs": self.runs[(page - 1) * 50 : page * 50]}
|
return 200, {"workflow_runs": self.runs[(page - 1) * 50 : page * 50]}
|
||||||
@@ -147,3 +159,116 @@ def test_image_build_jobs_are_not_posted(fake):
|
|||||||
)
|
)
|
||||||
bridge.post_statuses("r", SHA)
|
bridge.post_statuses("r", SHA)
|
||||||
assert f.posted == []
|
assert f.posted == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_non_blocking_jobs_are_not_posted_for_that_repo_only(fake, monkeypatch):
|
||||||
|
"""Grant ruled windy-pro's desktop/installer jobs non-blocking: they must not
|
||||||
|
reach GitHub for windy-pro, and the rule must not leak to other repos."""
|
||||||
|
monkeypatch.setattr(bridge, "NON_BLOCKING", {"windy-pro": {"ci/build-desktop"}})
|
||||||
|
runs = [_run(1, "ci.yml", "build-desktop", "failure"), _run(2, "ci.yml", "test", "success")]
|
||||||
|
f = fake(runs=runs)
|
||||||
|
bridge.post_statuses("windy-pro", SHA)
|
||||||
|
assert [p["context"] for p in f.posted] == ["windy-git/ci/test"]
|
||||||
|
f2 = fake(runs=runs)
|
||||||
|
bridge.post_statuses("windy-chat", SHA)
|
||||||
|
assert sorted(p["context"] for p in f2.posted) == [
|
||||||
|
"windy-git/ci/build-desktop",
|
||||||
|
"windy-git/ci/test",
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_default_non_blocking_is_grants_ruling():
|
||||||
|
assert bridge.NON_BLOCKING.get("windy-pro") == {
|
||||||
|
"ci/build-desktop",
|
||||||
|
"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
|
||||||
|
|||||||
202
api/tests/test_telemetry.py
Normal file
202
api/tests/test_telemetry.py
Normal 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
|
||||||
106
api/tests/test_webhooks_behavior.py
Normal file
106
api/tests/test_webhooks_behavior.py
Normal file
@@ -0,0 +1,106 @@
|
|||||||
|
"""G3.5 — the Eternitas webhook receiver, driven over HTTP (audit 2026-08-13).
|
||||||
|
|
||||||
|
The G3.5 invariants in test_invariants.py grep webhooks.py for strings; a
|
||||||
|
refactor that kept the strings and broke the behaviour would pass them all.
|
||||||
|
These send real requests through the real route (no DB: every case here stops
|
||||||
|
before the revocation handler) and assert what the receiver DOES.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import hashlib
|
||||||
|
import hmac
|
||||||
|
import json
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
import pytest
|
||||||
|
from fastapi import FastAPI
|
||||||
|
from fastapi.responses import JSONResponse
|
||||||
|
|
||||||
|
from api.app.config import Settings
|
||||||
|
from api.app.errors import RepairPointer
|
||||||
|
from api.app.routes import webhooks
|
||||||
|
|
||||||
|
SECRET = "s" * 64
|
||||||
|
URL = "/api/v1/webhooks/eternitas"
|
||||||
|
|
||||||
|
|
||||||
|
def _app(secret: str = SECRET) -> FastAPI:
|
||||||
|
app = FastAPI()
|
||||||
|
app.include_router(webhooks.router)
|
||||||
|
app.state.settings = Settings(eternitas_webhook_secret=secret)
|
||||||
|
|
||||||
|
@app.exception_handler(RepairPointer)
|
||||||
|
async def _h(_, exc: RepairPointer) -> JSONResponse:
|
||||||
|
return JSONResponse(status_code=exc.status_code, content=exc.detail)
|
||||||
|
|
||||||
|
return app
|
||||||
|
|
||||||
|
|
||||||
|
async def _post(body: bytes, headers: dict, secret: str = SECRET) -> httpx.Response:
|
||||||
|
transport = httpx.ASGITransport(app=_app(secret))
|
||||||
|
async with httpx.AsyncClient(transport=transport, base_url="http://t") as c:
|
||||||
|
return await c.post(
|
||||||
|
URL, content=body, headers={"content-type": "application/json", **headers}
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _sig(raw: bytes, secret: str = SECRET) -> str:
|
||||||
|
return hmac.new(secret.encode(), raw, hashlib.sha256).hexdigest()
|
||||||
|
|
||||||
|
|
||||||
|
# Deliberately odd spacing/key order: a receiver that re-serialises before
|
||||||
|
# hashing produces a different digest and must fail.
|
||||||
|
RAW = b'{"event":"windygit.selftest", "data": {"b": 2, "a": 1}}'
|
||||||
|
EVENT = {"x-eternitas-event": "windygit.selftest"}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_prefixed_and_bare_digests_are_both_accepted():
|
||||||
|
for header in (f"sha256={_sig(RAW)}", _sig(RAW)):
|
||||||
|
r = await _post(RAW, {**EVENT, "x-eternitas-signature": header})
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
assert r.json()["acted"] is False # unknown event: received, nothing done
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_digest_of_reserialised_json_is_refused():
|
||||||
|
reserialised = json.dumps(json.loads(RAW)).encode()
|
||||||
|
assert reserialised != RAW
|
||||||
|
r = await _post(RAW, {**EVENT, "x-eternitas-signature": f"sha256={_sig(reserialised)}"})
|
||||||
|
assert r.status_code == 401 and r.json()["code"] == "webhook_signature_invalid"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_forged_or_wrong_key_signature_is_refused():
|
||||||
|
for header in ("sha256=" + "0" * 64, f"sha256={_sig(RAW, 'other-secret')}", "garbage"):
|
||||||
|
r = await _post(RAW, {**EVENT, "x-eternitas-signature": header})
|
||||||
|
assert r.status_code == 401, header
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_signed_event_without_signature_is_refused():
|
||||||
|
r = await _post(RAW, EVENT)
|
||||||
|
assert r.status_code == 401
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_unset_secret_refuses_rather_than_accepts():
|
||||||
|
r = await _post(RAW, {**EVENT, "x-eternitas-signature": f"sha256={_sig(RAW)}"}, secret="")
|
||||||
|
assert r.status_code == 503 and r.json()["code"] == "webhook_secret_unset"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_revocation_with_bad_signature_never_reaches_the_handler():
|
||||||
|
body = b'{"event":"passport.revoked","passport":"ET26-TEST-GOOD"}'
|
||||||
|
r = await _post(
|
||||||
|
body,
|
||||||
|
{"x-eternitas-event": "passport.revoked", "x-eternitas-signature": "sha256=" + "f" * 64},
|
||||||
|
)
|
||||||
|
assert r.status_code == 401 # refused before any DB work
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_reachability_ping_acknowledges_but_never_acts():
|
||||||
|
r = await _post(b'{"anything": "at all"}', {"x-eternitas-event": "platform.test_ping"})
|
||||||
|
assert r.status_code == 200 and r.json()["acted"] is False
|
||||||
51
deploy/runner/egress.sh
Executable file
51
deploy/runner/egress.sh
Executable file
@@ -0,0 +1,51 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# CI egress filter (2026-09-23) — jobs reach the internet, never Grant's network.
|
||||||
|
#
|
||||||
|
# Measured before this existed: an ordinary (unprivileged) job container inside
|
||||||
|
# the CI dind could open SSH, Ollama, and every dev server on Veron
|
||||||
|
# (192.168.1.73:22/3000/3300/8080/11434) and anything else on the LAN, WireGuard
|
||||||
|
# or Tailscale. No container escape needed — a malicious npm/pip dependency in
|
||||||
|
# any first-party repo's CI could walk straight onto the fleet.
|
||||||
|
#
|
||||||
|
# All CI traffic leaves through the `windy-git-runner_jobs` bridge (dind NATs
|
||||||
|
# its job containers onto it). This script, run at boot and after any runner
|
||||||
|
# compose change, allows on that bridge:
|
||||||
|
# * traffic between the runners and dind (same bridge)
|
||||||
|
# * replies (ESTABLISHED/RELATED)
|
||||||
|
# * DNS (53) — Docker's embedded resolver forwards to the LAN router
|
||||||
|
# * everything public
|
||||||
|
# and drops: RFC1918, CGNAT/Tailscale (100.64/10), link-local, and ANY packet
|
||||||
|
# addressed to the host itself (INPUT), whatever interface IP it targets.
|
||||||
|
# Idempotent: owned chains are flushed and rebuilt; hooks are added once.
|
||||||
|
set -euo pipefail
|
||||||
|
|
||||||
|
NET=windy-git-runner_jobs
|
||||||
|
id=$(docker network inspect "$NET" --format '{{.Id}}')
|
||||||
|
BR="br-${id:0:12}"
|
||||||
|
ip link show "$BR" >/dev/null
|
||||||
|
|
||||||
|
iptables -N WG-CI-EGRESS 2>/dev/null || iptables -F WG-CI-EGRESS
|
||||||
|
iptables -A WG-CI-EGRESS -o "$BR" -j RETURN
|
||||||
|
iptables -A WG-CI-EGRESS -m conntrack --ctstate ESTABLISHED,RELATED -j RETURN
|
||||||
|
iptables -A WG-CI-EGRESS -p udp --dport 53 -j RETURN
|
||||||
|
iptables -A WG-CI-EGRESS -p tcp --dport 53 -j RETURN
|
||||||
|
for cidr in 10.0.0.0/8 172.16.0.0/12 192.168.0.0/16 100.64.0.0/10 169.254.0.0/16; do
|
||||||
|
iptables -A WG-CI-EGRESS -d "$cidr" -j DROP
|
||||||
|
done
|
||||||
|
iptables -A WG-CI-EGRESS -j RETURN
|
||||||
|
|
||||||
|
iptables -N WG-CI-INPUT 2>/dev/null || iptables -F WG-CI-INPUT
|
||||||
|
iptables -A WG-CI-INPUT -m conntrack --ctstate ESTABLISHED,RELATED -j RETURN
|
||||||
|
iptables -A WG-CI-INPUT -j DROP
|
||||||
|
|
||||||
|
# Hooks: remove any stale ones (the bridge name changes if the network is
|
||||||
|
# recreated), then add exactly one of each at the top.
|
||||||
|
for chain in DOCKER-USER INPUT; do
|
||||||
|
target=$([ "$chain" = INPUT ] && echo WG-CI-INPUT || echo WG-CI-EGRESS)
|
||||||
|
while read -r rule; do
|
||||||
|
iptables -D $chain ${rule#-A $chain }
|
||||||
|
done < <(iptables -S "$chain" | grep -- "-j $target" || true)
|
||||||
|
iptables -I "$chain" 1 -i "$BR" -j "$target"
|
||||||
|
done
|
||||||
|
|
||||||
|
echo "ci egress filter active on $BR ($NET)"
|
||||||
13
deploy/runner/windygit-ci-egress.service
Normal file
13
deploy/runner/windygit-ci-egress.service
Normal file
@@ -0,0 +1,13 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Windy Git - CI egress filter (jobs reach the internet, never the LAN/host)
|
||||||
|
After=docker.service
|
||||||
|
Requires=docker.service
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=oneshot
|
||||||
|
RemainAfterExit=yes
|
||||||
|
# The jobs network exists once the runner compose project is up; retry until it does.
|
||||||
|
ExecStart=/bin/bash -c 'for i in $(seq 1 60); do /srv/windygit/src/deploy/runner/egress.sh && exit 0; sleep 5; done; exit 1'
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=multi-user.target
|
||||||
13
deploy/systemd/windygit-backup.service
Normal file
13
deploy/systemd/windygit-backup.service
Normal file
@@ -0,0 +1,13 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Windy Git nightly backup (git bundles + windgit schema -> R2)
|
||||||
|
After=network-online.target docker.service
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=oneshot
|
||||||
|
WorkingDirectory=/srv/windygit/src
|
||||||
|
# The .env holds the R2 credentials. The script refuses to run without them
|
||||||
|
# rather than reporting a backup that did not happen.
|
||||||
|
EnvironmentFile=/srv/windygit/src/.env
|
||||||
|
ExecStart=/bin/bash /srv/windygit/src/scripts/backup.sh
|
||||||
|
Nice=10
|
||||||
|
IOSchedulingClass=idle
|
||||||
3
deploy/systemd/windygit-backup.service.d/windy-job.conf
Normal file
3
deploy/systemd/windygit-backup.service.d/windy-job.conf
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
[Service]
|
||||||
|
ExecStart=
|
||||||
|
ExecStart=/usr/local/bin/windy-job windygit-backup 26h --expect "ok — [0-9]+ repos" --owner 13 -- /bin/bash /srv/windygit/src/scripts/backup.sh
|
||||||
12
deploy/systemd/windygit-backup.timer
Normal file
12
deploy/systemd/windygit-backup.timer
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Nightly Windy Git backup
|
||||||
|
|
||||||
|
[Timer]
|
||||||
|
OnCalendar=*-*-* 04:17:00
|
||||||
|
# Grant's workstation is not always on at 04:17. Without this a missed window
|
||||||
|
# is simply skipped and the backup silently never runs.
|
||||||
|
Persistent=true
|
||||||
|
RandomizedDelaySec=600
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=timers.target
|
||||||
@@ -0,0 +1,3 @@
|
|||||||
|
[Service]
|
||||||
|
ExecStart=
|
||||||
|
ExecStart=/usr/local/bin/windy-job windygit-ci-prune 7h --expect "ci storage [0-9]+G" --owner 13 -- /srv/windygit/src/deploy/runner/prune.sh
|
||||||
12
deploy/systemd/windygit-sync.service
Normal file
12
deploy/systemd/windygit-sync.service
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Sync GitHub -> Windy Git (Phase 1: GitHub is the source of truth)
|
||||||
|
After=network-online.target docker.service
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=oneshot
|
||||||
|
WorkingDirectory=/srv/windygit/src
|
||||||
|
EnvironmentFile=/srv/windygit/src/.env
|
||||||
|
# GITHUB_TOKEN is set on the host only (root-only unit file / .env) — NEVER commit it.
|
||||||
|
Environment=GITHUB_OWNER=sneakyfree
|
||||||
|
ExecStart=/bin/bash /srv/windygit/src/scripts/sync_from_github.sh
|
||||||
|
Nice=10
|
||||||
3
deploy/systemd/windygit-sync.service.d/windy-job.conf
Normal file
3
deploy/systemd/windygit-sync.service.d/windy-job.conf
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
[Service]
|
||||||
|
ExecStart=
|
||||||
|
ExecStart=/usr/local/bin/windy-job windygit-sync 20m --expect "all repos in step with GitHub" --owner 13 -- /bin/bash /srv/windygit/src/scripts/sync_from_github.sh
|
||||||
10
deploy/systemd/windygit-sync.timer
Normal file
10
deploy/systemd/windygit-sync.timer
Normal file
@@ -0,0 +1,10 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Keep Windy Git in step with GitHub every 5 minutes
|
||||||
|
|
||||||
|
[Timer]
|
||||||
|
OnBootSec=3min
|
||||||
|
OnUnitActiveSec=5min
|
||||||
|
Persistent=true
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=timers.target
|
||||||
16
deploy/systemd/windygit-tunnel.service
Normal file
16
deploy/systemd/windygit-tunnel.service
Normal file
@@ -0,0 +1,16 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Windy Git - Cloudflare Tunnel (the only ingress; no inbound port is opened)
|
||||||
|
After=network-online.target
|
||||||
|
Wants=network-online.target
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=notify
|
||||||
|
ExecStart=/usr/bin/cloudflared --no-autoupdate --config /etc/cloudflared/config.yml tunnel run
|
||||||
|
Restart=always
|
||||||
|
RestartSec=5
|
||||||
|
# G1.4 - bounded, so a misbehaving ingress can never starve Grant's workstation.
|
||||||
|
MemoryMax=512M
|
||||||
|
CPUQuota=100%
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=multi-user.target
|
||||||
@@ -16,7 +16,11 @@ services:
|
|||||||
# I-12: baked at build time. A runtime COMMIT_SHA override is ignored.
|
# I-12: baked at build time. A runtime COMMIT_SHA override is ignored.
|
||||||
COMMIT_SHA: ${COMMIT_SHA_BUILD:-}
|
COMMIT_SHA: ${COMMIT_SHA_BUILD:-}
|
||||||
BUILT_AT: ${BUILT_AT:-}
|
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:
|
environment:
|
||||||
DATABASE_URL: postgresql+asyncpg://windygit:${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD}@db:5432/windygit
|
DATABASE_URL: postgresql+asyncpg://windygit:${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD}@db:5432/windygit
|
||||||
GITEA_BASE_URL: http://gitea:3000
|
GITEA_BASE_URL: http://gitea:3000
|
||||||
@@ -66,7 +70,13 @@ services:
|
|||||||
# G3.1 — a Windy account IS the account. Signing in with Windy provisions
|
# G3.1 — a Windy account IS the account. Signing in with Windy provisions
|
||||||
# the Gitea user on first arrival; nobody is asked to invent a second
|
# the Gitea user on first arrival; nobody is asked to invent a second
|
||||||
# identity for the same person, and no local password ever exists.
|
# identity for the same person, and no local password ever exists.
|
||||||
GITEA__oauth2_client__ENABLE_AUTO_REGISTRATION: "true"
|
# 🔴 OFF (2026-09-23). With it on, ANY stranger with a Windy Word account
|
||||||
|
# (public signup, not even email-verified) got a forge account on first
|
||||||
|
# sign-in — and the CI runners were instance-wide, so their workflows
|
||||||
|
# would run on Veron beside the R2 god token. Proven with a throwaway
|
||||||
|
# account, then closed. Opening the forge to non-Grant users is a §7
|
||||||
|
# Grant decision; until then new accounts are created deliberately.
|
||||||
|
GITEA__oauth2_client__ENABLE_AUTO_REGISTRATION: "false"
|
||||||
GITEA__oauth2_client__USERNAME: email
|
GITEA__oauth2_client__USERNAME: email
|
||||||
# 🔴 `login`, NOT `auto` (SSO #8). `auto` linked any hub login whose EMAIL
|
# 🔴 `login`, NOT `auto` (SSO #8). `auto` linked any hub login whose EMAIL
|
||||||
# matched an existing account — and windyadmin (SITE ADMIN) carries Grant's
|
# matched an existing account — and windyadmin (SITE ADMIN) carries Grant's
|
||||||
|
|||||||
@@ -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 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
|
the security-critical paths**, so verification persists instead of living in a
|
||||||
transcript.
|
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.
|
||||||
|
|||||||
@@ -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/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 |
|
| 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) |
|
| 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 |
|
| windy-cloud-sites | `POST /api/v1/sites/{id}/versions` | publish docs from a repo |
|
||||||
|
|
||||||
## Calls IN
|
## Calls IN
|
||||||
|
|||||||
@@ -1,6 +1,10 @@
|
|||||||
# RUNBOOK — Windy Git on Veron 1 (rung R0)
|
# RUNBOOK — Windy Git on Veron 1 (rung R0)
|
||||||
|
|
||||||
Host `Veron-1-5090`, WireGuard `10.10.0.6`, alias `wg-veron`. Passwordless sudo.
|
Host `Veron-1-5090`, WireGuard `10.10.0.6`, alias `wg-veron` (or `ts-veron`). Passwordless sudo.
|
||||||
|
|
||||||
|
**Checkouts (one-repo doctrine):** the ONE standing dev checkout is **OC5
|
||||||
|
`~/windy-git`** (platform repos live on OC5). `/srv/windygit/src` on Veron is the
|
||||||
|
*deploy* copy — it holds no local work. Nothing else should exist.
|
||||||
|
|
||||||
⛔ **Kit 0 is never a host for this service** (D-4). `api/app/main.py` refuses to
|
⛔ **Kit 0 is never a host for this service** (D-4). `api/app/main.py` refuses to
|
||||||
boot in production if it finds itself on `72.60.118.54`.
|
boot in production if it finds itself on `72.60.118.54`.
|
||||||
@@ -15,6 +19,9 @@ boot in production if it finds itself on `72.60.118.54`.
|
|||||||
| `/etc/cloudflared/config.yml` | tunnel ingress |
|
| `/etc/cloudflared/config.yml` | tunnel ingress |
|
||||||
| `/etc/cloudflared/windy-git.json` | tunnel credentials, mode 600 |
|
| `/etc/cloudflared/windy-git.json` | tunnel credentials, mode 600 |
|
||||||
| `/srv/windygit/src/.env` | secrets, mode 600, **never committed** |
|
| `/srv/windygit/src/.env` | secrets, mode 600, **never committed** |
|
||||||
|
| `/srv/windygit/git/gitea/conf/app.ini` | Gitea's persisted config — env-to-ini SETS but never UNSETS; edit here when removing a `GITEA__*` var |
|
||||||
|
| `/srv/windygit/sync/*.git` | bare staging copies the GitHub→Windy Git sync pushes from |
|
||||||
|
| `/srv/windygit/src/deploy/runner/.env` | `RUNNER_TOKEN` — a **windyadmin user-level** registration token (not instance-level; see CI) |
|
||||||
|
|
||||||
## Ports — all loopback, on purpose
|
## Ports — all loopback, on purpose
|
||||||
|
|
||||||
@@ -41,12 +48,15 @@ sudo systemctl status windygit-tunnel
|
|||||||
|
|
||||||
```bash
|
```bash
|
||||||
ssh wg-veron
|
ssh wg-veron
|
||||||
cd /srv/windygit/src && git pull
|
cd /srv/windygit/src && git fetch origin && git merge --ff-only origin/main # READ the output
|
||||||
export COMMIT_SHA_BUILD=$(git rev-parse HEAD) BUILT_AT=$(date -u +%Y-%m-%dT%H:%M:%SZ)
|
export COMMIT_SHA_BUILD=$(git rev-parse HEAD) BUILT_AT=$(date -u +%Y-%m-%dT%H:%M:%SZ)
|
||||||
sudo -E docker compose up -d --build
|
sudo -E docker compose up -d --build --no-deps api # API only: no forge restart
|
||||||
curl -s https://api.windygit.com/version # MUST equal git rev-parse HEAD
|
curl -s https://api.windygit.com/version # MUST equal git rev-parse HEAD
|
||||||
```
|
```
|
||||||
|
|
||||||
|
A Gitea config change (compose `GITEA__*`) needs `sudo docker compose up -d --no-deps gitea`
|
||||||
|
— a ~6 s forge outage; running CI jobs survive it. Check `app.ini` afterwards.
|
||||||
|
|
||||||
⚠️ **Never `git pull -q` in a deploy script.** `-q` hides *errors*, not just
|
⚠️ **Never `git pull -q` in a deploy script.** `-q` hides *errors*, not just
|
||||||
noise. On 2026-08-14 a divergent branch made `pull -q` fail silently and the
|
noise. On 2026-08-14 a divergent branch made `pull -q` fail silently and the
|
||||||
"deploy" ran for 20 minutes against stale code while reporting success. Use
|
"deploy" ran for 20 minutes against stale code while reporting success. Use
|
||||||
@@ -72,6 +82,41 @@ curl -sI https://app.windygit.com/ | head -1 # Gitea, 200
|
|||||||
sudo ss -tlnp | grep -E "3080|8600" # both must be 127.0.0.1
|
sudo ss -tlnp | grep -E "3080|8600" # both must be 127.0.0.1
|
||||||
```
|
```
|
||||||
|
|
||||||
|
## Timers (host systemd units — the sync timer is NOT in the repo)
|
||||||
|
|
||||||
|
| Unit | Cadence | Does |
|
||||||
|
|---|---|---|
|
||||||
|
| `windygit-sync.timer` | every 5 min (`OnUnitActiveSec`) | GitHub → Windy Git for `REPOS` in `scripts/sync_from_github.sh`, then `scripts/pr_status_bridge.py` (mirror PRs + GitHub commit statuses). A manual `systemctl start` RESETS the 5-min clock. |
|
||||||
|
| `windygit-backup.timer` | nightly | `git bundle` + pg_dump → R2, 30-day retention |
|
||||||
|
| `windygit-ci-prune.timer` | every 6 h | `deploy/runner/prune.sh` — CI dind storage, 60 GB cap |
|
||||||
|
| `windygit-tunnel.service` | always | the only ingress |
|
||||||
|
|
||||||
|
## CI (Gitea Actions) — see `docs/CUTOVER.md` for onboarding a repo
|
||||||
|
|
||||||
|
- **Six runners × capacity 1** (`deploy/runner/docker-compose.yml`), one shared
|
||||||
|
dind capped at 12 cores / 64 GB. Capacity >1 in one runner shares
|
||||||
|
`/root/.cache/act` between jobs and races (`lstat …: no such file`).
|
||||||
|
- **Runners are scoped to the `windyadmin` user** (`action_runner.owner_id=1`),
|
||||||
|
so only first-party repos run. A repo owned by anyone else — a plane-created
|
||||||
|
agent or `u-system` repo — gets NO runner. Re-registrations inherit this
|
||||||
|
because `RUNNER_TOKEN` is user-level.
|
||||||
|
- Job ceiling 90 min (`config.yaml` `runner.timeout`); a `config.yaml` change
|
||||||
|
needs each runner restarted **while idle** — `compose up -d` won't recreate it.
|
||||||
|
- `/actions/tasks` lists only PICKED-UP jobs. Queue truth is `action_run_job`
|
||||||
|
in the `gitea` DB: `sudo docker exec -i windy-git-db-1 psql -U windygit -d gitea`
|
||||||
|
(status 1 ok · 2 fail · 3 cancelled · 4 skipped · 5 waiting · 6 running).
|
||||||
|
- Job logs are in R2, not on disk. `GET /api/v1/repos/{o}/{r}/actions/jobs/{JOB_ID}/logs`
|
||||||
|
takes the `action_run_job` id, not the task id.
|
||||||
|
|
||||||
|
## Sign-in posture
|
||||||
|
|
||||||
|
- Windy SSO only: password + passkey forms OFF, `ACCOUNT_LINKING=login`,
|
||||||
|
**auto-registration OFF** — opening the forge to non-Grant users is a §7
|
||||||
|
Grant decision.
|
||||||
|
- **Break-glass:** `sudo docker exec -u git windy-git-gitea-1 gitea admin user generate-access-token --username windyadmin --token-name <name> --scopes <scopes> --raw`
|
||||||
|
(delete it after: `delete from access_token where name='<name>'` in the gitea DB —
|
||||||
|
Gitea refuses token management over token auth).
|
||||||
|
|
||||||
## Troubleshooting
|
## Troubleshooting
|
||||||
|
|
||||||
**A hostname returns 530 or won't resolve** — the tunnel is down. `sudo systemctl
|
**A hostname returns 530 or won't resolve** — the tunnel is down. `sudo systemctl
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ dependencies = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
[project.optional-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]
|
[tool.ruff]
|
||||||
line-length = 100
|
line-length = 100
|
||||||
|
|||||||
@@ -23,6 +23,23 @@ BUCKET="${R2_BUCKET_BACKUPS:-windy-git-backups}"
|
|||||||
KEEP_DAYS="${BACKUP_KEEP_DAYS:-30}"
|
KEEP_DAYS="${BACKUP_KEEP_DAYS:-30}"
|
||||||
FAILED=0
|
FAILED=0
|
||||||
|
|
||||||
|
# NEVER bundle these to R2 (orchestrator decision 2026-09-23). They carry
|
||||||
|
# credentials in plaintext — kit-army-config IS the lockbox, and the soul repos
|
||||||
|
# hold agent memory with keys in it — and these bundles are unencrypted, so
|
||||||
|
# anyone holding the R2 key could read every secret in the fleet. They are
|
||||||
|
# backed up ENCRYPTED elsewhere (Windy Drops lane, restic, restore-tested) and
|
||||||
|
# stay mirrored on Veron's own disk in Gitea. Extended globs, matched on name.
|
||||||
|
EXCLUDE="${BACKUP_EXCLUDE:-kit-army-config anima *-soul}"
|
||||||
|
excluded() {
|
||||||
|
local n=$1 pat pats
|
||||||
|
read -ra pats <<< "$EXCLUDE" # read never glob-expands; `for p in $EXCLUDE` would
|
||||||
|
for pat in "${pats[@]}"; do
|
||||||
|
# shellcheck disable=SC2053 # unquoted RHS: glob match is the point
|
||||||
|
[[ "$n" == $pat ]] && return 0
|
||||||
|
done
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
|
||||||
cleanup() { rm -rf "$WORK"; }
|
cleanup() { rm -rf "$WORK"; }
|
||||||
trap cleanup EXIT
|
trap cleanup EXIT
|
||||||
|
|
||||||
@@ -44,6 +61,10 @@ count=0
|
|||||||
for repo in "$GIT_ROOT"/*/*.git; do
|
for repo in "$GIT_ROOT"/*/*.git; do
|
||||||
owner="$(basename "$(dirname "$repo")")"
|
owner="$(basename "$(dirname "$repo")")"
|
||||||
name="$(basename "$repo" .git)"
|
name="$(basename "$repo" .git)"
|
||||||
|
if excluded "$name"; then
|
||||||
|
log "skip ${owner}/${name} (credential-bearing: never bundled to R2 in plaintext)"
|
||||||
|
continue
|
||||||
|
fi
|
||||||
out="$WORK/${owner}__${name}.bundle"
|
out="$WORK/${owner}__${name}.bundle"
|
||||||
|
|
||||||
# --all captures every ref, not just the default branch. A bundle of one
|
# --all captures every ref, not just the default branch. A bundle of one
|
||||||
|
|||||||
@@ -73,6 +73,9 @@ class Check:
|
|||||||
def _probe(c: Check) -> Result:
|
def _probe(c: Check) -> Result:
|
||||||
data = json.dumps(c.body).encode() if c.body else None
|
data = json.dumps(c.body).encode() if c.body else None
|
||||||
headers = {"User-Agent": "windy-git-canary/1.0", **c.headers}
|
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:
|
if data:
|
||||||
headers["Content-Type"] = "application/json"
|
headers["Content-Type"] = "application/json"
|
||||||
req = urllib.request.Request(c.url, data=data, method=c.method, headers=headers)
|
req = urllib.request.Request(c.url, data=data, method=c.method, headers=headers)
|
||||||
|
|||||||
10
scripts/cancel_unrunnable.sh
Executable file
10
scripts/cancel_unrunnable.sh
Executable file
@@ -0,0 +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
|
||||||
|
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)"
|
||||||
53
scripts/cancel_unrunnable.sql
Normal file
53
scripts/cancel_unrunnable.sql
Normal file
@@ -0,0 +1,53 @@
|
|||||||
|
-- Cancel CI jobs that can never run (called by scripts/cancel_unrunnable.sh).
|
||||||
|
--
|
||||||
|
-- A job whose runs-on names a label no Windy Git runner offers (ubuntu-latest,
|
||||||
|
-- macos-latest, windows-latest …) waits forever: Gitea evaluates a job's `if:`
|
||||||
|
-- only when a runner picks it, so even `if: false` / tag-only jobs sit in the
|
||||||
|
-- queue, invisible to /actions/tasks, and keep their run "waiting" for good.
|
||||||
|
-- After 30 minutes they are cancelled here; the run's status is then recomputed
|
||||||
|
-- (failure > still-active > cancelled > success), the same precedence Gitea uses.
|
||||||
|
-- Keep RUNNER_LABELS in step with deploy/runner/config.yaml.
|
||||||
|
BEGIN;
|
||||||
|
WITH dead AS (
|
||||||
|
UPDATE action_run_job j
|
||||||
|
SET status = 3, stopped = extract(epoch from now())::bigint, updated = extract(epoch from now())::bigint
|
||||||
|
WHERE j.status IN (5, 7)
|
||||||
|
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.id, j.run_id, j.name, j.runs_on, j.created
|
||||||
|
), runs AS (
|
||||||
|
UPDATE action_run r
|
||||||
|
SET status = CASE
|
||||||
|
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status = 2) THEN 2
|
||||||
|
WHEN EXISTS (SELECT 1 FROM action_run_job x WHERE x.run_id = r.id AND x.status IN (5, 6, 7)
|
||||||
|
AND x.id NOT IN (SELECT id FROM dead)) THEN r.status
|
||||||
|
ELSE 3 END,
|
||||||
|
stopped = CASE WHEN r.stopped = 0 THEN extract(epoch from now())::bigint ELSE r.stopped END
|
||||||
|
WHERE r.id IN (SELECT DISTINCT run_id FROM dead)
|
||||||
|
RETURNING r.id
|
||||||
|
)
|
||||||
|
-- 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;
|
||||||
@@ -229,13 +229,12 @@ def main() -> int:
|
|||||||
if not targets:
|
if not targets:
|
||||||
ap.error("name a repo, or pass --safe-batch / --list-candidates")
|
ap.error("name a repo, or pass --safe-batch / --list-candidates")
|
||||||
|
|
||||||
if "windy-pro" in targets:
|
# G11.5 RESOLVED 2026-09-23 (lane 8c, ~/windy-orchestra/WINDYPRO_CHECKOUTS.md):
|
||||||
sys.exit(
|
# a read-only audit of all 14 windy-pro checkouts on 5 machines found GitHub
|
||||||
"REFUSING windy-pro. Six checkouts exist, the build counter has forked "
|
# main is canonical (Kit 0 prod and Windy 0 sit exactly on it; the others are
|
||||||
"three ways (main 12 / overnight 34 / wave-44 56), and two sessions "
|
# stale, not divergent). Phase 1 keeps GitHub the source of truth anyway, so a
|
||||||
"recorded different HEADs hours apart. Resolve which is current and "
|
# writable Windy Git copy is CI only. Its six deploy/release workflows must be
|
||||||
"write it down BEFORE importing (G11.5)."
|
# disabled on import — see docs/CUTOVER.md.
|
||||||
)
|
|
||||||
|
|
||||||
if args.mirror:
|
if args.mirror:
|
||||||
print("mirror mode: repos will be read-only and will NOT run CI.\n")
|
print("mirror mode: repos will be read-only and will NOT run CI.\n")
|
||||||
|
|||||||
@@ -28,13 +28,17 @@ repo code, and no secret is handed to any repo.
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import base64
|
||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
import re
|
import re
|
||||||
import sys
|
import sys
|
||||||
|
import time
|
||||||
import urllib.error
|
import urllib.error
|
||||||
import urllib.request
|
import urllib.request
|
||||||
|
|
||||||
|
import yaml
|
||||||
|
|
||||||
GITEA = os.environ.get("BRIDGE_GITEA_URL", "http://localhost:3080").rstrip("/")
|
GITEA = os.environ.get("BRIDGE_GITEA_URL", "http://localhost:3080").rstrip("/")
|
||||||
PUBLIC = "https://app.windygit.com"
|
PUBLIC = "https://app.windygit.com"
|
||||||
GITEA_TOKEN = os.environ.get("GITEA_ADMIN_TOKEN", "")
|
GITEA_TOKEN = os.environ.get("GITEA_ADMIN_TOKEN", "")
|
||||||
@@ -46,7 +50,9 @@ WG_OWNER = os.environ.get("WINDYGIT_OWNER", "windyadmin")
|
|||||||
REPOS = os.environ.get(
|
REPOS = os.environ.get(
|
||||||
"BRIDGE_REPOS",
|
"BRIDGE_REPOS",
|
||||||
"windy-chat windy-mail windy-calendar Windy-Clone WindyCloud windy-search windy-connect"
|
"windy-chat windy-mail windy-calendar Windy-Clone WindyCloud windy-search windy-connect"
|
||||||
" windy-drops windy-code-web windy-code windy-traveler windy-registry eternitas",
|
" windy-drops windy-code-web windy-code windy-traveler windy-registry eternitas"
|
||||||
|
" windy-translate windytranslate-site windytraveler-site windy-hand"
|
||||||
|
" windy-cloud-sites windy-cloud-domains windy-cloud-vps windytalk windy-pro windy-mind",
|
||||||
).split()
|
).split()
|
||||||
|
|
||||||
# Gitea run status -> GitHub status state. `skipped` is deliberately absent: a
|
# Gitea run status -> GitHub status state. `skipped` is deliberately absent: a
|
||||||
@@ -69,6 +75,75 @@ MIRROR_TAG = "[GH#"
|
|||||||
# exists; that is a decision, recorded in docs/CUTOVER.md, not a failure.
|
# exists; that is a decision, recorded in docs/CUTOVER.md, not a failure.
|
||||||
NO_DAEMON_JOB = re.compile(r"docker", re.IGNORECASE)
|
NO_DAEMON_JOB = re.compile(r"docker", re.IGNORECASE)
|
||||||
|
|
||||||
|
# Jobs Grant ruled NON-BLOCKING (GRANT_DECISIONS_2026-09-23): still run on
|
||||||
|
# Windy Git and visible there, but not posted to GitHub, so they cannot turn a
|
||||||
|
# commit's combined status red. Format: "repo:workflow/job,workflow/job;repo2:..."
|
||||||
|
# windy-pro's desktop/installer jobs belong to Grant's desktop side (fixed from
|
||||||
|
# his Mac mini), not to any lane's merge gate.
|
||||||
|
NON_BLOCKING: dict[str, set[str]] = {}
|
||||||
|
for _entry in os.environ.get(
|
||||||
|
"BRIDGE_NON_BLOCKING", "windy-pro:ci/build-desktop,ci/test-installer,ci/reality-check"
|
||||||
|
).split(";"):
|
||||||
|
if ":" in _entry:
|
||||||
|
_repo, _jobs = _entry.split(":", 1)
|
||||||
|
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):
|
def _call(base: str, token_header: str, method: str, path: str, body=None):
|
||||||
req = urllib.request.Request(
|
req = urllib.request.Request(
|
||||||
@@ -83,12 +158,21 @@ def _call(base: str, token_header: str, method: str, path: str, body=None):
|
|||||||
"User-Agent": "windy-git-pr-bridge/1",
|
"User-Agent": "windy-git-pr-bridge/1",
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
try:
|
# Transport errors (TLS handshake timeout, reset) are retried: one GitHub
|
||||||
with urllib.request.urlopen(req, timeout=60) as r:
|
# blip used to fail the whole sync, flip its heartbeat to ok:false and page
|
||||||
raw = r.read()
|
# someone for nothing. HTTP errors are answers, not blips — never retried.
|
||||||
return r.status, (json.loads(raw) if raw else None)
|
for attempt in range(3):
|
||||||
except urllib.error.HTTPError as e:
|
try:
|
||||||
return e.code, None
|
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):
|
def gitea(method, path, body=None):
|
||||||
@@ -153,10 +237,13 @@ def post_statuses(repo: str, sha: str) -> None:
|
|||||||
for r in runs:
|
for r in runs:
|
||||||
if r["head_sha"] != sha or NO_DAEMON_JOB.search(r["name"]):
|
if r["head_sha"] != sha or NO_DAEMON_JOB.search(r["name"]):
|
||||||
continue
|
continue
|
||||||
|
if f"{r['workflow_id'].removesuffix('.yml')}/{r['name']}" in NON_BLOCKING.get(repo, ()):
|
||||||
|
continue
|
||||||
ctx = f"windy-git/{r['workflow_id'].removesuffix('.yml')}/{r['name']}"
|
ctx = f"windy-git/{r['workflow_id'].removesuffix('.yml')}/{r['name']}"
|
||||||
if ctx not in latest or r["id"] > latest[ctx]["id"]:
|
if ctx not in latest or r["id"] > latest[ctx]["id"]:
|
||||||
latest[ctx] = r
|
latest[ctx] = r
|
||||||
if not latest:
|
bad = invalid_workflows(repo, sha)
|
||||||
|
if not (latest or bad):
|
||||||
return
|
return
|
||||||
|
|
||||||
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
|
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
|
||||||
@@ -164,6 +251,21 @@ def post_statuses(repo: str, sha: str) -> None:
|
|||||||
for s in existing or []: # newest first
|
for s in existing or []: # newest first
|
||||||
current.setdefault(s["context"], s["state"])
|
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()):
|
for ctx, r in sorted(latest.items()):
|
||||||
state = STATE.get(r["status"])
|
state = STATE.get(r["status"])
|
||||||
if state is None or current.get(ctx) == state:
|
if state is None or current.get(ctx) == state:
|
||||||
|
|||||||
24
scripts/promote_to_ci.sh
Executable file
24
scripts/promote_to_ci.sh
Executable file
@@ -0,0 +1,24 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# promote_to_ci.sh <repo> [workflow-to-disable ...] — pull mirror -> writable CI repo.
|
||||||
|
# Run ON Veron as root. See docs/CUTOVER.md "Onboarding another private repo".
|
||||||
|
#
|
||||||
|
# ⚠️ It DELETES the mirror before importing (Gitea's migrate refuses an existing
|
||||||
|
# name). If the import then fails, the Windy Git copy is gone until you re-run —
|
||||||
|
# GitHub and the nightly R2 bundles still hold everything, but check first that
|
||||||
|
# scripts/import_from_github.py will accept the repo. (2026-09-23: windy-pro was
|
||||||
|
# deleted this way while the importer still refused it by name.)
|
||||||
|
set -euo pipefail
|
||||||
|
set -a; . /srv/windygit/src/.env; set +a
|
||||||
|
export IMPORT_GITEA_URL=http://localhost:3080
|
||||||
|
A=http://localhost:3080/api/v1; H="Authorization: token $GITEA_ADMIN_TOKEN"; r=$1; shift
|
||||||
|
info=$(curl -s -H "$H" $A/repos/windyadmin/$r)
|
||||||
|
m=$(echo "$info" | python3 -c 'import json,sys;print(json.load(sys.stdin).get("mirror"))')
|
||||||
|
if [ "$m" = True ]; then
|
||||||
|
curl -sf -o /dev/null -X DELETE -H "$H" $A/repos/windyadmin/$r
|
||||||
|
(cd /srv/windygit/src && python3 scripts/import_from_github.py "$r" | tail -1)
|
||||||
|
elif [ "$m" = False ]; then echo "$r already writable"; else echo "$r absent -> importing"; (cd /srv/windygit/src && python3 scripts/import_from_github.py "$r" | tail -1); fi
|
||||||
|
db=$(curl -s -H "$H" $A/repos/windyadmin/$r | python3 -c 'import json,sys;print(json.load(sys.stdin).get("default_branch","main"))')
|
||||||
|
for i in $(seq 1 120); do curl -sf -o /dev/null -H "$H" $A/repos/windyadmin/$r/branches/$db && break; sleep 5; done
|
||||||
|
for w in "$@"; do printf " disable %s: " "$w"; curl -s -o /dev/null -w '%{http_code}\n' -X PUT -H "$H" $A/repos/windyadmin/$r/actions/workflows/$w/disable; done
|
||||||
|
curl -s -H "$H" $A/repos/windyadmin/$r/actions/workflows | python3 -c 'import json,sys,os;print(" "+os.environ.get("R",""),[(w["path"].split("/")[-1],w["state"]) for w in json.load(sys.stdin).get("workflows",[])])'
|
||||||
|
echo " default=$db"
|
||||||
@@ -38,7 +38,17 @@ FAILED=0
|
|||||||
|
|
||||||
# Repos Windy Git tracks FROM GitHub. Remove a repo from this list at the moment
|
# Repos Windy Git tracks FROM GitHub. Remove a repo from this list at the moment
|
||||||
# it flips to Windy-Git-first, or the sync will fight its authors and win.
|
# it flips to Windy-Git-first, or the sync will fight its authors and win.
|
||||||
REPOS="${SYNC_REPOS:-windy-calendar windy-search windy-registry Windy-Clone WindyCloud windy-cloud-sites windy-mind eternitas windy-agent windy-git windy-chat windy-mail windy-connect windy-drops windy-code-web windy-code windy-traveler}"
|
REPOS="${SYNC_REPOS:-windy-calendar windy-search windy-registry Windy-Clone WindyCloud windy-cloud-sites windy-mind eternitas windy-agent windy-git windy-chat windy-mail windy-connect windy-drops windy-code-web windy-code windy-traveler windy-translate windytranslate-site windytraveler-site windy-hand windy-cloud-domains windy-cloud-vps windytalk windy-pro}"
|
||||||
|
|
||||||
|
# Repos whose TAGS must not reach Windy Git. A tag push fires `on: push: tags`
|
||||||
|
# workflows; windy-pro's build-electron is a matrix over ubuntu/macos/windows-
|
||||||
|
# latest, labels no runner here has, so every leg would queue forever (and
|
||||||
|
# queued jobs are invisible in /actions/tasks). Releases are built elsewhere.
|
||||||
|
NO_TAGS="${SYNC_NO_TAGS:-windy-pro}"
|
||||||
|
|
||||||
|
# `archive/*` branches never reach Windy Git (negative refspec, git >= 2.29).
|
||||||
|
# They are off-machine safety copies of unpushed work (one-repo doctrine), not
|
||||||
|
# work in progress: GitHub holds them, and CI time on them is waste.
|
||||||
|
|
||||||
mkdir -p "$WORK"
|
mkdir -p "$WORK"
|
||||||
log() { printf '[sync %s] %s\n' "$(date -u +%H:%M:%SZ)" "$*"; }
|
log() { printf '[sync %s] %s\n' "$(date -u +%H:%M:%SZ)" "$*"; }
|
||||||
@@ -63,18 +73,26 @@ for r in $REPOS; do
|
|||||||
|
|
||||||
if git --git-dir="$bare" push --quiet --force \
|
if git --git-dir="$bare" push --quiet --force \
|
||||||
"https://${WG_OWNER}:${GITEA_ADMIN_TOKEN}@${WG}/${WG_OWNER}/${r}.git" \
|
"https://${WG_OWNER}:${GITEA_ADMIN_TOKEN}@${WG}/${WG_OWNER}/${r}.git" \
|
||||||
'+refs/heads/*:refs/heads/*' '+refs/tags/*:refs/tags/*' 2>/dev/null; then
|
'+refs/heads/*:refs/heads/*' '^refs/heads/archive/*' $([[ " $NO_TAGS " == *" $r "* ]] || echo '+refs/tags/*:refs/tags/*') 2>/dev/null; then
|
||||||
log "$r ok (${before:0:7})"
|
log "$r ok (${before:0:7})"
|
||||||
else
|
else
|
||||||
log "FAILED push $r -> windy git"; FAILED=1
|
log "FAILED push $r -> windy git"; FAILED=1
|
||||||
fi
|
fi
|
||||||
done
|
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)"
|
||||||
|
|
||||||
# Private repos can't run GitHub Actions; mirror their open PRs here so CI
|
# 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.
|
# fires, and post the verdicts back to GitHub as commit statuses.
|
||||||
if ! python3 "$(dirname "$0")/pr_status_bridge.py"; then
|
if ! python3 "$(dirname "$0")/pr_status_bridge.py"; then
|
||||||
log "FAILED pr status bridge"; FAILED=1
|
log "FAILED pr status bridge"; FAILED=1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
# CI telemetry -> admin.windyword.ai (shapes declared with Windy Telemetry 40).
|
||||||
|
# Sends nothing until WINDYGIT_TELEMETRY_TOKEN is set; never fails the sync.
|
||||||
|
python3 "$(dirname "$0")/telemetry_emit.py" || log "telemetry emit failed (non-fatal)"
|
||||||
|
|
||||||
[[ "$FAILED" -ne 0 ]] && { log "COMPLETED WITH FAILURES"; exit 1; }
|
[[ "$FAILED" -ne 0 ]] && { log "COMPLETED WITH FAILURES"; exit 1; }
|
||||||
log "all repos in step with GitHub"
|
log "all repos in step with GitHub"
|
||||||
|
|||||||
276
scripts/telemetry_emit.py
Normal file
276
scripts/telemetry_emit.py
Normal file
@@ -0,0 +1,276 @@
|
|||||||
|
#!/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 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")
|
||||||
|
|
||||||
|
|
||||||
|
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,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
# 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},
|
||||||
|
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())
|
||||||
Reference in New Issue
Block a user