diff --git a/crud.py b/crud.py index 6eda100..a622d96 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,60 @@ 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). + + 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='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 + # 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(): await db.execute( @@ -1581,7 +1631,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). """