Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ CLAUDE.md.
Nothing is parsed before it is written to the log, and `rebuild` replays the
logs through the same code path a live run uses. If a parse is wrong, fix it and
rebuild; never edit the logs. `json.dumps(..., sort_keys=True)` in
`store.py:write_raw` is load-bearing - it is what lets two machines' logs be
`store.py:append_raw` is load-bearing - it is what lets two machines' logs be
merged with `sort -u`.

## Data-shape traps
Expand Down Expand Up @@ -175,6 +175,7 @@ merged with `sort -u`.
| Station access is labelled by hand into an append-only observation log, `lifts-data/survey/<CODE>.jsonl`, one line per fact with who, when and from what; a hand-maintained `stations.json` stays the failure mode. The graph replays the log, last line for a key wins, and says "another step-free way" only on a route every edge of which a person confirmed, so a page-seeded graph never says more than the prose. The site does not read it yet | `notes/step-free-graph.md` |
| The Metro Nation Dublin rail map is not a source: undefined "step-free", already behind the network, nothing it says survives one survey answer | `notes/step-free-graph.md` § What was learned |
| A delays site is a fourth repo, `baz8080/rail-delays`, reading `lifts-data`: not a second collector and not a poll target, because there is one endpoint, one response, and every delay notice is already logged. It carries its own decisions, and the collector here is not duplicated, extended or touched | `notes/delays-site.md` |
| The raw line is written before the database opens, and a database failure is exit 7. Replay orders each file by `fetched_at_utc`, so a `sort -u` merge is safe. The NTP wait is left for now | `notes/collector-review.md` |
| The access golden file pins the inputs it derives from, not just the outputs, so it guards code and nothing else. Corpus movement no longer fails it in either direction, a refreshed snapshot no longer fails it by name, and it runs without a `lifts-data` checkout instead of skipping. Reading `messages.text_raw` as if it were append-only reddened `main` three times in five days: the raw logs are append-only, the derived row is overwritten when Irish Rail rewords a live notice | `notes/station-access.md` § The golden file pins its inputs |

Decisions go in `notes/`, dated, with the rejected alternatives and their
Expand Down
2 changes: 1 addition & 1 deletion lift_access/fetch.py
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ def fetch_stations(log=print, attempts=3):


def write_snapshot(path, records, fetched_at=None):
"""One JSONL file, sorted, `sort_keys=True`, in the shape `store.write_raw` uses.
"""One JSONL file, sorted, `sort_keys=True`, in the shape `store.append_raw` uses.

Sorted so that two refreshes of unchanged data produce an identical file and
the scheduled job opens no PR.
Expand Down
66 changes: 61 additions & 5 deletions lift_status/alert.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

from __future__ import annotations

import contextlib
import hashlib
import json
import os
Expand All @@ -23,20 +24,27 @@
EXIT_UNREACHABLE = 3
EXIT_SCHEMA_DRIFT = 4
EXIT_STORAGE = 6
EXIT_DATABASE = 7

EXIT_MEANINGS = {
EXIT_OK: "success",
EXIT_AUTH: "API key rejected",
EXIT_UNREACHABLE: "messages API unreachable",
EXIT_SCHEMA_DRIFT: "API response shape changed",
EXIT_STORAGE: "data directory not writable",
EXIT_DATABASE: "database unusable; raw log still written",
}

BANNER_WIDTH = 78

# How long before an unchanged banner is worth pushing again.
ALERT_REPEAT_SECONDS = 24 * 60 * 60

# Consecutive clean runs (two hours at the 30-minute cadence) before the same
# fault counts as a new incident. One clean poll between two failures of a
# flapping API is not a recovery.
RECOVERED_AFTER_CLEAN_RUNS = 4


def banner(title: str, lines: list[str]) -> str:
bar = "!" * BANNER_WIDTH
Expand Down Expand Up @@ -153,6 +161,28 @@ def storage_banner(data_dir, problem: str) -> str:
)


def database_banner(data_dir, detail: str) -> str:
return banner(
"LIFT-STATUS: DATABASE UNUSABLE",
[
f"{detail}",
"",
"This run's response WAS written to the raw log, so nothing has been",
"lost yet, but the database is not being updated. By the error above:",
"",
" 'locked': something else holds it, usually 'lift rebuild' or",
" 'lift stats'. It clears when that finishes.",
"",
f" 'full' or 'No space': the SD card is full. Check: df -h {data_dir}",
"",
" 'malformed' or 'not a database': a power cut corrupted it. The raw",
" log rebuilds it:",
f" sudo mv {data_dir}/lift_status.db {data_dir}/lift_status.db.broken",
" sudo lift rebuild",
],
)


def _marker_path() -> Path:
state_dir = os.environ.get("LIFT_STATUS_DATA_DIR") or tempfile.gettempdir()
return Path(state_dir) / ".last-alert.json"
Expand Down Expand Up @@ -189,12 +219,37 @@ def _mark_delivered(message: str) -> None:
blip, which lands hardest at the only moment that matters: the first alert
of a collector that has stopped.
"""
_write_marker({"digest": _digest(message), "sent_at": time.time(), "clean_runs": 0})


def _write_marker(state: dict) -> None:
with contextlib.suppress(OSError):
_marker_path().write_text(json.dumps(state), encoding="utf-8")


def _restart_recovery() -> None:
with contextlib.suppress(OSError, ValueError, TypeError, AttributeError):
state = json.loads(_marker_path().read_text(encoding="utf-8"))
if state.get("clean_runs"):
_write_marker({**state, "clean_runs": 0})


def note_clean_run() -> None:
"""Close the repeat window once collection has stayed clean, so the same
fault coming back later is a new incident rather than sitting out the day."""
path = _marker_path()
try:
_marker_path().write_text(
json.dumps({"digest": _digest(message), "sent_at": time.time()}), encoding="utf-8"
)
except OSError:
pass
state = json.loads(path.read_text(encoding="utf-8"))
clean_runs = int(state.get("clean_runs", 0)) + 1
except FileNotFoundError:
return
except (OSError, ValueError, TypeError, AttributeError):
clean_runs = RECOVERED_AFTER_CLEAN_RUNS
if clean_runs >= RECOVERED_AFTER_CLEAN_RUNS:
with contextlib.suppress(OSError):
path.unlink(missing_ok=True)
else:
_write_marker({**state, "clean_runs": clean_runs})


def notify(message: str, dedup: bool = True) -> bool:
Expand All @@ -208,6 +263,7 @@ def notify(message: str, dedup: bool = True) -> bool:
if not url:
return False
if dedup and _suppressed(message):
_restart_recovery()
return False
try:
if "ntfy" in url:
Expand Down
10 changes: 6 additions & 4 deletions lift_status/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import time
import urllib.error
import urllib.request
import zlib

URL = "https://connect.irishrail.ie/realtime/messages?lang=en"

Expand Down Expand Up @@ -136,10 +137,11 @@ def _request(self) -> tuple[int, str]:
raise ApiError(f"{exc.code} from {self.url}: {body}", status=exc.code) from exc
except (urllib.error.URLError, TimeoutError, OSError) as exc:
raise TransientError(f"network failure for {self.url}: {exc}", status=None) from exc
# IncompleteRead/BadStatusLine are HTTPException, and UnicodeDecodeError
# from _decode is a ValueError - neither is an OSError, so without these
# two clauses they escape _run and the attempt is never logged at all.
except http.client.HTTPException as exc:
# IncompleteRead/BadStatusLine are HTTPException, a truncated or corrupt
# gzip body is EOFError or zlib.error, and UnicodeDecodeError from
# _decode is a ValueError - none is an OSError, so without these clauses
# they escape _run and the attempt is never logged at all.
except (http.client.HTTPException, EOFError, zlib.error) as exc:
raise TransientError(f"broken response from {self.url}: {exc!r}", status=None) from exc
except UnicodeDecodeError as exc:
raise ApiError(
Expand Down
76 changes: 53 additions & 23 deletions lift_status/poll.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
import contextlib
import fcntl
import json
import sqlite3
import sys
import uuid
from dataclasses import dataclass, field
Expand All @@ -25,7 +26,7 @@
from . import alert
from .client import ApiError, AuthError, MessagesClient, TransientError
from .parse import NOT_A_LIST, check_item_schema, parse_top_level
from .store import Store, utc_now_iso
from .store import Store, append_raw, utc_now_iso


@dataclass
Expand All @@ -40,6 +41,24 @@ class ApplyResult:
diff: dict | None = None


def classify_fetch_failure(http_status, body_text, network_error):
"""(outcome, detail, exit_code) when nothing was collected, else None."""
if http_status is None or http_status >= 400 or body_text is None:
# Either a true network-level failure (no response at all) or an HTTP
# error status. auth_error is split out because it needs a distinct,
# much louder alert; every other error status is folded into
# "unreachable" since the practical outcome is identical either way -
# nothing was collected this run.
#
# A sub-400 status does not imply a body: urllib does not follow a 300
# or 304, so those arrive as ApiError with no body, and a None body
# reaching json.loads() would crash past every alert path.
outcome = "auth_error" if http_status in (401, 403) else "unreachable"
exit_code = alert.EXIT_AUTH if outcome == "auth_error" else alert.EXIT_UNREACHABLE
return outcome, network_error or f"HTTP {http_status}", exit_code
return None


def apply_response(
store: Store, run_uuid: str, fetched_at: str, http_status, body_text, network_error
) -> ApplyResult:
Expand All @@ -57,20 +76,9 @@ def failed(outcome, detail, exit_code):
outcome=outcome, exit_code=exit_code, http_status=http_status, error_detail=detail
)

if http_status is None or http_status >= 400 or body_text is None:
# Either a true network-level failure (no response at all) or an HTTP
# error status. auth_error is split out because it needs a distinct,
# much louder alert; every other error status is folded into
# "unreachable" since the practical outcome is identical either way -
# nothing was collected this run.
#
# A sub-400 status does not imply a body: urllib does not follow a 300
# or 304, so those arrive as ApiError with no body, and a None body
# reaching json.loads() would crash past every alert path.
outcome = "auth_error" if http_status in (401, 403) else "unreachable"
exit_code = alert.EXIT_AUTH if outcome == "auth_error" else alert.EXIT_UNREACHABLE
detail = network_error or f"HTTP {http_status}"
return failed(outcome, detail, exit_code)
fetch_failure = classify_fetch_failure(http_status, body_text, network_error)
if fetch_failure:
return failed(*fetch_failure)

try:
parsed = parse_top_level(body_text)
Expand Down Expand Up @@ -171,15 +179,30 @@ def _run(data_dir: Path, client: MessagesClient) -> int:
body_text = None
network_error = repr(exc)

with Store(data_dir) as store:
store.write_raw(run_uuid, fetched_at, http_status, body_text, network_error)
result = apply_response(store, run_uuid, fetched_at, http_status, body_text, network_error)
try:
append_raw(data_dir, run_uuid, fetched_at, http_status, body_text, network_error)
except OSError as exc:
# check_writable's empty probe file passes on a full SD card.
problem = f"cannot append to the raw log: {exc}"
return alert.fail(alert.storage_banner(data_dir, problem), alert.EXIT_STORAGE)

if result.outcome == "auth_error":
banner = alert.auth_banner(client.masked_key, result.error_detail or "")
return alert.fail(banner, result.exit_code)
if result.outcome == "unreachable":
return alert.fail(alert.unreachable_banner(result.error_detail or ""), result.exit_code)
try:
with Store(data_dir) as store:
result = apply_response(
store, run_uuid, fetched_at, http_status, body_text, network_error
)
except (sqlite3.Error, OSError) as exc:
fetch_failure = classify_fetch_failure(http_status, body_text, network_error)
if fetch_failure:
# Nothing was collected, so the fetch is the news; the database
# alerts on the first run that has a response to lose.
return _fetch_failure_alert(client, *fetch_failure)
return alert.fail(alert.database_banner(data_dir, repr(exc)), alert.EXIT_DATABASE)

if result.outcome in ("auth_error", "unreachable"):
return _fetch_failure_alert(
client, result.outcome, result.error_detail or "", result.exit_code
)
if result.outcome in ("parse_error", "not_a_list"):
return alert.fail(alert.schema_root_banner(), result.exit_code)

Expand All @@ -192,9 +215,16 @@ def _run(data_dir: Path, client: MessagesClient) -> int:
)
if result.schema_drift_count:
return alert.fail(alert.schema_banner(result.drift_problems), result.exit_code)
alert.note_clean_run()
return alert.EXIT_OK


def _fetch_failure_alert(client: MessagesClient, outcome, detail, exit_code) -> int:
if outcome == "auth_error":
return alert.fail(alert.auth_banner(client.masked_key, detail), exit_code)
return alert.fail(alert.unreachable_banner(detail), exit_code)


def run_check(client: MessagesClient | None = None) -> int:
"""Validate connectivity and the API key without writing anything.

Expand Down
Loading
Loading