fix(cassettes): order ATM state events, and reconcile the bay set #44

Merged
padreug merged 2 commits from fix/cassette-state-reconcile into main 2026-09-22 20:30:45 +00:00
4 changed files with 144 additions and 85 deletions
Showing only changes of commit 27449e1d11 - Show all commits

fix(cassettes): order state events by created_at, not by one remembered id

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 <noreply@anthropic.com>
Padreug 2026-09-22 22:23:19 +02:00

108
crud.py
View file

@ -5,7 +5,7 @@
# through dca_machines.operator_user_id. See plan section "Identity & multi- # through dca_machines.operator_user_id. See plan section "Identity & multi-
# machine model". # machine model".
from datetime import datetime from datetime import datetime, timezone
from lnbits.db import Database from lnbits.db import Database
from lnbits.helpers import urlsafe_short_hash 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). # Cassette configs — operator-driven ATM cassette inventory (#29 v1.1).
# ============================================================================= # =============================================================================
# Row lifecycle per #29: # 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) # (consumer reading the ATM's one-shot bitspire-cassettes-state event)
# - Operator edit of denomination or count → update_cassette_config # - Operator edit of denomination or count → update_cassette_config
# (refuses to create new rows; the slot count is hardware-determined) # (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) # re-provisioning + new bootstrap event (not exposed in v1 here)
def _should_apply_bootstrap_state( def _as_unix(value) -> float | None:
existing_state_event_id: str | None, incoming_event_id: str """Normalise whatever the driver hands back for state_at to a unix float.
) -> bool:
"""Pure-function dedup gate for apply_bootstrap_state.
Returns False if any existing row for this machine already references SQLite stores these as integers and Postgres as timestamps, and an event's
the incoming event_id (relay re-delivery after restart). True otherwise. created_at arrives tz-aware, so comparing the raw values risks either a
TypeError (aware vs naive) or a silently wrong answer. Everything is
Extracted as a pure function so the dedup decision is unit-testable compared as seconds since the epoch instead.
without a database round-trip. The actual idempotency check in
apply_bootstrap_state fetches one existing row and passes its
state_event_id here.
""" """
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( def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool:
machine_id: str, position: int """Ordering gate for apply_reported_state.
) -> CassetteConfig | None:
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( return await db.fetchone(
"SELECT * FROM spirekeeper.cassette_configs " "SELECT * FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid AND position = :pos", "WHERE machine_id = :mid AND position = :pos",
@ -1497,7 +1528,7 @@ async def update_cassette_config(
) -> CassetteConfig | None: ) -> CassetteConfig | None:
"""Operator-driven row update: change denomination and/or count for a """Operator-driven row update: change denomination and/or count for a
single cassette slot. Refuses to create new rows — those only land via 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). lifecycle: hardware-determined slot count, not operator-creatable).
Returns None if the (machine_id, position) row doesn't exist. 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) return await get_cassette_config(machine_id, position)
async def apply_bootstrap_state( async def apply_reported_state(
machine_id: str, machine_id: str,
event_id: str, event_id: str,
event_created_at: datetime, event_created_at: datetime,
payload: PublishCassettesPayload, payload: PublishCassettesPayload,
) -> bool: ) -> bool:
"""Consume an ATM-published kind-30078 bitspire-cassettes-state:<m> event """Consume an ATM-published kind-30078 bitspire-cassettes-state:<m> 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 Returns True if the state was applied, False if the event was not newer
machine already references this event_id (idempotent on relay than what is already on file (see _should_apply_state_event).
re-delivery / restart).
Populates both the operator-believed columns (denomination, count, Populates both the operator-believed columns (denomination, count,
updated_at, updated_by='atm-bootstrap') AND the v2 reverse-channel updated_at, updated_by) and the reported columns (state_denomination,
columns (state_denomination, state_count, state_at, state_event_id) state_count, state_at, state_event_id), so the UI can show reported
so the operator's initial view matches the ATM's reported state. v2 against believed.
reconciliation UI will diverge them when continuous reverse-channel
events land + the operator subsequently edits.
""" """
existing_first: dict | None = await db.fetchone( oldest: dict | None = await db.fetchone(
"SELECT state_event_id FROM spirekeeper.cassette_configs " "SELECT state_at FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid LIMIT 1", "WHERE machine_id = :mid AND state_at IS NOT NULL "
"ORDER BY state_at ASC LIMIT 1",
{"mid": machine_id}, {"mid": machine_id},
) )
existing_event_id: str | None = None oldest_state_at = None
if existing_first is not None: if oldest is not None:
existing_event_id = ( oldest_state_at = (
existing_first.get("state_event_id") oldest.get("state_at")
if isinstance(existing_first, dict) if isinstance(oldest, dict)
else getattr(existing_first, "state_event_id", None) 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 return False
now = datetime.now() now = datetime.now()
for pos, row in payload.positions.items(): for pos, row in payload.positions.items():
await db.execute( await db.execute(
@ -1581,7 +1611,7 @@ async def apply_bootstrap_state(
"denom": row.denomination, "denom": row.denomination,
"count": row.count, "count": row.count,
"now": now, "now": now,
"by": "atm-bootstrap", "by": "atm-report",
"state_denom": row.denomination, "state_denom": row.denomination,
"state_count": row.count, "state_count": row.count,
"state_at": event_created_at, "state_at": event_created_at,

View file

@ -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:<atm_pubkey_hex> events # Subscribes to kind-30078 bitspire-cassettes-state:<atm_pubkey_hex> events
# published by each active machine's ATM on first boot (lamassu-next#56's # published by each active machine's ATM. Decrypts the NIP-44 v2 content with
# bootstrap publish path). Decrypts the NIP-44 v2 content with the operator's # the operator's privkey + ATM sender pubkey, validates as
# privkey + ATM sender pubkey, validates as PublishCassettesPayload, and # PublishCassettesPayload, and reconciles cassette_configs via
# upserts cassette_configs via apply_bootstrap_state. # apply_reported_state.
# #
# v1 = one-shot per machine (ATM's meta.bootstrapPublishedAt makes the # This is continuous, not one-shot: the ATM publishes on startup, after every
# publish idempotent on ATM-side restart; spirekeeper's apply_bootstrap_ # change to its bays, and on a heartbeat, and each event replaces the last.
# state dedups on state_event_id for relay re-delivery). # 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
# v2 (separate issue) = continuous reverse-channel consumer with a # an ATM published was already being applied. Ordering is enforced in
# last_state_created_at watermark for reconciliation UI. # apply_reported_state by created_at; it is not inferred from arrival order.
# #
# Implementation: polls nostrclient.router.NostrRouter.received_subscription_ # Implementation: polls nostrclient.router.NostrRouter.received_subscription_
# events keyed by our subscription_id. nostrclient's NostrRouter design is # 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) between retries (operator may install it later)
- inbound event fails sig-verify / decrypt / parse → log + skip - inbound event fails sig-verify / decrypt / parse → log + skip
the event, continue the loop the event, continue the loop
- apply_bootstrap_state errors → log + skip - apply_reported_state errors → log + skip
""" """
logger.info( logger.info(
"spirekeeper v2: cassette bootstrap consumer starting " "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 .cassette_transport import build_state_d_tags_for_machines
from .crud import ( from .crud import (
apply_bootstrap_state, apply_reported_state,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
list_all_active_machines, list_all_active_machines,
) )
@ -372,7 +372,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
await _handle_cassette_state_event( await _handle_cassette_state_event(
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_bootstrap_state, apply_reported_state,
) )
except Exception as exc: except Exception as exc:
logger.warning( logger.warning(
@ -386,7 +386,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
async def _handle_cassette_state_event( async def _handle_cassette_state_event(
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_bootstrap_state, apply_reported_state,
) -> None: ) -> None:
"""Verify signature, resolve the operator's signer, decrypt via the """Verify signature, resolve the operator's signer, decrypt via the
signer abstraction (bunker round-trip for RemoteBunkerSigner; direct 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) created_at_unix = event_obj.get("created_at", 0)
event_created_at = _datetime.fromtimestamp(int(created_at_unix), tz=_timezone.utc) 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 machine.id, event_id, event_created_at, payload
) )
if applied: if applied:

View file

@ -5,11 +5,11 @@ Covers the pure pieces that don't need a live DB:
- Pydantic validator behaviour on PublishCassettesPayload + the row / - Pydantic validator behaviour on PublishCassettesPayload + the row /
upsert models (position key coercion, integer ranges, multiple-same- upsert models (position key coercion, integer ranges, multiple-same-
denomination payloads, wire-format round-trip) denomination payloads, wire-format round-trip)
- _should_apply_bootstrap_state dedup helper (extracted from - _should_apply_state_event ordering gate (extracted from
apply_bootstrap_state so the relay-re-delivery decision is testable apply_reported_state so the decision is testable without a database
without a database round-trip) 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 machine ordering, etc.) follow the project convention from
test_deposit_currency.py: "Layer 2 is an endpoint-level behaviour better test_deposit_currency.py: "Layer 2 is an endpoint-level behaviour better
covered by an integration test against a running LNbits; tracked in #26 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. cassettes are valid.
""" """
from datetime import datetime, timedelta, timezone
import pytest import pytest
from ..crud import _should_apply_bootstrap_state from ..crud import _as_unix, _should_apply_state_event
from ..models import ( from ..models import (
CassettePayloadRow, CassettePayloadRow,
PublishCassettesPayload, 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, class TestShouldApplyStateEvent:
so this is sufficient for replay protection. v2 will need a """Ordering gate extracted from apply_reported_state so the decision is
`last_state_created_at` watermark in addition (per bitspire's testable without a DB.
`meta.lastKnownConfigCreatedAt` on the ATM side) — flagged in #29's
v2 forward-look section. 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): def test_applies_when_machine_has_no_state_yet(self):
assert _should_apply_bootstrap_state(None, "new-event-id") is True assert _should_apply_state_event(None, NOW) is True
def test_applies_when_existing_event_id_differs(self): def test_applies_when_strictly_newer(self):
assert _should_apply_bootstrap_state("old-event-id", "new-event-id") is True assert _should_apply_state_event(NOW, NOW + timedelta(seconds=1)) is True
def test_skips_when_existing_event_id_matches(self): def test_skips_an_exact_replay(self):
"""The same bootstrap event re-delivered after a relay reconnect """A re-delivered event carries the same stamp, so strict `>` covers
or spirekeeper restart should no-op, not re-upsert the same replay without needing to remember ids."""
rows (which would clobber any operator edits since).""" assert _should_apply_state_event(NOW, NOW) is False
assert _should_apply_bootstrap_state("same-event", "same-event") is False
def test_applies_when_existing_is_empty_string_and_incoming_is_id(self): def test_skips_an_event_that_arrives_late(self):
"""Defensive — a sentinel empty-string existing_state_event_id """The failure the id-only gate allowed: an older event landing after
shouldn't block a real incoming event from applying.""" a newer one used to overwrite it."""
assert _should_apply_bootstrap_state("", "real-event-id") is True 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

View file

@ -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). inside the test body sidesteps that without changing project config).
Full handler tests (the dispatch through verify_event → 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 DB; smoke-tested manually via the dev container per the project
convention (see test_deposit_currency.py rationale). convention (see test_deposit_currency.py rationale).
""" """