Compare commits

..

No commits in common. "d09c51f277c60d6cf8c157c2a088420a426f15e3" and "30870ebd16f4ef47c529c4f6705f287a942dd3dd" have entirely different histories.

10 changed files with 376 additions and 1317 deletions

View file

@ -17,14 +17,8 @@ publishes position-keyed cassette config to a target ATM via:
The ATM-side consumer (lamassu-next#56) subscribes by the d-tag + its own The ATM-side consumer (lamassu-next#56) subscribes by the d-tag + its own
npub, decrypts, validates, applies, hot-reloads HAL. npub, decrypts, validates, applies, hot-reloads HAL.
The operator → ATM direction carries OPERATIONS as of v2 (bitspire ADR-004): Reverse direction (ATM → operator, v1 = one-shot bootstrap on first boot,
refill, empty, recount, set_denomination, each with an id the machine dedups v2 = continuous reverse channel for reconciliation):
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 kind = 30078
tags = [ tags = [
@ -36,7 +30,7 @@ startup, after every change to its bays, and on a heartbeat):
This module owns the wire-format side of both directions. The consumer This module owns the wire-format side of both directions. The consumer
task (tasks.py) calls `decrypt_and_parse_state_event` per incoming event; task (tasks.py) calls `decrypt_and_parse_state_event` per incoming event;
the API endpoint (views_api.py) calls `publish_ops_to_atm` per operation. the API endpoint (views_api.py) calls `publish_to_atm` per operator submit.
The `<m>` placeholder semantics (load-bearing per the 2026-05-30T11:50Z The `<m>` placeholder semantics (load-bearing per the 2026-05-30T11:50Z
coord-log entry): always the ATM's hex pubkey, NEVER spirekeeper's coord-log entry): always the ATM's hex pubkey, NEVER spirekeeper's
@ -58,12 +52,7 @@ from lnbits.core.signers.base import (
) )
from lnbits.utils.nostr import normalize_public_key from lnbits.utils.nostr import normalize_public_key
from .models import ( from .models import Machine, PublishCassettesPayload
CassetteOp,
Machine,
PublishCassetteOpsPayload,
PublishCassettesPayload,
)
from .nip44 import Nip44Error from .nip44 import Nip44Error
from .nostr_publish import ( from .nostr_publish import (
NostrPublishError, NostrPublishError,
@ -86,7 +75,7 @@ __all__ = [
"RelayUnavailable", "RelayUnavailable",
"build_state_d_tags_for_machines", "build_state_d_tags_for_machines",
"decrypt_and_parse_state_event", "decrypt_and_parse_state_event",
"publish_ops_to_atm", "publish_to_atm",
] ]
_D_TAG_CONFIG_PREFIX = "bitspire-cassettes:" # operator → ATM _D_TAG_CONFIG_PREFIX = "bitspire-cassettes:" # operator → ATM
@ -163,41 +152,35 @@ def build_state_d_tags_for_machines(machines: list[Machine]) -> list[str]:
# ============================================================================= # =============================================================================
async def publish_ops_to_atm( async def publish_to_atm(
machine: Machine, machine: Machine,
ops: list[CassetteOp], payload: PublishCassettesPayload,
operator_user_id: str, operator_user_id: str,
) -> dict: ) -> dict:
"""Publish the operator's recent cassette OPERATIONS to the target ATM. """Build, encrypt, sign, and publish a kind-30078 cassette config event
from the operator to the target ATM.
The v2 wire (bitspire ADR-004). Replaces sending absolute counts, which Returns the signed event dict on success (caller may log event.id for
let a dashboard form loaded before a dispense silently discard that audit). Raises NostrPublishError subclasses (re-exported here as
dispense — the operator and the machine were both writing the same value CassetteTransportError, OperatorIdentityMissing, SignerUnavailable,
over a transport that never tells a writer it lost. RelayUnavailable) on hard failures.
`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) atm_pubkey_hex = _atm_hex_pubkey(machine)
payload = PublishCassetteOpsPayload(ops=ops)
signed = await publish_encrypted_kind_30078( signed = await publish_encrypted_kind_30078(
operator_user_id=operator_user_id, operator_user_id=operator_user_id,
recipient_pubkey_hex=atm_pubkey_hex, recipient_pubkey_hex=atm_pubkey_hex,
d_tag=_config_d_tag(atm_pubkey_hex), d_tag=_config_d_tag(atm_pubkey_hex),
payload=payload.to_wire_dict(), payload=payload.to_wire_dict(),
log_context=( log_context=(
f"cassette ops (machine={machine.id}, ops={[o.op_type for o in ops]})" f"cassette config (machine={machine.id}, "
f"positions={sorted(payload.positions.keys())})"
), ),
) )
return signed return signed
# ============================================================================= # =============================================================================
# Consume — ATM → operator (the machine's state reports) # Consume — ATM → operator (the bootstrap consumer task)
# ============================================================================= # =============================================================================

246
crud.py
View file

@ -12,11 +12,9 @@ from lnbits.helpers import urlsafe_short_hash
from .models import ( from .models import (
CassetteConfig, CassetteConfig,
CassetteOp,
ClientBalanceSummary, ClientBalanceSummary,
CommissionSplit, CommissionSplit,
CommissionSplitLeg, CommissionSplitLeg,
CreateCassetteOpData,
CreateDcaClientData, CreateDcaClientData,
CreateDcaPaymentData, CreateDcaPaymentData,
CreateDcaSettlementData, CreateDcaSettlementData,
@ -36,6 +34,7 @@ from .models import (
UpdateDepositStatusData, UpdateDepositStatusData,
UpdateMachineData, UpdateMachineData,
UpdateSuperConfigData, UpdateSuperConfigData,
UpsertCassetteConfigData,
UpsertDcaLpData, UpsertDcaLpData,
) )
@ -257,25 +256,6 @@ async def set_machine_unpaired(machine_id: str) -> Machine | None:
return await get_machine(machine_id) 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: async def delete_machine(machine_id: str) -> None:
await db.execute( await db.execute(
"DELETE FROM spirekeeper.dca_machines WHERE id = :id", "DELETE FROM spirekeeper.dca_machines WHERE id = :id",
@ -1464,11 +1444,10 @@ async def upsert_fleet_snapshot(
# Row lifecycle per #29: # Row lifecycle per #29:
# - First population for a (machine_id, position) pair → apply_reported_state # - First population for a (machine_id, position) pair → apply_reported_state
# (consumer reading the ATM's one-shot bitspire-cassettes-state event) # (consumer reading the ATM's one-shot bitspire-cassettes-state event)
# - The operator does NOT write these rows. It records operations # - Operator edit of denomination or count → update_cassette_config
# (cassette_ops) and the machine keeps the running count; these columns # (refuses to create new rows; the slot count is hardware-determined)
# hold what the machine last reported. # - Row creation/deletion for a new position → admin only, via ATM
# - Rows appear and disappear only as the machine reports its bay set — # re-provisioning + new bootstrap event (not exposed in v1 here)
# the slot count is hardware-determined.
def _as_unix(value) -> float | None: def _as_unix(value) -> float | None:
@ -1491,23 +1470,13 @@ def _as_unix(value) -> float | None:
return None return None
def _row_field(row, name): def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool:
"""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. """Ordering gate for apply_reported_state.
Applies only when the incoming event is strictly newer than the OLDEST Applies only when the incoming event is strictly newer than the OLDEST
state stamp on file for the machine. state stamp on file for the machine.
Three deliberate choices: Two deliberate choices:
- Compare created_at, not event ids. The old gate asked only whether the - 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: incoming id differed from one stored row, which is a one-event memory:
@ -1520,14 +1489,6 @@ def _should_apply_state_event(
some rows advanced and some not. Gating on the oldest means a partial 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 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. 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) oldest = _as_unix(oldest_state_at)
if oldest is None: if oldest is None:
@ -1535,11 +1496,7 @@ def _should_apply_state_event(
incoming = _as_unix(incoming_created_at) incoming = _as_unix(incoming_created_at)
if incoming is None: if incoming is None:
return False return False
if incoming != oldest:
return 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: async def get_cassette_config(machine_id: str, position: int) -> CassetteConfig | None:
@ -1562,6 +1519,38 @@ 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( async def apply_reported_state(
machine_id: str, machine_id: str,
event_id: str, event_id: str,
@ -1587,19 +1576,19 @@ async def apply_reported_state(
against believed. against believed.
""" """
oldest: dict | None = await db.fetchone( oldest: dict | None = await db.fetchone(
"SELECT state_at, state_seq FROM spirekeeper.cassette_configs " "SELECT state_at FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid AND state_at IS NOT NULL " "WHERE machine_id = :mid AND state_at IS NOT NULL "
"ORDER BY state_at ASC, state_seq ASC LIMIT 1", "ORDER BY state_at ASC LIMIT 1",
{"mid": machine_id}, {"mid": machine_id},
) )
oldest_state_at = None oldest_state_at = None
oldest_seq = None
if oldest is not None: if oldest is not None:
oldest_state_at = _row_field(oldest, "state_at") oldest_state_at = (
oldest_seq = _row_field(oldest, "state_seq") oldest.get("state_at")
if not _should_apply_state_event( if isinstance(oldest, dict)
oldest_state_at, event_created_at, oldest_seq, payload.seq else getattr(oldest, "state_at", None)
): )
if not _should_apply_state_event(oldest_state_at, event_created_at):
return False return False
# Drop bays the machine no longer reports, before writing the rest. A crash # Drop bays the machine no longer reports, before writing the rest. A crash
@ -1623,10 +1612,9 @@ async def apply_reported_state(
INSERT INTO spirekeeper.cassette_configs INSERT INTO spirekeeper.cassette_configs
(machine_id, position, denomination, count, updated_at, (machine_id, position, denomination, count, updated_at,
updated_by, state_denomination, state_count, state_at, updated_by, state_denomination, state_count, state_at,
state_event_id, state_seq) state_event_id)
VALUES (:mid, :pos, :denom, :count, :now, :by, 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 ON CONFLICT (machine_id, position) DO UPDATE SET
denomination = excluded.denomination, denomination = excluded.denomination,
count = excluded.count, count = excluded.count,
@ -1635,8 +1623,7 @@ async def apply_reported_state(
state_denomination = excluded.state_denomination, state_denomination = excluded.state_denomination,
state_count = excluded.state_count, state_count = excluded.state_count,
state_at = excluded.state_at, 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, "mid": machine_id,
@ -1649,139 +1636,6 @@ async def apply_reported_state(
"state_count": row.count, "state_count": row.count,
"state_at": event_created_at, "state_at": event_created_at,
"event_id": event_id, "event_id": event_id,
"state_seq": payload.seq,
}, },
) )
return True 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

View file

@ -860,84 +860,3 @@ async def m012_add_max_cash_in_sats(db):
await db.execute( await db.execute(
"ALTER TABLE spirekeeper.super_config ADD COLUMN max_cash_in_sats INTEGER" "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"
)

235
models.py
View file

@ -7,7 +7,7 @@
from datetime import datetime from datetime import datetime
from pydantic import BaseModel, root_validator, validator from pydantic import BaseModel, validator
# ============================================================================= # =============================================================================
# Machines — one row per bitSpire ATM, owned by exactly one operator. # Machines — one row per bitSpire ATM, owned by exactly one operator.
@ -63,10 +63,6 @@ class Machine(BaseModel):
# NIP-46 bunker pairing (S0 / #9). NULL until the spire is first paired. # NIP-46 bunker pairing (S0 / #9). NULL until the spire is first paired.
bunker_spire_key_name: str | None = None bunker_spire_key_name: str | None = None
paired_at: datetime | 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 created_at: datetime
updated_at: datetime updated_at: datetime
@ -677,9 +673,31 @@ class CassetteConfig(BaseModel):
state_count: int | None state_count: int | None
state_at: datetime | None state_at: datetime | None
state_event_id: str | 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 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): class CassettePayloadRow(BaseModel):
@ -703,41 +721,22 @@ class CassettePayloadRow(BaseModel):
class PublishCassettesPayload(BaseModel): class PublishCassettesPayload(BaseModel):
"""The decrypted content of the ATM → operator state document """The decrypted JSON content of a kind-30078 cassette event, both
(d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`). directions:
- operator → ATM (d-tag `bitspire-cassettes:<atm_pubkey_hex>`)
- ATM → operator (d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`)
It carried the operator → ATM direction too until v2 moved that to Wire shape: `{"positions": {"<pos_str>": {"denomination", "count"}}}`.
PublishCassetteOpsPayload. This is now the machine reporting what it JSON object keys are always strings; the validator coerces back to
holds, and the machine is the only writer of those counts. int on parse. The position key set MUST match what the receiver
already has (slot count is hardware-fixed; no add/remove from this
Wire shape: `{"positions": {"<pos_str>": {"denomination", "count"}}}` payload).
plus the optional fields below. JSON object keys are always strings; the
validator coerces back to int on parse.
No denomination-unique constraint: multiple same-denom cassettes are No denomination-unique constraint: multiple same-denom cassettes are
operationally valid (cash-out throughput on a popular denom). operationally valid (cash-out throughput on a popular denom).
The optional fields are all absent on older machines, so every one of them
defaults to a value meaning "this machine does not report that yet" rather
than to a value that would be wrong:
- `applied_ops`: operation ids the machine has applied. This is the
acknowledgement, and the only one an addressable event can carry — a
relay returns OK for an event it then discards, so the publisher is
never told anything. An empty list reads as "nothing acknowledged",
which is correct for a machine that has not yet applied any.
- `seq`: the machine's own monotonic counter, bumped on every local count
change. Regression detection independent of created_at, which is only
second-granular and can be forced by a bad clock.
- `counts_uncertain_since`: set when a dispense ended without the
dispenser reporting what it moved, so the counts above are the
machine's best guess rather than a measurement.
""" """
positions: dict[int, CassettePayloadRow] positions: dict[int, CassettePayloadRow]
applied_ops: list[str] = []
seq: int | None = None
counts_uncertain_since: int | None = None
@validator("positions", pre=True) @validator("positions", pre=True)
def coerce_string_keys_to_int(cls, v): def coerce_string_keys_to_int(cls, v):
@ -769,170 +768,6 @@ 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": "<uuid>", "at": 1790106060, "type": "refill",
# "position": 2, "bills": 100},
# {"id": "<uuid>", "at": 1790106061, "type": "empty", "position": 3},
# {"id": "<uuid>", "at": 1790106062, "type": "recount",
# "position": 1, "count": 37},
# {"id": "<uuid>", "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) # Fee-config Nostr payload — operator → ATM (aiolabs/satmachineadmin#39)
# ============================================================================= # =============================================================================

View file

@ -212,14 +212,15 @@ window.app = Vue.createApp({
loading: false, loading: false,
machine: null, machine: null,
settlements: [], settlements: [],
// Cassettes sub-tab state (v2, bitspire ADR-004) — see // Cassettes sub-tab state (#29 v1) — see openCassettePublishConfirm /
// openCassetteOpDialog / submitCassetteOp + the cassettes panel in // submitCassettePublish methods + the cassettes panel in
// templates/spirekeeper/index.html. Read-only by design: the // templates/spirekeeper/index.html.
// operator records operations, the machine owns the counts.
activeTab: 'settlements', activeTab: 'settlements',
cassettes: [], // machine-reported rows, not editable cassetteEdits: [], // editable working copy of cassette_configs rows
cassetteOps: [], // recent operations, newest first cassettesPristine: [], // last-known-clean snapshot for revert
cassettesLoading: false, cassettesLoading: false,
cassettesPublishing: false,
cassettesDirty: false,
cassettesError: null cassettesError: null
}, },
cassettesTable: { cassettesTable: {
@ -227,27 +228,14 @@ window.app = Vue.createApp({
{name: 'position', label: 'Bay', field: 'position', align: 'right'}, {name: 'position', label: 'Bay', field: 'position', align: 'right'},
{name: 'denomination', label: 'Denomination', field: 'denomination', align: 'right'}, {name: 'denomination', label: 'Denomination', field: 'denomination', align: 'right'},
{name: 'count', label: 'Count', field: 'count', align: 'right'}, {name: 'count', label: 'Count', field: 'count', align: 'right'},
{name: 'state_at', label: 'Machine reported', field: 'state_at', align: 'left'}, {name: 'state', label: 'ATM-reported', field: 'state_denomination', align: 'right'},
{name: 'actions', label: '', field: 'position', align: 'right'} {name: 'updated_at', label: 'Updated', field: 'updated_at', align: 'left'}
], ],
pagination: {rowsPerPage: 0} // hide pagination — cassette count is small pagination: {rowsPerPage: 0} // hide pagination — cassette count is small
}, },
cassetteOpDialog: { cassettePublishConfirm: {
show: false, 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: { partialDispenseDialog: {
show: false, show: false,
saving: false, saving: false,
@ -298,28 +286,6 @@ window.app = Vue.createApp({
}, },
computed: { 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() { superAnyFee() {
// Banner styling key — true when either directional super fee is // Banner styling key — true when either directional super fee is
// non-zero, so the banner reads as "active platform fee" instead // non-zero, so the banner reads as "active platform fee" instead
@ -977,8 +943,9 @@ window.app = Vue.createApp({
async viewMachine(machine) { async viewMachine(machine) {
this.machineDetail.machine = machine this.machineDetail.machine = machine
this.machineDetail.settlements = [] this.machineDetail.settlements = []
this.machineDetail.cassettes = [] this.machineDetail.cassetteEdits = []
this.machineDetail.cassetteOps = [] this.machineDetail.cassettesPristine = []
this.machineDetail.cassettesDirty = false
this.machineDetail.cassettesError = null this.machineDetail.cassettesError = null
this.machineDetail.activeTab = 'settlements' this.machineDetail.activeTab = 'settlements'
this.machineDetail.show = true this.machineDetail.show = true
@ -1005,25 +972,21 @@ window.app = Vue.createApp({
}, },
// ----------------------------------------------------------------- // -----------------------------------------------------------------
// Cassette inventory + operations (v2, bitspire ADR-004) // Cassette inventory (#29 v1)
// ----------------------------------------------------------------- // -----------------------------------------------------------------
async loadMachineCassettes() { async loadMachineCassettes() {
if (!this.machineDetail.machine) return if (!this.machineDetail.machine) return
this.machineDetail.cassettesLoading = true this.machineDetail.cassettesLoading = true
this.machineDetail.cassettesError = null this.machineDetail.cassettesError = null
const base = `${MACHINES_PATH}/${this.machineDetail.machine.id}`
try { try {
// The machine row is re-read too: counts_uncertain_since lives on it const {data} = await LNbits.api.request(
// and is set by the consumer as state events land, so the row the 'GET',
// machines table handed us goes stale while this dialog is open. `${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes`
const [machine, bays, ops] = await Promise.all([ )
LNbits.api.request('GET', base), const rows = (data || []).map(row => ({...row, _dirty: false}))
LNbits.api.request('GET', `${base}/cassettes`), this.machineDetail.cassetteEdits = rows
LNbits.api.request('GET', `${base}/cassettes/ops`) this.machineDetail.cassettesPristine = JSON.parse(JSON.stringify(rows))
]) this.machineDetail.cassettesDirty = false
if (machine.data) this.machineDetail.machine = machine.data
this.machineDetail.cassettes = bays.data || []
this.machineDetail.cassetteOps = ops.data || []
} catch (e) { } catch (e) {
this._notifyError(e, 'Failed to load cassettes') this._notifyError(e, 'Failed to load cassettes')
} finally { } finally {
@ -1031,91 +994,73 @@ window.app = Vue.createApp({
} }
}, },
cassetteOpIcon(opType) { markCassetteDirty(row) {
return ( // Find pristine match by position (the row identity) and compare;
{ // flip _dirty + overall dirty flag accordingly. Editable fields
refill: 'add_circle_outline', // are denomination + count; position is the immutable row key.
empty: 'remove_circle_outline', const pristine = this.machineDetail.cassettesPristine.find(
recount: 'fact_check', p => p.position === row.position
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)
}, },
cassetteOpSummary(op) { revertCassetteEdits() {
const fiat = (this.machineDetail.machine || {}).fiat_code || '' this.machineDetail.cassetteEdits = JSON.parse(
const bay = `Bay ${op.position}` JSON.stringify(this.machineDetail.cassettesPristine)
if (op.op_type === 'refill') return `${bay} — added ${op.bills} notes` )
if (op.op_type === 'empty') return `${bay} — emptied` this.machineDetail.cassettesDirty = false
if (op.op_type === 'recount') return `${bay} — recounted to ${op.count}` this.machineDetail.cassettesError = null
if (op.op_type === 'set_denomination') { },
return `${bay} — denomination set to ${op.denomination} ${fiat}`.trim()
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: { "<pos>": { 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)
} }
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 const payload = {positions}
d.error = null this.machineDetail.cassettesPublishing = true
try { try {
await LNbits.api.request( const {data} = await LNbits.api.request(
'POST', 'POST',
`${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes/ops`, `${MACHINES_PATH}/${this.machineDetail.machine.id}/cassettes/publish`,
null, null,
payload payload
) )
d.show = false 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
Quasar.Notify.create({ Quasar.Notify.create({
type: 'positive', type: 'positive',
message: 'Operation recorded and published to the ATM' message: 'Cassette config published to ATM'
}) })
await this.loadMachineCassettes()
} catch (e) { } catch (e) {
const detail = const detail =
(e && e.response && e.response.data && e.response.data.detail) || (e && e.response && e.response.data && e.response.data.detail) ||
'Could not record the operation' 'Publish failed'
// A 503 means the op IS recorded and will ride out with the next this.machineDetail.cassettesError = detail
// publish, so reload either way — the list should show it pending. this._notifyError(e, 'Publish failed')
d.error = detail
this._notifyError(e, 'Operation failed')
await this.loadMachineCassettes()
} finally { } finally {
d.saving = false this.machineDetail.cassettesPublishing = false
} }
}, },

View file

@ -338,8 +338,6 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
apply_reported_state, apply_reported_state,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
list_all_active_machines, list_all_active_machines,
mark_cassette_ops_acked,
set_machine_counts_uncertain,
) )
machines = await list_all_active_machines() machines = await list_all_active_machines()
@ -375,8 +373,6 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_reported_state, apply_reported_state,
mark_cassette_ops_acked,
set_machine_counts_uncertain,
) )
except Exception as exc: except Exception as exc:
logger.warning( logger.warning(
@ -387,63 +383,10 @@ async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
return filter_key return filter_key
async def _record_op_acknowledgements(
machine_id: str, payload, mark_cassette_ops_acked
) -> None:
"""Mark the operations a machine reports as applied.
Deliberately not gated on whether the state event advanced the counts. The
machine echoes its applied-op ids on EVERY state publish, so an event
carrying nothing new about the counts can still be the first one to tell us
an operation landed; gating on that would lose the acknowledgement.
This echo is the only acknowledgement the transport can carry. An
addressable event gives its publisher no failure signal at all — the relay
returns OK for an event it then discards — so without it the dashboard
could only ever show an operation as sent, never as delivered.
"""
if not payload.applied_ops:
return
newly_acked = await mark_cassette_ops_acked(machine_id, payload.applied_ops)
if newly_acked:
logger.info(
f"spirekeeper: machine {machine_id} acknowledged "
f"{newly_acked} cassette operation(s)"
)
async def _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( async def _handle_cassette_state_event(
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_reported_state, apply_reported_state,
mark_cassette_ops_acked,
set_machine_counts_uncertain,
) -> None: ) -> None:
"""Verify signature, resolve the operator's signer, decrypt via the """Verify signature, resolve the operator's signer, decrypt via the
signer abstraction (bunker round-trip for RemoteBunkerSigner; direct signer abstraction (bunker round-trip for RemoteBunkerSigner; direct
@ -544,17 +487,12 @@ async def _handle_cassette_state_event(
) )
if applied: if applied:
logger.info( logger.info(
f"spirekeeper: applied reported state event {event_id[:12]}... " f"spirekeeper: applied bootstrap state event {event_id[:12]}... "
f"to machine {machine.id} ({len(payload.positions)} cassettes)" f"to machine {machine.id} ({len(payload.positions)} cassettes)"
) )
else: else:
# Replay or an older event. Normal on relay reconnect. # Replay: event_id already on file. Normal on relay reconnect.
logger.debug( logger.debug(
f"spirekeeper: cassette state event {event_id[:12]}... " f"spirekeeper: cassette state event {event_id[:12]}... "
f"not newer than stored state for machine {machine.id} (no-op)" f"already applied to machine {machine.id} (replay 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)

View file

@ -1151,22 +1151,23 @@
<div class="col"> <div class="col">
<h6 class="q-my-none">Cassettes</h6> <h6 class="q-my-none">Cassettes</h6>
<p class="text-caption q-my-none" :style="{opacity: 0.7}"> <p class="text-caption q-my-none" :style="{opacity: 0.7}">
The machine owns these counts — it holds the notes. You Per-cassette count and physical bay position. Denomination
record what you <i>did</i> to a bay and it keeps the set is hardware-determined (re-provision via atm-tui to
running total. Bay count and denomination set are change). "Publish to ATM" encrypts + signs + sends the new
hardware-determined (re-provision via atm-tui to change config to the machine via Nostr.
the bays themselves).
</p> </p>
</div> </div>
<div class="col-auto"> <div class="col-auto">
<q-btn flat dense icon="refresh" label="Refresh" <q-btn flat dense icon="undo" label="Revert"
:loading="machineDetail.cassettesLoading" :disable="!machineDetail.cassettesDirty"
@click="loadMachineCassettes"> @click="revertCassetteEdits">
<q-tooltip>Re-read the machine's latest report</q-tooltip> <q-tooltip>Discard unsaved edits</q-tooltip>
</q-btn> </q-btn>
<q-btn color="primary" icon="add" label="Record operation" <q-btn color="primary" icon="cloud_upload"
:disable="!machineDetail.cassettes.length" label="Publish to ATM"
@click="openCassetteOpDialog()"></q-btn> :disable="!machineDetail.cassettesDirty"
:loading="machineDetail.cassettesPublishing"
@click="openCassettePublishConfirm"></q-btn>
</div> </div>
</div> </div>
@ -1178,108 +1179,64 @@
<span v-text="machineDetail.cassettesError"></span> <span v-text="machineDetail.cassettesError"></span>
</q-banner> </q-banner>
<q-banner v-if="machineDetail.machine <q-banner v-if="!machineDetail.cassetteEdits.length
&& machineDetail.machine.counts_uncertain_since"
class="bg-orange-1 text-grey-9 q-mb-md">
<template v-slot:avatar>
<q-icon name="help_outline" color="warning"></q-icon>
</template>
<b>The machine can't vouch for these counts.</b>
A dispense ended without a reliable count of what left the bay
(<span v-text="formatTime(machineDetail.machine.counts_uncertain_since)"></span>).
Open the bay, count it, and record a <b>Recount</b> — the
machine clears this on its own once you do.
</q-banner>
<q-banner v-if="!machineDetail.cassettes.length
&& !machineDetail.cassettesLoading" && !machineDetail.cassettesLoading"
class="bg-blue-1 text-grey-9"> class="bg-blue-1 text-grey-9">
<template v-slot:avatar> <template v-slot:avatar>
<q-icon name="hourglass_empty" color="blue"></q-icon> <q-icon name="hourglass_empty" color="blue"></q-icon>
</template> </template>
Waiting for the ATM's state event. Power on the ATM Waiting for the ATM's bootstrap state event. Power on the ATM
and confirm it has reached the configured relay; cassette and confirm it has reached the configured relay; cassette
rows will auto-populate on receipt. rows will auto-populate on receipt.
</q-banner> </q-banner>
<q-table v-if="machineDetail.cassettes.length" <q-table v-if="machineDetail.cassetteEdits.length"
dense flat dense flat
:rows="machineDetail.cassettes" :rows="machineDetail.cassetteEdits"
row-key="position" row-key="position"
:columns="cassettesTable.columns" :columns="cassettesTable.columns"
:pagination="cassettesTable.pagination" :pagination="cassettesTable.pagination"
hide-pagination> hide-pagination>
<template v-slot:body="props"> <template v-slot:body="props">
<q-tr :props="props"> <q-tr :props="props"
:style="props.row._dirty
? {boxShadow: 'inset 4px 0 0 0 #fdd835'}
: {}">
<q-td key="position" class="text-right"> <q-td key="position" class="text-right">
<b v-text="'Bay ' + props.row.position"></b> <b v-text="'Bay ' + props.row.position"></b>
</q-td> </q-td>
<q-td key="denomination" class="text-right"> <q-td key="denomination" class="text-right">
<span v-text="props.row.denomination"></span> <q-input v-model.number="props.row.denomination"
<span :style="{opacity: 0.6}" type="number" min="1" step="1" dense outlined
v-text="' ' + (machineDetail.machine.fiat_code || '')"></span> :suffix="machineDetail.machine.fiat_code || ''"
:style="{width: '140px', display: 'inline-block'}"
@update:model-value="markCassetteDirty(props.row)"></q-input>
</q-td> </q-td>
<q-td key="count" class="text-right"> <q-td key="count" class="text-right">
<b v-text="props.row.count"></b> <q-input v-model.number="props.row.count" type="number"
<span :style="{opacity: 0.6}"> notes</span> min="0" step="1" dense outlined
:style="{width: '120px', display: 'inline-block'}"
@update:model-value="markCassetteDirty(props.row)"></q-input>
</q-td> </q-td>
<q-td key="state_at"> <q-td key="state" class="text-right">
<span v-if="props.row.state_denomination !== null"
:style="{fontSize: '0.85em', opacity: 0.7}">
<span v-text="props.row.state_denomination"></span>
<span :style="{opacity: 0.6}"
v-text="' ' + (machineDetail.machine.fiat_code || '')"></span>
<span :style="{opacity: 0.6}"> · </span>
<span v-text="'×' + props.row.state_count"></span>
</span>
<span v-else :style="{opacity: 0.4}">—</span>
</q-td>
<q-td key="updated_at">
<span :style="{fontSize: '0.85em', opacity: 0.7}" <span :style="{fontSize: '0.85em', opacity: 0.7}"
v-text="props.row.state_at v-text="formatTime(props.row.updated_at)"></span>
? formatTime(props.row.state_at)
: 'never'"></span>
</q-td>
<q-td key="actions" class="text-right">
<q-btn flat dense size="sm" icon="edit_note"
@click="openCassetteOpDialog(props.row.position)">
<q-tooltip>Record an operation on this bay</q-tooltip>
</q-btn>
</q-td> </q-td>
</q-tr> </q-tr>
</template> </template>
</q-table> </q-table>
<div v-if="machineDetail.cassettes.length" class="q-mt-lg">
<div class="text-subtitle2 q-mb-xs">Recent operations</div>
<p class="text-caption q-mt-none q-mb-sm" :style="{opacity: 0.7}">
"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.
</p>
<q-banner v-if="!machineDetail.cassetteOps.length"
class="bg-grey-3 text-grey-9">
No operations recorded yet.
</q-banner>
<q-list v-else dense bordered separator>
<q-item v-for="op in machineDetail.cassetteOps" :key="op.id">
<q-item-section avatar>
<q-icon :name="cassetteOpIcon(op.op_type)"
:color="op.acked_at ? 'positive' : 'grey'"></q-icon>
</q-item-section>
<q-item-section>
<q-item-label v-text="cassetteOpSummary(op)"></q-item-label>
<q-item-label caption
v-text="formatTime(op.created_at)"></q-item-label>
</q-item-section>
<q-item-section side>
<q-chip dense size="sm"
:color="op.acked_at ? 'green-1' : 'orange-1'"
text-color="grey-9"
:label="op.acked_at ? 'Applied' : 'Pending'">
<q-tooltip v-if="op.acked_at">
Machine confirmed at
<span v-text="formatTime(op.acked_at)"></span>
</q-tooltip>
<q-tooltip v-else>
Sent; waiting for the machine to echo this id back
</q-tooltip>
</q-chip>
</q-item-section>
</q-item>
</q-list>
</div>
</q-tab-panel> </q-tab-panel>
</q-tab-panels> </q-tab-panels>
</q-card-section> </q-card-section>
@ -1287,82 +1244,53 @@
</q-dialog> </q-dialog>
<!-- =============================================================== --> <!-- =============================================================== -->
<!-- RECORD CASSETTE OPERATION DIALOG --> <!-- CASSETTE PUBLISH CONFIRM DIALOG -->
<!-- =============================================================== --> <!-- =============================================================== -->
<q-dialog v-model="cassetteOpDialog.show" persistent> <q-dialog v-model="cassettePublishConfirm.show" persistent>
<q-card :style="{width: '480px', maxWidth: '95vw'}"> <q-card :style="{width: '480px', maxWidth: '95vw'}">
<q-card-section class="row items-center q-pb-none"> <q-card-section class="row items-center q-pb-none">
<div class="text-h6">Record cassette operation</div> <div class="text-h6">Publish cassette config to ATM</div>
<q-space ></q-space> <q-space ></q-space>
<q-btn icon="close" flat round dense v-close-popup></q-btn> <q-btn icon="close" flat round dense v-close-popup></q-btn>
</q-card-section> </q-card-section>
<q-card-section> <q-card-section>
<q-banner class="bg-blue-1 text-grey-9 q-mb-md"> <q-banner class="bg-orange-1 text-grey-9 q-mb-md">
<template v-slot:avatar> <template v-slot:avatar>
<q-icon name="info" color="blue"></q-icon> <q-icon name="warning" color="warning"></q-icon>
</template> </template>
Record what you did to the bay. The machine applies it to the <b>This publish will overwrite the ATM's currently-tracked
count it already holds, so a dispense that happened while this counts.</b> If the ATM has dispensed cash since your last
dialog was open is kept, not overwritten. refill or count baseline, those decrements will be lost.
</q-banner> Publish only after a physical refill (a known total), not to
"tweak" counts mid-day. v2 reconciliation will replace this
<q-select v-model.number="cassetteOpDialog.position" modal with reconciled state display.
:options="cassetteBayOptions"
emit-value map-options
label="Bay" dense outlined
class="q-mb-md"></q-select>
<q-select v-model="cassetteOpDialog.op_type"
:options="cassetteOpTypeOptions"
emit-value map-options
label="Operation" dense outlined
class="q-mb-md"
@update:model-value="resetCassetteOpValue"></q-select>
<q-input v-if="cassetteOpDialog.op_type === 'refill'"
v-model.number="cassetteOpDialog.bills"
type="number" min="1" step="1" dense outlined
label="Notes added"
hint="How many notes you put IN — a delta, not a total."
></q-input>
<q-input v-if="cassetteOpDialog.op_type === 'recount'"
v-model.number="cassetteOpDialog.count"
type="number" min="0" step="1" dense outlined
label="Notes counted"
hint="What you physically counted in the bay, right now."
></q-input>
<q-input v-if="cassetteOpDialog.op_type === 'set_denomination'"
v-model.number="cassetteOpDialog.denomination"
type="number" min="1" step="1" dense outlined
label="Denomination"
:suffix="machineDetail.machine
? (machineDetail.machine.fiat_code || '') : ''"
hint="What is now loaded in that bay. The machine can't
know this — only you can."
></q-input>
<q-banner v-if="cassetteOpDialog.op_type === 'empty'"
class="bg-grey-3 text-grey-9">
The bay is now empty. Nothing else to enter.
</q-banner>
<q-banner v-if="cassetteOpDialog.error"
class="bg-red-1 text-grey-9 q-mt-md">
<template v-slot:avatar>
<q-icon name="warning" color="negative"></q-icon>
</template>
<span v-text="cassetteOpDialog.error"></span>
</q-banner> </q-banner>
<p class="q-mb-sm">Sending to ATM:</p>
<q-list dense bordered>
<q-item v-for="row in machineDetail.cassetteEdits"
:key="row.position">
<q-item-section>
<q-item-label>
<b v-text="'Bay ' + row.position"></b>
</q-item-label>
</q-item-section>
<q-item-section side>
<q-item-label caption>
<b v-text="row.denomination + ' ' +
(machineDetail.machine.fiat_code || '')"></b>
· count
<b v-text="row.count"></b>
</q-item-label>
</q-item-section>
</q-item>
</q-list>
</q-card-section> </q-card-section>
<q-card-actions align="right"> <q-card-actions align="right">
<q-btn flat label="Cancel" v-close-popup></q-btn> <q-btn flat label="Cancel" v-close-popup></q-btn>
<q-btn color="primary" <q-btn color="primary"
label="Record + publish" label="Publish to ATM"
:disable="!cassetteOpIsComplete" :loading="machineDetail.cassettesPublishing"
:loading="cassetteOpDialog.saving" @click="submitCassettePublish"></q-btn>
@click="submitCassetteOp"></q-btn>
</q-card-actions> </q-card-actions>
</q-card> </q-card>
</q-dialog> </q-dialog>

View file

@ -2,9 +2,9 @@
Tests for the v1.1 cassette-config layer (aiolabs/satmachineadmin#29). Tests for the v1.1 cassette-config layer (aiolabs/satmachineadmin#29).
Covers the pure pieces that don't need a live DB: Covers the pure pieces that don't need a live DB:
- Pydantic validator behaviour on PublishCassettesPayload + the row model - Pydantic validator behaviour on PublishCassettesPayload + the row /
(position key coercion, integer ranges, multiple-same-denomination upsert models (position key coercion, integer ranges, multiple-same-
payloads, wire-format round-trip) denomination payloads, wire-format round-trip)
- _should_apply_state_event ordering gate (extracted from - _should_apply_state_event ordering gate (extracted from
apply_reported_state so the decision is testable without a database apply_reported_state so the decision is testable without a database
round-trip) round-trip)
@ -29,6 +29,7 @@ from ..crud import _as_unix, _should_apply_state_event
from ..models import ( from ..models import (
CassettePayloadRow, CassettePayloadRow,
PublishCassettesPayload, PublishCassettesPayload,
UpsertCassetteConfigData,
) )
# ============================================================================= # =============================================================================
@ -148,6 +149,44 @@ class TestCassettePayloadRow:
CassettePayloadRow(denomination=20, count=-1) 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 # _should_apply_state_event — ordering gate
# ============================================================================= # =============================================================================
@ -208,31 +247,3 @@ class TestShouldApplyStateEvent:
"""Fail closed: an unparseable incoming stamp must not overwrite """Fail closed: an unparseable incoming stamp must not overwrite
state that is known-good.""" state that is known-good."""
assert _should_apply_state_event(NOW, "nonsense") is False 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

View file

@ -1,356 +0,0 @@
"""
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)]

View file

@ -26,7 +26,7 @@ from .cassette_transport import (
OperatorIdentityMissing, OperatorIdentityMissing,
RelayUnavailable, RelayUnavailable,
SignerUnavailable, SignerUnavailable,
publish_ops_to_atm, publish_to_atm,
) )
from .fee_transport import publish_fee_config from .fee_transport import publish_fee_config
from .pairing import ( from .pairing import (
@ -40,7 +40,6 @@ from .pairing import (
from .crud import ( from .crud import (
append_settlement_note, append_settlement_note,
count_completed_legs_for_settlement, count_completed_legs_for_settlement,
create_cassette_op,
create_dca_client, create_dca_client,
create_deposit, create_deposit,
create_machine, create_machine,
@ -48,7 +47,6 @@ from .crud import (
delete_deposit, delete_deposit,
delete_machine, delete_machine,
force_reset_stuck_settlement, force_reset_stuck_settlement,
get_cassette_ops_window,
get_client_balance_summary, get_client_balance_summary,
get_commission_splits, get_commission_splits,
get_dca_client, get_dca_client,
@ -68,12 +66,12 @@ from .crud import (
get_super_config, get_super_config,
list_all_active_machines, list_all_active_machines,
list_cassette_configs_for_machine, list_cassette_configs_for_machine,
list_cassette_ops,
lp_is_onboarded, lp_is_onboarded,
replace_commission_splits, replace_commission_splits,
reset_settlement_for_retry, reset_settlement_for_retry,
set_machine_pairing, set_machine_pairing,
set_machine_unpaired, set_machine_unpaired,
update_cassette_config,
update_dca_client, update_dca_client,
update_deposit, update_deposit,
update_deposit_status, update_deposit_status,
@ -88,10 +86,8 @@ from .distribution import (
from .models import ( from .models import (
AppendSettlementNoteData, AppendSettlementNoteData,
CassetteConfig, CassetteConfig,
CassetteOp,
ClientBalanceSummary, ClientBalanceSummary,
CommissionSplit, CommissionSplit,
CreateCassetteOpData,
CreateDcaClientData, CreateDcaClientData,
CreateDepositData, CreateDepositData,
CreateMachineData, CreateMachineData,
@ -102,6 +98,7 @@ from .models import (
Machine, Machine,
PairMachineData, PairMachineData,
PartialDispenseData, PartialDispenseData,
PublishCassettesPayload,
SetCommissionSplitsData, SetCommissionSplitsData,
SettleBalanceData, SettleBalanceData,
StuckSettlementsResponse, StuckSettlementsResponse,
@ -111,6 +108,7 @@ from .models import (
UpdateDepositStatusData, UpdateDepositStatusData,
UpdateMachineData, UpdateMachineData,
UpdateSuperConfigData, UpdateSuperConfigData,
UpsertCassetteConfigData,
) )
spirekeeper_api_router = APIRouter() spirekeeper_api_router = APIRouter()
@ -1104,21 +1102,16 @@ async def api_update_super_config(
# ============================================================================= # =============================================================================
# Cassettes — per-machine ATM inventory (bitspire ADR-004) # Cassette configs (#29 v1.1) — per-machine ATM cassette inventory
# ============================================================================= # =============================================================================
# GET /machines/{id}/cassettes — bays as the machine last reported them # v1.1 surface, paired with aiolabs/lamassu-next#56 ATM-side. Two endpoints:
# GET /machines/{id}/cassettes/ops — recent operations, with their acks # GET /machines/{id}/cassettes — list rows for the operator UI
# POST /machines/{id}/cassettes/ops — record one operation and publish # POST /machines/{id}/cassettes/publish — apply edits + publish kind-30078
# #
# The operator does not write counts. It records what it DID to a bay and the # Row creation (new (machine_id, position) pairs) is admin-only via the
# machine, which holds the notes, keeps the running total. The rows behind the # bootstrap consumer task — slot count is hardware-determined. Operator-
# first endpoint are the machine's report, not an operator draft. # side flow is edit-and-publish over the existing rows only; the editable
# # fields per row are denomination and count.
# 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( @spirekeeper_api_router.get(
@ -1136,96 +1129,105 @@ async def api_list_machine_cassettes(
return await list_cassette_configs_for_machine(machine_id) return await list_cassette_configs_for_machine(machine_id)
@spirekeeper_api_router.get(
"/api/v1/dca/machines/{machine_id}/cassettes/ops",
response_model=list[CassetteOp],
)
async def api_list_machine_cassette_ops(
machine_id: str,
user: User = Depends(check_user_exists),
) -> list[CassetteOp]:
"""Recent cassette operations for a machine, newest first.
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)
@spirekeeper_api_router.post( @spirekeeper_api_router.post(
"/api/v1/dca/machines/{machine_id}/cassettes/ops", "/api/v1/dca/machines/{machine_id}/cassettes/publish",
response_model=CassetteOp, response_model=list[CassetteConfig],
) )
async def api_create_machine_cassette_op( async def api_publish_machine_cassettes(
machine_id: str, machine_id: str,
data: CreateCassetteOpData, payload: PublishCassettesPayload,
user: User = Depends(check_user_exists), user: User = Depends(check_user_exists),
) -> CassetteOp: ) -> list[CassetteConfig]:
"""Record one cassette operation and publish the machine's recent window. """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:<atm_pubkey_hex> and
p=<atm_pubkey_hex>.
This replaces publishing absolute counts. The operator now records what it The `<m>` placeholder in the published d-tag is the ATM's hex pubkey
DID — a refill, an empty, a recount, a denomination change — and the from machine.machine_npub (canonicalised via normalize_public_key),
machine keeps the running total. Both sides used to write the same value NOT the internal dca_machines.id UUID — see #29 'machine_id semantics'
over a transport that never tells a writer it lost, so a dashboard form section and coord-log 2026-05-30T11:50Z load-bearing nudge.
loaded before a dispense silently discarded that dispense when published.
The op is recorded BEFORE the publish and is deliberately not rolled back Returns the fresh cassette_configs rows after the upserts so the UI
if the publish fails. It represents something that physically happened — can refresh its table from one round-trip.
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: Errors:
400 — machine not paired, so there is no ATM identity to publish to 400 — payload position set doesn't match the machine's stored set
400 — position is not a bay this machine has reported (operator publishing for a slot that doesn't exist on the
503 — signer offline, or relay/nostrclient unreachable. The operation is ATM; or the bootstrap hasn't landed yet so no rows exist)
still recorded and will be delivered with the next publish. 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
""" """
machine = await _machine_owned_by(machine_id, user.id) machine = await _machine_owned_by(machine_id, user.id)
if not machine.machine_npub: 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( raise HTTPException(
HTTPStatus.BAD_REQUEST, HTTPStatus.BAD_REQUEST,
"machine is not paired — pair it before recording cassette operations", "machine is not paired — pair it before publishing cassette config",
) )
existing = await list_cassette_configs_for_machine(machine_id) 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: if not existing:
raise HTTPException( raise HTTPException(
HTTPStatus.BAD_REQUEST, HTTPStatus.BAD_REQUEST,
( (
"No cassette rows for this machine yet — waiting for its state " "No cassette_configs rows exist for this machine yet — "
"event. Power on the ATM and confirm it has reached the " "waiting for the ATM's bootstrap state event. Power on the "
"configured relay; spirekeeper populates the bays on receipt." "ATM and confirm it has reached the configured relay; "
"spirekeeper will auto-populate cassette_configs on "
"receipt."
), ),
) )
known_positions = {row.position for row in existing} if existing_positions != incoming_positions:
if data.position not in known_positions: missing = existing_positions - incoming_positions
extra = incoming_positions - existing_positions
raise HTTPException( raise HTTPException(
HTTPStatus.BAD_REQUEST, HTTPStatus.BAD_REQUEST,
( (
f"position {data.position} is not a bay this machine has " "Payload position set doesn't match the machine's stored "
f"reported (has: {sorted(known_positions)}). Bay count is " f"set. Missing from payload: {sorted(missing)}; extra in "
"hardware-determined; re-provision via atm-tui to change it." f"payload: {sorted(extra)}. Slot count is hardware-fixed "
"— re-provision the ATM via atm-tui to add/remove physical "
"bays, then re-publish."
), ),
) )
op = await create_cassette_op(machine_id, data, created_by=user.id) # 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",
)
window = await get_cassette_ops_window(machine_id)
try: try:
await publish_ops_to_atm(machine, window, user.id) await publish_to_atm(machine, payload, user.id)
except OperatorIdentityMissing as exc: except OperatorIdentityMissing as exc:
raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc
except (SignerUnavailable, RelayUnavailable) as exc: except SignerUnavailable as exc:
raise HTTPException( raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc
HTTPStatus.SERVICE_UNAVAILABLE, except RelayUnavailable as exc:
f"{exc} — the operation was recorded and will be delivered with " raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc
"the next publish",
) from exc
except CassetteTransportError as exc: except CassetteTransportError as exc:
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
return op return await list_cassette_configs_for_machine(machine_id)