diff --git a/models.py b/models.py index 9112340..7a09093 100644 --- a/models.py +++ b/models.py @@ -721,22 +721,41 @@ class CassettePayloadRow(BaseModel): class PublishCassettesPayload(BaseModel): - """The decrypted JSON content of a kind-30078 cassette event, both - directions: - - operator → ATM (d-tag `bitspire-cassettes:`) - - ATM → operator (d-tag `bitspire-cassettes-state:`) + """The decrypted content of the ATM → operator state document + (d-tag `bitspire-cassettes-state:`). - Wire shape: `{"positions": {"": {"denomination", "count"}}}`. - JSON object keys are always strings; the validator coerces back to - int on parse. The position key set MUST match what the receiver - already has (slot count is hardware-fixed; no add/remove from this - payload). + It carried the operator → ATM direction too until v2 moved that to + PublishCassetteOpsPayload. This is now the machine reporting what it + holds, and the machine is the only writer of those counts. + + Wire shape: `{"positions": {"": {"denomination", "count"}}}` + plus the optional fields below. JSON object keys are always strings; the + validator coerces back to int on parse. No denomination-unique constraint: multiple same-denom cassettes are operationally valid (cash-out throughput on a popular denom). + + The optional fields are all absent on older machines, so every one of them + defaults to a value meaning "this machine does not report that yet" rather + than to a value that would be wrong: + + - `applied_ops`: operation ids the machine has applied. This is the + acknowledgement, and the only one an addressable event can carry — a + relay returns OK for an event it then discards, so the publisher is + never told anything. An empty list reads as "nothing acknowledged", + which is correct for a machine that has not yet applied any. + - `seq`: the machine's own monotonic counter, bumped on every local count + change. Regression detection independent of created_at, which is only + second-granular and can be forced by a bad clock. + - `counts_uncertain_since`: set when a dispense ended without the + dispenser reporting what it moved, so the counts above are the + machine's best guess rather than a measurement. """ positions: dict[int, CassettePayloadRow] + applied_ops: list[str] = [] + seq: int | None = None + counts_uncertain_since: int | None = None @validator("positions", pre=True) def coerce_string_keys_to_int(cls, v): diff --git a/tasks.py b/tasks.py index 61695e1..0d2f138 100644 --- a/tasks.py +++ b/tasks.py @@ -338,6 +338,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: apply_reported_state, get_machine_by_atm_pubkey_hex, list_all_active_machines, + mark_cassette_ops_acked, ) machines = await list_all_active_machines() @@ -373,6 +374,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: event_message, get_machine_by_atm_pubkey_hex, apply_reported_state, + mark_cassette_ops_acked, ) except Exception as exc: logger.warning( @@ -383,10 +385,36 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: return filter_key +async def _record_op_acknowledgements( + machine_id: str, payload, mark_cassette_ops_acked +) -> None: + """Mark the operations a machine reports as applied. + + Deliberately not gated on whether the state event advanced the counts. The + machine echoes its applied-op ids on EVERY state publish, so an event + carrying nothing new about the counts can still be the first one to tell us + an operation landed; gating on that would lose the acknowledgement. + + This echo is the only acknowledgement the transport can carry. An + addressable event gives its publisher no failure signal at all — the relay + returns OK for an event it then discards — so without it the dashboard + could only ever show an operation as sent, never as delivered. + """ + if not payload.applied_ops: + return + newly_acked = await mark_cassette_ops_acked(machine_id, payload.applied_ops) + if newly_acked: + logger.info( + f"spirekeeper: machine {machine_id} acknowledged " + f"{newly_acked} cassette operation(s)" + ) + + async def _handle_cassette_state_event( event_message, get_machine_by_atm_pubkey_hex, apply_reported_state, + mark_cassette_ops_acked, ) -> None: """Verify signature, resolve the operator's signer, decrypt via the signer abstraction (bunker round-trip for RemoteBunkerSigner; direct @@ -487,12 +515,16 @@ async def _handle_cassette_state_event( ) if applied: logger.info( - f"spirekeeper: applied bootstrap state event {event_id[:12]}... " + f"spirekeeper: applied reported state event {event_id[:12]}... " f"to machine {machine.id} ({len(payload.positions)} cassettes)" ) else: - # Replay: event_id already on file. Normal on relay reconnect. + # Replay or an older event. Normal on relay reconnect. logger.debug( f"spirekeeper: cassette state event {event_id[:12]}... " - f"already applied to machine {machine.id} (replay no-op)" + f"not newer than stored state for machine {machine.id} (no-op)" ) + + # Acknowledgement runs regardless of whether the counts were newer — see + # _record_op_acknowledgements for why. + await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked) diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index 7014a94..2ddeedd 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -15,7 +15,7 @@ from datetime import datetime, timezone import pytest from pydantic import ValidationError -from .. import cassette_transport +from .. import cassette_transport, tasks from ..crud import _should_ack_op from ..models import ( CASSETTE_OP_TYPES, @@ -23,6 +23,7 @@ from ..models import ( CreateCassetteOpData, Machine, PublishCassetteOpsPayload, + PublishCassettesPayload, ) AT = datetime.fromtimestamp(1790106060, timezone.utc) @@ -257,3 +258,62 @@ class TestPublishOpsToAtm: m = machine().copy(update={"machine_npub": npub}) await cassette_transport.publish_ops_to_atm(m, [], "op1") assert captured["d_tag"] == f"bitspire-cassettes:{ATM_HEX}" + + +# ============================================================================= +# _record_op_acknowledgements — the only ack this transport can carry +# ============================================================================= + + +def state_payload(**kw) -> PublishCassettesPayload: + base = {"positions": {"1": {"denomination": 50, "count": 24}}} + return PublishCassettesPayload(**{**base, **kw}) + + +class TestRecordOpAcknowledgements: + @pytest.mark.asyncio + async def test_marks_the_reported_ids(self): + calls = [] + + async def mark(machine_id, op_ids): + calls.append((machine_id, op_ids)) + return len(op_ids) + + await tasks._record_op_acknowledgements( + "m1", state_payload(applied_ops=["a", "b"]), mark + ) + assert calls == [("m1", ["a", "b"])] + + @pytest.mark.asyncio + async def test_does_nothing_when_the_machine_reports_none(self): + """An older machine sends no applied_ops at all, and a new one with + nothing applied sends an empty list. Neither should write.""" + calls = [] + + async def mark(machine_id, op_ids): + calls.append((machine_id, op_ids)) + return 0 + + await tasks._record_op_acknowledgements("m1", state_payload(), mark) + await tasks._record_op_acknowledgements( + "m1", state_payload(applied_ops=[]), mark + ) + assert calls == [] + + @pytest.mark.asyncio + async def test_acks_even_when_the_counts_were_not_newer(self): + """The machine echoes its applied ids on every publish, including + heartbeats that carry nothing new about the counts. One of those can + still be the first event to tell us an operation landed, so the ack + must not depend on the state having advanced.""" + seen = [] + + async def mark(machine_id, op_ids): + seen.extend(op_ids) + return len(op_ids) + + # Same positions as already on file — a pure heartbeat. + await tasks._record_op_acknowledgements( + "m1", state_payload(applied_ops=["late-ack"]), mark + ) + assert seen == ["late-ack"]