Publish cassette operations instead of counts #46
10 changed files with 1317 additions and 376 deletions
|
|
@ -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 `<m>` 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)
|
||||
# =============================================================================
|
||||
|
||||
|
||||
|
|
|
|||
248
crud.py
248
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
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
)
|
||||
|
|
|
|||
235
models.py
235
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_pubkey_hex>`)
|
||||
- ATM → operator (d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`)
|
||||
"""The decrypted content of the ATM → operator state document
|
||||
(d-tag `bitspire-cassettes-state:<atm_pubkey_hex>`).
|
||||
|
||||
Wire shape: `{"positions": {"<pos_str>": {"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": {"<pos_str>": {"denomination", "count"}}}`
|
||||
plus the optional fields below. JSON object keys are always strings; the
|
||||
validator coerces back to int on parse.
|
||||
|
||||
No denomination-unique constraint: multiple same-denom cassettes are
|
||||
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": "<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)
|
||||
# =============================================================================
|
||||
|
|
|
|||
|
|
@ -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: { "<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)
|
||||
}
|
||||
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
|
||||
}
|
||||
},
|
||||
|
||||
|
|
|
|||
68
tasks.py
68
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)
|
||||
|
|
|
|||
|
|
@ -1151,23 +1151,22 @@
|
|||
<div class="col">
|
||||
<h6 class="q-my-none">Cassettes</h6>
|
||||
<p class="text-caption q-my-none" :style="{opacity: 0.7}">
|
||||
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 <i>did</i> 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).
|
||||
</p>
|
||||
</div>
|
||||
<div class="col-auto">
|
||||
<q-btn flat dense icon="undo" label="Revert"
|
||||
:disable="!machineDetail.cassettesDirty"
|
||||
@click="revertCassetteEdits">
|
||||
<q-tooltip>Discard unsaved edits</q-tooltip>
|
||||
<q-btn flat dense icon="refresh" label="Refresh"
|
||||
:loading="machineDetail.cassettesLoading"
|
||||
@click="loadMachineCassettes">
|
||||
<q-tooltip>Re-read the machine's latest report</q-tooltip>
|
||||
</q-btn>
|
||||
<q-btn color="primary" icon="cloud_upload"
|
||||
label="Publish to ATM"
|
||||
:disable="!machineDetail.cassettesDirty"
|
||||
:loading="machineDetail.cassettesPublishing"
|
||||
@click="openCassettePublishConfirm"></q-btn>
|
||||
<q-btn color="primary" icon="add" label="Record operation"
|
||||
:disable="!machineDetail.cassettes.length"
|
||||
@click="openCassetteOpDialog()"></q-btn>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
|
|
@ -1179,64 +1178,108 @@
|
|||
<span v-text="machineDetail.cassettesError"></span>
|
||||
</q-banner>
|
||||
|
||||
<q-banner v-if="!machineDetail.cassetteEdits.length
|
||||
<q-banner v-if="machineDetail.machine
|
||||
&& 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"
|
||||
class="bg-blue-1 text-grey-9">
|
||||
<template v-slot:avatar>
|
||||
<q-icon name="hourglass_empty" color="blue"></q-icon>
|
||||
</template>
|
||||
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.
|
||||
</q-banner>
|
||||
|
||||
<q-table v-if="machineDetail.cassetteEdits.length"
|
||||
<q-table v-if="machineDetail.cassettes.length"
|
||||
dense flat
|
||||
:rows="machineDetail.cassetteEdits"
|
||||
:rows="machineDetail.cassettes"
|
||||
row-key="position"
|
||||
:columns="cassettesTable.columns"
|
||||
:pagination="cassettesTable.pagination"
|
||||
hide-pagination>
|
||||
<template v-slot:body="props">
|
||||
<q-tr :props="props"
|
||||
:style="props.row._dirty
|
||||
? {boxShadow: 'inset 4px 0 0 0 #fdd835'}
|
||||
: {}">
|
||||
<q-tr :props="props">
|
||||
<q-td key="position" class="text-right">
|
||||
<b v-text="'Bay ' + props.row.position"></b>
|
||||
</q-td>
|
||||
<q-td key="denomination" class="text-right">
|
||||
<q-input v-model.number="props.row.denomination"
|
||||
type="number" min="1" step="1" dense outlined
|
||||
:suffix="machineDetail.machine.fiat_code || ''"
|
||||
:style="{width: '140px', display: 'inline-block'}"
|
||||
@update:model-value="markCassetteDirty(props.row)"></q-input>
|
||||
<span v-text="props.row.denomination"></span>
|
||||
<span :style="{opacity: 0.6}"
|
||||
v-text="' ' + (machineDetail.machine.fiat_code || '')"></span>
|
||||
</q-td>
|
||||
<q-td key="count" class="text-right">
|
||||
<q-input v-model.number="props.row.count" type="number"
|
||||
min="0" step="1" dense outlined
|
||||
:style="{width: '120px', display: 'inline-block'}"
|
||||
@update:model-value="markCassetteDirty(props.row)"></q-input>
|
||||
<b v-text="props.row.count"></b>
|
||||
<span :style="{opacity: 0.6}"> notes</span>
|
||||
</q-td>
|
||||
<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">
|
||||
<q-td key="state_at">
|
||||
<span :style="{fontSize: '0.85em', opacity: 0.7}"
|
||||
v-text="formatTime(props.row.updated_at)"></span>
|
||||
v-text="props.row.state_at
|
||||
? 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-tr>
|
||||
</template>
|
||||
</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-panels>
|
||||
</q-card-section>
|
||||
|
|
@ -1244,53 +1287,82 @@
|
|||
</q-dialog>
|
||||
|
||||
<!-- =============================================================== -->
|
||||
<!-- CASSETTE PUBLISH CONFIRM DIALOG -->
|
||||
<!-- RECORD CASSETTE OPERATION DIALOG -->
|
||||
<!-- =============================================================== -->
|
||||
<q-dialog v-model="cassettePublishConfirm.show" persistent>
|
||||
<q-dialog v-model="cassetteOpDialog.show" persistent>
|
||||
<q-card :style="{width: '480px', maxWidth: '95vw'}">
|
||||
<q-card-section class="row items-center q-pb-none">
|
||||
<div class="text-h6">Publish cassette config to ATM</div>
|
||||
<div class="text-h6">Record cassette operation</div>
|
||||
<q-space ></q-space>
|
||||
<q-btn icon="close" flat round dense v-close-popup></q-btn>
|
||||
</q-card-section>
|
||||
<q-card-section>
|
||||
<q-banner class="bg-orange-1 text-grey-9 q-mb-md">
|
||||
<q-banner class="bg-blue-1 text-grey-9 q-mb-md">
|
||||
<template v-slot:avatar>
|
||||
<q-icon name="warning" color="warning"></q-icon>
|
||||
<q-icon name="info" color="blue"></q-icon>
|
||||
</template>
|
||||
<b>This publish will overwrite the ATM's currently-tracked
|
||||
counts.</b> 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.
|
||||
</q-banner>
|
||||
|
||||
<q-select v-model.number="cassetteOpDialog.position"
|
||||
: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>
|
||||
<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-actions align="right">
|
||||
<q-btn flat label="Cancel" v-close-popup></q-btn>
|
||||
<q-btn color="primary"
|
||||
label="Publish to ATM"
|
||||
:loading="machineDetail.cassettesPublishing"
|
||||
@click="submitCassettePublish"></q-btn>
|
||||
label="Record + publish"
|
||||
:disable="!cassetteOpIsComplete"
|
||||
:loading="cassetteOpDialog.saving"
|
||||
@click="submitCassetteOp"></q-btn>
|
||||
</q-card-actions>
|
||||
</q-card>
|
||||
</q-dialog>
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
356
tests/test_cassette_ops.py
Normal file
356
tests/test_cassette_ops.py
Normal file
|
|
@ -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)]
|
||||
158
views_api.py
158
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:<atm_pubkey_hex> and
|
||||
p=<atm_pubkey_hex>.
|
||||
) -> list[CassetteOp]:
|
||||
"""Recent cassette operations for a machine, newest first.
|
||||
|
||||
The `<m>` 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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue