Compare commits
3 Commits
c83f808a60
...
telemetry-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e3b69fa759 | ||
|
|
4acf50d9ef | ||
|
|
b7a7e94df0 |
@@ -120,6 +120,9 @@ class Telemetry:
|
||||
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:
|
||||
@@ -147,6 +150,7 @@ class Telemetry:
|
||||
}
|
||||
)
|
||||
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:
|
||||
@@ -188,6 +192,8 @@ class Telemetry:
|
||||
"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)
|
||||
@@ -199,7 +205,7 @@ class Telemetry:
|
||||
self._reset_window()
|
||||
|
||||
# ---- sending ------------------------------------------------------------
|
||||
def _post(self, batch: list[dict]) -> int:
|
||||
def _post(self, batch: list[dict]) -> tuple[int, dict]:
|
||||
req = urllib.request.Request(
|
||||
self.url,
|
||||
data=json.dumps({"events": batch}).encode(),
|
||||
@@ -211,19 +217,32 @@ class Telemetry:
|
||||
},
|
||||
)
|
||||
with urllib.request.urlopen(req, timeout=20) as r:
|
||||
return r.status
|
||||
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 = await asyncio.to_thread(self._post, batch)
|
||||
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."""
|
||||
|
||||
@@ -8,6 +8,7 @@ no status at all.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import importlib.util
|
||||
from pathlib import Path
|
||||
|
||||
@@ -35,12 +36,23 @@ def _run(i, wf, job, status, sha=SHA, n=1):
|
||||
|
||||
|
||||
class Fake:
|
||||
def __init__(self, runs=(), statuses=(), gh_prs=(), wg_prs=()):
|
||||
def __init__(self, runs=(), statuses=(), gh_prs=(), wg_prs=(), workflows=None):
|
||||
self.runs, self.statuses = list(runs), list(statuses)
|
||||
self.workflows = workflows or {} # {path: yaml text} at every commit
|
||||
self.gh_prs, self.wg_prs = list(gh_prs), list(wg_prs)
|
||||
self.posted, self.opened, self.closed = [], [], []
|
||||
|
||||
def gitea(self, method, path, body=None):
|
||||
if "/contents/" in path:
|
||||
want = path.split("/contents/", 1)[1].split("?", 1)[0]
|
||||
if want in self.workflows:
|
||||
return 200, {"content": base64.b64encode(self.workflows[want].encode()).decode()}
|
||||
files = [
|
||||
{"type": "file", "name": k.rsplit("/", 1)[1], "path": k}
|
||||
for k in self.workflows
|
||||
if k.rsplit("/", 1)[0] == want
|
||||
]
|
||||
return (200, files) if files else (404, None)
|
||||
if "/actions/tasks" in path:
|
||||
page = int(path.rsplit("page=", 1)[1])
|
||||
return 200, {"workflow_runs": self.runs[(page - 1) * 50 : page * 50]}
|
||||
@@ -209,3 +221,54 @@ def test_transport_blips_are_retried_but_http_errors_are_not(monkeypatch):
|
||||
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
|
||||
|
||||
@@ -160,3 +160,43 @@ def test_synthetic_is_forwarded_downstream_only_for_synthetic_requests():
|
||||
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
|
||||
|
||||
@@ -31,7 +31,7 @@ dependencies = [
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
dev = ["pytest>=8.3", "pytest-asyncio>=0.24", "ruff>=0.7", "mypy>=1.13"]
|
||||
dev = ["pytest>=8.3", "pyyaml>=6.0", "pytest-asyncio>=0.24", "ruff>=0.7", "mypy>=1.13"]
|
||||
|
||||
[tool.ruff]
|
||||
line-length = 100
|
||||
|
||||
@@ -28,6 +28,7 @@ repo code, and no secret is handed to any repo.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
@@ -36,6 +37,8 @@ import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
import yaml
|
||||
|
||||
GITEA = os.environ.get("BRIDGE_GITEA_URL", "http://localhost:3080").rstrip("/")
|
||||
PUBLIC = "https://app.windygit.com"
|
||||
GITEA_TOKEN = os.environ.get("GITEA_ADMIN_TOKEN", "")
|
||||
@@ -86,6 +89,62 @@ for _entry in os.environ.get(
|
||||
NON_BLOCKING[_repo.strip()] = {j.strip() for j in _jobs.split(",") if j.strip()}
|
||||
|
||||
|
||||
# Gitea reads the FIRST of these dirs that has workflow files at a commit (1.24).
|
||||
WORKFLOW_DIRS = (".gitea/workflows", ".github/workflows")
|
||||
|
||||
|
||||
def workflow_problem(text: str) -> str | None:
|
||||
"""Why Gitea would drop this workflow file, or None if it looks runnable.
|
||||
|
||||
Gitea skips an invalid workflow with one log line and fires no run at all,
|
||||
so on GitHub the PR just shows nothing, and people wait for CI that is never
|
||||
coming. These are the shapes we have actually hit, not a full schema.
|
||||
"""
|
||||
try:
|
||||
doc = yaml.safe_load(text)
|
||||
except yaml.YAMLError as e:
|
||||
mark = getattr(e, "problem_mark", None)
|
||||
return f"invalid YAML at line {mark.line + 1}" if mark else "invalid YAML"
|
||||
if not isinstance(doc, dict):
|
||||
return "not a YAML mapping"
|
||||
if "on" not in doc and True not in doc: # YAML 1.1 reads a bare `on` as True
|
||||
return "no `on:` trigger"
|
||||
jobs = doc.get("jobs")
|
||||
if not isinstance(jobs, dict) or not jobs:
|
||||
return "no `jobs:`"
|
||||
for name, job in jobs.items():
|
||||
if not isinstance(job, dict):
|
||||
return f"job `{name}` is not a mapping"
|
||||
if "runs-on" not in job and "uses" not in job:
|
||||
return f"job `{name}` has no `runs-on:`"
|
||||
return None
|
||||
|
||||
|
||||
def invalid_workflows(repo: str, sha: str) -> dict[str, tuple[str, str]]:
|
||||
"""{context: (path, problem)} for each workflow file at `sha` that won't run."""
|
||||
for d in WORKFLOW_DIRS:
|
||||
st, entries = gitea("GET", f"/repos/{WG_OWNER}/{repo}/contents/{d}?ref={sha}")
|
||||
if st == 404:
|
||||
continue
|
||||
if st != 200:
|
||||
raise RuntimeError(f"{repo}: Windy Git {d}@{sha[:7]} -> {st}")
|
||||
files = [e for e in entries or [] if e.get("type") == "file"
|
||||
and e["name"].endswith((".yml", ".yaml"))]
|
||||
if not files:
|
||||
continue
|
||||
bad = {}
|
||||
for e in files:
|
||||
st, f = gitea("GET", f"/repos/{WG_OWNER}/{repo}/contents/{e['path']}?ref={sha}")
|
||||
if st != 200:
|
||||
raise RuntimeError(f"{repo}: Windy Git {e['path']}@{sha[:7]} -> {st}")
|
||||
problem = workflow_problem(base64.b64decode(f["content"]).decode("utf-8", "replace"))
|
||||
if problem:
|
||||
stem = re.sub(r"\.ya?ml$", "", e["name"])
|
||||
bad[f"windy-git/{stem}/workflow"] = (e["path"], problem)
|
||||
return bad
|
||||
return {}
|
||||
|
||||
|
||||
def _call(base: str, token_header: str, method: str, path: str, body=None):
|
||||
req = urllib.request.Request(
|
||||
base + path,
|
||||
@@ -183,7 +242,8 @@ def post_statuses(repo: str, sha: str) -> None:
|
||||
ctx = f"windy-git/{r['workflow_id'].removesuffix('.yml')}/{r['name']}"
|
||||
if ctx not in latest or r["id"] > latest[ctx]["id"]:
|
||||
latest[ctx] = r
|
||||
if not latest:
|
||||
bad = invalid_workflows(repo, sha)
|
||||
if not (latest or bad):
|
||||
return
|
||||
|
||||
st, existing = github("GET", f"/repos/{GH_OWNER}/{repo}/commits/{sha}/statuses?per_page=100")
|
||||
@@ -191,6 +251,21 @@ def post_statuses(repo: str, sha: str) -> None:
|
||||
for s in existing or []: # newest first
|
||||
current.setdefault(s["context"], s["state"])
|
||||
|
||||
for ctx, (path, problem) in sorted(bad.items()):
|
||||
if current.get(ctx) == "error":
|
||||
continue
|
||||
st, _ = github(
|
||||
"POST",
|
||||
f"/repos/{GH_OWNER}/{repo}/statuses/{sha}",
|
||||
{
|
||||
"state": "error",
|
||||
"context": ctx,
|
||||
"description": f"Windy Git ignored this workflow, no CI ran: {problem}"[:140],
|
||||
"target_url": f"{PUBLIC}/{WG_OWNER}/{repo}/src/commit/{sha}/{path}",
|
||||
},
|
||||
)
|
||||
print(f" {repo}@{sha[:7]} {ctx} = error ({problem}) -> {st}")
|
||||
|
||||
for ctx, r in sorted(latest.items()):
|
||||
state = STATE.get(r["status"])
|
||||
if state is None or current.get(ctx) == state:
|
||||
|
||||
@@ -161,6 +161,12 @@ def main() -> int:
|
||||
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
|
||||
@@ -216,6 +222,7 @@ def main() -> int:
|
||||
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,
|
||||
@@ -229,11 +236,20 @@ def main() -> int:
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as r:
|
||||
body = r.read()[:300]
|
||||
body = r.read()
|
||||
if r.status >= 300:
|
||||
raise urllib.error.HTTPError(
|
||||
INGEST, r.status, body.decode(errors="replace"), None, None
|
||||
INGEST, r.status, body[:300].decode(errors="replace"), None, None
|
||||
)
|
||||
try:
|
||||
resp = json.loads(body or b"{}")
|
||||
except ValueError:
|
||||
resp = {}
|
||||
q = resp.get("quarantined") if isinstance(resp, dict) else None
|
||||
if isinstance(q, int) and q > 0:
|
||||
quarantined += q
|
||||
reasons = "; ".join(map(str, resp.get("rejections") or [])) or "no reason given"
|
||||
print(f"[telemetry] WARNING {q} row(s) QUARANTINED by the ledger: {reasons}")
|
||||
except urllib.error.HTTPError as e:
|
||||
print(f"[telemetry] FAILED ingest HTTP {e.code}: {e.read()[:200]!r}")
|
||||
return 1 # state NOT advanced: the same rows retry next run
|
||||
@@ -246,7 +262,11 @@ def main() -> int:
|
||||
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}, 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
|
||||
|
||||
Reference in New Issue
Block a user