diff --git a/CLAUDE.md b/CLAUDE.md index e8a0186..c305dcc 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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 @@ -175,6 +175,7 @@ merged with `sort -u`. | Station access is labelled by hand into an append-only observation log, `lifts-data/survey/.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 diff --git a/lift_access/fetch.py b/lift_access/fetch.py index 2c4e971..e3bda0c 100644 --- a/lift_access/fetch.py +++ b/lift_access/fetch.py @@ -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. diff --git a/lift_status/alert.py b/lift_status/alert.py index 7985276..b3f3fce 100644 --- a/lift_status/alert.py +++ b/lift_status/alert.py @@ -9,6 +9,7 @@ from __future__ import annotations +import contextlib import hashlib import json import os @@ -23,6 +24,7 @@ EXIT_UNREACHABLE = 3 EXIT_SCHEMA_DRIFT = 4 EXIT_STORAGE = 6 +EXIT_DATABASE = 7 EXIT_MEANINGS = { EXIT_OK: "success", @@ -30,6 +32,7 @@ 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 @@ -37,6 +40,11 @@ # 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 @@ -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" @@ -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: @@ -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: diff --git a/lift_status/client.py b/lift_status/client.py index e0ae816..c924b15 100644 --- a/lift_status/client.py +++ b/lift_status/client.py @@ -19,6 +19,7 @@ import time import urllib.error import urllib.request +import zlib URL = "https://connect.irishrail.ie/realtime/messages?lang=en" @@ -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( diff --git a/lift_status/poll.py b/lift_status/poll.py index 5c618d1..9e81ff2 100644 --- a/lift_status/poll.py +++ b/lift_status/poll.py @@ -17,6 +17,7 @@ import contextlib import fcntl import json +import sqlite3 import sys import uuid from dataclasses import dataclass, field @@ -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 @@ -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: @@ -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) @@ -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) @@ -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. diff --git a/lift_status/store.py b/lift_status/store.py index e1a6889..edec629 100644 --- a/lift_status/store.py +++ b/lift_status/store.py @@ -97,6 +97,43 @@ def utc_now_iso() -> str: return datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ") +def append_raw(data_dir, run_uuid, fetched_at, http_status, body, network_error) -> None: + """Append one line for this run attempt. Called before any parsing, and + before the database is opened, so a database that cannot be opened never + costs the response. + + Fsynced: a run happens once per 30 minutes, so the extra syscall cost + is irrelevant, and this is the durability point the rest of the design + depends on - a crash or power loss right after this call must not lose + the response. + """ + raw_dir = Path(data_dir) / RAW_DIRNAME + raw_dir.mkdir(parents=True, exist_ok=True) + date_part = fetched_at[:10].replace("-", "") + path = raw_dir / f"messages-{date_part}.jsonl" + line = json.dumps( + { + "run_uuid": run_uuid, + "fetched_at_utc": fetched_at, + "http_status": http_status, + "body": body, + "network_error": network_error, + }, + sort_keys=True, + ) + with path.open("a+b") as f: + # A power cut mid-append leaves a last line with no newline, and the + # next record appended onto it would be lost with the fragment. + end = f.seek(0, os.SEEK_END) + if end: + f.seek(end - 1) + if f.read(1) != b"\n": + f.write(b"\n") + f.write((line + "\n").encode("utf-8")) + f.flush() + os.fsync(f.fileno()) + + class Store: def __init__(self, data_dir): self.data_dir = Path(data_dir) @@ -118,59 +155,36 @@ def __exit__(self, exc_type, exc, tb) -> None: # -- raw JSONL log ----------------------------------------------------- - def write_raw(self, run_uuid, fetched_at, http_status, body, network_error) -> None: - """Append one line for this run attempt. Called before any parsing. - - Fsynced: a run happens once per 30 minutes, so the extra syscall cost - is irrelevant, and this is the durability point the rest of the design - depends on - a crash or power loss right after this call must not lose - the response. - """ - date_part = fetched_at[:10].replace("-", "") - path = self.raw_dir / f"messages-{date_part}.jsonl" - line = json.dumps( - { - "run_uuid": run_uuid, - "fetched_at_utc": fetched_at, - "http_status": http_status, - "body": body, - "network_error": network_error, - }, - sort_keys=True, - ) - with path.open("a", encoding="utf-8") as f: - f.write(line + "\n") - f.flush() - os.fsync(f.fileno()) - def iter_raw_lines(self): - """Yield every recorded run attempt, oldest file first, append order - within a file. Never sorts by the embedded timestamp - a Pi's clock can - jump (e.g. before NTP sync after a reboot), and replay must follow the - order runs actually happened in, not a timestamp that might be wrong. + """Yield every recorded run attempt, oldest file first, and in + `fetched_at_utc` order within a file, so two collectors' logs merged + with `sort -u` (which orders lines by their text) replay as they + happened. The sort is stable, so runs stamped in the same second keep + their order in the file. An undecodable line is skipped and counted in self.raw_decode_errors - rather than aborting the replay: write_raw's append is not atomic, so a + rather than aborting the replay: append_raw is not atomic, so a power cut leaves a truncated last line, and one bad line must not make every good line behind it unreplayable. """ self.raw_decode_errors = 0 for path in sorted(self.raw_dir.glob("messages-*.jsonl")): + records = [] with path.open("r", encoding="utf-8") as f: for lineno, line in enumerate(f, 1): line = line.strip() if not line: continue try: - record = json.loads(line) + records.append(json.loads(line)) except json.JSONDecodeError as exc: self.raw_decode_errors += 1 print( f"warning: skipping unreadable line {path.name}:{lineno}: {exc}", file=sys.stderr, ) - continue - yield record + records.sort(key=lambda r: r.get("fetched_at_utc") or "") + yield from records # -- runs ---------------------------------------------------------- @@ -349,7 +363,7 @@ def diff_and_update_messages(self, run_id: int, observed_at: str, items: list) - grace = max(1, int(raw_grace)) except ValueError: # A typo in the env file must not stop collection dead here, after - # write_raw and before any alert path. + # the raw append and before any alert path. print( f"warning: LIFT_STATUS_GRACE_MISSES={raw_grace!r} is not a number; " f"using {DEFAULT_GRACE_MISSES}", diff --git a/notes/collector-review.md b/notes/collector-review.md new file mode 100644 index 0000000..e9566da --- /dev/null +++ b/notes/collector-review.md @@ -0,0 +1,82 @@ +# The collector's first review + +`lift_status/` and `scripts/` were written before PRs here went through a review +agent, and barely changed after: `parse.py`, `client.py` and `__main__.py` not at +all, `poll.py` by 16 lines. They are also the only code whose mistakes a +rebuild cannot undo, since a response that never reached the raw log is gone. +So they were reviewed once, as whole files, on 2026-09-24. The rest of the +repository was left alone: `lift_access/` and most of `lift_site/` have been +reviewed diff by diff since 2026-08-26, and the real-corpus and golden tests pin +what they publish. + +Ten findings. Eight were fixed, one was not a real path, and one is left for now. + +## Fixed + +- **The database could cost the response.** `Store()` opened SQLite and ran the + schema before the raw append, so a database left corrupt by a power cut, or + locked past the 5s busy timeout, stopped every poll from reaching the log, + with a traceback and no alert. `append_raw` now writes the line before the + database is opened, and a database failure is its own alert and exit code 7, + saying the response was kept and what to do for a lock, a full card or + corruption. A fetch that failed keeps its own alert even when the database + is broken too, since then there was no response to keep. +- **A full SD card was a traceback, not the storage alert.** `check_writable` + touches an empty file, which needs no data block and passes on a full disk. + An `OSError` from the append now goes to the storage banner. +- **A cut-off last line took the next good one with it.** A power cut mid-append + leaves no trailing newline, and the next record was appended onto the + fragment, so replay dropped both. The append now finishes the broken line + first, and only the fragment is lost. +- **A truncated or corrupt gzip body escaped the client.** `gzip.decompress` + raises `EOFError` or `zlib.error`, neither an `OSError`, so they bypassed the + retry, the raw line and the alert. They are now a `TransientError`. +- **A second outage within a day of the first was silent.** The alert dedup + marker outlived a recovery, so the same banner after a clean stretch was + suppressed for up to 24h. The marker now counts consecutive clean runs and + goes after four, two hours at the 30-minute cadence; a suppressed failure + resets the count. Clearing it on the first clean run was tried first and + rejected in review: a flapping API would then alert on every failure, which + is what the window exists to stop. +- **systemd could kill a poll before it logged anything.** `TimeoutStartSec=60` + was below the client's own worst case (three attempts of a 15s connect and a + 15s read, plus backoff and DNS). It is 300 now. +- **The backup could hang forever.** A oneshot has no start timeout by default, + and ssh had no keepalive, so a push stalled on a half-open connection held the + unit "activating" and turned every later firing into a no-op. ssh now has + `ConnectTimeout` and `ServerAliveInterval`, and the unit has a 15-minute cap. + The script traps TERM, because dash skips the EXIT trap on a signal it does + not trap, and the cap would otherwise end the backup without an alert. + +- **A `sort -u` merge replayed out of order.** `sort_keys=True` is there so two + collectors' logs can be merged with `sort -u`, which is also how git's + conflict on a shared day file gets resolved. But `sort -u` orders lines by + their first key, `body`, and replay followed line order, so a merged file + replayed out of time order. Replay now sorts each file by `fetched_at_utc`, + stably. The cost is the clock jump the old line-order rule was for: after a + power cut, fake-hwclock restores the last hourly save, so a catch-up poll can + be stamped up to an hour before runs already in the same file, and a rebuild + applies it before them where the live run applied it after. Line order only + ever covered part of that, since a stamp that crosses midnight already lands + in the wrong day's file, and it cannot survive a merge at all. On 2026-09-24 all 2,224 real lines were already in + time order within their files, so no rebuild moved. The owner's call: + merging logs has to work. + +## Not a real path + +- **The backup merges `origin` into the tree the poller appends to, without the + poll lock.** A merge only rewrites files that changed upstream, and nothing + but the Pi writes `raw/`: the stations workflow touches `stations/` and + `survey/`, through PRs. Taking the lock for a whole fetch and push would + instead make a poll skip. + +## Left for now + +- **`time-sync.target` does not wait for NTP.** It is reached as soon as + timesyncd starts unless `systemd-time-wait-sync.service` is enabled, and the + install does not enable it, so a catch-up poll after a reboot can be stamped + with fake-hwclock's time. Enabling the wait service risks a poll that never + runs if NTP is unreachable, which is the silent failure this collector exists + to avoid. Left as it is by the owner on 2026-09-24. If it is taken up, the + shape is a poll that checks whether the clock is synced and still writes the + line, flagged, rather than one that waits or skips. diff --git a/notes/delays-site.md b/notes/delays-site.md index 2b1bbc8..099162a 100644 --- a/notes/delays-site.md +++ b/notes/delays-site.md @@ -18,7 +18,7 @@ and 2026-09-11, over 1620 successful runs. The question "a second collector, or a delay poll target added to `lift_status`" has a shorter answer than it looks: **there is nothing to collect**. One endpoint returns the whole feed in one response - `client.py` says so and the log proves -it - and `store.write_raw` writes that response verbatim before anything reads +it - and `store.append_raw` writes that response verbatim before anything reads it. Every delay notice the other site will ever show is already on disk, back to the first poll. diff --git a/notes/station-access.md b/notes/station-access.md index 4175891..5f611df 100644 --- a/notes/station-access.md +++ b/notes/station-access.md @@ -753,7 +753,7 @@ nothing may ever need to. ## The snapshot `lifts-data/stations/irishrail-.jsonl` holds every payload **verbatim**, -one per line, `sort_keys=True`, the shape `store.write_raw` uses. 7.8 MB plain, +one per line, `sort_keys=True`, the shape `store.append_raw` uses. 7.8 MB plain, which git stores at about 2 MB and which greps and diffs. Never edited; the derivation is always recomputed from it. diff --git a/scripts/backup-to-git.sh b/scripts/backup-to-git.sh index 5e7d3e6..ea2125c 100755 --- a/scripts/backup-to-git.sh +++ b/scripts/backup-to-git.sh @@ -37,6 +37,9 @@ Nothing new is offsite. Check: exit "$status" } trap on_exit EXIT +# dash skips the EXIT trap on a signal it does not trap, and TERM is how the +# unit's TimeoutStartSec ends a stalled push. +trap 'exit 143' TERM INT cd "$DATA_DIR" || { notify "lift-status backup: $DATA_DIR does not exist. Nothing is being backed up." diff --git a/scripts/systemd/lift-status-backup.service b/scripts/systemd/lift-status-backup.service index 245d7a9..22980c0 100644 --- a/scripts/systemd/lift-status-backup.service +++ b/scripts/systemd/lift-status-backup.service @@ -16,7 +16,13 @@ Environment=LIFT_STATUS_DATA_DIR=/var/lib/lift-status # inside the tree being backed up. # Quoted: systemd splits an unquoted Environment= value on whitespace into # separate assignments, which would set GIT_SSH_COMMAND to just "ssh". -Environment="GIT_SSH_COMMAND=ssh -i /etc/lift-status-deploy-key -o IdentitiesOnly=yes -o UserKnownHostsFile=/etc/lift-status-known_hosts -o StrictHostKeyChecking=accept-new" +Environment="GIT_SSH_COMMAND=ssh -i /etc/lift-status-deploy-key -o IdentitiesOnly=yes -o UserKnownHostsFile=/etc/lift-status-known_hosts -o StrictHostKeyChecking=accept-new -o ConnectTimeout=30 -o ServerAliveInterval=15 -o ServerAliveCountMax=4" + +# A oneshot has no start timeout by default, so a push stalled on a half-open +# connection would hold the unit "activating" and turn every later firing into +# a no-op. The ssh keepalives above make the stall fail into the script's own +# alert; this is the backstop, and the script alerts on its TERM too. +TimeoutStartSec=15min EnvironmentFile=-/etc/lift-status.env ExecStart=/opt/lift-status/scripts/backup-to-git.sh diff --git a/scripts/systemd/lift-status.service b/scripts/systemd/lift-status.service index 793321c..6c5ccd9 100644 --- a/scripts/systemd/lift-status.service +++ b/scripts/systemd/lift-status.service @@ -26,7 +26,11 @@ EnvironmentFile=-/etc/lift-status.env WorkingDirectory=/opt/lift-status ExecStart=/usr/bin/python3 -m lift_status poll -TimeoutStartSec=60 +# Above the client's own worst case - three attempts, each a 15s connect and a +# 15s read, plus backoff and DNS - so a slow API ends in a logged run and an +# alert rather than a SIGTERM that writes neither. Still well inside the +# 30-minute interval. +TimeoutStartSec=300 # It reads its own code and writes one state directory. Nothing else. NoNewPrivileges=true diff --git a/tests/test_alert.py b/tests/test_alert.py index 07a9c73..54e93b9 100644 --- a/tests/test_alert.py +++ b/tests/test_alert.py @@ -64,6 +64,23 @@ def test_the_retry_after_a_failure_is_delivered_and_then_suppresses(self): self.assertTrue(self._send()) self.assertTrue(alert._suppressed(BANNER)) + def _clean_runs(self, n): + for _ in range(n): + alert.note_clean_run() + + def test_a_recovery_that_holds_makes_the_same_fault_a_new_alert(self): + self._send() + self._clean_runs(alert.RECOVERED_AFTER_CLEAN_RUNS - 1) + self.assertTrue(self.marker.exists()) + self._clean_runs(1) + self.assertFalse(alert._suppressed(BANNER)) + + def test_a_fault_flapping_between_clean_polls_stays_one_alert(self): + self._send() + for _ in range(3): + self._clean_runs(alert.RECOVERED_AFTER_CLEAN_RUNS - 1) + self.assertFalse(self._send()) + def test_a_different_banner_is_never_suppressed(self): self._send() self.assertFalse(alert._suppressed("lift-status: the disk is full")) diff --git a/tests/test_client.py b/tests/test_client.py index 1e4d6f5..0b43c16 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -107,6 +107,18 @@ def test_gzip_response_is_decoded_regardless_of_accept_encoding(self): status, body = client.get_messages_raw() self.assertEqual(body, "[]") + def test_a_truncated_or_corrupt_gzip_body_is_a_transient_error(self): + import gzip + + whole = gzip.compress(b"[]" * 1000) + for broken in (whole[: len(whole) // 2], whole[:10] + b"x" * 20 + whole[30:]): + resp = FakeResponse(broken, headers={"Content-Encoding": "gzip"}) + client = MessagesClient(retries=2, sleep=lambda s: None) + with mock.patch("urllib.request.urlopen", return_value=resp) as urlopen: + with self.assertRaises(TransientError): + client.get_messages_raw() + self.assertEqual(urlopen.call_count, 2) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_poll.py b/tests/test_poll.py index ebfdfa0..2f3f79d 100644 --- a/tests/test_poll.py +++ b/tests/test_poll.py @@ -1,3 +1,4 @@ +import errno import fcntl import json import os @@ -16,6 +17,10 @@ class PollTestCase(unittest.TestCase): def setUp(self): self._tmp = tempfile.TemporaryDirectory() self.data_dir = Path(self._tmp.name) + # The alert marker lives under this variable, and a clean run deletes it. + env = mock.patch.dict("os.environ", {"LIFT_STATUS_DATA_DIR": str(self.data_dir)}) + env.start() + self.addCleanup(env.stop) def tearDown(self): self._tmp.cleanup() @@ -186,6 +191,49 @@ def test_unwritable_data_dir_exits_storage_code(self): os.chmod(self.data_dir, 0o700) +class TestTheRawLineDoesNotDependOnTheDatabase(PollTestCase): + def _raw_lines(self): + return [ + json.loads(line) + for path in (self.data_dir / "raw").glob("*.jsonl") + for line in path.read_text(encoding="utf-8").splitlines() + ] + + def test_a_corrupt_database_still_logs_the_response_and_alerts(self): + (self.data_dir / "lift_status.db").write_bytes(b"not a database, after a power cut") + body = json.dumps([make_item()]) + code = poll.run_poll(self.data_dir, client=FakeClient([(200, body)])) + self.assertEqual(code, alert.EXIT_DATABASE) + self.assertEqual([r["body"] for r in self._raw_lines()], [body]) + + def test_a_failed_fetch_keeps_its_own_alert_when_the_database_is_broken_too(self): + (self.data_dir / "lift_status.db").write_bytes(b"not a database, after a power cut") + code = poll.run_poll(self.data_dir, client=FakeClient([AuthError("401", status=401)])) + self.assertEqual(code, alert.EXIT_AUTH) + self.assertEqual(len(self._raw_lines()), 1) + + def test_a_full_disk_at_the_append_is_a_storage_alert(self): + full = OSError(errno.ENOSPC, "No space left on device") + with mock.patch.object(poll, "append_raw", side_effect=full): + code = poll.run_poll(self.data_dir, client=FakeClient([(200, "[]")])) + self.assertEqual(code, alert.EXIT_STORAGE) + self.assertEqual(self._run_count(), 0) + + +class TestACleanRunClosesTheRepeatWindow(PollTestCase): + def test_clean_runs_clear_the_marker_and_a_failed_one_does_not(self): + marker = self.data_dir / ".last-alert.json" + marker.write_text("{}", encoding="utf-8") + poll.run_poll(self.data_dir, client=FakeClient([TransientError("down")])) + self.assertTrue(marker.exists()) + clean = [(200, "[]")] * alert.RECOVERED_AFTER_CLEAN_RUNS + for response in clean[:-1]: + poll.run_poll(self.data_dir, client=FakeClient([response])) + self.assertTrue(marker.exists()) + poll.run_poll(self.data_dir, client=FakeClient(clean[-1:])) + self.assertFalse(marker.exists()) + + class TestMissingApiKey(PollTestCase): """No key configured is a config error, distinct from a rejected key: it must not fetch and must not record a run.""" diff --git a/tests/test_rebuild.py b/tests/test_rebuild.py index 80d9c57..172a547 100644 --- a/tests/test_rebuild.py +++ b/tests/test_rebuild.py @@ -14,7 +14,7 @@ from lift_status import poll from lift_status.client import AuthError -from lift_status.store import Store +from lift_status.store import Store, append_raw from tests.helpers import FakeClient, make_item STATION_A = make_item(head="Station A - Lift out of order", codes=["AAA"]) @@ -118,6 +118,33 @@ def test_synthetic_history_round_trips_through_rebuild(self): self.assertEqual(before["runs"], after["runs"]) self.assertEqual(before["unidentifiable"], after["unidentifiable"]) + def test_two_collectors_logs_merged_with_sort_u_rebuild_the_same_history(self): + seen = json.dumps([STATION_A, STATION_B]) + only_a = json.dumps([STATION_A]) + history = [ + ("r1", "2026-08-08T12:00:00Z", 200, seen, None), + ("r2", "2026-08-08T12:30:00Z", None, None, "TransientError('down')"), + ("r3", "2026-08-08T13:00:00Z", 200, only_a, None), + ("r4", "2026-08-08T13:30:00Z", 200, only_a, None), + ("r5", "2026-08-08T14:00:00Z", 200, seen, None), + ] + for run in history: + append_raw(self.data_dir, *run) + self.assertEqual(poll.run_rebuild(self.data_dir), 0) + before = _snapshot(self.data_dir) + [b] = [m for m in before["messages"].values() if m["head"] == STATION_B["head"]] + self.assertEqual(b["reopen_count"], 1) + + # The second machine logged the same runs; git's conflict, resolved the + # way CLAUDE.md says, is both copies through sort -u. + path = self.data_dir / "raw" / "messages-20260808.jsonl" + lines = path.read_text(encoding="utf-8").splitlines() + path.write_text("\n".join(sorted(set(lines + lines))) + "\n", encoding="utf-8") + self.assertNotEqual(path.read_text(encoding="utf-8").splitlines(), lines) + + self.assertEqual(poll.run_rebuild(self.data_dir), 0) + self.assertEqual(_snapshot(self.data_dir), before) + def test_rebuild_with_no_raw_logs_is_a_harmless_noop(self): # A fresh install has nothing to replay and nothing to lose. code = poll.run_rebuild(self.data_dir) diff --git a/tests/test_store.py b/tests/test_store.py index 130b890..b5f230e 100644 --- a/tests/test_store.py +++ b/tests/test_store.py @@ -5,7 +5,7 @@ import unittest from pathlib import Path -from lift_status.store import Store, utc_now_iso +from lift_status.store import Store, append_raw, utc_now_iso from tests.helpers import make_item _run_counter = itertools.count() @@ -49,7 +49,7 @@ def _listings(self, key): class TestWriteRaw(StoreTestCase): def test_writes_one_jsonl_line(self): - self.store.write_raw("run-1", "2026-08-08T12:00:00Z", 200, "[]", None) + append_raw(self.data_dir, "run-1", "2026-08-08T12:00:00Z", 200, "[]", None) path = self.data_dir / "raw" / "messages-20260808.jsonl" self.assertTrue(path.exists()) line = json.loads(path.read_text().strip()) @@ -59,12 +59,22 @@ def test_writes_one_jsonl_line(self): self.assertIsNone(line["network_error"]) def test_appends_multiple_lines_same_day(self): - self.store.write_raw("run-1", "2026-08-08T12:00:00Z", 200, "[]", None) - self.store.write_raw("run-2", "2026-08-08T12:30:00Z", 200, "[]", None) + append_raw(self.data_dir, "run-1", "2026-08-08T12:00:00Z", 200, "[]", None) + append_raw(self.data_dir, "run-2", "2026-08-08T12:30:00Z", 200, "[]", None) path = self.data_dir / "raw" / "messages-20260808.jsonl" lines = path.read_text().strip().splitlines() self.assertEqual(len(lines), 2) + def test_a_line_cut_off_by_a_power_cut_does_not_take_the_next_one_with_it(self): + path = self.data_dir / "raw" / "messages-20260808.jsonl" + append_raw(self.data_dir, "run-1", "2026-08-08T12:00:00Z", 200, "[]", None) + with path.open("a", encoding="utf-8") as f: + f.write('{"body": "[{\\"head') + append_raw(self.data_dir, "run-3", "2026-08-08T13:00:00Z", 200, "[]", None) + replayed = [r["run_uuid"] for r in self.store.iter_raw_lines()] + self.assertEqual(replayed, ["run-1", "run-3"]) + self.assertEqual(self.store.raw_decode_errors, 1) + class TestDiffAndUpdateMessages(StoreTestCase): def test_new_message_is_inserted_open(self):