diff --git a/DNA_STRAND_MASTER_PLAN.md b/DNA_STRAND_MASTER_PLAN.md index f422ba0..15bf121 100644 --- a/DNA_STRAND_MASTER_PLAN.md +++ b/DNA_STRAND_MASTER_PLAN.md @@ -81,7 +81,7 @@ Numbered because code cites them. Changing one requires an ADR that names it. 1. **I-1 · Gitea is a component, never a merged tree.** Our code lives in our services and calls Gitea's REST API. Any patch to Gitea source lives in `patches/` as a numbered, rebasable diff with a one-line justification, and `make check` fails if `patches/` grows past **3** files without an ADR. 2. **I-2 · The membrane is ENUMERATED.** - **Calls out:** `windy-cloud` kernel `GET /api/v1/storage/objects` + `HEAD` (read user objects to version them) · `windy-cloud` `POST /api/v1/storage/quota/check` · `eternitas` `GET /api/v1/trust/{passport}` (band + allowed_actions) · `eternitas` `GET /api/v1/registry/{passport}/integrity` · `account-server` OIDC discovery + JWKS · `windy-cloud-sites` `POST /api/v1/sites/{id}/versions` (publish docs from a repo). + **Calls out:** `windy-cloud` kernel `GET /api/v1/storage/objects` + `HEAD` (read user objects to version them) · `windy-cloud` `POST /api/v1/storage/quota/check` · `eternitas` `GET /api/v1/trust/{passport}` (band + allowed_actions) · `eternitas` `GET /api/v1/registry/{passport}/integrity` · `account-server` OIDC discovery + JWKS · `windy-cloud-sites` `POST /api/v1/sites/{id}/versions` (publish docs from a repo). · windy-admin ledger `POST https://admin.windyword.ai/v1/events` (field telemetry, 2026-09-23: `ci.run`, `ci.job_cancelled`, `service.boot`, `service.health`, `forge.auth.failed` — shapes declared with the ledger owner first; no content, no passports, no emails) **Calls in:** `POST /internal/repo-from-folder` (Cloud portal: git-enable a folder) · `POST /internal/mirror-status` (ops). **Events out:** `repo.created`, `repo.pushed`, `release.published`, `model.published`, `ci.completed`. **Events in:** `passport.revoked` (fail-closed), `storage.quota.exceeded`, `identity.created`. diff --git a/api/app/config.py b/api/app/config.py index 1008b45..da32d66 100644 --- a/api/app/config.py +++ b/api/app/config.py @@ -66,6 +66,12 @@ class Settings(BaseSettings): # Flip to True once the hub emits aud on every access token. hub_require_aud: bool = False + # ---- field telemetry (admin.windyword.ai ledger) ---------------------- + # Unset token = nothing sent, nothing buffered. The token lives in the + # root-only /etc/windygit/telemetry.env on Veron, never in the repo. + windygit_telemetry_token: str = "" + telemetry_ingest_url: str = "https://admin.windyword.ai/v1/events" + # Internal callers (the Cloud portal calling /internal/*). A first-class # caller class, not a bypass: unset means service calls are REFUSED. service_token: str = "" diff --git a/api/app/main.py b/api/app/main.py index 7bf946e..c59aa7a 100644 --- a/api/app/main.py +++ b/api/app/main.py @@ -6,11 +6,13 @@ component and is reached only over its REST API (D-2 / I-1). from __future__ import annotations +import asyncio import logging import socket +import time from contextlib import asynccontextmanager -from fastapi import FastAPI +from fastapi import FastAPI, Request from fastapi.exceptions import RequestValidationError from fastapi.responses import JSONResponse from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine @@ -25,6 +27,7 @@ from api.app.providers.registry import ( R2Provider, ) from api.app.routes import health, repos, webhooks +from api.app.telemetry import Telemetry, caller_class logging.basicConfig( level=logging.INFO, @@ -98,8 +101,23 @@ async def lifespan(app: FastAPI): # systemd Restart=always, plus the runbook's `systemctl status`. ] + telemetry = Telemetry( + settings.telemetry_ingest_url, + settings.windygit_telemetry_token, + environment=settings.environment, + commit_sha=info.commit_sha, + version=info.version, + ) + app.state.telemetry = telemetry + telemetry.boot() + await telemetry.flush() + task = asyncio.create_task(telemetry.run()) if telemetry.enabled else None + yield + if task is not None: + task.cancel() + await telemetry.flush() if engine is not None: await engine.dispose() @@ -119,8 +137,39 @@ app.include_router(repos.router) app.include_router(webhooks.router) +@app.middleware("http") +async def _count_requests(request: Request, call_next): + """Heartbeat counts (requests, 4xx/5xx, refusals, p95). Never raises.""" + start = time.perf_counter() + response = await call_next(request) + tel = getattr(request.app.state, "telemetry", None) + if tel is not None: + tel.record_request( + response.status_code, + (time.perf_counter() - start) * 1000, + refused=getattr(request.state, "refused", False), + ) + return response + + @app.exception_handler(RepairPointer) -async def _repair_pointer_handler(_, exc: RepairPointer) -> JSONResponse: +async def _repair_pointer_handler(request: Request, exc: RepairPointer) -> JSONResponse: + tel = getattr(request.app.state, "telemetry", None) + detail = exc.detail if isinstance(exc.detail, dict) else {} + code = detail.get("code") + if tel is not None and code in tel.auth_codes: + # A refusal is a failure row (field-visibility rule 1). The caller is + # unauthenticated by definition, so: system actor, no actor_id, and + # the route TEMPLATE, never the concrete path. + request.state.refused = True + route = request.scope.get("route") + tel.auth_failed( + code=code, + http_status=exc.status_code, + caller=caller_class(request.headers), + route=getattr(route, "path", None), + upstream_status=getattr(exc, "upstream_status", None), + ) return JSONResponse(status_code=exc.status_code, content=exc.detail) diff --git a/api/app/telemetry.py b/api/app/telemetry.py new file mode 100644 index 0000000..dcea56f --- /dev/null +++ b/api/app/telemetry.py @@ -0,0 +1,212 @@ +"""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 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") + + +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] = [] + + # ---- 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: + 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, + ) -> None: + if code not in AUTH_CODES: + return + meta: dict = {"code": code, "http_status": int(http_status), "caller": caller} + 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, + } + 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]) -> int: + 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: + return r.status + + async def flush(self) -> None: + if not self.enabled or not self.buffer: + return + batch = self.buffer[:500] + try: + status = await asyncio.to_thread(self._post, batch) + 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)] + + 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() diff --git a/api/tests/test_telemetry.py b/api/tests/test_telemetry.py new file mode 100644 index 0000000..5abb8c1 --- /dev/null +++ b/api/tests/test_telemetry.py @@ -0,0 +1,142 @@ +"""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"} + + +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" + ) diff --git a/docker-compose.yml b/docker-compose.yml index 3db8c4e..04dc001 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -16,7 +16,11 @@ services: # I-12: baked at build time. A runtime COMMIT_SHA override is ignored. COMMIT_SHA: ${COMMIT_SHA_BUILD:-} BUILT_AT: ${BUILT_AT:-} - env_file: [.env] + env_file: + - .env + # WINDYGIT_TELEMETRY_TOKEN (root-only on Veron). Optional: no file = no telemetry. + - path: /etc/windygit/telemetry.env + required: false environment: DATABASE_URL: postgresql+asyncpg://windygit:${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD}@db:5432/windygit GITEA_BASE_URL: http://gitea:3000 diff --git a/docs/MEMBRANE.v1.md b/docs/MEMBRANE.v1.md index 96fc779..63923f8 100644 --- a/docs/MEMBRANE.v1.md +++ b/docs/MEMBRANE.v1.md @@ -15,6 +15,7 @@ Mirrored into `windy-cloud` and `eternitas` on change. | eternitas | `GET /api/v1/trust/{passport}` | band + allowed_actions | | eternitas | `GET /api/v1/registry/{passport}/integrity` | ⚠️ note the path — `windy-registry` calls `/api/v1/passports/{p}/status`, which 404s, which is why the integrity index has never been populated | | account-server | OIDC discovery + JWKS | human identity (G3.1) | +| windy-admin ledger | `POST /v1/events` (admin.windyword.ai) | field telemetry: `ci.run`, `ci.job_cancelled`, `service.boot`, `service.health`, `forge.auth.failed`. Shapes are declared with the ledger owner BEFORE shipping (the server quarantines undeclared keys). Codes, counts, route templates only | | windy-cloud-sites | `POST /api/v1/sites/{id}/versions` | publish docs from a repo | ## Calls IN