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/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/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/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}`")
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/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)
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/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)
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
+
+