From c96996adf01f685edf8ad355c5a1ccfca52d9483 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 06:50:28 +0000 Subject: [PATCH 1/3] Review the collector: the raw line no longer waits on the database lift_status/ and scripts/ predate review on this repository and are the only code whose mistakes a rebuild cannot undo, so they were reviewed whole, once. Seven of ten findings are fixed here; notes/collector-review.md has all ten, including the one that was not a real path and the two left open as decisions. - Write the raw line before opening SQLite. A database left corrupt by a power cut stopped every poll from reaching the log, with no alert; it is now exit 7 with a banner saying the response was kept. - Route an OSError from the append to the storage alert: the empty probe file in check_writable passes on a full SD card. - Finish a line a power cut left without a newline before appending, so the fragment no longer takes the next good record with it. - Treat a truncated or corrupt gzip body (EOFError, zlib.error) as a transient error instead of letting it escape unlogged. - Clear the alert dedup marker on a clean run, so the same fault coming back within a day alerts again. - Raise the poll unit's TimeoutStartSec above the client's own worst case, and bound the backup unit and its ssh so a stalled push cannot hang it. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01AnrcmsrTHYCgnVGyqBqji4 --- CLAUDE.md | 3 +- lift_status/alert.py | 26 +++++++++ lift_status/client.py | 10 ++-- lift_status/poll.py | 21 +++++-- lift_status/store.py | 64 +++++++++++++-------- notes/collector-review.md | 67 ++++++++++++++++++++++ scripts/systemd/lift-status-backup.service | 8 ++- scripts/systemd/lift-status.service | 6 +- tests/test_alert.py | 6 ++ tests/test_client.py | 12 ++++ tests/test_poll.py | 38 ++++++++++++ tests/test_store.py | 10 ++++ 12 files changed, 237 insertions(+), 34 deletions(-) create mode 100644 notes/collector-review.md diff --git a/CLAUDE.md b/CLAUDE.md index e8a0186..1b3d6ae 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 is opened, so a corrupt or locked database costs the derived rows and never the response, and alerts as exit 7. The collector and its units were reviewed whole once; two findings are open by choice | `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_status/alert.py b/lift_status/alert.py index 7985276..ad1437e 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 @@ -153,6 +156,22 @@ 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. A database left", + "corrupt by a power cut is rebuilt from the raw log:", + "", + 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" @@ -197,6 +216,13 @@ def _mark_delivered(message: str) -> None: pass +def clear() -> None: + """Close the repeat window on a clean run, so the same fault coming back + later is a new incident and alerts again rather than sitting out the day.""" + with contextlib.suppress(OSError): + _marker_path().unlink(missing_ok=True) + + def notify(message: str, dedup: bool = True) -> bool: """Push to LIFT_STATUS_ALERT_WEBHOOK. Returns whether it was delivered. 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..a900706 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 @@ -171,9 +172,20 @@ 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) + + 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: + return alert.fail(alert.database_banner(data_dir, repr(exc)), alert.EXIT_DATABASE) if result.outcome == "auth_error": banner = alert.auth_banner(client.masked_key, result.error_detail or "") @@ -192,6 +204,7 @@ 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.clear() return alert.EXIT_OK diff --git a/lift_status/store.py b/lift_status/store.py index e1a6889..00b8955 100644 --- a/lift_status/store.py +++ b/lift_status/store.py @@ -97,6 +97,46 @@ 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("ab") 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. + if f.tell() and _last_byte(path) != b"\n": + f.write(b"\n") + f.write((line + "\n").encode("utf-8")) + f.flush() + os.fsync(f.fileno()) + + +def _last_byte(path: Path) -> bytes: + with path.open("rb") as f: + f.seek(-1, os.SEEK_END) + return f.read(1) + + class Store: def __init__(self, data_dir): self.data_dir = Path(data_dir) @@ -119,29 +159,7 @@ 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()) + append_raw(self.data_dir, run_uuid, fetched_at, http_status, body, network_error) def iter_raw_lines(self): """Yield every recorded run attempt, oldest file first, append order diff --git a/notes/collector-review.md b/notes/collector-review.md new file mode 100644 index 0000000..5062c57 --- /dev/null +++ b/notes/collector-review.md @@ -0,0 +1,67 @@ +# 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. Seven were fixed, one was not a real path, and two are open. + +## Fixed + +- **The database could cost the response.** `Store()` opened SQLite and ran the + schema before `write_raw`, 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. +- **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 run was + suppressed for up to 24h. A clean run now clears it. The cost is that a + fault flapping at poll granularity alerts on each return; each poll already + retries three times, so that is a real outage each time. +- **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. + +## 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. + +## Open + +- **`sort -u` and replay order disagree.** CLAUDE.md says `sort_keys=True` lets + two machines' logs be merged with `sort -u`. Replay follows line order within + a file, and `sort -u` orders lines by their first key, `body`, so a merged + file replays out of time order. It has never happened (one collector), and + the fix is a choice: sort by `fetched_at_utc` within a file on replay, which + `iter_raw_lines` currently refuses because of clock jumps, or merge by + timestamp rather than `sort -u`. +- **`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, so it is not changed blind. diff --git a/scripts/systemd/lift-status-backup.service b/scripts/systemd/lift-status-backup.service index 245d7a9..3823a35 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. +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..daa0b77 100644 --- a/tests/test_alert.py +++ b/tests/test_alert.py @@ -64,6 +64,12 @@ def test_the_retry_after_a_failure_is_delivered_and_then_suppresses(self): self.assertTrue(self._send()) self.assertTrue(alert._suppressed(BANNER)) + def test_a_clean_run_in_between_makes_the_same_fault_a_new_alert(self): + self._send() + alert.clear() + self.assertFalse(alert._suppressed(BANNER)) + alert.clear() + 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..85b775a 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,39 @@ 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_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_a_successful_run_clears_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()) + poll.run_poll(self.data_dir, client=FakeClient([(200, "[]")])) + 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_store.py b/tests/test_store.py index 130b890..6f954cd 100644 --- a/tests/test_store.py +++ b/tests/test_store.py @@ -65,6 +65,16 @@ def test_appends_multiple_lines_same_day(self): 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" + self.store.write_raw("run-1", "2026-08-08T12:00:00Z", 200, "[]", None) + with path.open("a", encoding="utf-8") as f: + f.write('{"body": "[{\\"head') + self.store.write_raw("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): From e1cffd24e012f9e88b9f5114e07e1b0620754ef3 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 08:18:07 +0000 Subject: [PATCH 2/3] Replay each raw file in fetched_at order, so merged logs rebuild the same sort_keys=True exists so two collectors' logs can be merged with sort -u, which is also how git's conflict on a shared day file would be 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: closures, miss counts and listing stretches would all have come out wrong. iter_raw_lines now sorts each file by fetched_at_utc, stably. The old line-order rule guarded against clock jumps, but only partly, since a pre-NTP stamp already lands in the wrong day's file. All 2,224 real lines are already in time order within their files, and rebuild then stats gives identical output before and after. The NTP wait stays as it is, by the owner's call, and the note says so. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01AnrcmsrTHYCgnVGyqBqji4 --- CLAUDE.md | 2 +- lift_status/store.py | 18 ++++++++++-------- notes/collector-review.md | 26 ++++++++++++++++---------- tests/test_rebuild.py | 28 ++++++++++++++++++++++++++++ 4 files changed, 55 insertions(+), 19 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 1b3d6ae..a78da90 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -175,7 +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 is opened, so a corrupt or locked database costs the derived rows and never the response, and alerts as exit 7. The collector and its units were reviewed whole once; two findings are open by choice | `notes/collector-review.md` | +| The raw line is written before the database is opened, so a corrupt or locked database costs the derived rows and never the response, and alerts as exit 7. Replay orders each file by `fetched_at_utc`, so a `sort -u` merge of two collectors' logs rebuilds the same history. The collector and its units were reviewed whole once; 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_status/store.py b/lift_status/store.py index 00b8955..8b6a800 100644 --- a/lift_status/store.py +++ b/lift_status/store.py @@ -162,33 +162,35 @@ def write_raw(self, run_uuid, fetched_at, http_status, body, network_error) -> N append_raw(self.data_dir, run_uuid, fetched_at, http_status, body, network_error) 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 ---------------------------------------------------------- diff --git a/notes/collector-review.md b/notes/collector-review.md index 5062c57..7af63dc 100644 --- a/notes/collector-review.md +++ b/notes/collector-review.md @@ -9,7 +9,7 @@ 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. Seven were fixed, one was not a real path, and two are open. +Ten findings. Eight were fixed, one was not a real path, and one is left for now. ## Fixed @@ -42,6 +42,17 @@ Ten findings. Seven were fixed, one was not a real path, and two are open. 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. +- **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 line-order rule was there for clock jumps, and it only ever + covered part of them: a pre-NTP stamp already lands in the wrong day's file, + which is ordered by name. 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 @@ -50,18 +61,13 @@ Ten findings. Seven were fixed, one was not a real path, and two are open. `survey/`, through PRs. Taking the lock for a whole fetch and push would instead make a poll skip. -## Open +## Left for now -- **`sort -u` and replay order disagree.** CLAUDE.md says `sort_keys=True` lets - two machines' logs be merged with `sort -u`. Replay follows line order within - a file, and `sort -u` orders lines by their first key, `body`, so a merged - file replays out of time order. It has never happened (one collector), and - the fix is a choice: sort by `fetched_at_utc` within a file on replay, which - `iter_raw_lines` currently refuses because of clock jumps, or merge by - timestamp rather than `sort -u`. - **`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, so it is not changed blind. + 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/tests/test_rebuild.py b/tests/test_rebuild.py index 80d9c57..1d0847d 100644 --- a/tests/test_rebuild.py +++ b/tests/test_rebuild.py @@ -118,6 +118,34 @@ 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), + ] + with Store(self.data_dir) as store: + for run in history: + store.write_raw(*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) From d0ce1b94815cdc29367f6d76be6439804155fe02 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 08:24:17 +0000 Subject: [PATCH 3/3] Act on the review of the collector PR - A fetch that failed keeps its own alert when the database is broken too. The database banner said the response was kept when there was no response, and hid a rejected key behind exit 7. - The database banner says what to do for a lock, a full card and corruption, instead of telling every error to move the database aside. - The alert marker clears after four consecutive clean runs, not the first. Clearing on one clean poll made a flapping API alert on every failure, which is what the repeat window exists to stop. - backup-to-git.sh traps TERM: dash skips the EXIT trap on a signal it does not trap, so the new 15-minute cap ended a stalled push silently. - Store.write_raw is gone, so the only way to write the raw line is the one that does not open the database first; the newline check reads through the append handle. - The replay-order note states the power-cut case it gives up, and the CLAUDE.md row is a pointer again. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01AnrcmsrTHYCgnVGyqBqji4 --- CLAUDE.md | 2 +- lift_access/fetch.py | 2 +- lift_status/alert.py | 58 ++++++++++++++++------ lift_status/poll.py | 57 +++++++++++++-------- lift_status/store.py | 20 +++----- notes/collector-review.md | 27 ++++++---- notes/delays-site.md | 2 +- notes/station-access.md | 2 +- scripts/backup-to-git.sh | 3 ++ scripts/systemd/lift-status-backup.service | 2 +- tests/test_alert.py | 17 +++++-- tests/test_poll.py | 14 +++++- tests/test_rebuild.py | 7 ++- tests/test_store.py | 12 ++--- 14 files changed, 149 insertions(+), 76 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index a78da90..c305dcc 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -175,7 +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 is opened, so a corrupt or locked database costs the derived rows and never the response, and alerts as exit 7. Replay orders each file by `fetched_at_utc`, so a `sort -u` merge of two collectors' logs rebuilds the same history. The collector and its units were reviewed whole once; the NTP wait is left for now | `notes/collector-review.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 ad1437e..b3f3fce 100644 --- a/lift_status/alert.py +++ b/lift_status/alert.py @@ -40,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 @@ -163,11 +168,17 @@ def database_banner(data_dir, detail: str) -> str: 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. A database left", - "corrupt by a power cut is rebuilt from the raw log:", + "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}", "", - f" sudo mv {data_dir}/lift_status.db {data_dir}/lift_status.db.broken", - " sudo lift rebuild", + " '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", ], ) @@ -208,19 +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. """ - try: - _marker_path().write_text( - json.dumps({"digest": _digest(message), "sent_at": time.time()}), encoding="utf-8" - ) - except OSError: - pass + _write_marker({"digest": _digest(message), "sent_at": time.time(), "clean_runs": 0}) -def clear() -> None: - """Close the repeat window on a clean run, so the same fault coming back - later is a new incident and alerts again rather than sitting out the day.""" +def _write_marker(state: dict) -> None: with contextlib.suppress(OSError): - _marker_path().unlink(missing_ok=True) + _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: + 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: @@ -234,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/poll.py b/lift_status/poll.py index a900706..9e81ff2 100644 --- a/lift_status/poll.py +++ b/lift_status/poll.py @@ -41,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: @@ -58,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) @@ -185,13 +192,17 @@ def _run(data_dir: Path, client: MessagesClient) -> int: 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 == "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) + 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) @@ -204,10 +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.clear() + 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 8b6a800..edec629 100644 --- a/lift_status/store.py +++ b/lift_status/store.py @@ -121,22 +121,19 @@ def append_raw(data_dir, run_uuid, fetched_at, http_status, body, network_error) }, sort_keys=True, ) - with path.open("ab") as f: + 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. - if f.tell() and _last_byte(path) != b"\n": - f.write(b"\n") + 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()) -def _last_byte(path: Path) -> bytes: - with path.open("rb") as f: - f.seek(-1, os.SEEK_END) - return f.read(1) - - class Store: def __init__(self, data_dir): self.data_dir = Path(data_dir) @@ -158,9 +155,6 @@ 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_raw(self.data_dir, run_uuid, fetched_at, http_status, body, network_error) - def iter_raw_lines(self): """Yield every recorded run attempt, oldest file first, and in `fetched_at_utc` order within a file, so two collectors' logs merged @@ -369,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 index 7af63dc..e9566da 100644 --- a/notes/collector-review.md +++ b/notes/collector-review.md @@ -14,11 +14,13 @@ 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 `write_raw`, so a database left corrupt by a power cut, or + 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. + 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. @@ -30,10 +32,12 @@ Ten findings. Eight were fixed, one was not a real path, and one is left for now 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 run was - suppressed for up to 24h. A clean run now clears it. The cost is that a - fault flapping at poll granularity alerts on each return; each poll already - retries three times, so that is a real outage each time. + 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. @@ -41,15 +45,20 @@ Ten findings. Eight were fixed, one was not a real path, and one is left for now 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 line-order rule was there for clock jumps, and it only ever - covered part of them: a pre-NTP stamp already lands in the wrong day's file, - which is ordered by name. On 2026-09-24 all 2,224 real lines were already in + 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. 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 3823a35..22980c0 100644 --- a/scripts/systemd/lift-status-backup.service +++ b/scripts/systemd/lift-status-backup.service @@ -21,7 +21,7 @@ Environment="GIT_SSH_COMMAND=ssh -i /etc/lift-status-deploy-key -o IdentitiesOnl # 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. +# alert; this is the backstop, and the script alerts on its TERM too. TimeoutStartSec=15min EnvironmentFile=-/etc/lift-status.env diff --git a/tests/test_alert.py b/tests/test_alert.py index daa0b77..54e93b9 100644 --- a/tests/test_alert.py +++ b/tests/test_alert.py @@ -64,11 +64,22 @@ def test_the_retry_after_a_failure_is_delivered_and_then_suppresses(self): self.assertTrue(self._send()) self.assertTrue(alert._suppressed(BANNER)) - def test_a_clean_run_in_between_makes_the_same_fault_a_new_alert(self): + 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() - alert.clear() + self._clean_runs(alert.RECOVERED_AFTER_CLEAN_RUNS - 1) + self.assertTrue(self.marker.exists()) + self._clean_runs(1) self.assertFalse(alert._suppressed(BANNER)) - alert.clear() + + 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() diff --git a/tests/test_poll.py b/tests/test_poll.py index 85b775a..2f3f79d 100644 --- a/tests/test_poll.py +++ b/tests/test_poll.py @@ -206,6 +206,12 @@ def test_a_corrupt_database_still_logs_the_response_and_alerts(self): 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): @@ -215,12 +221,16 @@ def test_a_full_disk_at_the_append_is_a_storage_alert(self): class TestACleanRunClosesTheRepeatWindow(PollTestCase): - def test_a_successful_run_clears_the_marker_and_a_failed_one_does_not(self): + 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()) - poll.run_poll(self.data_dir, client=FakeClient([(200, "[]")])) + 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()) diff --git a/tests/test_rebuild.py b/tests/test_rebuild.py index 1d0847d..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"]) @@ -128,9 +128,8 @@ def test_two_collectors_logs_merged_with_sort_u_rebuild_the_same_history(self): ("r4", "2026-08-08T13:30:00Z", 200, only_a, None), ("r5", "2026-08-08T14:00:00Z", 200, seen, None), ] - with Store(self.data_dir) as store: - for run in history: - store.write_raw(*run) + 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"]] diff --git a/tests/test_store.py b/tests/test_store.py index 6f954cd..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,18 +59,18 @@ 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" - 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) with path.open("a", encoding="utf-8") as f: f.write('{"body": "[{\\"head') - self.store.write_raw("run-3", "2026-08-08T13:00:00Z", 200, "[]", None) + 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)