From b8a5e6352af0f596a2459db1ea074242e8a58b28 Mon Sep 17 00:00:00 2001 From: Padreug Date: Sat, 10 Oct 2026 21:51:51 +0200 Subject: [PATCH 1/4] feat(schema): dispense outcome on settlements, dispense_reports, cash-out hold mirror (ADR-005) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit m016: an append-only `dispense_reports` table (lamassu-server's cash_out_actions shape — one row per report the machine sent, so a retry, a late report and a remediation report stay distinct); dispense_confirmed / dispense_error / dispense_error_code / dispense_raw_code / dispense_error_class / dispense_reported_at / dispensed_fiat_cents on dca_settlements; cash_out_held_since / _reason / _code on dca_machines beside counts_uncertain_since. Settlement lifecycle gains awaiting_dispense (cash_out at insert — paid, waiting for the machine's report), partial_pending (some notes out, value short; held whole until the operator records the resolution) and cash_owed (nothing out; legs never run). dispense_unreported is derived by the worklist, not stored. resume_cash_out joins CASSETTE_OP_TYPES as a machine-wide op: position 0, no position on the wire, no bay fields. It rides the operator channel the machine already consumes; the machine honours it only if stamped after the hold began. A recount releases the hold too. crud: get_settlement_by_txid (the join the machine's extra.txid already provides), apply_dispense_outcome (copies the report onto the settlement and finally writes bills_json / cassettes_json with what actually came out), the dispense_reports accessors incl. adopting a report that arrived before its payment, set_machine_cash_out_hold, and the three new worklist buckets. --- crud.py | 250 +++++++++++++++++++++++++++++++++++++++++++++++++- migrations.py | 79 ++++++++++++++++ models.py | 200 +++++++++++++++++++++++++++++++++++++--- 3 files changed, 511 insertions(+), 18 deletions(-) diff --git a/crud.py b/crud.py index c4bffd5..1858cb6 100644 --- a/crud.py +++ b/crud.py @@ -27,6 +27,8 @@ from .models import ( DcaLpPreferences, DcaPayment, DcaSettlement, + DispenseReport, + DispenseReportIn, Machine, PublishCassettesPayload, SuperConfig, @@ -257,6 +259,32 @@ async def set_machine_unpaired(machine_id: str) -> Machine | None: return await get_machine(machine_id) +async def set_machine_cash_out_hold( + machine_id: str, + since: datetime | None, + reason: str | None, + code: str | None, +) -> None: + """Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row. + + Written on every state event, including when it is None: the machine + clearing the hold — after an operator recount or resume_cash_out — is as + important as it setting one. `updated_at` is left alone for the same + reason as counts_uncertain_since: this is the machine reporting about + itself on a heartbeat, not an operator editing the machine. + """ + await db.execute( + """ + UPDATE spirekeeper.dca_machines + SET cash_out_held_since = :since, + cash_out_held_reason = :reason, + cash_out_held_code = :code + WHERE id = :id + """, + {"id": machine_id, "since": since, "reason": reason, "code": code}, + ) + + 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. @@ -711,6 +739,171 @@ async def create_settlement_idempotent( return await get_settlement(settlement_id) +async def get_settlement_by_txid( + machine_id: str, bitspire_txid: str +) -> DcaSettlement | None: + """The settlement a machine report refers to. `bitspire_txid` comes from + the invoice's extra.txid, stamped by the machine at create_invoice time, + so it is the natural join for a report that names the same txid.""" + return await db.fetchone( + "SELECT * FROM spirekeeper.dca_settlements " + "WHERE machine_id = :mid AND bitspire_txid = :txid", + {"mid": machine_id, "txid": bitspire_txid}, + DcaSettlement, + ) + + +async def apply_dispense_outcome( + settlement_id: str, + report: DispenseReportIn, + new_status: str, + reported_at: datetime, +) -> DcaSettlement | None: + """Copy the machine's report onto the settlement and move it (ADR-005 §1). + + Fills the never-before-written bills_json / cassettes_json with what + actually came out, not what was provisioned. `error_message` is left to + the distribution path; the dispense error lives in its own columns. + """ + import json as _json + + await db.execute( + """ + UPDATE spirekeeper.dca_settlements + SET status = :status, + dispense_confirmed = :confirmed, + dispense_error = :error, + dispense_error_code = :error_code, + dispense_raw_code = :raw_code, + dispense_error_class = :error_class, + dispense_reported_at = :reported_at, + dispensed_fiat_cents = :dispensed_fiat_cents, + bills_json = :bills_json, + cassettes_json = :cassettes_json, + processing_claim = NULL + WHERE id = :id + """, + { + "id": settlement_id, + "status": new_status, + "confirmed": report.dispense_confirmed, + "error": report.error, + "error_code": report.error_code, + "raw_code": report.raw_code, + "error_class": report.error_class, + "reported_at": reported_at, + "dispensed_fiat_cents": report.dispensed_fiat_cents, + "bills_json": _json.dumps([b.dict() for b in report.bills]), + "cassettes_json": _json.dumps([c.dict() for c in report.cassettes]), + }, + ) + return await get_settlement(settlement_id) + + +# --------------------------------------------------------------------------- +# Dispense reports (ADR-005 §2) — append-only +# --------------------------------------------------------------------------- + + +async def get_dispense_report( + machine_id: str, txid: str, reported_at: int +) -> DispenseReport | None: + """A report is identified by (machine, txid, at): the machine resends the + same report until acked, and a byte-identical resend must not grow the + log. A remediation report for the same txid carries a later `at`.""" + return await db.fetchone( + "SELECT * FROM spirekeeper.dispense_reports " + "WHERE machine_id = :mid AND txid = :txid AND reported_at = :at", + {"mid": machine_id, "txid": txid, "at": datetime.fromtimestamp(reported_at)}, + DispenseReport, + ) + + +async def insert_dispense_report( + machine_id: str, settlement_id: str | None, report: DispenseReportIn +) -> DispenseReport: + import json as _json + + report_id = urlsafe_short_hash() + await db.execute( + """ + INSERT INTO spirekeeper.dispense_reports + (id, machine_id, settlement_id, txid, payment_hash, dispense_confirmed, + error, error_code, raw_code, error_class, fiat_cents, currency, + bills_json, cassettes_json, counts_uncertain, remediates_txid, + reported_at, received_at) + VALUES (:id, :machine_id, :settlement_id, :txid, :payment_hash, + :confirmed, :error, :error_code, :raw_code, :error_class, + :fiat_cents, :currency, :bills_json, :cassettes_json, + :counts_uncertain, :remediates_txid, :reported_at, :received_at) + """, + { + "id": report_id, + "machine_id": machine_id, + "settlement_id": settlement_id, + "txid": report.txid, + "payment_hash": report.payment_hash, + "confirmed": report.dispense_confirmed, + "error": report.error, + "error_code": report.error_code, + "raw_code": report.raw_code, + "error_class": report.error_class, + "fiat_cents": report.fiat_cents, + "currency": report.currency, + "bills_json": _json.dumps([b.dict() for b in report.bills]), + "cassettes_json": _json.dumps([c.dict() for c in report.cassettes]), + "counts_uncertain": report.counts_uncertain, + "remediates_txid": report.remediates_txid, + "reported_at": datetime.fromtimestamp(report.at), + "received_at": datetime.now(), + }, + ) + row = await db.fetchone( + "SELECT * FROM spirekeeper.dispense_reports WHERE id = :id", + {"id": report_id}, + DispenseReport, + ) + assert row is not None, "Newly inserted dispense report couldn't be retrieved" + return row + + +async def link_dispense_reports_to_settlement( + machine_id: str, txid: str, settlement_id: str +) -> int: + """A report can arrive before its payment lands (hold invoices settle + after the dispense; the invoice listener can lag). When the settlement is + finally inserted, adopt the orphan rows.""" + result = await db.execute( + "UPDATE spirekeeper.dispense_reports SET settlement_id = :sid " + "WHERE machine_id = :mid AND txid = :txid AND settlement_id IS NULL", + {"sid": settlement_id, "mid": machine_id, "txid": txid}, + ) + return getattr(result, "rowcount", 0) or 0 + + +async def get_latest_unlinked_dispense_report( + machine_id: str, txid: str +) -> DispenseReport | None: + return await db.fetchone( + "SELECT * FROM spirekeeper.dispense_reports " + "WHERE machine_id = :mid AND txid = :txid AND settlement_id IS NULL " + "ORDER BY reported_at DESC LIMIT 1", + {"mid": machine_id, "txid": txid}, + DispenseReport, + ) + + +async def get_dispense_reports_for_settlement( + settlement_id: str, +) -> list[DispenseReport]: + return await db.fetchall( + "SELECT * FROM spirekeeper.dispense_reports WHERE settlement_id = :sid " + "ORDER BY reported_at ASC", + {"sid": settlement_id}, + DispenseReport, + ) + + async def get_settlement(settlement_id: str) -> DcaSettlement | None: return await db.fetchone( "SELECT * FROM spirekeeper.dca_settlements WHERE id = :id", @@ -752,7 +945,14 @@ async def get_stuck_settlements_for_operator( ) -> dict: """Operator worklist of settlements that didn't process cleanly. - Returns a dict with four keyed lists: + Returns a dict with seven keyed lists. The first three are ADR-005 §6 — + the only ones whose meaning is "a customer is owed money": + - 'cash_owed': the machine reported nothing dispensed; legs never ran. + - 'partial_pending': some notes out, value short; held whole until the + operator records the resolution. + - 'dispense_unreported': awaiting_dispense older than the threshold — + the machine never reported (crashed, offline, or an old build). + Then the original four: - 'rejected': any status='rejected' (Nostr attribution cross-check failed — signer didn't match the machine identity). Distinct from 'errored' because retry is wrong: the row was misrouted, @@ -817,7 +1017,48 @@ async def get_stuck_settlements_for_operator( {"uid": operator_user_id, "threshold": threshold_at}, DcaSettlement, ) + # ADR-005 §6 — the owed-cash buckets. cash_owed / partial_pending are + # stored statuses; dispense_unreported is derived: a cash-out that landed + # and never heard from its machine within the threshold. + cash_owed = await db.fetchall( + """ + SELECT s.* + FROM spirekeeper.dca_settlements s + JOIN spirekeeper.dca_machines m ON m.id = s.machine_id + WHERE m.operator_user_id = :uid AND s.status = 'cash_owed' + ORDER BY s.created_at DESC + """, + {"uid": operator_user_id}, + DcaSettlement, + ) + partial_pending = await db.fetchall( + """ + SELECT s.* + FROM spirekeeper.dca_settlements s + JOIN spirekeeper.dca_machines m ON m.id = s.machine_id + WHERE m.operator_user_id = :uid AND s.status = 'partial_pending' + ORDER BY s.created_at DESC + """, + {"uid": operator_user_id}, + DcaSettlement, + ) + dispense_unreported = await db.fetchall( + """ + SELECT s.* + FROM spirekeeper.dca_settlements s + JOIN spirekeeper.dca_machines m ON m.id = s.machine_id + WHERE m.operator_user_id = :uid + AND s.status = 'awaiting_dispense' + AND s.created_at < :threshold + ORDER BY s.created_at ASC + """, + {"uid": operator_user_id, "threshold": threshold_at}, + DcaSettlement, + ) return { + "cash_owed": cash_owed, + "partial_pending": partial_pending, + "dispense_unreported": dispense_unreported, "rejected": rejected, "errored": errored, "stuck_pending": stuck_pending, @@ -870,9 +1111,10 @@ async def mark_settlement_status( status: str, error_message: str | None = None, ) -> DcaSettlement | None: - """Status: 'pending' | 'processing' | 'processed' | 'partial' | - 'refunded' | 'errored'. Clears processing_claim on terminal states so a - fresh claim attempt won't see a stale token.""" + """Status: 'awaiting_dispense' | 'pending' | 'processing' | 'processed' | + 'partial_pending' | 'cash_owed' | 'partial' | 'refunded' | 'errored'. + Clears processing_claim on terminal states so a fresh claim attempt won't + see a stale token.""" await db.execute( """ UPDATE spirekeeper.dca_settlements diff --git a/migrations.py b/migrations.py index f96fe53..ed3d6f1 100644 --- a/migrations.py +++ b/migrations.py @@ -941,3 +941,82 @@ async def m015_add_cassette_state_seq(db): await db.execute( "ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER" ) + + +async def m016_dispense_outcome(db): + """The dispense outcome becomes a first-class fact (bitspire ADR-005 §2). + + Until now a cash-out settlement was captured the instant the payment + landed: `_handle_payment` spawned distribution in the same breath, so by + the time a dispenser jammed two seconds later the legs were already paid + and the dashboard honestly reported `processed`. The machine now reports + every cash-out's outcome over a `report_dispense` RPC and the settlement + waits for it (`awaiting_dispense`) before anything moves. + + `dispense_reports` is append-only, one row per report the machine sent — + lamassu-server's `cash_out_actions` shape — so a retry, a late report and + a remediation report are all visible as distinct rows. `settlement_id` is + NULL for a report whose payment this server never saw. + + The settlement carries the three lamassu fields (dispense_confirmed, + error, error_code) plus raw_code / error_class and the fiat value that + actually left the machine, so the partial-dispense dialog can be + pre-filled with the hardware's own number instead of a typed one. + + The machine's cash-out hold is mirrored onto its registry row beside + counts_uncertain_since: a latched machine refuses cash-out until an + operator recounts or publishes `resume_cash_out`, and the dashboard needs + to show that and offer the button. + """ + await db.execute( + f""" + CREATE TABLE IF NOT EXISTS spirekeeper.dispense_reports ( + id TEXT PRIMARY KEY, + machine_id TEXT NOT NULL, + settlement_id TEXT, + txid TEXT NOT NULL, + payment_hash TEXT, + dispense_confirmed BOOLEAN NOT NULL, + error TEXT, + error_code TEXT, + raw_code TEXT, + error_class TEXT, + fiat_cents INTEGER NOT NULL, + currency TEXT NOT NULL, + bills_json TEXT NOT NULL, + cassettes_json TEXT NOT NULL, + counts_uncertain BOOLEAN NOT NULL DEFAULT false, + remediates_txid TEXT, + reported_at TIMESTAMP NOT NULL, + received_at TIMESTAMP NOT NULL DEFAULT {db.timestamp_now} + ); + """ + ) + await db.execute( + "CREATE INDEX IF NOT EXISTS dispense_reports_txid_idx " + "ON dispense_reports (machine_id, txid)" + ) + await db.execute( + "CREATE INDEX IF NOT EXISTS dispense_reports_settlement_idx " + "ON dispense_reports (settlement_id)" + ) + for col, typ in ( + ("dispense_confirmed", "BOOLEAN"), + ("dispense_error", "TEXT"), + ("dispense_error_code", "TEXT"), + ("dispense_raw_code", "TEXT"), + ("dispense_error_class", "TEXT"), + ("dispense_reported_at", "TIMESTAMP"), + ("dispensed_fiat_cents", "INTEGER"), + ): + await db.execute( + f"ALTER TABLE spirekeeper.dca_settlements ADD COLUMN {col} {typ}" + ) + for col, typ in ( + ("cash_out_held_since", "TIMESTAMP"), + ("cash_out_held_reason", "TEXT"), + ("cash_out_held_code", "TEXT"), + ): + await db.execute( + f"ALTER TABLE spirekeeper.dca_machines ADD COLUMN {col} {typ}" + ) diff --git a/models.py b/models.py index dc6f533..f1422fc 100644 --- a/models.py +++ b/models.py @@ -67,6 +67,13 @@ class Machine(BaseModel): # 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 + # ADR-005 §5: the machine has latched cash-out off after a terminal + # dispenser fault. Mirrored from its state document. Cleared when the + # machine reports the hold released (an operator recount or a + # resume_cash_out op). The dashboard shows it and offers the button. + cash_out_held_since: datetime | None = None + cash_out_held_reason: str | None = None + cash_out_held_code: str | None = None created_at: datetime updated_at: datetime @@ -311,19 +318,39 @@ class DcaSettlement(BaseModel): fee_mismatch_sats: int | None = None bills_json: str | None cassettes_json: str | None - # 'pending' (default at insert) + # Lifecycle (bitspire ADR-005 §1 — payment is authorization, the + # machine's dispense report is capture; distribution waits for capture): + # 'awaiting_dispense' (cash_out at insert: paid, waiting for the report) + # 'pending' (cash_in at insert; cash_out once dispense_confirmed) # 'processing' (claim taken by distribution processor) # 'processed' (all legs paid) - # 'partial' (operator marked partial-dispense after the fact) + # 'partial_pending' (report says some notes out, value short — holds + # EVERYTHING until the operator records how the shortfall + # was resolved; then one distribution at the final amount) + # 'cash_owed' (report says nothing out — legs never run, funds stay in + # the machine wallet, the customer is owed; worklist first) + # 'partial' (operator confirmed a partial amount; distributed scaled) # 'refunded' (operator-initiated refund) # 'errored' (operational distribution failure — retry path applies) # 'rejected' (Nostr attribution cross-check failed at land time; # never went near distribution. error_message holds the # reason. Retry is wrong — investigate the machine.) + # 'dispense_unreported' is NOT stored: the worklist derives it from + # awaiting_dispense rows older than its threshold. status: str error_message: str | None processed_at: datetime | None created_at: datetime + # ADR-005 §2 — copied from the machine's report. The three lamassu + # fields plus raw_code / error_class; dispensed_fiat_cents is what the + # hardware says physically left, which pre-fills partial-dispense. + dispense_confirmed: bool | None = None + dispense_error: str | None = None + dispense_error_code: str | None = None + dispense_raw_code: str | None = None + dispense_error_class: str | None = None + dispense_reported_at: datetime | None = None + dispensed_fiat_cents: int | None = None # Append-only audit memo. Populated when an operator triggers an in-place # adjustment (partial-dispense, manual reconciliation override). Each # entry timestamped + records original values so the overwrite is @@ -336,6 +363,110 @@ class DcaSettlement(BaseModel): processing_claim: str | None = None +# ============================================================================= +# Dispense outcome (bitspire ADR-005 §2) — machine → spirekeeper `report_dispense` +# ============================================================================= + + +class DispenseReportBill(BaseModel): + denomination: int + requested: int + dispensed: int + rejected: int + + +class DispenseReportCassette(BaseModel): + position: int + denomination: int + provisioned: int + dispensed: int + rejected: int + + +class DispenseReportIn(BaseModel): + """The RPC body as the machine sends it. Mirrors @bitSpire/lnbits + DispenseReportBody. Field names follow lamassu-server's cash_out_txs / + cash_out_actions (dispense_confirmed, error, error_code). Idempotent on + (txid, at): the machine resends until acked; a byte-identical resend is + acknowledged without a new row.""" + + txid: str + payment_hash: str | None = None + tx_type: str = "cash_out" + dispense_confirmed: bool + error: str | None = None + error_code: str | None = None + raw_code: str | None = None + error_class: str | None = None # 'terminal' | 'recoverable' | 'inventory' + fiat_cents: int + currency: str + bills: list[DispenseReportBill] = [] + cassettes: list[DispenseReportCassette] = [] + counts_uncertain: bool = False + remediates_txid: str | None = None + at: int + + @validator("txid") + def _txid_present(cls, v): + if not v or not v.strip(): + raise ValueError("txid is required") + return v.strip() + + @validator("tx_type") + def _cash_out_only(cls, v): + if v != "cash_out": + raise ValueError("report_dispense is for cash_out transactions only") + return v + + @validator("error_class") + def _known_class(cls, v): + if v is not None and v not in ("terminal", "recoverable", "inventory"): + raise ValueError( + f"error_class must be terminal|recoverable|inventory, got {v!r}" + ) + return v + + @validator("fiat_cents") + def _fiat_non_negative(cls, v): + if v < 0: + raise ValueError("fiat_cents must be >= 0") + return v + + @property + def dispensed_fiat_cents(self) -> int: + """What the hardware says physically left, in cents.""" + return sum(b.denomination * b.dispensed for b in self.bills) * 100 + + @property + def total_dispensed_notes(self) -> int: + return sum(b.dispensed for b in self.bills) + + +class DispenseReport(BaseModel): + """One stored report (append-only — a retry, a late report and a + remediation report are distinct rows). `settlement_id` is NULL when the + payment this report refers to was never seen by this server.""" + + id: str + machine_id: str + settlement_id: str | None + txid: str + payment_hash: str | None + dispense_confirmed: bool + error: str | None + error_code: str | None + raw_code: str | None + error_class: str | None + fiat_cents: int + currency: str + bills_json: str + cassettes_json: str + counts_uncertain: bool = False + remediates_txid: str | None = None + reported_at: datetime + received_at: datetime + + # ============================================================================= # Commission splits — operator-defined remainder allocation per machine. # ============================================================================= @@ -570,6 +701,13 @@ class StuckSettlementsResponse(BaseModel): """ threshold_minutes: int + # ADR-005 §6 — the only buckets whose meaning is "a customer is owed + # money"; rendered first. + cash_owed: list = [] # list[DcaSettlement] + partial_pending: list = [] + # awaiting_dispense older than the threshold: the machine never reported. + # A machine on an old build lands here too — that is the upgrade path. + dispense_unreported: list = [] rejected: list # list[DcaSettlement] errored: list stuck_pending: list @@ -738,6 +876,11 @@ class PublishCassettesPayload(BaseModel): applied_ops: list[str] = [] seq: int | None = None counts_uncertain_since: int | None = None + # ADR-005 §5 — additive like counts_uncertain_since. Present while the + # machine refuses cash-out after a terminal dispenser fault. + cash_out_held_since: int | None = None + cash_out_held_reason: str | None = None + cash_out_held_code: str | None = None @validator("positions", pre=True) def coerce_string_keys_to_int(cls, v): @@ -803,7 +946,21 @@ class PublishCassettesPayload(BaseModel): # landed after the same problem. -CASSETTE_OP_TYPES = ("refill", "empty", "recount", "set_denomination") +CASSETTE_OP_TYPES = ( + "refill", + "empty", + "recount", + "set_denomination", + "resume_cash_out", +) + +# ADR-005 §5: `resume_cash_out` is not a cassette operation — it releases the +# machine's cash-out hold without touching a bay — but it rides the same +# operator event, with the same id/at dedup shape, because that channel is the +# one the machine already consumes. It is machine-wide, so its position is 0 +# and its wire form carries no position at all. A recount also releases the +# hold (same "operator at the open machine" gesture). +MACHINE_WIDE_OP_TYPES = ("resume_cash_out",) class CassetteOp(BaseModel): @@ -837,11 +994,17 @@ class CassetteOp(BaseModel): 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 + @root_validator(skip_on_failure=True) + def _position_matches_scope(cls, values): + pos, typ = values.get("position"), values.get("op_type") + if typ in MACHINE_WIDE_OP_TYPES: + if pos != 0: + raise ValueError( + f"{typ} is machine-wide; position must be 0, got {pos}" + ) + elif pos is None or pos <= 0: + raise ValueError(f"position must be > 0, got {pos}") + return values def to_wire_dict(self) -> dict: """The published form. Drops the fields this op_type does not use, so @@ -850,8 +1013,10 @@ class CassetteOp(BaseModel): "id": self.id, "at": int(self.created_at.timestamp()), "type": self.op_type, - "position": self.position, } + if self.op_type in MACHINE_WIDE_OP_TYPES: + return out + out["position"] = self.position if self.op_type == "refill": out["bills"] = self.bills elif self.op_type == "recount": @@ -880,11 +1045,17 @@ class CreateCassetteOpData(BaseModel): 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 + @root_validator(skip_on_failure=True) + def _position_matches_scope(cls, values): + pos, typ = values.get("position"), values.get("op_type") + if typ in MACHINE_WIDE_OP_TYPES: + if pos != 0: + raise ValueError( + f"{typ} is machine-wide; position must be 0, got {pos}" + ) + elif pos is None or pos <= 0: + raise ValueError(f"position must be > 0, got {pos}") + return values @validator("bills") def _bills_positive(cls, v): @@ -911,6 +1082,7 @@ class CreateCassetteOpData(BaseModel): "recount": "count", "set_denomination": "denomination", "empty": None, + "resume_cash_out": None, }[values.get("op_type")] if required is not None and values.get(required) is None: raise ValueError(f"{values['op_type']} requires `{required}`") From 44c2afa5bb4b14936318198ca4411ab711a6cb79 Mon Sep 17 00:00:00 2001 From: Padreug Date: Sat, 10 Oct 2026 21:51:51 +0200 Subject: [PATCH 2/4] =?UTF-8?q?feat(transport):=20report=5Fdispense=20RPC?= =?UTF-8?q?=20=E2=80=94=20capture=20a=20cash-out=20on=20the=20machine's=20?= =?UTF-8?q?report,=20not=20on=20payment=20(ADR-005=20=C2=A71=E2=80=93?= =?UTF-8?q?=C2=A72)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The structural fix for bitspire#122. _handle_payment used to spawn process_settlement the instant a cash_out payment landed — before the machine had begun to dispense — so a jam two seconds later found the legs already paid and `processed` was the honest answer. Payment is now the authorization; the machine's report is the capture. A cash_out lands as awaiting_dispense and is not distributed. The new `report_dispense` handler (identity from the VERIFIED transport sender, same as create_withdraw / get_machine_config) stores every report append-only and moves the settlement: dispense_confirmed → pending and distribution runs; some notes out → partial_pending, held whole (ADR-005 Decision 1, one distribution when the shortfall is resolved); nothing out → cash_owed, first on the worklist. A report naming remediates_txid moves the owed settlement it names to pending in full. Already-captured settlements are recorded but never moved — a report cannot un-pay legs. A byte-identical resend is acked without a new row. Both orders of arrival are handled: a report that precedes its payment (hold invoices settle after the dispense; the invoice listener can lag) is stored unlinked and adopted when the settlement is inserted, through the same transition. counts_uncertain on a report mirrors onto the machine immediately rather than at the next heartbeat. The state-event consumer mirrors cash_out_held_* onto dca_machines, including clearing it. Soft-fails like the other RPCs: without register_rpc the settlements sit in awaiting_dispense and surface as dispense_unreported — the honest state. --- __init__.py | 6 + dispense_transport.py | 299 ++++++++++++++++++++++++++++++++++++++++++ tasks.py | 65 ++++++++- 3 files changed, 368 insertions(+), 2 deletions(-) create mode 100644 dispense_transport.py diff --git a/__init__.py b/__init__.py index 08ea43a..5bea74e 100644 --- a/__init__.py +++ b/__init__.py @@ -6,6 +6,7 @@ from loguru import logger from .cashin_transport import register_create_withdraw_rpc from .crud import db +from .dispense_transport import register_dispense_report_rpc from .machine_config_transport import register_machine_config_rpc from .nostr_transport_roster import register_with_lnbits as register_roster_with_lnbits from .tasks import wait_for_cassette_state_events, wait_for_paid_invoices @@ -68,6 +69,11 @@ def spirekeeper_start(): # config over the transport, leaving "awaiting configuration" with no # per-machine env provisioning. Soft-fails if register_rpc isn't exposed. register_machine_config_rpc() + # Dispense outcome capture (bitspire ADR-005 §2 / #122): register the + # report_dispense RPC. A cash-out settlement now waits in awaiting_dispense + # until the machine reports; the success report is what distributes it, + # a failure report puts the customer on the owed-cash worklist. + register_dispense_report_rpc() __all__ = [ diff --git a/dispense_transport.py b/dispense_transport.py new file mode 100644 index 0000000..2b13efa --- /dev/null +++ b/dispense_transport.py @@ -0,0 +1,299 @@ +""" +Dispense outcome capture: the `report_dispense` nostr-transport RPC +(bitspire ADR-005 §1-§2, aiolabs/bitspire#122). + +A cash-out used to be captured the instant its payment landed: `_handle_payment` +spawned distribution in the same breath, so when a dispenser jammed two seconds +later the legs were already paid and `processed` was the honest answer. Payment +is now the authorization and the machine's dispense report is the capture. A +`cash_out` settlement lands as `awaiting_dispense` and this handler moves it: + + dispense_confirmed → pending → distribution runs + some notes out, value short → partial_pending (held whole — ADR-005 + Decision 1: one distribution when the + shortfall is resolved) + nothing out → cash_owed (legs never run; the customer + is owed; first on the worklist) + report carries remediates_txid → the owed/partial settlement it names + goes to pending and distributes in full + +Every report is stored append-only in `dispense_reports` (lamassu-server's +`cash_out_actions` shape), including the success ones — the success report is +what captures. Identity is the VERIFIED transport sender, never the body, same +as `create_withdraw` / `get_machine_config`. Idempotent on (txid, at): the +machine resends until it gets an OK, and a byte-identical resend is acked +without a new row or a second transition. + +A report can also arrive BEFORE its payment lands (hold invoices settle after +the dispense; the invoice listener can lag). It is stored with no settlement; +`_handle_payment` adopts it when the settlement is inserted and applies the +same transition. See `apply_report_to_settlement`. +""" + +from __future__ import annotations + +import asyncio +from datetime import datetime + +from loguru import logger + +from .crud import ( + apply_dispense_outcome, + get_dispense_report, + get_machine_by_atm_pubkey_hex, + get_settlement_by_txid, + insert_dispense_report, + link_dispense_reports_to_settlement, + set_machine_counts_uncertain, +) +from .models import DcaSettlement, DispenseReport, DispenseReportIn, Machine + +_RPC_NAME = "report_dispense" + +# Statuses a first report may move. Anything else (processed, pending, +# processing, errored, rejected, partial, refunded) is recorded but not moved — +# a report cannot un-pay legs, and a late report for an already-captured sale +# is information, not an instruction. +_CAPTURABLE = ("awaiting_dispense",) +# Statuses a remediation report may close out. +_OWED = ("cash_owed", "partial_pending") + +# Strong references to in-flight distribution tasks, same reason as tasks.py. +_inflight: set[asyncio.Task] = set() + + +def _outcome_status(report: DispenseReportIn) -> str: + if report.dispense_confirmed: + return "pending" + if report.total_dispensed_notes > 0: + return "partial_pending" + return "cash_owed" + + +def _spawn_distribution(settlement_id: str) -> None: + # Lazy import: distribution imports crud, and crud is what this module + # already depends on; importing at module load would make a cycle. + from .distribution import process_settlement + + task = asyncio.create_task(process_settlement(settlement_id)) + _inflight.add(task) + task.add_done_callback(_inflight.discard) + + +async def apply_report_to_settlement( + settlement: DcaSettlement, + report: DispenseReportIn, + machine: Machine, +) -> str: + """Move `settlement` according to `report`. Returns the resulting status. + + Shared by the RPC handler (report after payment) and `_handle_payment` + (payment after report), so both orders of arrival take the same path. + """ + reported_at = datetime.fromtimestamp(report.at) + + if report.remediates_txid: + # Handled by the caller against the settlement the remediation names; + # for the remediation's OWN txid there is nothing to capture. + return settlement.status + + if settlement.status not in _CAPTURABLE: + logger.info( + f"spirekeeper: report_dispense for settlement {settlement.id} in status " + f"{settlement.status!r} — recorded, not moved (txid={report.txid})" + ) + return settlement.status + + new_status = _outcome_status(report) + updated = await apply_dispense_outcome( + settlement.id, report, new_status, reported_at + ) + status = updated.status if updated else new_status + + if new_status == "pending": + _spawn_distribution(settlement.id) + logger.info( + f"spirekeeper: dispense CONFIRMED for settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}, " + f"{report.fiat_cents / 100:.2f} {report.currency}) — distributing" + ) + elif new_status == "partial_pending": + logger.warning( + f"spirekeeper: PARTIAL dispense for settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}): " + f"{report.dispensed_fiat_cents / 100:.2f} of " + f"{report.fiat_cents / 100:.2f} " + f"{report.currency} left the machine; {report.error_code or 'no error'} " + f"{report.raw_code or ''}. Held until the operator resolves the shortfall." + ) + else: + logger.error( + f"spirekeeper: CASH OWED — settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}): customer paid " + f"{report.fiat_cents / 100:.2f} {report.currency}, nothing dispensed; " + f"{report.error_code or 'no error'} {report.raw_code or ''}: " + f"{report.error or ''}" + ) + return status + + +async def _apply_remediation( + machine: Machine, report: DispenseReportIn +) -> tuple[str | None, str | None]: + """A manual dispense closed out an earlier failed txid. Returns + (settlement_id, status) of the remediated settlement, or (None, None).""" + assert report.remediates_txid + target = await get_settlement_by_txid(machine.id, report.remediates_txid) + if target is None: + logger.warning( + f"spirekeeper: remediation report {report.txid} names txid " + f"{report.remediates_txid} with no settlement on this server" + ) + return None, None + if target.status not in _OWED: + logger.info( + f"spirekeeper: remediation report {report.txid} for settlement " + f"{target.id} in status {target.status!r} — recorded, not moved" + ) + return target.id, target.status + if not report.dispense_confirmed: + logger.warning( + f"spirekeeper: remediation report {report.txid} for settlement " + f"{target.id} did not itself confirm — settlement stays {target.status}" + ) + return target.id, target.status + # The customer is whole: the sale is the full amount. Keep the ORIGINAL + # report's columns on the settlement (that is what happened at the sale); + # the remediation is its own dispense_reports row. + from .crud import mark_settlement_status + + await mark_settlement_status(target.id, "pending", None) + _spawn_distribution(target.id) + logger.info( + f"spirekeeper: settlement {target.id} remediated by manual dispense " + f"{report.txid} — distributing in full" + ) + return target.id, "pending" + + +async def handle_report_dispense(auth, request) -> dict: + """nostr-transport RPC handler. `auth` is the roster-resolved auth context + (unused — the machine is identified from the signature); `request` is a + NostrRpcRequest with `body` and `sender_pubkey` (verified). + + Returns `{txid, received, settlement_status}`; raises ValueError (→ transport + ERROR reply) for an unpaired sender or a malformed body. The machine treats + anything but OK as "resend later", so a malformed report is retried — which + is right: the bug is on one side or the other and the row must not be lost. + """ + sender = (request.sender_pubkey or "").lower() + if not sender: + raise ValueError("missing verified sender_pubkey") + machine = await get_machine_by_atm_pubkey_hex(sender) + if machine is None: + raise ValueError("sender pubkey is not a paired machine") + + try: + report = DispenseReportIn(**(request.body or {})) + except Exception as exc: # pydantic ValidationError, TypeError + raise ValueError(f"invalid report_dispense body: {exc}") from exc + + # Idempotency: the machine resends until acked. + existing = await get_dispense_report(machine.id, report.txid, report.at) + if existing is not None: + settlement = ( + await get_settlement_by_txid(machine.id, report.txid) + if existing.settlement_id + else None + ) + return { + "txid": report.txid, + "received": True, + "settlement_status": settlement.status if settlement else None, + "duplicate": True, + } + + if report.counts_uncertain: + # The state document carries this too; mirroring it here means the + # dashboard learns at report time rather than at the next heartbeat. + await set_machine_counts_uncertain(machine.id, datetime.now()) + + settlement = await get_settlement_by_txid(machine.id, report.txid) + stored: DispenseReport = await insert_dispense_report( + machine.id, settlement.id if settlement else None, report + ) + + status: str | None + if report.remediates_txid: + _, status = await _apply_remediation(machine, report) + elif settlement is None: + # Payment not landed yet (or never will). Kept unlinked; + # _handle_payment adopts it when the settlement is inserted. + logger.warning( + f"spirekeeper: report_dispense {report.txid} from machine {machine.id} " + f"has no settlement yet (confirmed={report.dispense_confirmed}) — " + f"stored unlinked, will attach when the payment lands" + ) + status = None + else: + status = await apply_report_to_settlement(settlement, report, machine) + + logger.info( + f"spirekeeper: report_dispense stored id={stored.id} machine={machine.id} " + f"txid={report.txid} confirmed={report.dispense_confirmed} → {status}" + ) + return {"txid": report.txid, "received": True, "settlement_status": status} + + +async def adopt_unlinked_report( + settlement: DcaSettlement, machine: Machine, report_row: DispenseReport +) -> str: + """`_handle_payment` found a report that arrived before the payment: link + it and apply the same transition the live handler would have.""" + import json as _json + + await link_dispense_reports_to_settlement( + machine.id, report_row.txid, settlement.id + ) + report = DispenseReportIn( + txid=report_row.txid, + payment_hash=report_row.payment_hash, + dispense_confirmed=report_row.dispense_confirmed, + error=report_row.error, + error_code=report_row.error_code, + raw_code=report_row.raw_code, + error_class=report_row.error_class, + fiat_cents=report_row.fiat_cents, + currency=report_row.currency, + bills=_json.loads(report_row.bills_json or "[]"), + cassettes=_json.loads(report_row.cassettes_json or "[]"), + counts_uncertain=report_row.counts_uncertain, + remediates_txid=report_row.remediates_txid, + at=int(report_row.reported_at.timestamp()), + ) + logger.info( + f"spirekeeper: adopting early dispense report {report_row.id} for " + f"settlement {settlement.id} (report preceded the payment)" + ) + return await apply_report_to_settlement(settlement, report, machine) + + +def register_dispense_report_rpc() -> None: + """Register `report_dispense` with the lnbits nostr transport. Soft-fails + if the transport doesn't expose `register_rpc` (older lnbits) — then + cash-out settlements wait in awaiting_dispense and surface on the + worklist as dispense_unreported, which is the honest state.""" + try: + from lnbits.core.services.nostr_transport.dispatcher import ( # type: ignore + AUTH_ACCOUNT, + register_rpc, + ) + except ImportError: + logger.warning( + "spirekeeper: nostr-transport register_rpc unavailable; " + "'report_dispense' not registered (ADR-005 capture disabled — " + "cash-out settlements will sit in awaiting_dispense)" + ) + return + register_rpc(_RPC_NAME, handle_report_dispense, AUTH_ACCOUNT) + logger.info("spirekeeper: registered nostr-transport RPC 'report_dispense'") diff --git a/tasks.py b/tasks.py index 6e5ac55..7001bae 100644 --- a/tasks.py +++ b/tasks.py @@ -169,8 +169,17 @@ async def _handle_payment(payment: Payment) -> None: if isinstance(nostr_event_id, str) and nostr_event_id: data.bitspire_event_id = nostr_event_id - # 3) Insert + distribute. - settlement = await create_settlement_idempotent(data, initial_status="pending") + # 3) Insert. ADR-005 §1: payment is authorization, the machine's dispense + # report is capture. A cash_out waits in `awaiting_dispense` for that + # report (handled in dispense_transport); distribution runs only once it + # says dispense_confirmed. A cash_in has no dispense and proceeds as + # before. Before this gate the legs were paid sub-second, before the + # machine had even begun to dispense — which is how a jam read + # `processed` on 2026-10-09 (bitspire#122). + is_cash_out = data.tx_type == "cash_out" + settlement = await create_settlement_idempotent( + data, initial_status="awaiting_dispense" if is_cash_out else "pending" + ) if settlement is None: logger.error( f"spirekeeper: failed to insert settlement for " @@ -185,6 +194,10 @@ async def _handle_payment(payment: Payment) -> None: f"(super_fee={data.platform_fee_sats} " f"operator_fee={data.operator_fee_sats})" ) + if is_cash_out: + await _await_dispense_or_adopt(settlement, machine, data) + return + # Spawn distribution on a background task so the LNbits invoice queue # (shared across all extensions) keeps draining while we move sats. # Concurrency-safe: process_settlement uses claim_settlement_for_processing @@ -196,6 +209,23 @@ async def _handle_payment(payment: Payment) -> None: task.add_done_callback(_inflight_distributions.discard) +async def _await_dispense_or_adopt( + settlement, machine: Machine, data: CreateDcaSettlementData +) -> None: + """A cash_out waits for the machine's dispense report (ADR-005 §1) — unless + the report is already here. Under hold invoices the payment settles AFTER + the dispense, and the invoice listener can lag the transport, so an orphan + report for this txid is adopted and applied now.""" + if settlement.status != "awaiting_dispense" or not data.bitspire_txid: + return + from .crud import get_latest_unlinked_dispense_report + from .dispense_transport import adopt_unlinked_report + + early = await get_latest_unlinked_dispense_report(machine.id, data.bitspire_txid) + if early is not None: + await adopt_unlinked_report(settlement, machine, early) + + async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) -> None: """Insert a minimal `dca_settlements` row with `status='rejected'` and the exception message for operator forensics. @@ -438,12 +468,38 @@ async def _record_counts_uncertainty( await set_machine_counts_uncertain(machine_id, since) +async def _record_cash_out_hold( + machine_id: str, payload, set_machine_cash_out_hold +) -> None: + """Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row. + + Written on every state event, including when absent, because the machine + clearing the hold — after an operator recount or resume_cash_out — matters + exactly as much as it setting one. + """ + from datetime import datetime as _datetime + from datetime import timezone as _timezone + + since = None + if payload.cash_out_held_since is not None: + since = _datetime.fromtimestamp( + int(payload.cash_out_held_since), tz=_timezone.utc + ) + await set_machine_cash_out_hold( + machine_id, + since, + payload.cash_out_held_reason if since else None, + payload.cash_out_held_code if since else None, + ) + + 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, + set_machine_cash_out_hold=None, ) -> None: """Verify signature, resolve the operator's signer, decrypt via the signer abstraction (bunker round-trip for RemoteBunkerSigner; direct @@ -558,3 +614,8 @@ async def _handle_cassette_state_event( # _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) + if set_machine_cash_out_hold is None: + from .crud import set_machine_cash_out_hold as _default_set_hold + + set_machine_cash_out_hold = _default_set_hold + await _record_cash_out_hold(machine.id, payload, set_machine_cash_out_hold) From fdba2d363f3bb9b5a9f57838f3bd5e8948967fa8 Mon Sep 17 00:00:00 2001 From: Padreug Date: Sat, 10 Oct 2026 21:51:52 +0200 Subject: [PATCH 3/4] =?UTF-8?q?feat(dashboard):=20owed-cash=20worklist=20b?= =?UTF-8?q?uckets,=20prefilled=20partial=20dispense,=20resume=20cash-out?= =?UTF-8?q?=20(ADR-005=20=C2=A75=E2=80=93=C2=A76)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three buckets render first on the worklist — cash_owed, partial_pending, dispense_unreported (awaiting_dispense older than the threshold) — the only ones whose meaning is "a customer is owed money". partial_pending rows open the partial-dispense dialog pre-filled from the machine's report: the fraction from dispensed_fiat_cents / fiat_amount and the dispenser's error in the note, so the operator confirms a number the hardware produced rather than typing one. Machine detail shows a held-cash-out banner (code, time, reason) with a Resume button; POST /machines/{id}/resume-cash-out records a resume_cash_out op and publishes the window. The machine clears the hold on receipt and the banner clears on its next state report. --- static/js/index.js | 82 ++++++++++++++++++++++++++++++- templates/spirekeeper/index.html | 25 ++++++++++ views_api.py | 84 ++++++++++++++++++++++++++++---- 3 files changed, 179 insertions(+), 12 deletions(-) diff --git a/static/js/index.js b/static/js/index.js index e8a31e2..54b6014 100644 --- a/static/js/index.js +++ b/static/js/index.js @@ -70,6 +70,10 @@ window.app = Vue.createApp({ // Worklist (P9g) worklist: { + // ADR-005 §6 — owed-cash buckets first + cash_owed: [], + partial_pending: [], + dispense_unreported: [], rejected: [], errored: [], stuck_pending: [], @@ -372,6 +376,33 @@ window.app = Vue.createApp({ }, worklistBuckets() { return [ + { + key: 'cash_owed', + label: + 'Cash owed — customer paid, machine dispensed nothing. ' + + 'Legs never ran; funds are in the machine wallet.', + icon: 'money_off', + color: 'negative', + rows: this.worklist.cash_owed + }, + { + key: 'partial_pending', + label: + 'Partial dispense — some notes out, value short. Held whole ' + + 'until you record how the shortfall was resolved.', + icon: 'call_split', + color: 'deep-orange', + rows: this.worklist.partial_pending + }, + { + key: 'dispense_unreported', + label: + 'Unreported — cash-out paid, machine never reported the dispense. ' + + 'Check the machine; an old build lands here too.', + icon: 'help_outline', + color: 'amber', + rows: this.worklist.dispense_unreported + }, { key: 'rejected', label: 'Rejected — Nostr attribution failed; investigate machine', @@ -560,6 +591,9 @@ window.app = Vue.createApp({ try { const {data} = await LNbits.api.request('GET', STUCK_PATH) this.worklistCount = + (data?.cash_owed?.length || 0) + + (data?.partial_pending?.length || 0) + + (data?.dispense_unreported?.length || 0) + (data?.rejected?.length || 0) + (data?.errored?.length || 0) + (data?.stuck_pending?.length || 0) + @@ -575,11 +609,17 @@ window.app = Vue.createApp({ const {data} = await LNbits.api.request( 'GET', `${STUCK_PATH}?threshold_minutes=${this.worklistThreshold}` ) + this.worklist.cash_owed = data?.cash_owed || [] + this.worklist.partial_pending = data?.partial_pending || [] + this.worklist.dispense_unreported = data?.dispense_unreported || [] this.worklist.rejected = data?.rejected || [] this.worklist.errored = data?.errored || [] this.worklist.stuck_pending = data?.stuck_pending || [] this.worklist.stuck_processing = data?.stuck_processing || [] this.worklist.totalCount = + this.worklist.cash_owed.length + + this.worklist.partial_pending.length + + this.worklist.dispense_unreported.length + this.worklist.rejected.length + this.worklist.errored.length + this.worklist.stuck_pending.length + @@ -1208,12 +1248,50 @@ window.app = Vue.createApp({ openPartialDispense(settlement) { this.partialDispenseDialog.settlement = settlement this.partialDispenseDialog.mode = 'fraction' - this.partialDispenseDialog.dispensed_fraction = null + // ADR-005: pre-fill from the machine's report — the hardware's own count + // of what left — so the operator confirms a number rather than typing one. + const dispensedCents = settlement.dispensed_fiat_cents + const fiat = Number(settlement.fiat_amount) + this.partialDispenseDialog.dispensed_fraction = + dispensedCents != null && fiat > 0 + ? Math.round((dispensedCents / 100 / fiat) * 10000) / 10000 + : null this.partialDispenseDialog.dispensed_sats = null - this.partialDispenseDialog.notes = '' + this.partialDispenseDialog.notes = settlement.dispense_error + ? `Machine reported: ${settlement.dispense_error_code || ''} ${settlement.dispense_raw_code || ''} — ${settlement.dispense_error}`.trim() + : '' this.partialDispenseDialog.show = true }, + // ADR-005 §5 — release a machine's cash-out hold after a terminal + // dispenser fault, when the jam was cleared without a recount. + confirmResumeCashOut(machine) { + Quasar.Dialog.create({ + title: 'Resume cash-out?', + message: + 'The machine latched cash-out off after a dispenser fault' + + (machine.cash_out_held_code ? ` (${machine.cash_out_held_code})` : '') + + '. Only do this after the transport path has been physically cleared. ' + + 'A recount releases the hold too, and also fixes the bay count.', + cancel: true, + persistent: true + }).onOk(async () => { + try { + await LNbits.api.request( + 'POST', + `/spirekeeper/api/v1/dca/machines/${machine.id}/resume-cash-out` + ) + Quasar.Notify.create({ + type: 'positive', + message: 'Resume published — the machine clears the hold on receipt' + }) + if (this.machineDetail && this.machineDetail.machine) await this.reloadMachineDetail() + } catch (e) { + this._notifyError(e, 'Resume cash-out failed') + } + }) + }, + async submitPartialDispense() { const d = this.partialDispenseDialog const body = {notes: d.notes || null} diff --git a/templates/spirekeeper/index.html b/templates/spirekeeper/index.html index ad01503..77d9612 100644 --- a/templates/spirekeeper/index.html +++ b/templates/spirekeeper/index.html @@ -661,6 +661,12 @@ @click="viewMachineFromWorklist(props.row)"> Open machine detail + + Record the resolution (pre-filled from the machine's report) + + + + Cash-out is held. + The machine latched cash-out off after a terminal dispenser fault + ( + at ): + . + Clear the transport path, then either record a Recount + (which also fixes the count) or release it here. + + + diff --git a/views_api.py b/views_api.py index 33fb2de..7d6bf9f 100644 --- a/views_api.py +++ b/views_api.py @@ -19,6 +19,7 @@ from lnbits.core.services.nsec_bunker import ( ) from lnbits.decorators import check_super_user, check_user_exists from lnbits.utils.nostr import normalize_public_key +from loguru import logger from .calculations import MAX_FEE_FRACTION_PER_DIRECTION from .cassette_transport import ( @@ -28,15 +29,6 @@ from .cassette_transport import ( SignerUnavailable, publish_ops_to_atm, ) -from .fee_transport import publish_fee_config -from .pairing import ( - PairResult, - PairingError, - RevokeResult, - default_relay_endpoint, - pair_spire, - revoke_spire, -) from .crud import ( append_settlement_note, count_completed_legs_for_settlement, @@ -85,6 +77,7 @@ from .distribution import ( process_settlement, settle_lp_balance, ) +from .fee_transport import publish_fee_config from .models import ( AppendSettlementNoteData, CassetteConfig, @@ -112,6 +105,14 @@ from .models import ( UpdateMachineData, UpdateSuperConfigData, ) +from .pairing import ( + PairingError, + PairResult, + RevokeResult, + default_relay_endpoint, + pair_spire, + revoke_spire, +) spirekeeper_api_router = APIRouter() @@ -768,7 +769,13 @@ async def api_list_stuck_settlements( ) -> StuckSettlementsResponse: """Operator worklist of settlements that didn't process cleanly. - Returns four lists: + Returns seven lists. The first three (ADR-005 §6) mean a customer is owed + money and render first: + - cash_owed: the machine reported nothing dispensed; nothing moved + - partial_pending: some notes out, value short; held until resolved + - dispense_unreported: cash-out landed, machine never reported within + the threshold + Then: - rejected: Nostr attribution cross-check failed — signer didn't match the machine identity. Investigate; do not retry. - errored: distribution ran and failed; retry endpoint handles these @@ -783,6 +790,9 @@ async def api_list_stuck_settlements( buckets = await get_stuck_settlements_for_operator(user.id, threshold_minutes) return StuckSettlementsResponse( threshold_minutes=threshold_minutes, + cash_owed=buckets["cash_owed"], + partial_pending=buckets["partial_pending"], + dispense_unreported=buckets["dispense_unreported"], rejected=buckets["rejected"], errored=buckets["errored"], stuck_pending=buckets["stuck_pending"], @@ -1229,3 +1239,57 @@ async def api_create_machine_cassette_op( raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc return op +@spirekeeper_api_router.post( + "/api/v1/dca/machines/{machine_id}/resume-cash-out", + response_model=CassetteOp, +) +async def api_resume_cash_out( + machine_id: str, + user: User = Depends(check_user_exists), +) -> CassetteOp: + """Release a machine's cash-out hold (bitspire ADR-005 §5). + + After a terminal dispenser fault the machine refuses cash-out until an + operator has been to it. A `recount` releases the hold as a side effect; + this is for the case where the jam was cleared without touching a bay + count. Recorded as a machine-wide op (position 0) and published on the + same operator channel as the cassette ops — the machine honours it only if + it is stamped after the hold began, so a re-delivered old resume cannot + clear a newer fault. + + Errors mirror the cassette-op endpoint: 400 unpaired, 503 signer/relay + unavailable (the op is recorded and rides out with the next publish). + """ + machine = await _machine_owned_by(machine_id, user.id) + if not machine.machine_npub: + raise HTTPException( + HTTPStatus.BAD_REQUEST, + "machine is not paired — there is no ATM identity to publish to", + ) + if machine.cash_out_held_since is None: + logger.info( + f"spirekeeper: resume_cash_out for machine {machine_id} with no hold " + "on file — publishing anyway (the machine is the authority)" + ) + + op = await create_cassette_op( + machine_id, + CreateCassetteOpData(position=0, op_type="resume_cash_out"), + created_by=user.id, + ) + window = await get_cassette_ops_window(machine_id) + try: + await publish_ops_to_atm(machine, window, user.id) + except OperatorIdentityMissing as exc: + raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc + except (SignerUnavailable, RelayUnavailable) as exc: + raise HTTPException( + HTTPStatus.SERVICE_UNAVAILABLE, + f"{exc} — the resume 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 op + + From 79413dc06d8aa69815723ca4e638b27b30670a63 Mon Sep 17 00:00:00 2001 From: Padreug Date: Sat, 10 Oct 2026 21:51:52 +0200 Subject: [PATCH 4/4] =?UTF-8?q?test:=20dispense=20outcome=20capture=20?= =?UTF-8?q?=E2=80=94=20handler=20transitions,=20payment=20gate,=20resume?= =?UTF-8?q?=20op,=20hold=20mirror?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Twenty tests in the project's style (asyncio.run, monkeypatched crud, no DB): confirmed → pending + distribution; nothing out → cash_owed; some out → partial_pending with nothing spawned; already-captured recorded not moved; identical resend acked without a row; report-before-payment stored unlinked and adopted when the payment lands; remediation moves the owed settlement; a remediation that did not confirm leaves it; unpaired sender / malformed body refused. The payment gate: cash_out → awaiting_dispense with no distribution, cash_in unchanged. The resume_cash_out op's position rule and wire shape, the worklist model's new buckets, and the state-document hold mirror set-and-clear. The existing nulls-never-reach-the-wire test learns the machine-wide op. --- tests/test_cassette_ops.py | 2 + tests/test_dispense_outcome.py | 542 +++++++++++++++++++++++++++++++++ 2 files changed, 544 insertions(+) create mode 100644 tests/test_dispense_outcome.py diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index 150d493..35ca48c 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -124,6 +124,8 @@ class TestWireShape: "recount": {"count": 1}, "set_denomination": {"denomination": 1}, "empty": {}, + # machine-wide (ADR-005 §5): position 0, no position on the wire + "resume_cash_out": {"position": 0}, }[op_type] wire = op(op_type=op_type, **kw).to_wire_dict() assert None not in wire.values() diff --git a/tests/test_dispense_outcome.py b/tests/test_dispense_outcome.py new file mode 100644 index 0000000..363dc5d --- /dev/null +++ b/tests/test_dispense_outcome.py @@ -0,0 +1,542 @@ +""" +Tests for dispense-outcome capture (bitspire ADR-005 §1-§2, #122). + +Covers the pure pieces and the handler's transitions with the crud layer +monkeypatched (no DB), in the project's established style: asyncio.run inside +the test body, SimpleNamespace for request/payment shapes. + + - models: DispenseReportIn validation + derived numbers; resume_cash_out as a + machine-wide op (position 0, no position on the wire); the three new + worklist buckets default empty. + - handler: confirmed → pending + distribution spawned; nothing out → + cash_owed; some out → partial_pending; already-captured settlements are + recorded but not moved; a byte-identical resend is acked without a new + row; an orphan report (payment not landed) is stored unlinked; a + remediation report moves the owed settlement to pending; unpaired sender + and bad bodies are refused. + - gate: _handle_payment inserts cash_out as awaiting_dispense and does NOT + spawn distribution; cash_in is unchanged; an early report is adopted. + - consumer: the state document's cash_out_held_* mirrors onto the machine, + including clearing it. +""" + +import asyncio +from datetime import datetime, timezone +from types import SimpleNamespace +from typing import Any + +import pytest +from pydantic import ValidationError + +from .. import crud as crud_mod +from .. import dispense_transport, tasks +from ..dispense_transport import ( + _outcome_status, + handle_report_dispense, +) +from ..models import ( + CASSETTE_OP_TYPES, + CassetteOp, + CreateCassetteOpData, + CreateDcaSettlementData, + DcaSettlement, + DispenseReportIn, + Machine, + PublishCassettesPayload, + StuckSettlementsResponse, +) + +_NOW = datetime(2026, 10, 9, 7, 2, 33) +_ATM_HEX = "df2003343784b69cb813b2a4fd231f83ae81133279251c735414f9909baa7ac6" +_TXID = "tx_mv0madw6_wdhtea1v" +_HASH = "6f216df32c36" + "0" * 52 + + +def _machine(**over) -> Machine: + base: dict[str, Any] = { + "id": "m1", + "operator_user_id": "op1", + "machine_npub": _ATM_HEX, + "wallet_id": "w1", + "name": "sintra", + "location": None, + "fiat_code": "EUR", + "is_active": True, + "created_at": _NOW, + "updated_at": _NOW, + } + base.update(over) + return Machine(**base) + + +def _settlement(status="awaiting_dispense", **over) -> DcaSettlement: + base: dict[str, Any] = { + "id": "s1", + "machine_id": "m1", + "payment_hash": _HASH, + "bitspire_event_id": None, + "bitspire_txid": _TXID, + "wire_sats": 54440, + "fiat_amount": 40.0, + "fiat_code": "EUR", + "exchange_rate": 1361.0, + "principal_sats": 54440, + "fee_sats": 0, + "platform_fee_sats": 0, + "operator_fee_sats": 0, + "tx_type": "cash_out", + "bills_json": None, + "cassettes_json": None, + "status": status, + "error_message": None, + "processed_at": None, + "created_at": _NOW, + } + base.update(over) + return DcaSettlement(**base) + + +def _report(**over) -> dict: + """The wire body for sintra's 2026-10-09 jam, as a dict.""" + body: dict[str, Any] = { + "txid": _TXID, + "payment_hash": _HASH, + "tx_type": "cash_out", + "dispense_confirmed": False, + "error": "Note stopped at the cassette exit", + "error_code": "F56DispenseError", + "raw_code": "78 42", + "error_class": "terminal", + "fiat_cents": 4000, + "currency": "EUR", + "bills": [{"denomination": 20, "requested": 2, "dispensed": 0, "rejected": 0}], + "cassettes": [ + { + "position": 1, + "denomination": 50, + "provisioned": 0, + "dispensed": 0, + "rejected": 0, + }, + { + "position": 2, + "denomination": 20, + "provisioned": 2, + "dispensed": 0, + "rejected": 0, + }, + ], + "counts_uncertain": True, + "at": 1791529353, + } + body.update(over) + return body + + +def _confirmed() -> dict: + return _report( + dispense_confirmed=True, + error=None, + error_code=None, + raw_code=None, + error_class=None, + bills=[{"denomination": 20, "requested": 2, "dispensed": 2, "rejected": 0}], + counts_uncertain=False, + ) + + +def _partial() -> dict: + return _report( + bills=[{"denomination": 20, "requested": 2, "dispensed": 1, "rejected": 1}], + counts_uncertain=False, + ) + + +# --------------------------------------------------------------------------- +# Models +# --------------------------------------------------------------------------- + + +class TestDispenseReportIn: + def test_derived_numbers(self): + r = DispenseReportIn(**_partial()) + assert r.dispensed_fiat_cents == 2000 + assert r.total_dispensed_notes == 1 + assert DispenseReportIn(**_confirmed()).dispensed_fiat_cents == 4000 + assert DispenseReportIn(**_report()).total_dispensed_notes == 0 + + def test_outcome_routing(self): + assert _outcome_status(DispenseReportIn(**_confirmed())) == "pending" + assert _outcome_status(DispenseReportIn(**_partial())) == "partial_pending" + assert _outcome_status(DispenseReportIn(**_report())) == "cash_owed" + + def test_rejects_cash_in_and_unknown_class(self): + with pytest.raises(ValidationError): + DispenseReportIn(**_report(tx_type="cash_in")) + with pytest.raises(ValidationError): + DispenseReportIn(**_report(error_class="weird")) + with pytest.raises(ValidationError): + DispenseReportIn(**_report(txid=" ")) + with pytest.raises(ValidationError): + DispenseReportIn(**_report(fiat_cents=-1)) + + +class TestResumeCashOutOp: + def test_is_a_known_type_and_machine_wide(self): + assert "resume_cash_out" in CASSETTE_OP_TYPES + op = CassetteOp( + id="r1", + machine_id="m1", + position=0, + op_type="resume_cash_out", + created_at=_NOW, + ) + wire = op.to_wire_dict() + assert wire == { + "id": "r1", + "at": int(_NOW.timestamp()), + "type": "resume_cash_out", + } + assert "position" not in wire + + def test_position_must_be_zero_for_resume_and_positive_otherwise(self): + with pytest.raises(ValidationError): + CassetteOp( + id="r1", + machine_id="m1", + position=2, + op_type="resume_cash_out", + created_at=_NOW, + ) + with pytest.raises(ValidationError): + CreateCassetteOpData(position=1, op_type="resume_cash_out") + with pytest.raises(ValidationError): + CreateCassetteOpData(position=0, op_type="refill", bills=5) + CreateCassetteOpData(position=0, op_type="resume_cash_out") # ok + + def test_resume_carries_no_bay_fields(self): + with pytest.raises(ValidationError): + CreateCassetteOpData(position=0, op_type="resume_cash_out", count=3) + + +class TestWorklistModel: + def test_new_buckets_default_empty(self): + r = StuckSettlementsResponse( + threshold_minutes=30, + rejected=[], + errored=[], + stuck_pending=[], + stuck_processing=[], + ) + assert ( + r.cash_owed == [] + and r.partial_pending == [] + and r.dispense_unreported == [] + ) + + +# --------------------------------------------------------------------------- +# Handler +# --------------------------------------------------------------------------- + + +class _Wired: + """Monkeypatched crud layer for handle_report_dispense.""" + + def __init__(self, monkeypatch, *, machine, settlement, existing_report=None): + self.inserted = [] + self.applied = [] + self.spawned = [] + self.uncertain = [] + self.statuses = [] + self.settlement = settlement + + async def get_machine(_hex): + return machine + + async def get_report(_mid, _txid, _at): + return existing_report + + async def get_settlement(_mid, txid): + if self.settlement is not None and self.settlement.bitspire_txid == txid: + return self.settlement + return None + + async def insert(mid, sid, report): + self.inserted.append((mid, sid, report)) + return SimpleNamespace(id="rep1", settlement_id=sid, txid=report.txid) + + async def apply(sid, report, new_status, reported_at): + self.applied.append((sid, new_status, reported_at)) + return ( + self.settlement.copy(update={"status": new_status}) + if self.settlement + else None + ) + + async def set_uncertain(mid, since): + self.uncertain.append((mid, since)) + + async def mark_status(sid, status, _msg): + self.statuses.append((sid, status)) + return None + + monkeypatch.setattr( + dispense_transport, "get_machine_by_atm_pubkey_hex", get_machine + ) + monkeypatch.setattr(dispense_transport, "get_dispense_report", get_report) + monkeypatch.setattr( + dispense_transport, "get_settlement_by_txid", get_settlement + ) + monkeypatch.setattr(dispense_transport, "insert_dispense_report", insert) + monkeypatch.setattr(dispense_transport, "apply_dispense_outcome", apply) + monkeypatch.setattr( + dispense_transport, "set_machine_counts_uncertain", set_uncertain + ) + monkeypatch.setattr( + dispense_transport, "_spawn_distribution", self.spawned.append + ) + monkeypatch.setattr(crud_mod, "mark_settlement_status", mark_status) + + +def _req(body, sender=_ATM_HEX): + return SimpleNamespace(body=body, sender_pubkey=sender, event_id="ev1") + + +class TestHandleReportDispense: + def test_confirmed_captures_and_distributes(self, monkeypatch): + w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement()) + out = asyncio.run(handle_report_dispense(None, _req(_confirmed()))) + assert out["received"] is True + assert out["settlement_status"] == "pending" + assert w.applied == [("s1", "pending", datetime.fromtimestamp(1791529353))] + assert w.spawned == ["s1"] + assert w.inserted[0][1] == "s1" # linked to the settlement + assert w.uncertain == [] + + def test_nothing_out_is_cash_owed_and_nothing_moves(self, monkeypatch): + w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement()) + out = asyncio.run(handle_report_dispense(None, _req(_report()))) + assert out["settlement_status"] == "cash_owed" + assert w.applied[0][1] == "cash_owed" + assert w.spawned == [] + # counts_uncertain on the report mirrors onto the machine immediately + assert len(w.uncertain) == 1 and w.uncertain[0][0] == "m1" + + def test_some_out_is_partial_pending_held_whole(self, monkeypatch): + w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement()) + out = asyncio.run(handle_report_dispense(None, _req(_partial()))) + assert out["settlement_status"] == "partial_pending" + assert w.spawned == [] # ADR-005 Decision 1: one distribution, when final + + def test_already_captured_settlement_is_recorded_not_moved(self, monkeypatch): + w = _Wired( + monkeypatch, machine=_machine(), settlement=_settlement(status="processed") + ) + out = asyncio.run(handle_report_dispense(None, _req(_report()))) + assert out["settlement_status"] == "processed" + assert w.applied == [] and w.spawned == [] + assert len(w.inserted) == 1 # the row still lands — it is information + + def test_identical_resend_is_acked_without_a_new_row(self, monkeypatch): + existing = SimpleNamespace(id="rep0", settlement_id="s1", txid=_TXID) + w = _Wired( + monkeypatch, + machine=_machine(), + settlement=_settlement(status="cash_owed"), + existing_report=existing, + ) + out = asyncio.run(handle_report_dispense(None, _req(_report()))) + assert out["received"] is True and out.get("duplicate") is True + assert out["settlement_status"] == "cash_owed" + assert w.inserted == [] and w.applied == [] + + def test_report_before_payment_is_stored_unlinked(self, monkeypatch): + w = _Wired(monkeypatch, machine=_machine(), settlement=None) + out = asyncio.run(handle_report_dispense(None, _req(_confirmed()))) + assert out["settlement_status"] is None + assert w.inserted[0][1] is None + assert w.spawned == [] + + def test_remediation_moves_the_owed_settlement_to_pending(self, monkeypatch): + owed = _settlement(status="cash_owed") + w = _Wired(monkeypatch, machine=_machine(), settlement=owed) + body = _confirmed() + body.update(txid="manual-1", remediates_txid=_TXID, at=1791530000) + out = asyncio.run(handle_report_dispense(None, _req(body))) + assert out["settlement_status"] == "pending" + assert w.statuses == [("s1", "pending")] + assert w.spawned == ["s1"] + assert w.applied == [] # the ORIGINAL report's columns stay on the settlement + + def test_remediation_that_did_not_confirm_leaves_it_owed(self, monkeypatch): + w = _Wired( + monkeypatch, machine=_machine(), settlement=_settlement(status="cash_owed") + ) + body = _report() + body.update(txid="manual-2", remediates_txid=_TXID, at=1791530001) + out = asyncio.run(handle_report_dispense(None, _req(body))) + assert out["settlement_status"] == "cash_owed" + assert w.statuses == [] and w.spawned == [] + + def test_unpaired_sender_and_bad_body_are_refused(self, monkeypatch): + _Wired(monkeypatch, machine=None, settlement=None) + with pytest.raises(ValueError, match="not a paired machine"): + asyncio.run(handle_report_dispense(None, _req(_report()))) + _Wired(monkeypatch, machine=_machine(), settlement=_settlement()) + with pytest.raises(ValueError, match="invalid report_dispense body"): + asyncio.run(handle_report_dispense(None, _req({"txid": _TXID}))) + with pytest.raises(ValueError, match="sender_pubkey"): + asyncio.run(handle_report_dispense(None, _req(_report(), sender=""))) + + +# --------------------------------------------------------------------------- +# Gate in _handle_payment +# --------------------------------------------------------------------------- + + +def _payment(is_in=True): + return SimpleNamespace( + success=True, + wallet_id="w1", + extra={ + "source": "bitspire", + "type": "cash_out" if is_in else "cash_in", + "txid": _TXID, + }, + is_in=is_in, + sat=54440 if is_in else -54440, + payment_hash=_HASH, + ) + + +def _data(tx_type) -> CreateDcaSettlementData: + return CreateDcaSettlementData( + machine_id="m1", + payment_hash=_HASH, + bitspire_txid=_TXID, + wire_sats=54440, + fiat_amount=40.0, + fiat_code="EUR", + exchange_rate=1361.0, + principal_sats=54440, + fee_sats=0, + platform_fee_sats=0, + operator_fee_sats=0, + tx_type=tx_type, + ) + + +class _GateWired: + def __init__(self, monkeypatch, *, tx_type, early_report=None): + self.created = [] + self.spawned = [] + self.adopted = [] + + async def get_machine(_wid): + return _machine() + + def attribution(_machine, _extra): + return None + + async def get_super(): + return SimpleNamespace(id="default") + + def parse(**_kw): + return _data(tx_type) + + async def create(data, initial_status, error_message=None): + self.created.append((data.tx_type, initial_status)) + return _settlement(status=initial_status, tx_type=data.tx_type) + + async def process(sid): + self.spawned.append(sid) + + async def early(_mid, _txid): + return early_report + + async def adopt(settlement, machine, row): + self.adopted.append((settlement.id, row.id)) + return "pending" + + monkeypatch.setattr(tasks, "get_active_machine_by_wallet_id", get_machine) + monkeypatch.setattr(tasks, "assert_nostr_attribution", attribution) + monkeypatch.setattr(tasks, "get_super_config", get_super) + monkeypatch.setattr(tasks, "parse_settlement", parse) + monkeypatch.setattr(tasks, "create_settlement_idempotent", create) + monkeypatch.setattr(tasks, "process_settlement", process) + monkeypatch.setattr(crud_mod, "get_latest_unlinked_dispense_report", early) + monkeypatch.setattr(dispense_transport, "adopt_unlinked_report", adopt) + + +async def _drain(): + # let any create_task'd distribution run + await asyncio.sleep(0) + + +class TestPaymentGate: + def test_cash_out_lands_awaiting_dispense_and_does_not_distribute( + self, monkeypatch + ): + w = _GateWired(monkeypatch, tx_type="cash_out") + + async def run(): + await tasks._handle_payment(_payment(is_in=True)) + await _drain() + + asyncio.run(run()) + assert w.created == [("cash_out", "awaiting_dispense")] + assert w.spawned == [] + assert w.adopted == [] + + def test_cash_in_is_unchanged(self, monkeypatch): + w = _GateWired(monkeypatch, tx_type="cash_in") + + async def run(): + await tasks._handle_payment(_payment(is_in=False)) + await _drain() + + asyncio.run(run()) + assert w.created == [("cash_in", "pending")] + assert w.spawned == ["s1"] + + def test_early_report_is_adopted_when_the_payment_lands(self, monkeypatch): + early = SimpleNamespace(id="rep-early", txid=_TXID) + w = _GateWired(monkeypatch, tx_type="cash_out", early_report=early) + + async def run(): + await tasks._handle_payment(_payment(is_in=True)) + await _drain() + + asyncio.run(run()) + assert w.adopted == [("s1", "rep-early")] + assert w.spawned == [] # adoption decides; the fake adopt did not spawn + + +# --------------------------------------------------------------------------- +# Consumer mirror +# --------------------------------------------------------------------------- + + +class TestCashOutHoldMirror: + def test_sets_and_clears_from_the_state_document(self): + calls = [] + + async def setter(mid, since, reason, code): + calls.append((mid, since, reason, code)) + + held = PublishCassettesPayload( + positions={"2": {"denomination": 20, "count": 54}}, + cash_out_held_since=1791529353, + cash_out_held_reason="Note stopped at the cassette exit", + cash_out_held_code="78 42", + ) + clear = PublishCassettesPayload( + positions={"2": {"denomination": 20, "count": 54}} + ) + asyncio.run(tasks._record_cash_out_hold("m1", held, setter)) + asyncio.run(tasks._record_cash_out_hold("m1", clear, setter)) + assert calls[0][0] == "m1" + assert calls[0][1] == datetime.fromtimestamp(1791529353, tz=timezone.utc) + assert calls[0][2:] == ("Note stopped at the cassette exit", "78 42") + assert calls[1] == ("m1", None, None, None)