From d8190375a6650eb3852899f144dcae6c6442c4ba Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 09:50:28 +0200 Subject: [PATCH 1/8] feat(cassettes): schema and models for operator operations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit First piece of the v2 wire (bitspire ADR-004). The operator stops publishing counts and starts publishing what it DID; the machine, which holds the notes, keeps the running total. A value with one writer cannot be clobbered, which is the whole point: the absolute-count wire let a form loaded before a dispense discard that dispense when published, and nothing in an addressable event can tell the loser it lost. m013 adds cassette_ops, append-only. The id is minted here and is the idempotency key the machine dedups on, because a delta applied twice is wrong and addressable events are re-delivered on reconnect. acked_at is set when the machine reports that id back, which is the only acknowledgement this transport can carry. The models enforce that an op carries exactly the one field its type means, so an instance is publishable by construction — the same contract FeeConfigPayload has — and nulls never reach the wire for the machine to disambiguate. recount is the only absolute, deliberately: it is what an operator opening a bay and counting actually does, and it stays auditable as its own act rather than looking like a stale form. Vocabulary mirrors lamassu-server's cash_unit_operation_type. Co-Authored-By: Claude Fable 5.1 --- migrations.py | 41 +++++++++ models.py | 166 ++++++++++++++++++++++++++++++++++++- tests/test_cassette_ops.py | 142 +++++++++++++++++++++++++++++++ 3 files changed, 348 insertions(+), 1 deletion(-) create mode 100644 tests/test_cassette_ops.py diff --git a/migrations.py b/migrations.py index 6ede0eb..52aebf3 100644 --- a/migrations.py +++ b/migrations.py @@ -860,3 +860,44 @@ async def m012_add_max_cash_in_sats(db): await db.execute( "ALTER TABLE spirekeeper.super_config ADD COLUMN max_cash_in_sats INTEGER" ) + + +async def m013_add_cassette_ops(db): + """Cassette operations — the operator→ATM v2 wire (aiolabs/bitspire ADR-004). + + Until now the operator published absolute counts and the ATM applied them + outright. Both sides wrote the same value over a transport that never tells + a writer it lost, so a dashboard form loaded before a dispense would + silently discard that dispense when published. The fix is to stop the + operator writing counts at all: it publishes OPERATIONS and the machine, + which holds the physical notes, owns the running count. + + Each row here is one operator intent — a refill, an empty, a recount, a + denomination change. `id` is minted here and is the idempotency key the ATM + dedups on, because a delta applied twice is wrong and addressable events + are re-delivered on reconnect. `acked_at` is set when the machine reports + the id back in its state document, which is the only acknowledgement this + transport can carry. + + Kept append-only on purpose: the published window is a slice of this table, + and an operation the machine has not yet acknowledged must stay publishable. + """ + await db.execute(f""" + CREATE TABLE IF NOT EXISTS spirekeeper.cassette_ops ( + id TEXT PRIMARY KEY, + machine_id TEXT NOT NULL, + position INTEGER NOT NULL, + op_type TEXT NOT NULL, + bills INTEGER, + count INTEGER, + denomination INTEGER, + created_at TIMESTAMP NOT NULL DEFAULT {db.timestamp_now}, + created_by TEXT, + acked_at TIMESTAMP + ); + """) + # The publisher reads the most recent N for a machine on every publish. + await db.execute( + "CREATE INDEX IF NOT EXISTS cassette_ops_machine_idx " + "ON cassette_ops (machine_id, created_at DESC)" + ) diff --git a/models.py b/models.py index 5a4e529..9112340 100644 --- a/models.py +++ b/models.py @@ -7,7 +7,7 @@ from datetime import datetime -from pydantic import BaseModel, validator +from pydantic import BaseModel, root_validator, validator # ============================================================================= # Machines — one row per bitSpire ATM, owned by exactly one operator. @@ -768,6 +768,170 @@ class PublishCassettesPayload(BaseModel): } +# ============================================================================= +# Cassette operations — operator → ATM v2 (bitspire ADR-004) +# ============================================================================= +# The operator no longer publishes counts. It publishes what it DID, and the +# machine — which holds the notes — keeps the running total. A value with one +# writer cannot be clobbered, which is the point: the old absolute-count wire +# let a form loaded before a dispense discard that dispense when published, and +# nothing in an addressable event can tell the loser it lost. +# +# Wire shape (kind-30078 content, NIP-44 v2 encrypted, schema_version 2): +# { +# "schema_version": 2, +# "ops": [ +# {"id": "", "at": 1790106060, "type": "refill", +# "position": 2, "bills": 100}, +# {"id": "", "at": 1790106061, "type": "empty", "position": 3}, +# {"id": "", "at": 1790106062, "type": "recount", +# "position": 1, "count": 37}, +# {"id": "", "at": 1790106063, "type": "set_denomination", +# "position": 1, "denomination": 50} +# ] +# } +# +# `ops` is a WINDOW of recent operations, not just the newest. An event the +# machine missed is carried again by the next one, so the channel heals itself +# without the operator noticing. `id` is the idempotency key: deltas are not +# idempotent and addressable events are re-delivered on reconnect, so the +# machine records what it applied and ignores repeats. +# +# The vocabulary mirrors lamassu-server's cash_unit_operation_type +# (refill / empty / count-change), which is where the ancestor of this fleet +# landed after the same problem. + + +CASSETTE_OP_TYPES = ("refill", "empty", "recount", "set_denomination") + + +class CassetteOp(BaseModel): + """One operator intent against one bay, as stored and as published. + + Exactly one of bills/count/denomination is meaningful, decided by op_type: + - refill → bills, the number of notes ADDED (a delta) + - empty → none; the bay was emptied + - recount → count, an absolute the operator physically counted + - set_denomination → denomination, what is now loaded in that bay + + `recount` is the only absolute, and deliberately so: it is what an operator + opening a bay and counting actually does, and it is auditable as a distinct + act rather than being indistinguishable from a stale form. + """ + + id: str + machine_id: str + position: int + op_type: str + bills: int | None = None + count: int | None = None + denomination: int | None = None + created_at: datetime + created_by: str | None = None + acked_at: datetime | None = None + + @validator("op_type") + def _known_op_type(cls, v): + if v not in CASSETTE_OP_TYPES: + raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}") + return v + + @validator("position") + def _position_positive(cls, v): + if v <= 0: + raise ValueError(f"position must be > 0, got {v}") + return v + + def to_wire_dict(self) -> dict: + """The published form. Drops the fields this op_type does not use, so + the machine never has to guess which of three nullable columns applies.""" + out: dict = { + "id": self.id, + "at": int(self.created_at.timestamp()), + "type": self.op_type, + "position": self.position, + } + if self.op_type == "refill": + out["bills"] = self.bills + elif self.op_type == "recount": + out["count"] = self.count + elif self.op_type == "set_denomination": + out["denomination"] = self.denomination + return out + + +class CreateCassetteOpData(BaseModel): + """Operator submits one operation from the dashboard. + + Validated per type here rather than at the endpoint so an instance is + always publishable, matching FeeConfigPayload's contract. + """ + + position: int + op_type: str + bills: int | None = None + count: int | None = None + denomination: int | None = None + + @validator("op_type") + def _known_op_type(cls, v): + if v not in CASSETTE_OP_TYPES: + raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}") + return v + + @validator("position") + def _position_positive(cls, v): + if v <= 0: + raise ValueError(f"position must be > 0, got {v}") + return v + + @validator("bills") + def _bills_positive(cls, v): + if v is not None and v <= 0: + raise ValueError("bills must be > 0 (a refill adds notes)") + return v + + @validator("count") + def _count_non_negative(cls, v): + if v is not None and v < 0: + raise ValueError("count must be >= 0") + return v + + @validator("denomination") + def _denomination_positive(cls, v): + if v is not None and v <= 0: + raise ValueError("denomination must be > 0") + return v + + @root_validator(skip_on_failure=True) + def _field_matches_type(cls, values): + required = { + "refill": "bills", + "recount": "count", + "set_denomination": "denomination", + "empty": None, + }[values.get("op_type")] + if required is not None and values.get(required) is None: + raise ValueError(f"{values['op_type']} requires `{required}`") + for field in ("bills", "count", "denomination"): + if field != required and values.get(field) is not None: + raise ValueError(f"{values['op_type']} must not carry `{field}`") + return values + + +class PublishCassetteOpsPayload(BaseModel): + """The decrypted content of a v2 operator → ATM cassette event.""" + + schema_version: int = 2 + ops: list[CassetteOp] + + def to_wire_dict(self) -> dict: + return { + "schema_version": self.schema_version, + "ops": [op.to_wire_dict() for op in self.ops], + } + + # ============================================================================= # Fee-config Nostr payload — operator → ATM (aiolabs/satmachineadmin#39) # ============================================================================= diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py new file mode 100644 index 0000000..8187fa2 --- /dev/null +++ b/tests/test_cassette_ops.py @@ -0,0 +1,142 @@ +""" +Tests for the v2 cassette-operations models (bitspire ADR-004). + +The operator no longer publishes counts; it publishes operations and the +machine keeps the running total. These cover the pure pieces: per-type field +validation, and the wire shape the publisher ships. + +A CreateCassetteOpData instance is meant to be publishable by construction — +same contract as FeeConfigPayload — so the type/field agreement is enforced in +the model rather than at the endpoint. +""" + +from datetime import datetime, timezone + +import pytest +from pydantic import ValidationError + +from ..models import ( + CASSETTE_OP_TYPES, + CassetteOp, + CreateCassetteOpData, + PublishCassetteOpsPayload, +) + +AT = datetime.fromtimestamp(1790106060, timezone.utc) + + +def op(**kw) -> CassetteOp: + base = {"id": "op-1", "machine_id": "m1", "position": 2, "created_at": AT} + return CassetteOp(**{**base, **kw}) + + +class TestCreateCassetteOpData: + def test_accepts_one_of_each_type(self): + CreateCassetteOpData(position=2, op_type="refill", bills=100) + CreateCassetteOpData(position=3, op_type="empty") + CreateCassetteOpData(position=1, op_type="recount", count=37) + CreateCassetteOpData(position=1, op_type="set_denomination", denomination=50) + + @pytest.mark.parametrize( + "kwargs", + [ + {"position": 2, "op_type": "refill"}, + {"position": 1, "op_type": "recount"}, + {"position": 1, "op_type": "set_denomination"}, + ], + ) + def test_rejects_a_type_missing_its_field(self, kwargs): + with pytest.raises(ValidationError): + CreateCassetteOpData(**kwargs) + + @pytest.mark.parametrize( + "kwargs", + [ + {"position": 2, "op_type": "refill", "bills": 1, "count": 5}, + {"position": 3, "op_type": "empty", "bills": 1}, + {"position": 1, "op_type": "recount", "count": 1, "denomination": 50}, + ], + ) + def test_rejects_a_type_carrying_a_foreign_field(self, kwargs): + """An op that carries two meanings is ambiguous on the wire, and the + machine would have to guess which one to apply.""" + with pytest.raises(ValidationError): + CreateCassetteOpData(**kwargs) + + def test_rejects_a_refill_of_zero_or_fewer_notes(self): + """A refill is a delta that adds notes. Zero is a no-op an operator + did not mean, and negative is a withdrawal wearing a refill's name.""" + for bills in (0, -5): + with pytest.raises(ValidationError): + CreateCassetteOpData(position=2, op_type="refill", bills=bills) + + def test_allows_a_recount_to_zero(self): + """Distinct from refill: counting a bay and finding it empty is a real + and important observation.""" + assert CreateCassetteOpData(position=2, op_type="recount", count=0).count == 0 + + def test_rejects_a_negative_recount_and_a_non_positive_denomination(self): + with pytest.raises(ValidationError): + CreateCassetteOpData(position=2, op_type="recount", count=-1) + with pytest.raises(ValidationError): + CreateCassetteOpData(position=2, op_type="set_denomination", denomination=0) + + def test_rejects_an_unknown_type_and_a_non_positive_position(self): + with pytest.raises(ValidationError): + CreateCassetteOpData(position=1, op_type="drain") + with pytest.raises(ValidationError): + CreateCassetteOpData(position=0, op_type="empty") + + +class TestWireShape: + def test_each_type_ships_only_its_own_field(self): + assert op(op_type="refill", bills=100).to_wire_dict() == { + "id": "op-1", + "at": 1790106060, + "type": "refill", + "position": 2, + "bills": 100, + } + assert op(op_type="empty").to_wire_dict() == { + "id": "op-1", + "at": 1790106060, + "type": "empty", + "position": 2, + } + assert op(op_type="recount", count=37).to_wire_dict()["count"] == 37 + assert ( + op(op_type="set_denomination", denomination=50).to_wire_dict()[ + "denomination" + ] + == 50 + ) + + def test_nulls_never_reach_the_wire(self): + """The row has three nullable columns and one op only ever means one + of them. Shipping the other two as null would make the machine guess.""" + for op_type in CASSETTE_OP_TYPES: + kw = { + "refill": {"bills": 1}, + "recount": {"count": 1}, + "set_denomination": {"denomination": 1}, + "empty": {}, + }[op_type] + wire = op(op_type=op_type, **kw).to_wire_dict() + assert None not in wire.values() + + def test_payload_declares_v2_and_preserves_order(self): + ops = [ + op(id="a", op_type="refill", bills=1), + op(id="b", op_type="empty"), + ] + wire = PublishCassetteOpsPayload(ops=ops).to_wire_dict() + assert wire["schema_version"] == 2 + assert [o["id"] for o in wire["ops"]] == ["a", "b"] + + def test_an_empty_window_is_representable(self): + """A machine with no operator history still gets a well-formed + payload rather than the publisher having to special-case it.""" + assert PublishCassetteOpsPayload(ops=[]).to_wire_dict() == { + "schema_version": 2, + "ops": [], + } From b2157db219e84b9831393ac57e3e496c983bca5c Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:33:46 +0200 Subject: [PATCH 2/8] feat(cassettes): record and read operator operations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Append-only, and deliberately nothing here writes cassette_configs. That table now holds only what the machine has reported; letting an operation write it would put back the second writer this whole design exists to remove. get_cassette_ops_window returns oldest-first because order is meaning: a recount followed by a refill is not the same as the reverse. It takes the most recent N and reverses, so the window slides without the machine ever seeing them out of sequence. The window is what makes the channel self-healing, so it has to cover a plausible outage rather than just the newest change — a machine that missed one event still sees the operation in the next. _should_ack_op is extracted pure, the same way the state-event gate is, because three of its rules are easy to get wrong and none need a database to test: an unknown id closes out nothing, one machine must never be able to ack another machine's operation, and the FIRST acknowledgement is the one worth keeping. That last one matters because the machine echoes a window, so every id comes back many times over; overwriting would keep sliding the timestamp forward and lose when the operation actually landed. Co-Authored-By: Claude Fable 5.1 --- crud.py | 134 +++++++++++++++++++++++++++++++++++++ tests/test_cassette_ops.py | 26 +++++++ 2 files changed, 160 insertions(+) diff --git a/crud.py b/crud.py index a622d96..bc5cb0a 100644 --- a/crud.py +++ b/crud.py @@ -12,9 +12,11 @@ from lnbits.helpers import urlsafe_short_hash from .models import ( CassetteConfig, + CassetteOp, ClientBalanceSummary, CommissionSplit, CommissionSplitLeg, + CreateCassetteOpData, CreateDcaClientData, CreateDcaPaymentData, CreateDcaSettlementData, @@ -1639,3 +1641,135 @@ async def apply_reported_state( }, ) return True + + +# --------------------------------------------------------------------------- +# Cassette operations — the v2 operator → ATM wire (bitspire ADR-004) +# --------------------------------------------------------------------------- +# Append-only. The operator records what it DID to a bay; the machine keeps the +# running count. Nothing here ever updates cassette_configs: that table now +# holds only what the machine has reported, and letting an operation write it +# would reintroduce the second writer this design exists to remove. + +# How many operations ride along in each published event. The window is what +# makes the channel self-healing — a machine that missed one event still sees +# the operation in the next — so it needs to cover a plausible outage, not just +# the newest change. Twenty is roughly a fortnight of refills on a busy fleet +# and still a small payload. +CASSETTE_OPS_WINDOW = 20 + + +async def create_cassette_op( + machine_id: str, data: CreateCassetteOpData, created_by: str | None +) -> CassetteOp: + """Record one operator intent against one bay. + + The id minted here is the idempotency key: it travels on the wire, the + machine records the ones it has applied, and a re-delivered event is + therefore free rather than double-counted. + """ + op_id = urlsafe_short_hash() + await db.execute( + """ + INSERT INTO spirekeeper.cassette_ops + (id, machine_id, position, op_type, bills, count, denomination, + created_at, created_by) + VALUES (:id, :machine_id, :position, :op_type, :bills, :count, + :denomination, :created_at, :created_by) + """, + { + "id": op_id, + "machine_id": machine_id, + "position": data.position, + "op_type": data.op_type, + "bills": data.bills, + "count": data.count, + "denomination": data.denomination, + "created_at": datetime.now(), + "created_by": created_by, + }, + ) + op = await get_cassette_op(op_id) + assert op is not None, "Newly recorded cassette op couldn't be retrieved" + return op + + +async def get_cassette_op(op_id: str) -> CassetteOp | None: + return await db.fetchone( + "SELECT * FROM spirekeeper.cassette_ops WHERE id = :id", + {"id": op_id}, + CassetteOp, + ) + + +async def get_cassette_ops_window( + machine_id: str, limit: int = CASSETTE_OPS_WINDOW +) -> list[CassetteOp]: + """The operations a publish carries, oldest first. + + Selected newest-first to take the most recent `limit`, then reversed so the + machine applies them in the order they happened. That ordering matters: + a recount followed by a refill is not the same as the reverse. + """ + rows = await db.fetchall( + "SELECT * FROM spirekeeper.cassette_ops WHERE machine_id = :mid " + "ORDER BY created_at DESC, id DESC LIMIT :limit", + {"mid": machine_id, "limit": limit}, + CassetteOp, + ) + return list(reversed(rows)) + + +async def list_cassette_ops( + machine_id: str, limit: int = CASSETTE_OPS_WINDOW +) -> list[CassetteOp]: + """Newest-first, for the dashboard's history panel.""" + return await db.fetchall( + "SELECT * FROM spirekeeper.cassette_ops WHERE machine_id = :mid " + "ORDER BY created_at DESC, id DESC LIMIT :limit", + {"mid": machine_id, "limit": limit}, + CassetteOp, + ) + + +def _should_ack_op(existing: CassetteOp | None, machine_id: str) -> bool: + """Pure decision behind mark_cassette_ops_acked, extracted so it is + testable without a database — same approach as the state-event gate. + + Three reasons not to ack, and each matters: + - the id is unknown, so there is nothing to close out; + - it belongs to another machine, and one machine reporting an id must + never be able to close out another machine's operation; + - it is already acked, and the first acknowledgement is the interesting + one. The machine echoes a WINDOW, so every id arrives many times over; + overwriting would keep sliding the timestamp forward and lose when the + operation actually landed. + """ + if existing is None: + return False + if existing.machine_id != machine_id: + return False + return existing.acked_at is None + + +async def mark_cassette_ops_acked(machine_id: str, op_ids: list[str]) -> int: + """Record that the machine reported these operation ids as applied. + + Returns how many rows moved from unacked to acked. See _should_ack_op for + which reports are ignored and why. + """ + if not op_ids: + return 0 + acked = 0 + now = datetime.now() + for op_id in op_ids: + existing = await get_cassette_op(op_id) + if not _should_ack_op(existing, machine_id): + continue + await db.execute( + "UPDATE spirekeeper.cassette_ops SET acked_at = :now " + "WHERE id = :id AND machine_id = :mid", + {"now": now, "id": op_id, "mid": machine_id}, + ) + acked += 1 + return acked diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index 8187fa2..ab6cd88 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -15,6 +15,7 @@ from datetime import datetime, timezone import pytest from pydantic import ValidationError +from ..crud import _should_ack_op from ..models import ( CASSETTE_OP_TYPES, CassetteOp, @@ -140,3 +141,28 @@ class TestWireShape: "schema_version": 2, "ops": [], } + + +class TestShouldAckOp: + """The pure decision behind mark_cassette_ops_acked. + + The machine echoes a WINDOW of applied ids on every state publish, so the + same id arrives repeatedly and from a machine that may not own it. + """ + + def test_acks_an_unacked_op_for_the_reporting_machine(self): + assert _should_ack_op(op(op_type="empty"), "m1") is True + + def test_ignores_an_unknown_id(self): + assert _should_ack_op(None, "m1") is False + + def test_ignores_an_op_belonging_to_another_machine(self): + """Ids are unique, but a report from one machine must never close out + another machine's operation.""" + assert _should_ack_op(op(op_type="empty", machine_id="m2"), "m1") is False + + def test_keeps_the_first_acknowledgement(self): + """Every subsequent window carries the id again. Re-acking would slide + the timestamp forward and lose when the operation actually landed.""" + already = op(op_type="empty", acked_at=AT) + assert _should_ack_op(already, "m1") is False From 90bd43d6dae200f671e13c5f92bc748820d2b27f Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:35:42 +0200 Subject: [PATCH 3/8] feat(cassettes): publish operations to the ATM MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The v2 operator to ATM wire. Same kind-30078 document and the same d-tag the counts wire used, because the machine subscribes by that tag and the document is addressable, so v2 replaces v1 in place. Sends a WINDOW of recent operations, oldest-first, not just the newest change. Each publish replaces the last, so a machine that was offline for one of them would otherwise never see that operation again; carrying the recent history means the channel heals itself without anyone noticing it broke. Re-delivery costs nothing because every op carries an id the machine dedups on. Tests pin the contract rather than the implementation: the d-tag, that the payload declares v2 and carries no positions key, that window order survives the publisher untouched, that an empty window still ships a well-formed document so a machine can tell "no operations" from "operator still on v1", and that an npub entered in the UI is normalised to hex — get that last one wrong and the machine's subscription filter silently never matches. Additive. The endpoints still publish counts until the next commit, so the tree is not left half-switched. Co-Authored-By: Claude Fable 5.1 --- cassette_transport.py | 50 +++++++++++++++++++-- tests/test_cassette_ops.py | 91 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 138 insertions(+), 3 deletions(-) diff --git a/cassette_transport.py b/cassette_transport.py index e6517a1..f64ddec 100644 --- a/cassette_transport.py +++ b/cassette_transport.py @@ -17,8 +17,14 @@ publishes position-keyed cassette config to a target ATM via: The ATM-side consumer (lamassu-next#56) subscribes by the d-tag + its own npub, decrypts, validates, applies, hot-reloads HAL. -Reverse direction (ATM → operator, v1 = one-shot bootstrap on first boot, -v2 = continuous reverse channel for reconciliation): +The operator → ATM direction carries OPERATIONS as of v2 (bitspire ADR-004): +refill, empty, recount, set_denomination, each with an id the machine dedups +on. It used to carry absolute counts, which meant the operator and the machine +both wrote the same value over a transport that never tells a writer it lost — +so a form loaded before a dispense discarded that dispense when published. + +Reverse direction (ATM → operator, continuous: the machine publishes on +startup, after every change to its bays, and on a heartbeat): kind = 30078 tags = [ @@ -52,7 +58,12 @@ from lnbits.core.signers.base import ( ) from lnbits.utils.nostr import normalize_public_key -from .models import Machine, PublishCassettesPayload +from .models import ( + CassetteOp, + Machine, + PublishCassetteOpsPayload, + PublishCassettesPayload, +) from .nip44 import Nip44Error from .nostr_publish import ( NostrPublishError, @@ -179,6 +190,39 @@ async def publish_to_atm( return signed +async def publish_ops_to_atm( + machine: Machine, + ops: list[CassetteOp], + operator_user_id: str, +) -> dict: + """Publish the operator's recent cassette OPERATIONS to the target ATM. + + The v2 wire (bitspire ADR-004). Replaces sending absolute counts, which + let a dashboard form loaded before a dispense silently discard that + dispense — the operator and the machine were both writing the same value + over a transport that never tells a writer it lost. + + `ops` is a WINDOW, oldest-first, not just the newest change. The event is + addressable, so each publish replaces the last, and a machine that was + offline for one of them would otherwise never see that operation again. + Carrying the recent history means the channel heals itself without anyone + noticing it broke. Re-delivery is harmless because each op carries an id + the machine dedups on. + """ + atm_pubkey_hex = _atm_hex_pubkey(machine) + payload = PublishCassetteOpsPayload(ops=ops) + signed = await publish_encrypted_kind_30078( + operator_user_id=operator_user_id, + recipient_pubkey_hex=atm_pubkey_hex, + d_tag=_config_d_tag(atm_pubkey_hex), + payload=payload.to_wire_dict(), + log_context=( + f"cassette ops (machine={machine.id}, ops={[o.op_type for o in ops]})" + ), + ) + return signed + + # ============================================================================= # Consume — ATM → operator (the bootstrap consumer task) # ============================================================================= diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index ab6cd88..7014a94 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -15,11 +15,13 @@ from datetime import datetime, timezone import pytest from pydantic import ValidationError +from .. import cassette_transport from ..crud import _should_ack_op from ..models import ( CASSETTE_OP_TYPES, CassetteOp, CreateCassetteOpData, + Machine, PublishCassetteOpsPayload, ) @@ -166,3 +168,92 @@ class TestShouldAckOp: the timestamp forward and lose when the operation actually landed.""" already = op(op_type="empty", acked_at=AT) assert _should_ack_op(already, "m1") is False + + +# ============================================================================= +# publish_ops_to_atm — the v2 wire contract +# ============================================================================= + +ATM_HEX = "df2003343784b69cb813b2a4fd231f83ae81133279251c735414f9909baa7ac6" + + +def machine() -> Machine: + return Machine( + id="m1", + operator_user_id="op1", + machine_npub=ATM_HEX, + wallet_id="w1", + name="Cinderella", + location=None, + fiat_code="EUR", + is_active=True, + created_at=AT, + updated_at=AT, + ) + + +@pytest.fixture +def captured(monkeypatch): + """Capture what the transport would publish, without a relay or signer.""" + seen: dict = {} + + async def fake_publish(**kwargs): + seen.update(kwargs) + return {"id": "event-id"} + + monkeypatch.setattr( + cassette_transport, "publish_encrypted_kind_30078", fake_publish + ) + return seen + + +class TestPublishOpsToAtm: + @pytest.mark.asyncio + async def test_publishes_v2_ops_to_the_config_d_tag(self, captured): + ops = [ + op(id="a", op_type="refill", bills=100), + op(id="b", op_type="empty", position=3), + ] + await cassette_transport.publish_ops_to_atm(machine(), ops, "op1") + + # Same d-tag as the counts wire it replaces: the machine subscribes by + # this tag, and the document is addressable, so v2 replaces v1 in place. + assert captured["d_tag"] == f"bitspire-cassettes:{ATM_HEX}" + assert captured["recipient_pubkey_hex"] == ATM_HEX + assert captured["operator_user_id"] == "op1" + + payload = captured["payload"] + assert payload["schema_version"] == 2 + assert [o["id"] for o in payload["ops"]] == ["a", "b"] + assert payload["ops"][0]["bills"] == 100 + assert "positions" not in payload + + @pytest.mark.asyncio + async def test_preserves_window_order(self, captured): + """Order is meaning: a recount then a refill is not the same as the + reverse, so the publisher must not re-sort what crud handed it.""" + ops = [ + op(id="first", op_type="recount", count=10), + op(id="second", op_type="refill", bills=5), + ] + await cassette_transport.publish_ops_to_atm(machine(), ops, "op1") + assert [o["id"] for o in captured["payload"]["ops"]] == ["first", "second"] + + @pytest.mark.asyncio + async def test_publishes_an_empty_window_rather_than_skipping(self, captured): + """A machine with no operator history still gets a well-formed v2 + document, so it can tell 'no operations' from 'operator still on v1'.""" + await cassette_transport.publish_ops_to_atm(machine(), [], "op1") + assert captured["payload"] == {"schema_version": 2, "ops": []} + + @pytest.mark.asyncio + async def test_accepts_an_npub_and_publishes_hex(self, captured): + """Operators enter either form in the UI; the d-tag is always hex, or + the machine's subscription filter silently never matches.""" + import bech32 + + data = bech32.convertbits(bytes.fromhex(ATM_HEX), 8, 5) + npub = bech32.bech32_encode("npub", data) + 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}" From 3d8368bcc44ce0b76655654fbae3ac571eb0de05 Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:38:12 +0200 Subject: [PATCH 4/8] feat(cassettes): consume the machine's operation acknowledgements MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- models.py | 37 +++++++++++++++++------ tasks.py | 38 +++++++++++++++++++++-- tests/test_cassette_ops.py | 62 +++++++++++++++++++++++++++++++++++++- 3 files changed, 124 insertions(+), 13 deletions(-) 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"] From 2ac3e2064e6ef8936101b6e19ee873c145f0ffbf Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:44:18 +0200 Subject: [PATCH 5/8] feat(cassettes): swap the count-publish endpoint for operation endpoints MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The operator can no longer write a count. POST .../cassettes/ops records one operation — refill, empty, recount, set_denomination — and publishes the machine's recent window; GET .../cassettes/ops lists them newest first with acked_at, so the dashboard can tell a delivered operation from one merely sent. POST .../cassettes/publish is gone, along with update_cassette_config and UpsertCassetteConfigData. Nothing in the operator can now set a count, which is the point: a value with one writer cannot be clobbered. Under the old endpoint a dashboard form loaded before a dispense silently discarded that dispense on publish, and neither side could detect it — addressable events order by created_at at second granularity and a relay returns OK for an event it then drops, so the losing writer is never told. The op is recorded before the publish and is deliberately not rolled back when the publish fails. It records something that physically happened; notes went into a bay whether or not a relay was reachable. The window carries recent operations rather than just the newest, so an op that missed its own publish rides out with the next one. Validation rejects an unpaired machine and a position the machine has not reported. Bay count stays hardware-determined. --- cassette_transport.py | 33 +------ crud.py | 42 ++------- models.py | 25 ------ tests/test_cassette_configs.py | 45 +--------- views_api.py | 158 ++++++++++++++++----------------- 5 files changed, 89 insertions(+), 214 deletions(-) diff --git a/cassette_transport.py b/cassette_transport.py index f64ddec..bc84523 100644 --- a/cassette_transport.py +++ b/cassette_transport.py @@ -36,7 +36,7 @@ startup, after every change to its bays, and on a heartbeat): This module owns the wire-format side of both directions. The consumer task (tasks.py) calls `decrypt_and_parse_state_event` per incoming event; -the API endpoint (views_api.py) calls `publish_to_atm` per operator submit. +the API endpoint (views_api.py) calls `publish_ops_to_atm` per operation. The `` placeholder semantics (load-bearing per the 2026-05-30T11:50Z coord-log entry): always the ATM's hex pubkey, NEVER spirekeeper's @@ -86,7 +86,7 @@ __all__ = [ "RelayUnavailable", "build_state_d_tags_for_machines", "decrypt_and_parse_state_event", - "publish_to_atm", + "publish_ops_to_atm", ] _D_TAG_CONFIG_PREFIX = "bitspire-cassettes:" # operator → ATM @@ -163,33 +163,6 @@ def build_state_d_tags_for_machines(machines: list[Machine]) -> list[str]: # ============================================================================= -async def publish_to_atm( - machine: Machine, - payload: PublishCassettesPayload, - operator_user_id: str, -) -> dict: - """Build, encrypt, sign, and publish a kind-30078 cassette config event - from the operator to the target ATM. - - Returns the signed event dict on success (caller may log event.id for - audit). Raises NostrPublishError subclasses (re-exported here as - CassetteTransportError, OperatorIdentityMissing, SignerUnavailable, - RelayUnavailable) on hard failures. - """ - atm_pubkey_hex = _atm_hex_pubkey(machine) - signed = await publish_encrypted_kind_30078( - operator_user_id=operator_user_id, - recipient_pubkey_hex=atm_pubkey_hex, - d_tag=_config_d_tag(atm_pubkey_hex), - payload=payload.to_wire_dict(), - log_context=( - f"cassette config (machine={machine.id}, " - f"positions={sorted(payload.positions.keys())})" - ), - ) - return signed - - async def publish_ops_to_atm( machine: Machine, ops: list[CassetteOp], @@ -224,7 +197,7 @@ async def publish_ops_to_atm( # ============================================================================= -# Consume — ATM → operator (the bootstrap consumer task) +# Consume — ATM → operator (the machine's state reports) # ============================================================================= diff --git a/crud.py b/crud.py index bc5cb0a..e83b9af 100644 --- a/crud.py +++ b/crud.py @@ -36,7 +36,6 @@ from .models import ( UpdateDepositStatusData, UpdateMachineData, UpdateSuperConfigData, - UpsertCassetteConfigData, UpsertDcaLpData, ) @@ -1446,10 +1445,11 @@ async def upsert_fleet_snapshot( # Row lifecycle per #29: # - 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) -# - Row creation/deletion for a new position → admin only, via ATM -# re-provisioning + new bootstrap event (not exposed in v1 here) +# - The operator does NOT write these rows. It records operations +# (cassette_ops) and the machine keeps the running count; these columns +# hold what the machine last reported. +# - Rows appear and disappear only as the machine reports its bay set — +# the slot count is hardware-determined. def _as_unix(value) -> float | None: @@ -1521,38 +1521,6 @@ async def list_cassette_configs_for_machine( ) -async def update_cassette_config( - machine_id: str, - position: int, - data: UpsertCassetteConfigData, - *, - updated_by: str | None = None, -) -> 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_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. - """ - existing = await get_cassette_config(machine_id, position) - if existing is None: - return None - update_data: dict = {k: v for k, v in data.dict().items() if v is not None} - if not update_data: - return existing - update_data["updated_at"] = datetime.now() - update_data["updated_by"] = updated_by - set_clause = ", ".join(f"{k} = :{k}" for k in update_data) - update_data["mid"] = machine_id - update_data["pos"] = position - await db.execute( - f"UPDATE spirekeeper.cassette_configs SET {set_clause} " - "WHERE machine_id = :mid AND position = :pos", - update_data, - ) - return await get_cassette_config(machine_id, position) - - async def apply_reported_state( machine_id: str, event_id: str, diff --git a/models.py b/models.py index 7a09093..f2a1f3d 100644 --- a/models.py +++ b/models.py @@ -675,31 +675,6 @@ class CassetteConfig(BaseModel): state_event_id: str | None -class UpsertCassetteConfigData(BaseModel): - """Operator edits a single cassette row's denomination or count from - the dashboard. Both fields optional; pass only those changed. - Position is not edited — it's the row's identity (hardware bay).""" - - denomination: int | None = None - count: int | None = None - - @validator("denomination") - def denomination_positive(cls, v): - if v is None: - return v - if v <= 0: - raise ValueError("denomination must be > 0") - return v - - @validator("count") - def count_non_negative(cls, v): - if v is None: - return v - if v < 0: - raise ValueError("count must be >= 0") - return v - - class CassettePayloadRow(BaseModel): """One position's payload values in the wire-format `{"positions": {"": {"denomination", "count"}}}`.""" diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index 242aa13..b6adf2d 100644 --- a/tests/test_cassette_configs.py +++ b/tests/test_cassette_configs.py @@ -2,9 +2,9 @@ Tests for the v1.1 cassette-config layer (aiolabs/satmachineadmin#29). 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) + - Pydantic validator behaviour on PublishCassettesPayload + the row model + (position key coercion, integer ranges, multiple-same-denomination + payloads, wire-format round-trip) - _should_apply_state_event ordering gate (extracted from apply_reported_state so the decision is testable without a database round-trip) @@ -29,7 +29,6 @@ from ..crud import _as_unix, _should_apply_state_event from ..models import ( CassettePayloadRow, PublishCassettesPayload, - UpsertCassetteConfigData, ) # ============================================================================= @@ -149,44 +148,6 @@ class TestCassettePayloadRow: CassettePayloadRow(denomination=20, count=-1) -# ============================================================================= -# UpsertCassetteConfigData — operator-edit form -# ============================================================================= - - -class TestUpsertCassetteConfigData: - """Operator-driven row edit. Both fields optional; same int constraints - as the wire-format row but applied independently per-edit. Position is - NOT editable — it's the row's identity (the hardware bay number).""" - - def test_partial_update_count_only(self): - d = UpsertCassetteConfigData(count=80) - assert d.count == 80 - assert d.denomination is None - - def test_partial_update_denomination_only(self): - """v1.1 operational case: operator records a cartridge swap at - refill — slot 1 was $20, dispatcher replaced with $50.""" - d = UpsertCassetteConfigData(denomination=50) - assert d.denomination == 50 - assert d.count is None - - def test_empty_update_is_legal(self): - """An empty UpsertCassetteConfigData parses fine; the CRUD short- - circuits a no-op on empty payload (no SQL emitted).""" - d = UpsertCassetteConfigData() - assert d.count is None - assert d.denomination is None - - def test_rejects_negative_count(self): - with pytest.raises(ValueError): - UpsertCassetteConfigData(count=-1) - - def test_rejects_non_positive_denomination(self): - with pytest.raises(ValueError): - UpsertCassetteConfigData(denomination=0) - - # ============================================================================= # _should_apply_state_event — ordering gate # ============================================================================= diff --git a/views_api.py b/views_api.py index 73b2d0b..33fb2de 100644 --- a/views_api.py +++ b/views_api.py @@ -26,7 +26,7 @@ from .cassette_transport import ( OperatorIdentityMissing, RelayUnavailable, SignerUnavailable, - publish_to_atm, + publish_ops_to_atm, ) from .fee_transport import publish_fee_config from .pairing import ( @@ -40,6 +40,7 @@ from .pairing import ( from .crud import ( append_settlement_note, count_completed_legs_for_settlement, + create_cassette_op, create_dca_client, create_deposit, create_machine, @@ -47,6 +48,7 @@ from .crud import ( delete_deposit, delete_machine, force_reset_stuck_settlement, + get_cassette_ops_window, get_client_balance_summary, get_commission_splits, get_dca_client, @@ -66,12 +68,12 @@ from .crud import ( get_super_config, list_all_active_machines, list_cassette_configs_for_machine, + list_cassette_ops, lp_is_onboarded, replace_commission_splits, reset_settlement_for_retry, set_machine_pairing, set_machine_unpaired, - update_cassette_config, update_dca_client, update_deposit, update_deposit_status, @@ -86,8 +88,10 @@ from .distribution import ( from .models import ( AppendSettlementNoteData, CassetteConfig, + CassetteOp, ClientBalanceSummary, CommissionSplit, + CreateCassetteOpData, CreateDcaClientData, CreateDepositData, CreateMachineData, @@ -98,7 +102,6 @@ from .models import ( Machine, PairMachineData, PartialDispenseData, - PublishCassettesPayload, SetCommissionSplitsData, SettleBalanceData, StuckSettlementsResponse, @@ -108,7 +111,6 @@ from .models import ( UpdateDepositStatusData, UpdateMachineData, UpdateSuperConfigData, - UpsertCassetteConfigData, ) spirekeeper_api_router = APIRouter() @@ -1102,16 +1104,21 @@ async def api_update_super_config( # ============================================================================= -# Cassette configs (#29 v1.1) — per-machine ATM cassette inventory +# Cassettes — per-machine ATM inventory (bitspire ADR-004) # ============================================================================= -# v1.1 surface, paired with aiolabs/lamassu-next#56 ATM-side. Two endpoints: -# GET /machines/{id}/cassettes — list rows for the operator UI -# POST /machines/{id}/cassettes/publish — apply edits + publish kind-30078 +# GET /machines/{id}/cassettes — bays as the machine last reported them +# GET /machines/{id}/cassettes/ops — recent operations, with their acks +# POST /machines/{id}/cassettes/ops — record one operation and publish # -# Row creation (new (machine_id, position) pairs) is admin-only via the -# bootstrap consumer task — slot count is hardware-determined. Operator- -# side flow is edit-and-publish over the existing rows only; the editable -# fields per row are denomination and count. +# The operator does not write counts. It records what it DID to a bay and the +# machine, which holds the notes, keeps the running total. The rows behind the +# first endpoint are the machine's report, not an operator draft. +# +# This replaced an edit-and-publish form over absolute counts. Both sides wrote +# the same value across a transport that never tells a writer it lost, so a +# form loaded before a dispense discarded that dispense when published — and +# nothing could detect it afterwards. Bay count stays hardware-determined: +# rows appear and disappear only as the machine reports them. @spirekeeper_api_router.get( @@ -1129,105 +1136,96 @@ async def api_list_machine_cassettes( return await list_cassette_configs_for_machine(machine_id) -@spirekeeper_api_router.post( - "/api/v1/dca/machines/{machine_id}/cassettes/publish", - response_model=list[CassetteConfig], +@spirekeeper_api_router.get( + "/api/v1/dca/machines/{machine_id}/cassettes/ops", + response_model=list[CassetteOp], ) -async def api_publish_machine_cassettes( +async def api_list_machine_cassette_ops( machine_id: str, - payload: PublishCassettesPayload, user: User = Depends(check_user_exists), -) -> list[CassetteConfig]: - """Operator submits the full per-machine cassette state for publish to - the ATM. Validates the position set matches what's currently in - cassette_configs for the machine (slot count is hardware-fixed), - upserts each row, then encrypts + signs + publishes a kind-30078 - event tagged with d=bitspire-cassettes: and - p=. +) -> list[CassetteOp]: + """Recent cassette operations for a machine, newest first. - The `` placeholder in the published d-tag is the ATM's hex pubkey - from machine.machine_npub (canonicalised via normalize_public_key), - NOT the internal dca_machines.id UUID — see #29 'machine_id semantics' - section and coord-log 2026-05-30T11:50Z load-bearing nudge. + Each carries acked_at: null until the machine has reported that id back in + its state document, which is how the dashboard distinguishes an operation + that has been delivered from one that has merely been sent. + """ + await _machine_owned_by(machine_id, user.id) + return await list_cassette_ops(machine_id) - Returns the fresh cassette_configs rows after the upserts so the UI - can refresh its table from one round-trip. + +@spirekeeper_api_router.post( + "/api/v1/dca/machines/{machine_id}/cassettes/ops", + response_model=CassetteOp, +) +async def api_create_machine_cassette_op( + machine_id: str, + data: CreateCassetteOpData, + user: User = Depends(check_user_exists), +) -> CassetteOp: + """Record one cassette operation and publish the machine's recent window. + + This replaces publishing absolute counts. The operator now records what it + DID — a refill, an empty, a recount, a denomination change — and the + machine keeps the running total. Both sides used to write the same value + over a transport that never tells a writer it lost, so a dashboard form + loaded before a dispense silently discarded that dispense when published. + + The op is recorded BEFORE the publish and is deliberately not rolled back + if the publish fails. It represents something that physically happened — + notes went into a bay — and that stays true whether or not a relay was + reachable. Because each publish carries a window of recent operations + rather than just the newest, an op that missed its own publish rides out + with the next one. Errors: - 400 — payload position set doesn't match the machine's stored set - (operator publishing for a slot that doesn't exist on the - ATM; or the bootstrap hasn't landed yet so no rows exist) - 400 — operator hasn't onboarded a Nostr identity - 503 — signer offline / client-side-only, or nostrclient extension - not installed on this LNbits instance - 500 — anything else from the publish path + 400 — machine not paired, so there is no ATM identity to publish to + 400 — position is not a bay this machine has reported + 503 — signer offline, or relay/nostrclient unreachable. The operation is + still recorded and will be delivered with the next publish. """ machine = await _machine_owned_by(machine_id, user.id) if not machine.machine_npub: - # Unpaired machine (machine_npub None — nullable since #29/m011) has no - # ATM identity to publish a cassette config to. Fail fast with a clean - # 400 instead of crashing publish_to_atm's normalize_public_key(None). raise HTTPException( HTTPStatus.BAD_REQUEST, - "machine is not paired — pair it before publishing cassette config", + "machine is not paired — pair it before recording cassette operations", ) existing = await list_cassette_configs_for_machine(machine_id) - existing_positions = {row.position for row in existing} - incoming_positions = set(payload.positions.keys()) - if not existing: raise HTTPException( HTTPStatus.BAD_REQUEST, ( - "No cassette_configs rows exist for this machine yet — " - "waiting for the ATM's bootstrap state event. Power on the " - "ATM and confirm it has reached the configured relay; " - "spirekeeper will auto-populate cassette_configs on " - "receipt." + "No cassette rows for this machine yet — waiting for its state " + "event. Power on the ATM and confirm it has reached the " + "configured relay; spirekeeper populates the bays on receipt." ), ) - if existing_positions != incoming_positions: - missing = existing_positions - incoming_positions - extra = incoming_positions - existing_positions + known_positions = {row.position for row in existing} + if data.position not in known_positions: raise HTTPException( HTTPStatus.BAD_REQUEST, ( - "Payload position set doesn't match the machine's stored " - f"set. Missing from payload: {sorted(missing)}; extra in " - f"payload: {sorted(extra)}. Slot count is hardware-fixed " - "— re-provision the ATM via atm-tui to add/remove physical " - "bays, then re-publish." + f"position {data.position} is not a bay this machine has " + f"reported (has: {sorted(known_positions)}). Bay count is " + "hardware-determined; re-provision via atm-tui to change it." ), ) - # Apply each per-row edit so the operator-believed state on - # spirekeeper reflects the published payload, even if the ATM - # ack lands later (v2). updated_by audit-stamps the operator user id. - for pos, row in payload.positions.items(): - updated = await update_cassette_config( - machine_id, - pos, - UpsertCassetteConfigData(denomination=row.denomination, count=row.count), - updated_by=user.id, - ) - if updated is None: - # Defensive — we just validated the row exists, but a - # concurrent delete could land between. Surface as 500. - raise HTTPException( - HTTPStatus.INTERNAL_SERVER_ERROR, - f"cassette row for position {pos} disappeared mid-publish", - ) + op = await create_cassette_op(machine_id, data, created_by=user.id) + window = await get_cassette_ops_window(machine_id) try: - await publish_to_atm(machine, payload, user.id) + await publish_ops_to_atm(machine, window, user.id) except OperatorIdentityMissing as exc: raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc - except SignerUnavailable as exc: - raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc - except RelayUnavailable as exc: - raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc + except (SignerUnavailable, RelayUnavailable) as exc: + raise HTTPException( + HTTPStatus.SERVICE_UNAVAILABLE, + f"{exc} — the operation was recorded and will be delivered with " + "the next publish", + ) from exc except CassetteTransportError as exc: raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc - return await list_cassette_configs_for_machine(machine_id) + return op From c76a1bb12506a90efcdd97c70f2ad7fd703fd759 Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:47:30 +0200 Subject: [PATCH 6/8] feat(cassettes): persist the machine's counts-uncertain marker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The machine has been publishing counts_uncertain_since since v1 of the state document and spirekeeper has been parsing it into a field nobody read. That defeats the point of the marker: it exists so a human opens the bay and recounts. A dispenser can throw, or time out, after notes have physically moved. The machine cannot know how many left, so rather than decrement a number it would be guessing at, it stamps the moment (bitspire ADR-004, decision 3). m014 gives that stamp a home on dca_machines, and the consumer mirrors it on every state event. Stored on the machine rather than the bay because the uncertainty is about the dispense as a whole; a multi-bay dispense that fails midway gives no reliable way to attribute it to one position. Written through even when the machine reports None. The machine clearing the marker is as important as setting it — the operator recounted, the bay is trustworthy again — and a banner that never goes away is a banner nobody reads. --- crud.py | 19 +++++++++++++++++++ migrations.py | 21 +++++++++++++++++++++ models.py | 4 ++++ tasks.py | 32 +++++++++++++++++++++++++++++++- tests/test_cassette_ops.py | 37 +++++++++++++++++++++++++++++++++++++ 5 files changed, 112 insertions(+), 1 deletion(-) diff --git a/crud.py b/crud.py index e83b9af..8508e9a 100644 --- a/crud.py +++ b/crud.py @@ -257,6 +257,25 @@ async def set_machine_unpaired(machine_id: str) -> Machine | None: return await get_machine(machine_id) +async def set_machine_counts_uncertain(machine_id: str, since: datetime | None) -> None: + """Record (or clear) the machine's own "I can't vouch for these counts" + marker, straight from its state document. + + `updated_at` is deliberately left alone. This is the machine reporting + about itself on a five-minute heartbeat, not an operator editing the + machine, and touching the timestamp on every heartbeat would make the + registry look perpetually just-modified. + """ + await db.execute( + """ + UPDATE spirekeeper.dca_machines + SET counts_uncertain_since = :since + WHERE id = :id + """, + {"since": since, "id": machine_id}, + ) + + async def delete_machine(machine_id: str) -> None: await db.execute( "DELETE FROM spirekeeper.dca_machines WHERE id = :id", diff --git a/migrations.py b/migrations.py index 52aebf3..8f5ddcf 100644 --- a/migrations.py +++ b/migrations.py @@ -901,3 +901,24 @@ async def m013_add_cassette_ops(db): "CREATE INDEX IF NOT EXISTS cassette_ops_machine_idx " "ON cassette_ops (machine_id, created_at DESC)" ) + + +async def m014_add_counts_uncertain_since(db): + """Surface the machine's "I don't know what left the bay" marker. + + A dispenser can throw, or time out, after notes have physically moved. The + machine cannot know how many left, so rather than decrement a number it + would be guessing at, it stamps the moment and reports it (bitspire + ADR-004, decision 3). It has been publishing this field since v1 of the + state document and the operator has been discarding it, which defeats the + point: the marker exists to tell a human to open the bay and recount. + + Stored on the machine, not the bay, because the uncertainty is about the + dispense as a whole — a multi-bay dispense that fails midway leaves no + reliable way to attribute it to one position. Cleared to NULL by the + machine's own report once it is confident again. + """ + await db.execute( + "ALTER TABLE spirekeeper.dca_machines " + "ADD COLUMN counts_uncertain_since TIMESTAMP" + ) diff --git a/models.py b/models.py index f2a1f3d..1a0658c 100644 --- a/models.py +++ b/models.py @@ -63,6 +63,10 @@ class Machine(BaseModel): # NIP-46 bunker pairing (S0 / #9). NULL until the spire is first paired. bunker_spire_key_name: str | None = None paired_at: datetime | None = None + # Set when the machine reports that a dispense ended without a reliable + # count of what left the bay; cleared by the machine's own report. The + # dashboard turns this into a prompt to open the bay and recount. + counts_uncertain_since: datetime | None = None created_at: datetime updated_at: datetime diff --git a/tasks.py b/tasks.py index 0d2f138..6e5ac55 100644 --- a/tasks.py +++ b/tasks.py @@ -339,6 +339,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: get_machine_by_atm_pubkey_hex, list_all_active_machines, mark_cassette_ops_acked, + set_machine_counts_uncertain, ) machines = await list_all_active_machines() @@ -375,6 +376,7 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str: get_machine_by_atm_pubkey_hex, apply_reported_state, mark_cassette_ops_acked, + set_machine_counts_uncertain, ) except Exception as exc: logger.warning( @@ -410,11 +412,38 @@ async def _record_op_acknowledgements( ) +async def _record_counts_uncertainty( + machine_id: str, payload, set_machine_counts_uncertain +) -> None: + """Mirror the machine's counts-uncertain marker onto its registry row. + + The machine sets this when a dispense ended without a reliable count of + what physically left the bay — a dispenser throw, or a timeout. It cannot + know how many notes moved, so it says so instead of decrementing a number + it would be guessing at. + + Written on every state event, including when it is None, because the + machine clearing the marker is exactly as important as setting it: the + operator has recounted, the bay is trustworthy again, and a banner that + never goes away is a banner nobody reads. + """ + from datetime import datetime as _datetime + from datetime import timezone as _timezone + + since = None + if payload.counts_uncertain_since is not None: + since = _datetime.fromtimestamp( + int(payload.counts_uncertain_since), tz=_timezone.utc + ) + await set_machine_counts_uncertain(machine_id, since) + + async def _handle_cassette_state_event( event_message, get_machine_by_atm_pubkey_hex, apply_reported_state, mark_cassette_ops_acked, + set_machine_counts_uncertain, ) -> None: """Verify signature, resolve the operator's signer, decrypt via the signer abstraction (bunker round-trip for RemoteBunkerSigner; direct @@ -526,5 +555,6 @@ async def _handle_cassette_state_event( ) # Acknowledgement runs regardless of whether the counts were newer — see - # _record_op_acknowledgements for why. + # _record_op_acknowledgements for why. Same for the uncertainty marker. await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked) + await _record_counts_uncertainty(machine.id, payload, set_machine_counts_uncertain) diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index 2ddeedd..150d493 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -317,3 +317,40 @@ class TestRecordOpAcknowledgements: "m1", state_payload(applied_ops=["late-ack"]), mark ) assert seen == ["late-ack"] + + +# ============================================================================= +# _record_counts_uncertainty — the machine saying "don't trust these counts" +# ============================================================================= + + +class TestRecordCountsUncertainty: + @pytest.mark.asyncio + async def test_stores_the_reported_moment_as_utc(self): + calls = [] + + async def setter(machine_id, since): + calls.append((machine_id, since)) + + await tasks._record_counts_uncertainty( + "m1", state_payload(counts_uncertain_since=1790110546), setter + ) + assert len(calls) == 1 + machine_id, since = calls[0] + assert machine_id == "m1" + assert since is not None + assert since.tzinfo is not None + assert int(since.timestamp()) == 1790110546 + + @pytest.mark.asyncio + async def test_clears_the_marker_when_the_machine_is_confident_again(self): + """A banner that never goes away is a banner nobody reads. The + machine dropping the field is how the operator learns the recount + took, so None must be written through rather than skipped.""" + calls = [] + + async def setter(machine_id, since): + calls.append((machine_id, since)) + + await tasks._record_counts_uncertainty("m1", state_payload(), setter) + assert calls == [("m1", None)] From 79a4f8329347c32aee3414fbb0b0e21572e2ff16 Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:50:08 +0200 Subject: [PATCH 7/8] feat(cassettes): record operations from the dashboard instead of counts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The cassettes tab no longer has editable count fields, because there is no longer an endpoint that would accept them. The bays render read-only from the machine's own report, and a Record-operation dialog captures what the operator did: a refill in notes added, an empty, a recount, a denomination change. This removes the failure the tab used to invite. A form loaded before a dispense held a count that was already wrong, and publishing it overwrote the dispense with no error on either side. Recording a delta instead means a dispense that happened while the dialog was open is kept rather than discarded, and a recount is now an explicit act — what an operator opening a bay and counting actually does — rather than being indistinguishable from a stale form. A recent-operations list shows each one as Applied or Pending from acked_at, which is the machine echoing the id back. Pending needs no retry button: every publish carries the recent window, so an operation that missed its own publish keeps being re-offered until it lands, and saying so in the panel is more useful than a button that would do nothing new. The counts-uncertain banner tells the operator when the machine cannot vouch for its own numbers and asks for the recount that clears it. The machine row is re-read on every cassette refresh, since that flag is set by the consumer while the dialog is open. --- static/js/index.js | 203 ++++++++++++++++++---------- templates/spirekeeper/index.html | 224 ++++++++++++++++++++----------- 2 files changed, 277 insertions(+), 150 deletions(-) diff --git a/static/js/index.js b/static/js/index.js index 6c722b3..e8a31e2 100644 --- a/static/js/index.js +++ b/static/js/index.js @@ -212,15 +212,14 @@ window.app = Vue.createApp({ loading: false, machine: null, settlements: [], - // Cassettes sub-tab state (#29 v1) — see openCassettePublishConfirm / - // submitCassettePublish methods + the cassettes panel in - // templates/spirekeeper/index.html. + // Cassettes sub-tab state (v2, bitspire ADR-004) — see + // openCassetteOpDialog / submitCassetteOp + the cassettes panel in + // templates/spirekeeper/index.html. Read-only by design: the + // operator records operations, the machine owns the counts. activeTab: 'settlements', - cassetteEdits: [], // editable working copy of cassette_configs rows - cassettesPristine: [], // last-known-clean snapshot for revert + cassettes: [], // machine-reported rows, not editable + cassetteOps: [], // recent operations, newest first cassettesLoading: false, - cassettesPublishing: false, - cassettesDirty: false, cassettesError: null }, cassettesTable: { @@ -228,14 +227,27 @@ window.app = Vue.createApp({ {name: 'position', label: 'Bay', field: 'position', align: 'right'}, {name: 'denomination', label: 'Denomination', field: 'denomination', align: 'right'}, {name: 'count', label: 'Count', field: 'count', align: 'right'}, - {name: 'state', label: 'ATM-reported', field: 'state_denomination', align: 'right'}, - {name: 'updated_at', label: 'Updated', field: 'updated_at', align: 'left'} + {name: 'state_at', label: 'Machine reported', field: 'state_at', align: 'left'}, + {name: 'actions', label: '', field: 'position', align: 'right'} ], pagination: {rowsPerPage: 0} // hide pagination — cassette count is small }, - cassettePublishConfirm: { - show: false + cassetteOpDialog: { + show: false, + saving: false, + error: null, + position: null, + op_type: 'refill', + bills: null, + count: null, + denomination: null }, + cassetteOpTypeOptions: [ + {label: 'Refill — notes added', value: 'refill'}, + {label: 'Empty — bay emptied', value: 'empty'}, + {label: 'Recount — notes counted', value: 'recount'}, + {label: 'Set denomination', value: 'set_denomination'} + ], partialDispenseDialog: { show: false, saving: false, @@ -286,6 +298,28 @@ window.app = Vue.createApp({ }, computed: { + cassetteBayOptions() { + return this.machineDetail.cassettes.map(row => ({ + label: `Bay ${row.position} — ${row.denomination} ${ + (this.machineDetail.machine || {}).fiat_code || '' + } ×${row.count}`, + value: row.position + })) + }, + + cassetteOpIsComplete() { + // Mirrors CreateCassetteOpData's root validator: exactly one value + // field, decided by the type. Enforced here only to keep the button + // honest — the server rejects a malformed op regardless. + const d = this.cassetteOpDialog + if (!d.position) return false + if (d.op_type === 'empty') return true + if (d.op_type === 'refill') return Number(d.bills) > 0 + if (d.op_type === 'recount') return d.count !== null && Number(d.count) >= 0 + if (d.op_type === 'set_denomination') return Number(d.denomination) > 0 + return false + }, + superAnyFee() { // Banner styling key — true when either directional super fee is // non-zero, so the banner reads as "active platform fee" instead @@ -943,9 +977,8 @@ window.app = Vue.createApp({ async viewMachine(machine) { this.machineDetail.machine = machine this.machineDetail.settlements = [] - this.machineDetail.cassetteEdits = [] - this.machineDetail.cassettesPristine = [] - this.machineDetail.cassettesDirty = false + this.machineDetail.cassettes = [] + this.machineDetail.cassetteOps = [] this.machineDetail.cassettesError = null this.machineDetail.activeTab = 'settlements' this.machineDetail.show = true @@ -972,21 +1005,25 @@ window.app = Vue.createApp({ }, // ----------------------------------------------------------------- - // Cassette inventory (#29 v1) + // Cassette inventory + operations (v2, bitspire ADR-004) // ----------------------------------------------------------------- async loadMachineCassettes() { if (!this.machineDetail.machine) return this.machineDetail.cassettesLoading = true this.machineDetail.cassettesError = null + const base = `${MACHINES_PATH}/${this.machineDetail.machine.id}` try { - const {data} = await LNbits.api.request( - 'GET', - `${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes` - ) - const rows = (data || []).map(row => ({...row, _dirty: false})) - this.machineDetail.cassetteEdits = rows - this.machineDetail.cassettesPristine = JSON.parse(JSON.stringify(rows)) - this.machineDetail.cassettesDirty = false + // The machine row is re-read too: counts_uncertain_since lives on it + // and is set by the consumer as state events land, so the row the + // machines table handed us goes stale while this dialog is open. + const [machine, bays, ops] = await Promise.all([ + LNbits.api.request('GET', base), + LNbits.api.request('GET', `${base}/cassettes`), + LNbits.api.request('GET', `${base}/cassettes/ops`) + ]) + if (machine.data) this.machineDetail.machine = machine.data + this.machineDetail.cassettes = bays.data || [] + this.machineDetail.cassetteOps = ops.data || [] } catch (e) { this._notifyError(e, 'Failed to load cassettes') } finally { @@ -994,73 +1031,91 @@ window.app = Vue.createApp({ } }, - markCassetteDirty(row) { - // Find pristine match by position (the row identity) and compare; - // flip _dirty + overall dirty flag accordingly. Editable fields - // are denomination + count; position is the immutable row key. - const pristine = this.machineDetail.cassettesPristine.find( - p => p.position === row.position + cassetteOpIcon(opType) { + return ( + { + refill: 'add_circle_outline', + empty: 'remove_circle_outline', + recount: 'fact_check', + set_denomination: 'sell' + }[opType] || 'help_outline' ) - row._dirty = - !pristine || - Number(row.denomination) !== Number(pristine.denomination) || - Number(row.count) !== Number(pristine.count) - this.machineDetail.cassettesDirty = - this.machineDetail.cassetteEdits.some(r => r._dirty) }, - revertCassetteEdits() { - this.machineDetail.cassetteEdits = JSON.parse( - JSON.stringify(this.machineDetail.cassettesPristine) - ) - this.machineDetail.cassettesDirty = false - this.machineDetail.cassettesError = null - }, - - openCassettePublishConfirm() { - if (!this.machineDetail.cassettesDirty) return - this.machineDetail.cassettesError = null - this.cassettePublishConfirm.show = true - }, - - async submitCassettePublish() { - // Build the PublishCassettesPayload shape (v1.1, position-keyed): - // { positions: { "": { denomination, count }, ... } } - // The API enforces the position set matches what's stored — - // since we only edit existing rows, this should always pass. - const positions = {} - for (const row of this.machineDetail.cassetteEdits) { - positions[String(row.position)] = { - denomination: Number(row.denomination), - count: Number(row.count) - } + cassetteOpSummary(op) { + const fiat = (this.machineDetail.machine || {}).fiat_code || '' + const bay = `Bay ${op.position}` + if (op.op_type === 'refill') return `${bay} — added ${op.bills} notes` + if (op.op_type === 'empty') return `${bay} — emptied` + if (op.op_type === 'recount') return `${bay} — recounted to ${op.count}` + if (op.op_type === 'set_denomination') { + return `${bay} — denomination set to ${op.denomination} ${fiat}`.trim() } - const payload = {positions} - this.machineDetail.cassettesPublishing = true + return `${bay} — ${op.op_type}` + }, + + openCassetteOpDialog(position) { + const bays = this.machineDetail.cassettes + if (!bays.length) return + Object.assign(this.cassetteOpDialog, { + show: true, + saving: false, + error: null, + position: position || bays[0].position, + op_type: 'refill', + bills: null, + count: null, + denomination: null + }) + }, + + resetCassetteOpValue() { + // Each type carries exactly one value field and the server rejects an + // op that carries a foreign one. Clearing on every type switch means a + // number typed under the previous type can't ride along invisibly. + Object.assign(this.cassetteOpDialog, { + bills: null, + count: null, + denomination: null, + error: null + }) + }, + + async submitCassetteOp() { + const d = this.cassetteOpDialog + if (!this.cassetteOpIsComplete) return + const payload = {position: Number(d.position), op_type: d.op_type} + if (d.op_type === 'refill') payload.bills = Number(d.bills) + if (d.op_type === 'recount') payload.count = Number(d.count) + if (d.op_type === 'set_denomination') { + payload.denomination = Number(d.denomination) + } + d.saving = true + d.error = null try { - const {data} = await LNbits.api.request( + await LNbits.api.request( 'POST', - `${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes/publish`, + `${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes/ops`, null, payload ) - const fresh = (data || []).map(r => ({...r, _dirty: false})) - this.machineDetail.cassetteEdits = fresh - this.machineDetail.cassettesPristine = JSON.parse(JSON.stringify(fresh)) - this.machineDetail.cassettesDirty = false - this.cassettePublishConfirm.show = false + d.show = false Quasar.Notify.create({ type: 'positive', - message: 'Cassette config published to ATM' + message: 'Operation recorded and published to the ATM' }) + await this.loadMachineCassettes() } catch (e) { const detail = (e && e.response && e.response.data && e.response.data.detail) || - 'Publish failed' - this.machineDetail.cassettesError = detail - this._notifyError(e, 'Publish failed') + 'Could not record the operation' + // A 503 means the op IS recorded and will ride out with the next + // publish, so reload either way — the list should show it pending. + d.error = detail + this._notifyError(e, 'Operation failed') + await this.loadMachineCassettes() } finally { - this.machineDetail.cassettesPublishing = false + d.saving = false } }, diff --git a/templates/spirekeeper/index.html b/templates/spirekeeper/index.html index fd26d1c..ad01503 100644 --- a/templates/spirekeeper/index.html +++ b/templates/spirekeeper/index.html @@ -1151,23 +1151,22 @@
Cassettes

- Per-cassette count and physical bay position. Denomination - set is hardware-determined (re-provision via atm-tui to - change). "Publish to ATM" encrypts + signs + sends the new - config to the machine via Nostr. + The machine owns these counts — it holds the notes. You + record what you did to a bay and it keeps the + running total. Bay count and denomination set are + hardware-determined (re-provision via atm-tui to change + the bays themselves).

- - Discard unsaved edits + + Re-read the machine's latest report - +
@@ -1179,64 +1178,108 @@ - + + The machine can't vouch for these counts. + A dispense ended without a reliable count of what left the bay + (). + Open the bay, count it, and record a Recount — the + machine clears this on its own once you do. + + + - Waiting for the ATM's bootstrap state event. Power on the ATM + Waiting for the ATM's state event. Power on the ATM and confirm it has reached the configured relay; cassette rows will auto-populate on receipt. - +
+
Recent operations
+

+ "Pending" means sent but not yet echoed back by the machine. + Every publish carries the recent window, so a pending + operation keeps being re-offered until it lands — there is + nothing to retry by hand. +

+ + No operations recorded yet. + + + + + + + + + + + + + + Machine confirmed at + + + + Sent; waiting for the machine to echo this id back + + + + + +
+ @@ -1244,53 +1287,82 @@ - + - + -
Publish cassette config to ATM
+
Record cassette operation
- + - This publish will overwrite the ATM's currently-tracked - counts. If the ATM has dispensed cash since your last - refill or count baseline, those decrements will be lost. - Publish only after a physical refill (a known total), not to - "tweak" counts mid-day. v2 reconciliation will replace this - modal with reconciled state display. + Record what you did to the bay. The machine applies it to the + count it already holds, so a dispense that happened while this + dialog was open is kept, not overwritten. + + + + + + + + + + + + + + The bay is now empty. Nothing else to enter. + + + + + -

Sending to ATM:

- - - - - - - - - - - · count - - - - -
+ label="Record + publish" + :disable="!cassetteOpIsComplete" + :loading="cassetteOpDialog.saving" + @click="submitCassetteOp">
From 5f60b3fe31ca14b3e9b0cb014ddc08617d303eb6 Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:58:04 +0200 Subject: [PATCH 8/8] fix(cassettes): break the same-second tie with the machine's counter The ordering gate compares created_at, which NIP-01 defines at one-second granularity. A dispense and the publish that follows it land inside one second routinely, so the report was dropped and the operator kept the pre-dispense count until the next heartbeat five minutes later. The machine bumps a counter on every local change to a bay count and carries it in its state document. m015 stores it per row, and the gate consults it only when the stamps are equal, where created_at carries no information at all. Only on equality, deliberately. A machine whose state.db was replaced restarts its counter at zero while its wall clock keeps moving forward; gating on the counter across different stamps would lock that machine out for good. Equal stamps with no counter on either side stay closed, which costs one heartbeat and risks nothing. --- crud.py | 53 +++++++++++++++++++++++++--------- migrations.py | 19 ++++++++++++ models.py | 3 ++ tests/test_cassette_configs.py | 28 ++++++++++++++++++ 4 files changed, 89 insertions(+), 14 deletions(-) diff --git a/crud.py b/crud.py index 8508e9a..c4bffd5 100644 --- a/crud.py +++ b/crud.py @@ -1491,13 +1491,23 @@ def _as_unix(value) -> float | None: return None -def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: +def _row_field(row, name): + """Read one column from a row the data layer may hand back as either a + mapping or an object, depending on the driver.""" + if isinstance(row, dict): + return row.get(name) + return getattr(row, name, None) + + +def _should_apply_state_event( + oldest_state_at, incoming_created_at, oldest_seq=None, incoming_seq=None +) -> 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: + Three 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: @@ -1510,6 +1520,14 @@ def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: 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. + - `seq` breaks a same-second tie, and only that. NIP-01 stamps at + one-second granularity, so a dispense and the publish that follows it + can share a stamp and the later report would be dropped. The machine + bumps seq on every local count change, so a higher seq at an equal stamp + is strictly newer. It is consulted ONLY on equality: a machine whose + state.db was replaced restarts its counter at zero, and gating on seq + across different stamps would lock that machine out permanently while + its wall clock kept moving forward. """ oldest = _as_unix(oldest_state_at) if oldest is None: @@ -1517,7 +1535,11 @@ def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: incoming = _as_unix(incoming_created_at) if incoming is None: return False - return incoming > oldest + if incoming != oldest: + return incoming > oldest + if incoming_seq is None or oldest_seq is None: + return False + return incoming_seq > oldest_seq async def get_cassette_config(machine_id: str, position: int) -> CassetteConfig | None: @@ -1565,19 +1587,19 @@ async def apply_reported_state( against believed. """ oldest: dict | None = await db.fetchone( - "SELECT state_at FROM spirekeeper.cassette_configs " + "SELECT state_at, state_seq FROM spirekeeper.cassette_configs " "WHERE machine_id = :mid AND state_at IS NOT NULL " - "ORDER BY state_at ASC LIMIT 1", + "ORDER BY state_at ASC, state_seq ASC LIMIT 1", {"mid": machine_id}, ) oldest_state_at = None + oldest_seq = 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_state_event(oldest_state_at, event_created_at): + oldest_state_at = _row_field(oldest, "state_at") + oldest_seq = _row_field(oldest, "state_seq") + if not _should_apply_state_event( + oldest_state_at, event_created_at, oldest_seq, payload.seq + ): return False # Drop bays the machine no longer reports, before writing the rest. A crash @@ -1601,9 +1623,10 @@ async def apply_reported_state( INSERT INTO spirekeeper.cassette_configs (machine_id, position, denomination, count, updated_at, updated_by, state_denomination, state_count, state_at, - state_event_id) + state_event_id, state_seq) VALUES (:mid, :pos, :denom, :count, :now, :by, - :state_denom, :state_count, :state_at, :event_id) + :state_denom, :state_count, :state_at, :event_id, + :state_seq) ON CONFLICT (machine_id, position) DO UPDATE SET denomination = excluded.denomination, count = excluded.count, @@ -1612,7 +1635,8 @@ async def apply_reported_state( state_denomination = excluded.state_denomination, state_count = excluded.state_count, state_at = excluded.state_at, - state_event_id = excluded.state_event_id + state_event_id = excluded.state_event_id, + state_seq = excluded.state_seq """, { "mid": machine_id, @@ -1625,6 +1649,7 @@ async def apply_reported_state( "state_count": row.count, "state_at": event_created_at, "event_id": event_id, + "state_seq": payload.seq, }, ) return True diff --git a/migrations.py b/migrations.py index 8f5ddcf..f96fe53 100644 --- a/migrations.py +++ b/migrations.py @@ -922,3 +922,22 @@ async def m014_add_counts_uncertain_since(db): "ALTER TABLE spirekeeper.dca_machines " "ADD COLUMN counts_uncertain_since TIMESTAMP" ) + + +async def m015_add_cassette_state_seq(db): + """Break the same-second tie in the state-event ordering gate. + + The gate compares `created_at`, which NIP-01 defines at one-second + granularity — so two reports from the same second are indistinguishable to + it, and the later one is dropped. A dispense and the publish that follows + it land inside one second routinely. + + The machine bumps `seq` on every local change to a bay count, whatever + caused it, and carries it in the state document. Stored per row alongside + `state_at` and consulted ONLY when the stamps are equal, so a machine whose + state.db was replaced — seq back to zero, wall clock still moving forward — + is not locked out by its own counter. + """ + await db.execute( + "ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER" + ) diff --git a/models.py b/models.py index 1a0658c..dc6f533 100644 --- a/models.py +++ b/models.py @@ -677,6 +677,9 @@ class CassetteConfig(BaseModel): state_count: int | None state_at: datetime | None state_event_id: str | None + # The machine's own counter, bumped on every local count change. Breaks a + # same-second tie in the ordering gate, where created_at cannot. + state_seq: int | None = None class CassettePayloadRow(BaseModel): diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index b6adf2d..bf597d9 100644 --- a/tests/test_cassette_configs.py +++ b/tests/test_cassette_configs.py @@ -208,3 +208,31 @@ class TestShouldApplyStateEvent: """Fail closed: an unparseable incoming stamp must not overwrite state that is known-good.""" assert _should_apply_state_event(NOW, "nonsense") is False + + def test_seq_breaks_a_same_second_tie(self): + """NIP-01 stamps at one-second granularity, so a dispense and the + publish that follows it share a stamp routinely. Without the counter + the later report is dropped and the operator keeps the pre-dispense + count until the next heartbeat.""" + assert _should_apply_state_event(NOW, NOW, 4, 5) is True + assert _should_apply_state_event(NOW, NOW, 5, 5) is False + assert _should_apply_state_event(NOW, NOW, 5, 4) is False + + def test_seq_is_ignored_when_the_stamps_differ(self): + """A machine whose state.db was replaced restarts its counter at zero + while its wall clock keeps moving forward. Gating on the counter + across different stamps would lock that machine out for good.""" + assert ( + _should_apply_state_event(NOW, NOW + timedelta(seconds=1), 900, 0) is True + ) + assert ( + _should_apply_state_event(NOW, NOW - timedelta(seconds=1), 0, 900) is False + ) + + def test_same_second_without_a_counter_stays_closed(self): + """An older machine sends no seq at all. Equal stamps then carry no + evidence the report is newer, and fail-closed is the safe read: the + machine republishes on its heartbeat a second later.""" + assert _should_apply_state_event(NOW, NOW, None, None) is False + assert _should_apply_state_event(NOW, NOW, None, 5) is False + assert _should_apply_state_event(NOW, NOW, 5, None) is False