diff --git a/cassette_transport.py b/cassette_transport.py index e6517a1..bc84523 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 = [ @@ -30,7 +36,7 @@ v2 = continuous reverse channel for reconciliation): 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 @@ -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, @@ -75,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 @@ -152,35 +163,41 @@ def build_state_d_tags_for_machines(machines: list[Machine]) -> list[str]: # ============================================================================= -async def publish_to_atm( +async def publish_ops_to_atm( machine: Machine, - payload: PublishCassettesPayload, + ops: list[CassetteOp], operator_user_id: str, ) -> dict: - """Build, encrypt, sign, and publish a kind-30078 cassette config event - from the operator to the target ATM. + """Publish the operator's recent cassette OPERATIONS 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. + 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 config (machine={machine.id}, " - f"positions={sorted(payload.positions.keys())})" + f"cassette ops (machine={machine.id}, ops={[o.op_type for o in ops]})" ), ) return signed # ============================================================================= -# Consume — ATM → operator (the bootstrap consumer task) +# Consume — ATM → operator (the machine's state reports) # ============================================================================= diff --git a/crud.py b/crud.py index a622d96..c4bffd5 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, @@ -34,7 +36,6 @@ from .models import ( UpdateDepositStatusData, UpdateMachineData, UpdateSuperConfigData, - UpsertCassetteConfigData, UpsertDcaLpData, ) @@ -256,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", @@ -1444,10 +1464,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: @@ -1470,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: @@ -1489,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: @@ -1496,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: @@ -1519,38 +1562,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, @@ -1576,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 @@ -1612,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, @@ -1623,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, @@ -1636,6 +1649,139 @@ async def apply_reported_state( "state_count": row.count, "state_at": event_created_at, "event_id": event_id, + "state_seq": payload.seq, }, ) 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/migrations.py b/migrations.py index 6ede0eb..f96fe53 100644 --- a/migrations.py +++ b/migrations.py @@ -860,3 +860,84 @@ 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)" + ) + + +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" + ) + + +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 5a4e529..dc6f533 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. @@ -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 @@ -673,31 +677,9 @@ class CassetteConfig(BaseModel): state_count: int | None state_at: datetime | None 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 + # 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): @@ -721,22 +703,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): @@ -768,6 +769,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/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/tasks.py b/tasks.py index 61695e1..6e5ac55 100644 --- a/tasks.py +++ b/tasks.py @@ -338,6 +338,8 @@ 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, + set_machine_counts_uncertain, ) machines = await list_all_active_machines() @@ -373,6 +375,8 @@ 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, + set_machine_counts_uncertain, ) except Exception as exc: logger.warning( @@ -383,10 +387,63 @@ 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 _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 @@ -487,12 +544,17 @@ 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. 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/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">
diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index 242aa13..bf597d9 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 # ============================================================================= @@ -247,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 diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py new file mode 100644 index 0000000..150d493 --- /dev/null +++ b/tests/test_cassette_ops.py @@ -0,0 +1,356 @@ +""" +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 .. import cassette_transport, tasks +from ..crud import _should_ack_op +from ..models import ( + CASSETTE_OP_TYPES, + CassetteOp, + CreateCassetteOpData, + Machine, + PublishCassetteOpsPayload, + PublishCassettesPayload, +) + +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": [], + } + + +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 + + +# ============================================================================= +# 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}" + + +# ============================================================================= +# _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"] + + +# ============================================================================= +# _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)] 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