From 27449e1d11241c1e0b5242e79cdd72b8deaea2c3 Mon Sep 17 00:00:00 2001 From: Padreug Date: Tue, 22 Sep 2026 22:23:19 +0200 Subject: [PATCH 1/2] fix(cassettes): order state events by created_at, not by one remembered id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The gate on the ATM-state consumer compared the incoming event id against the id stored on a single arbitrary row (SELECT ... LIMIT 1, no ORDER BY). That is a one-event memory, not a watermark: a re-delivered A, B, A applied three times. Worse, created_at was parsed, written to state_at and then never compared, so an event arriving late overwrote newer state — nothing in the path ever looked at the clock. Events are now applied only when strictly newer than the OLDEST state stamp on file. Strict '>' subsumes replay dedup, since a replay carries the same stamp. Oldest rather than newest is deliberate: every execute in this data layer commits on its own, so a multi-row apply cannot be made atomic here, and gating on the oldest means a crash mid-apply is re-applied on the next event instead of being mistaken for a complete one. The ATM republishes on a heartbeat, so it converges. Stamps are compared as unix floats because SQLite returns integers, Postgres returns timestamps and the incoming value is tz-aware; comparing raw would either raise or quietly mislead. An unparseable incoming stamp fails closed. Also renames apply_bootstrap_state to apply_reported_state and corrects the module comments. There has never been a once-per-machine guard, so calling it a one-shot bootstrap consumer described something the code did not do. Refs #43 Co-Authored-By: Claude Fable 5.1 --- crud.py | 108 ++++++++++++++++---------- tasks.py | 32 ++++---- tests/test_cassette_configs.py | 85 +++++++++++++------- tests/test_cassette_state_consumer.py | 2 +- 4 files changed, 143 insertions(+), 84 deletions(-) diff --git a/crud.py b/crud.py index 6eda100..f0f7f37 100644 --- a/crud.py +++ b/crud.py @@ -5,7 +5,7 @@ # through dca_machines.operator_user_id. See plan section "Identity & multi- # machine model". -from datetime import datetime +from datetime import datetime, timezone from lnbits.db import Database from lnbits.helpers import urlsafe_short_hash @@ -1442,7 +1442,7 @@ async def upsert_fleet_snapshot( # Cassette configs — operator-driven ATM cassette inventory (#29 v1.1). # ============================================================================= # Row lifecycle per #29: -# - First population for a (machine_id, position) pair → apply_bootstrap_state +# - First population for a (machine_id, position) pair → apply_reported_state # (consumer reading the ATM's one-shot bitspire-cassettes-state event) # - Operator edit of denomination or count → update_cassette_config # (refuses to create new rows; the slot count is hardware-determined) @@ -1450,25 +1450,56 @@ async def upsert_fleet_snapshot( # re-provisioning + new bootstrap event (not exposed in v1 here) -def _should_apply_bootstrap_state( - existing_state_event_id: str | None, incoming_event_id: str -) -> bool: - """Pure-function dedup gate for apply_bootstrap_state. +def _as_unix(value) -> float | None: + """Normalise whatever the driver hands back for state_at to a unix float. - Returns False if any existing row for this machine already references - the incoming event_id (relay re-delivery after restart). True otherwise. - - Extracted as a pure function so the dedup decision is unit-testable - without a database round-trip. The actual idempotency check in - apply_bootstrap_state fetches one existing row and passes its - state_event_id here. + SQLite stores these as integers and Postgres as timestamps, and an event's + created_at arrives tz-aware, so comparing the raw values risks either a + TypeError (aware vs naive) or a silently wrong answer. Everything is + compared as seconds since the epoch instead. """ - return existing_state_event_id != incoming_event_id + if value is None: + return None + if isinstance(value, datetime): + if value.tzinfo is None: + value = value.replace(tzinfo=timezone.utc) + return value.timestamp() + try: + return float(value) + except (TypeError, ValueError): + return None -async def get_cassette_config( - machine_id: str, position: int -) -> CassetteConfig | None: +def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: + """Ordering gate for apply_reported_state. + + Applies only when the incoming event is strictly newer than the OLDEST + state stamp on file for the machine. + + Two deliberate choices: + + - Compare created_at, not event ids. The old gate asked only whether the + incoming id differed from one stored row, which is a one-event memory: + a re-delivered A, B, A applied three times, and an event that arrived + late overwrote newer state because nothing ever looked at the clock. + Strict `>` also subsumes exact-replay dedup, since a replay carries the + same stamp. + - Oldest, not newest. Every execute in this data layer commits on its own, + so a multi-row apply cannot be made atomic here; a crash mid-loop leaves + some rows advanced and some not. Gating on the oldest means a partial + apply is re-applied on the next event rather than being mistaken for a + complete one, and the ATM republishes on a heartbeat, so it converges. + """ + oldest = _as_unix(oldest_state_at) + if oldest is None: + return True + incoming = _as_unix(incoming_created_at) + if incoming is None: + return False + return incoming > oldest + + +async def get_cassette_config(machine_id: str, position: int) -> CassetteConfig | None: return await db.fetchone( "SELECT * FROM spirekeeper.cassette_configs " "WHERE machine_id = :mid AND position = :pos", @@ -1497,7 +1528,7 @@ async def update_cassette_config( ) -> CassetteConfig | None: """Operator-driven row update: change denomination and/or count for a single cassette slot. Refuses to create new rows — those only land via - apply_bootstrap_state() consuming an ATM bootstrap event (per #29 row + apply_reported_state() consuming an ATM bootstrap event (per #29 row lifecycle: hardware-determined slot count, not operator-creatable). Returns None if the (machine_id, position) row doesn't exist. """ @@ -1520,41 +1551,40 @@ async def update_cassette_config( return await get_cassette_config(machine_id, position) -async def apply_bootstrap_state( +async def apply_reported_state( machine_id: str, event_id: str, event_created_at: datetime, payload: PublishCassettesPayload, ) -> bool: """Consume an ATM-published kind-30078 bitspire-cassettes-state: event - and upsert one cassette_configs row per position in the payload. + and reconcile cassette_configs for the machine against it. - Returns True if the upsert ran; False if any existing row for this - machine already references this event_id (idempotent on relay - re-delivery / restart). + Returns True if the state was applied, False if the event was not newer + than what is already on file (see _should_apply_state_event). Populates both the operator-believed columns (denomination, count, - updated_at, updated_by='atm-bootstrap') AND the v2 reverse-channel - columns (state_denomination, state_count, state_at, state_event_id) - so the operator's initial view matches the ATM's reported state. v2 - reconciliation UI will diverge them when continuous reverse-channel - events land + the operator subsequently edits. + updated_at, updated_by) and the reported columns (state_denomination, + state_count, state_at, state_event_id), so the UI can show reported + against believed. """ - existing_first: dict | None = await db.fetchone( - "SELECT state_event_id FROM spirekeeper.cassette_configs " - "WHERE machine_id = :mid LIMIT 1", + oldest: dict | None = await db.fetchone( + "SELECT state_at FROM spirekeeper.cassette_configs " + "WHERE machine_id = :mid AND state_at IS NOT NULL " + "ORDER BY state_at ASC LIMIT 1", {"mid": machine_id}, ) - existing_event_id: str | None = None - if existing_first is not None: - existing_event_id = ( - existing_first.get("state_event_id") - if isinstance(existing_first, dict) - else getattr(existing_first, "state_event_id", None) + oldest_state_at = None + if oldest is not None: + oldest_state_at = ( + oldest.get("state_at") + if isinstance(oldest, dict) + else getattr(oldest, "state_at", None) ) - if not _should_apply_bootstrap_state(existing_event_id, event_id): + if not _should_apply_state_event(oldest_state_at, event_created_at): return False + now = datetime.now() for pos, row in payload.positions.items(): await db.execute( @@ -1581,7 +1611,7 @@ async def apply_bootstrap_state( "denom": row.denomination, "count": row.count, "now": now, - "by": "atm-bootstrap", + "by": "atm-report", "state_denom": row.denomination, "state_count": row.count, "state_at": event_created_at, diff --git a/tasks.py b/tasks.py index e3865d2..61695e1 100644 --- a/tasks.py +++ b/tasks.py @@ -244,20 +244,20 @@ async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) - # ============================================================================= -# Cassette bootstrap consumer (#29 v1) +# Cassette state consumer (#29) # ============================================================================= # Subscribes to kind-30078 bitspire-cassettes-state: events -# published by each active machine's ATM on first boot (lamassu-next#56's -# bootstrap publish path). Decrypts the NIP-44 v2 content with the operator's -# privkey + ATM sender pubkey, validates as PublishCassettesPayload, and -# upserts cassette_configs via apply_bootstrap_state. +# published by each active machine's ATM. Decrypts the NIP-44 v2 content with +# the operator's privkey + ATM sender pubkey, validates as +# PublishCassettesPayload, and reconciles cassette_configs via +# apply_reported_state. # -# v1 = one-shot per machine (ATM's meta.bootstrapPublishedAt makes the -# publish idempotent on ATM-side restart; spirekeeper's apply_bootstrap_ -# state dedups on state_event_id for relay re-delivery). -# -# v2 (separate issue) = continuous reverse-channel consumer with a -# last_state_created_at watermark for reconciliation UI. +# This is continuous, not one-shot: the ATM publishes on startup, after every +# change to its bays, and on a heartbeat, and each event replaces the last. +# The comments here used to call it a "v1 bootstrap" consumer, which was +# misleading — there has never been a once-per-machine guard, so every event +# an ATM published was already being applied. Ordering is enforced in +# apply_reported_state by created_at; it is not inferred from arrival order. # # Implementation: polls nostrclient.router.NostrRouter.received_subscription_ # events keyed by our subscription_id. nostrclient's NostrRouter design is @@ -285,7 +285,7 @@ async def wait_for_cassette_state_events() -> None: between retries (operator may install it later) - inbound event fails sig-verify / decrypt / parse → log + skip the event, continue the loop - - apply_bootstrap_state errors → log + skip + - apply_reported_state errors → log + skip """ logger.info( "spirekeeper v2: cassette bootstrap consumer starting " @@ -335,7 +335,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: from .cassette_transport import build_state_d_tags_for_machines from .crud import ( - apply_bootstrap_state, + apply_reported_state, get_machine_by_atm_pubkey_hex, list_all_active_machines, ) @@ -372,7 +372,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: await _handle_cassette_state_event( event_message, get_machine_by_atm_pubkey_hex, - apply_bootstrap_state, + apply_reported_state, ) except Exception as exc: logger.warning( @@ -386,7 +386,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: async def _handle_cassette_state_event( event_message, get_machine_by_atm_pubkey_hex, - apply_bootstrap_state, + apply_reported_state, ) -> None: """Verify signature, resolve the operator's signer, decrypt via the signer abstraction (bunker round-trip for RemoteBunkerSigner; direct @@ -482,7 +482,7 @@ async def _handle_cassette_state_event( created_at_unix = event_obj.get("created_at", 0) event_created_at = _datetime.fromtimestamp(int(created_at_unix), tz=_timezone.utc) - applied = await apply_bootstrap_state( + applied = await apply_reported_state( machine.id, event_id, event_created_at, payload ) if applied: diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index 703c3bb..242aa13 100644 --- a/tests/test_cassette_configs.py +++ b/tests/test_cassette_configs.py @@ -5,11 +5,11 @@ Covers the pure pieces that don't need a live DB: - Pydantic validator behaviour on PublishCassettesPayload + the row / upsert models (position key coercion, integer ranges, multiple-same- denomination payloads, wire-format round-trip) - - _should_apply_bootstrap_state dedup helper (extracted from - apply_bootstrap_state so the relay-re-delivery decision is testable - without a database round-trip) + - _should_apply_state_event ordering gate (extracted from + apply_reported_state so the decision is testable without a database + round-trip) -DB-touching tests (apply_bootstrap_state actually upserting, list-by- +DB-touching tests (apply_reported_state actually reconciling, list-by- machine ordering, etc.) follow the project convention from test_deposit_currency.py: "Layer 2 is an endpoint-level behaviour better covered by an integration test against a running LNbits; tracked in #26 @@ -21,9 +21,11 @@ denomination + count are operator-editable per row, multiple same-denom cassettes are valid. """ +from datetime import datetime, timedelta, timezone + import pytest -from ..crud import _should_apply_bootstrap_state +from ..crud import _as_unix, _should_apply_state_event from ..models import ( CassettePayloadRow, PublishCassettesPayload, @@ -186,35 +188,62 @@ class TestUpsertCassetteConfigData: # ============================================================================= -# _should_apply_bootstrap_state — relay re-delivery dedup +# _should_apply_state_event — ordering gate # ============================================================================= +NOW = datetime(2026, 9, 22, 12, 0, 0, tzinfo=timezone.utc) -class TestShouldApplyBootstrapState: - """Pure-function dedup gate extracted from apply_bootstrap_state so the - decision is testable without a DB. Logic: apply if-and-only-if the - existing row's state_event_id differs from the incoming event_id. - In v1.1 the ATM publishes the bootstrap event exactly once per machine, - so this is sufficient for replay protection. v2 will need a - `last_state_created_at` watermark in addition (per bitspire's - `meta.lastKnownConfigCreatedAt` on the ATM side) — flagged in #29's - v2 forward-look section. +class TestShouldApplyStateEvent: + """Ordering gate extracted from apply_reported_state so the decision is + testable without a DB. + + It compares created_at against the OLDEST state stamp on file. The gate it + replaced compared event ids against a single arbitrary row, which is a + one-event memory: a re-delivered A, B, A applied three times, and a late + event overwrote newer state because nothing looked at the clock. """ - def test_applies_when_no_existing_row(self): - assert _should_apply_bootstrap_state(None, "new-event-id") is True + def test_applies_when_machine_has_no_state_yet(self): + assert _should_apply_state_event(None, NOW) is True - def test_applies_when_existing_event_id_differs(self): - assert _should_apply_bootstrap_state("old-event-id", "new-event-id") is True + def test_applies_when_strictly_newer(self): + assert _should_apply_state_event(NOW, NOW + timedelta(seconds=1)) is True - def test_skips_when_existing_event_id_matches(self): - """The same bootstrap event re-delivered after a relay reconnect - or spirekeeper restart should no-op, not re-upsert the same - rows (which would clobber any operator edits since).""" - assert _should_apply_bootstrap_state("same-event", "same-event") is False + def test_skips_an_exact_replay(self): + """A re-delivered event carries the same stamp, so strict `>` covers + replay without needing to remember ids.""" + assert _should_apply_state_event(NOW, NOW) is False - def test_applies_when_existing_is_empty_string_and_incoming_is_id(self): - """Defensive — a sentinel empty-string existing_state_event_id - shouldn't block a real incoming event from applying.""" - assert _should_apply_bootstrap_state("", "real-event-id") is True + def test_skips_an_event_that_arrives_late(self): + """The failure the id-only gate allowed: an older event landing after + a newer one used to overwrite it.""" + assert _should_apply_state_event(NOW, NOW - timedelta(seconds=30)) is False + + def test_reapplies_over_a_partial_write(self): + """Every execute in this data layer commits on its own, so a crash + mid-apply can leave some rows advanced and some not. Gating on the + oldest stamp means the next event re-applies rather than treating the + partial write as complete.""" + partially_applied_oldest = NOW - timedelta(seconds=300) + assert _should_apply_state_event(partially_applied_oldest, NOW) is True + + def test_handles_naive_and_epoch_stamps(self): + """SQLite hands these back as integers and Postgres as timestamps, and + the incoming stamp is tz-aware. Comparing raw would raise or silently + mislead.""" + assert _as_unix(NOW) == NOW.timestamp() + assert _as_unix(NOW.replace(tzinfo=None)) == NOW.timestamp() + assert _as_unix(int(NOW.timestamp())) == float(int(NOW.timestamp())) + assert _as_unix(None) is None + assert _as_unix("not-a-time") is None + assert _should_apply_state_event(int(NOW.timestamp()), NOW) is False + assert ( + _should_apply_state_event(int(NOW.timestamp()), NOW + timedelta(seconds=5)) + is True + ) + + def test_skips_when_incoming_stamp_is_unusable(self): + """Fail closed: an unparseable incoming stamp must not overwrite + state that is known-good.""" + assert _should_apply_state_event(NOW, "nonsense") is False diff --git a/tests/test_cassette_state_consumer.py b/tests/test_cassette_state_consumer.py index a0840bc..d7cdd8e 100644 --- a/tests/test_cassette_state_consumer.py +++ b/tests/test_cassette_state_consumer.py @@ -22,7 +22,7 @@ demonstrates the project lacks an asyncio plugin in CI; using asyncio.run inside the test body sidesteps that without changing project config). Full handler tests (the dispatch through verify_event → -get_machine_by_atm_pubkey_hex → apply_bootstrap_state) need a live LNbits +get_machine_by_atm_pubkey_hex → apply_reported_state) need a live LNbits DB; smoke-tested manually via the dev container per the project convention (see test_deposit_currency.py rationale). """ From 106da5b46bf7f1323d93612a30841623779cad01 Mon Sep 17 00:00:00 2001 From: Padreug Date: Tue, 22 Sep 2026 22:23:40 +0200 Subject: [PATCH 2/2] fix(cassettes): drop bays the machine no longer reports MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit apply_reported_state upserted the positions in the payload and left every other row untouched. When a machine's bay count shrank, the stale row stayed: the dashboard showed a mix of old and new bays, and publish validation then rejected every operator payload for a position-set mismatch. The error text told the operator to fix it with atm-tui, which could not propagate either, so the only way out was DELETE FROM by hand. The payload is the machine's full bay set and the machine owns that layout, so positions absent from it are now deleted. The delete runs before the upserts: since this data layer commits per statement, a crash between the two leaves rows missing rather than stale, and the next heartbeat re-inserts them — the safe direction of the two. Refs #43 Co-Authored-By: Claude Fable 5.1 --- crud.py | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/crud.py b/crud.py index f0f7f37..a622d96 100644 --- a/crud.py +++ b/crud.py @@ -1563,6 +1563,13 @@ async def apply_reported_state( Returns True if the state was applied, False if the event was not newer than what is already on file (see _should_apply_state_event). + The payload is the machine's full bay set, and the machine owns that + layout — the bay count is hardware-determined. So positions absent from + the payload are DELETED here. Leaving them behind was half the problem a + shrinking bay count caused: the stale row stayed, the dashboard showed a + mix of old and new, and every later publish was rejected for a position + mismatch with no way out but hand-written SQL. + Populates both the operator-believed columns (denomination, count, updated_at, updated_by) and the reported columns (state_denomination, state_count, state_at, state_event_id), so the UI can show reported @@ -1584,6 +1591,19 @@ async def apply_reported_state( if not _should_apply_state_event(oldest_state_at, event_created_at): return False + # Drop bays the machine no longer reports, before writing the rest. A crash + # between the two leaves rows missing rather than stale, and the next + # heartbeat re-inserts them — the safe direction. + keep = sorted(payload.positions.keys()) + placeholders = ", ".join(f":p{i}" for i in range(len(keep))) + delete_values: dict = {"mid": machine_id} + for i, pos in enumerate(keep): + delete_values[f"p{i}"] = pos + await db.execute( + "DELETE FROM spirekeeper.cassette_configs " + f"WHERE machine_id = :mid AND position NOT IN ({placeholders})", + delete_values, + ) now = datetime.now() for pos, row in payload.positions.items():