Publish cassette operations instead of counts #46

Merged
padreug merged 8 commits from feat/cassette-ops-publisher into main 2026-09-23 21:16:34 +00:00
3 changed files with 124 additions and 13 deletions
Showing only changes of commit 3d8368bcc4 - Show all commits

feat(cassettes): consume the machine's operation acknowledgements

The machine echoes the operation ids it has applied in its state
document, and this records them. That echo is the only acknowledgement
this transport can carry: an addressable event gives its publisher no
failure signal at all, since the relay returns OK for an event it then
discards. Without it the dashboard could only ever show an operation as
sent, never as delivered.

Deliberately not gated on whether the state event advanced the counts.
The machine echoes its applied ids on every publish, heartbeats included,
so an event carrying nothing new about the counts can still be the first
one to tell us an operation landed.

The consumer goes in before the producer on purpose. The machine does not
send applied_ops yet, and every new field on the state payload defaults to
a value meaning "this machine does not report that yet" rather than to
one that would be wrong — an absent list reads as nothing acknowledged,
which is exactly right for a machine that has applied nothing.

Also picks up seq and counts_uncertain_since, which the machine already
publishes and this side was dropping on the floor.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Padreug 2026-09-23 12:38:12 +02:00

View file

@ -721,22 +721,41 @@ class CassettePayloadRow(BaseModel):
class PublishCassettesPayload(BaseModel): class PublishCassettesPayload(BaseModel):
"""The decrypted JSON content of a kind-30078 cassette event, both """The decrypted content of the ATM → operator state document
directions: (d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`).
- operator → ATM (d-tag `bitspire-cassettes:<atm_pubkey_hex>`)
- ATM → operator (d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`)
Wire shape: `{"positions": {"<pos_str>": {"denomination", "count"}}}`. It carried the operator → ATM direction too until v2 moved that to
JSON object keys are always strings; the validator coerces back to PublishCassetteOpsPayload. This is now the machine reporting what it
int on parse. The position key set MUST match what the receiver holds, and the machine is the only writer of those counts.
already has (slot count is hardware-fixed; no add/remove from this
payload). Wire shape: `{"positions": {"<pos_str>": {"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 No denomination-unique constraint: multiple same-denom cassettes are
operationally valid (cash-out throughput on a popular denom). 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] positions: dict[int, CassettePayloadRow]
applied_ops: list[str] = []
seq: int | None = None
counts_uncertain_since: int | None = None
@validator("positions", pre=True) @validator("positions", pre=True)
def coerce_string_keys_to_int(cls, v): def coerce_string_keys_to_int(cls, v):

View file

@ -338,6 +338,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
apply_reported_state, apply_reported_state,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
list_all_active_machines, list_all_active_machines,
mark_cassette_ops_acked,
) )
machines = await list_all_active_machines() machines = await list_all_active_machines()
@ -373,6 +374,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_reported_state, apply_reported_state,
mark_cassette_ops_acked,
) )
except Exception as exc: except Exception as exc:
logger.warning( logger.warning(
@ -383,10 +385,36 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
return filter_key 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( async def _handle_cassette_state_event(
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_reported_state, apply_reported_state,
mark_cassette_ops_acked,
) -> 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
@ -487,12 +515,16 @@ async def _handle_cassette_state_event(
) )
if applied: if applied:
logger.info( 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)" f"to machine {machine.id} ({len(payload.positions)} cassettes)"
) )
else: else:
# Replay: event_id already on file. Normal on relay reconnect. # Replay or an older event. Normal on relay reconnect.
logger.debug( logger.debug(
f"spirekeeper: cassette state event {event_id[:12]}... " 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)

View file

@ -15,7 +15,7 @@ from datetime import datetime, timezone
import pytest import pytest
from pydantic import ValidationError from pydantic import ValidationError
from .. import cassette_transport from .. import cassette_transport, tasks
from ..crud import _should_ack_op from ..crud import _should_ack_op
from ..models import ( from ..models import (
CASSETTE_OP_TYPES, CASSETTE_OP_TYPES,
@ -23,6 +23,7 @@ from ..models import (
CreateCassetteOpData, CreateCassetteOpData,
Machine, Machine,
PublishCassetteOpsPayload, PublishCassetteOpsPayload,
PublishCassettesPayload,
) )
AT = datetime.fromtimestamp(1790106060, timezone.utc) AT = datetime.fromtimestamp(1790106060, timezone.utc)
@ -257,3 +258,62 @@ class TestPublishOpsToAtm:
m = machine().copy(update={"machine_npub": npub}) m = machine().copy(update={"machine_npub": npub})
await cassette_transport.publish_ops_to_atm(m, [], "op1") await cassette_transport.publish_ops_to_atm(m, [], "op1")
assert captured["d_tag"] == f"bitspire-cassettes:{ATM_HEX}" 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"]