Compare commits

...

3 commits

Author SHA1 Message Date
30870ebd16 Merge pull request 'fix(cassettes): order ATM state events, and reconcile the bay set' (#44) from fix/cassette-state-reconcile into main
Some checks failed
ci.yml / Merge pull request 'fix(cassettes): order ATM state events, and reconcile the bay set' (#44) from fix/cassette-state-reconcile into main (push) Failing after 0s
/ release (push) Has been cancelled
/ pullrequest (push) Has been cancelled
Reviewed-on: #44
2026-09-22 20:30:44 +00:00
106da5b46b fix(cassettes): drop bays the machine no longer reports
Some checks failed
ci.yml / fix(cassettes): drop bays the machine no longer reports (pull_request) Failing after 0s
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 <noreply@anthropic.com>
2026-09-22 22:23:40 +02:00
27449e1d11 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>
2026-09-22 22:23:19 +02:00
4 changed files with 164 additions and 85 deletions

128
crud.py
View file

@ -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:<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
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,

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
# 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:

View file

@ -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

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).
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).
"""