Compare commits

..

No commits in common. "79413dc06d8aa69815723ca4e638b27b30670a63" and "d09c51f277c60d6cf8c157c2a088420a426f15e3" have entirely different histories.

11 changed files with 32 additions and 1602 deletions

View file

@ -6,7 +6,6 @@ from loguru import logger
from .cashin_transport import register_create_withdraw_rpc from .cashin_transport import register_create_withdraw_rpc
from .crud import db from .crud import db
from .dispense_transport import register_dispense_report_rpc
from .machine_config_transport import register_machine_config_rpc from .machine_config_transport import register_machine_config_rpc
from .nostr_transport_roster import register_with_lnbits as register_roster_with_lnbits 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 from .tasks import wait_for_cassette_state_events, wait_for_paid_invoices
@ -69,11 +68,6 @@ def spirekeeper_start():
# config over the transport, leaving "awaiting configuration" with no # config over the transport, leaving "awaiting configuration" with no
# per-machine env provisioning. Soft-fails if register_rpc isn't exposed. # per-machine env provisioning. Soft-fails if register_rpc isn't exposed.
register_machine_config_rpc() 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__ = [ __all__ = [

250
crud.py
View file

@ -27,8 +27,6 @@ from .models import (
DcaLpPreferences, DcaLpPreferences,
DcaPayment, DcaPayment,
DcaSettlement, DcaSettlement,
DispenseReport,
DispenseReportIn,
Machine, Machine,
PublishCassettesPayload, PublishCassettesPayload,
SuperConfig, SuperConfig,
@ -259,32 +257,6 @@ async def set_machine_unpaired(machine_id: str) -> Machine | None:
return await get_machine(machine_id) return await get_machine(machine_id)
async def set_machine_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: 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" """Record (or clear) the machine's own "I can't vouch for these counts"
marker, straight from its state document. marker, straight from its state document.
@ -739,171 +711,6 @@ async def create_settlement_idempotent(
return await get_settlement(settlement_id) 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: async def get_settlement(settlement_id: str) -> DcaSettlement | None:
return await db.fetchone( return await db.fetchone(
"SELECT * FROM spirekeeper.dca_settlements WHERE id = :id", "SELECT * FROM spirekeeper.dca_settlements WHERE id = :id",
@ -945,14 +752,7 @@ async def get_stuck_settlements_for_operator(
) -> dict: ) -> dict:
"""Operator worklist of settlements that didn't process cleanly. """Operator worklist of settlements that didn't process cleanly.
Returns a dict with seven keyed lists. The first three are ADR-005 §6 — Returns a dict with four keyed lists:
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 - 'rejected': any status='rejected' (Nostr attribution cross-check
failed — signer didn't match the machine identity). Distinct failed — signer didn't match the machine identity). Distinct
from 'errored' because retry is wrong: the row was misrouted, from 'errored' because retry is wrong: the row was misrouted,
@ -1017,48 +817,7 @@ async def get_stuck_settlements_for_operator(
{"uid": operator_user_id, "threshold": threshold_at}, {"uid": operator_user_id, "threshold": threshold_at},
DcaSettlement, 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 { return {
"cash_owed": cash_owed,
"partial_pending": partial_pending,
"dispense_unreported": dispense_unreported,
"rejected": rejected, "rejected": rejected,
"errored": errored, "errored": errored,
"stuck_pending": stuck_pending, "stuck_pending": stuck_pending,
@ -1111,10 +870,9 @@ async def mark_settlement_status(
status: str, status: str,
error_message: str | None = None, error_message: str | None = None,
) -> DcaSettlement | None: ) -> DcaSettlement | None:
"""Status: 'awaiting_dispense' | 'pending' | 'processing' | 'processed' | """Status: 'pending' | 'processing' | 'processed' | 'partial' |
'partial_pending' | 'cash_owed' | 'partial' | 'refunded' | 'errored'. 'refunded' | 'errored'. Clears processing_claim on terminal states so a
Clears processing_claim on terminal states so a fresh claim attempt won't fresh claim attempt won't see a stale token."""
see a stale token."""
await db.execute( await db.execute(
""" """
UPDATE spirekeeper.dca_settlements UPDATE spirekeeper.dca_settlements

View file

@ -1,299 +0,0 @@
"""
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'")

View file

@ -941,82 +941,3 @@ async def m015_add_cassette_state_seq(db):
await db.execute( await db.execute(
"ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER" "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}"
)

200
models.py
View file

@ -67,13 +67,6 @@ class Machine(BaseModel):
# count of what left the bay; cleared by the machine's own report. The # 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. # dashboard turns this into a prompt to open the bay and recount.
counts_uncertain_since: datetime | None = None 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 created_at: datetime
updated_at: datetime updated_at: datetime
@ -318,39 +311,19 @@ class DcaSettlement(BaseModel):
fee_mismatch_sats: int | None = None fee_mismatch_sats: int | None = None
bills_json: str | None bills_json: str | None
cassettes_json: str | None cassettes_json: str | None
# Lifecycle (bitspire ADR-005 §1 — payment is authorization, the # 'pending' (default at insert)
# 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) # 'processing' (claim taken by distribution processor)
# 'processed' (all legs paid) # 'processed' (all legs paid)
# 'partial_pending' (report says some notes out, value short — holds # 'partial' (operator marked partial-dispense after the fact)
# 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) # 'refunded' (operator-initiated refund)
# 'errored' (operational distribution failure — retry path applies) # 'errored' (operational distribution failure — retry path applies)
# 'rejected' (Nostr attribution cross-check failed at land time; # 'rejected' (Nostr attribution cross-check failed at land time;
# never went near distribution. error_message holds the # never went near distribution. error_message holds the
# reason. Retry is wrong — investigate the machine.) # 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 status: str
error_message: str | None error_message: str | None
processed_at: datetime | None processed_at: datetime | None
created_at: datetime 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 # Append-only audit memo. Populated when an operator triggers an in-place
# adjustment (partial-dispense, manual reconciliation override). Each # adjustment (partial-dispense, manual reconciliation override). Each
# entry timestamped + records original values so the overwrite is # entry timestamped + records original values so the overwrite is
@ -363,110 +336,6 @@ class DcaSettlement(BaseModel):
processing_claim: str | None = None 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. # Commission splits — operator-defined remainder allocation per machine.
# ============================================================================= # =============================================================================
@ -701,13 +570,6 @@ class StuckSettlementsResponse(BaseModel):
""" """
threshold_minutes: int 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] rejected: list # list[DcaSettlement]
errored: list errored: list
stuck_pending: list stuck_pending: list
@ -876,11 +738,6 @@ class PublishCassettesPayload(BaseModel):
applied_ops: list[str] = [] applied_ops: list[str] = []
seq: int | None = None seq: int | None = None
counts_uncertain_since: 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) @validator("positions", pre=True)
def coerce_string_keys_to_int(cls, v): def coerce_string_keys_to_int(cls, v):
@ -946,21 +803,7 @@ class PublishCassettesPayload(BaseModel):
# landed after the same problem. # landed after the same problem.
CASSETTE_OP_TYPES = ( CASSETTE_OP_TYPES = ("refill", "empty", "recount", "set_denomination")
"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): class CassetteOp(BaseModel):
@ -994,17 +837,11 @@ class CassetteOp(BaseModel):
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}") raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
return v return v
@root_validator(skip_on_failure=True) @validator("position")
def _position_matches_scope(cls, values): def _position_positive(cls, v):
pos, typ = values.get("position"), values.get("op_type") if v <= 0:
if typ in MACHINE_WIDE_OP_TYPES: raise ValueError(f"position must be > 0, got {v}")
if pos != 0: return v
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: def to_wire_dict(self) -> dict:
"""The published form. Drops the fields this op_type does not use, so """The published form. Drops the fields this op_type does not use, so
@ -1013,10 +850,8 @@ class CassetteOp(BaseModel):
"id": self.id, "id": self.id,
"at": int(self.created_at.timestamp()), "at": int(self.created_at.timestamp()),
"type": self.op_type, "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": if self.op_type == "refill":
out["bills"] = self.bills out["bills"] = self.bills
elif self.op_type == "recount": elif self.op_type == "recount":
@ -1045,17 +880,11 @@ class CreateCassetteOpData(BaseModel):
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}") raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
return v return v
@root_validator(skip_on_failure=True) @validator("position")
def _position_matches_scope(cls, values): def _position_positive(cls, v):
pos, typ = values.get("position"), values.get("op_type") if v <= 0:
if typ in MACHINE_WIDE_OP_TYPES: raise ValueError(f"position must be > 0, got {v}")
if pos != 0: return v
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") @validator("bills")
def _bills_positive(cls, v): def _bills_positive(cls, v):
@ -1082,7 +911,6 @@ class CreateCassetteOpData(BaseModel):
"recount": "count", "recount": "count",
"set_denomination": "denomination", "set_denomination": "denomination",
"empty": None, "empty": None,
"resume_cash_out": None,
}[values.get("op_type")] }[values.get("op_type")]
if required is not None and values.get(required) is None: if required is not None and values.get(required) is None:
raise ValueError(f"{values['op_type']} requires `{required}`") raise ValueError(f"{values['op_type']} requires `{required}`")

View file

@ -70,10 +70,6 @@ window.app = Vue.createApp({
// Worklist (P9g) // Worklist (P9g)
worklist: { worklist: {
// ADR-005 §6 — owed-cash buckets first
cash_owed: [],
partial_pending: [],
dispense_unreported: [],
rejected: [], rejected: [],
errored: [], errored: [],
stuck_pending: [], stuck_pending: [],
@ -376,33 +372,6 @@ window.app = Vue.createApp({
}, },
worklistBuckets() { worklistBuckets() {
return [ 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', key: 'rejected',
label: 'Rejected — Nostr attribution failed; investigate machine', label: 'Rejected — Nostr attribution failed; investigate machine',
@ -591,9 +560,6 @@ window.app = Vue.createApp({
try { try {
const {data} = await LNbits.api.request('GET', STUCK_PATH) const {data} = await LNbits.api.request('GET', STUCK_PATH)
this.worklistCount = this.worklistCount =
(data?.cash_owed?.length || 0) +
(data?.partial_pending?.length || 0) +
(data?.dispense_unreported?.length || 0) +
(data?.rejected?.length || 0) + (data?.rejected?.length || 0) +
(data?.errored?.length || 0) + (data?.errored?.length || 0) +
(data?.stuck_pending?.length || 0) + (data?.stuck_pending?.length || 0) +
@ -609,17 +575,11 @@ window.app = Vue.createApp({
const {data} = await LNbits.api.request( const {data} = await LNbits.api.request(
'GET', `${STUCK_PATH}?threshold_minutes=${this.worklistThreshold}` '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.rejected = data?.rejected || []
this.worklist.errored = data?.errored || [] this.worklist.errored = data?.errored || []
this.worklist.stuck_pending = data?.stuck_pending || [] this.worklist.stuck_pending = data?.stuck_pending || []
this.worklist.stuck_processing = data?.stuck_processing || [] this.worklist.stuck_processing = data?.stuck_processing || []
this.worklist.totalCount = this.worklist.totalCount =
this.worklist.cash_owed.length +
this.worklist.partial_pending.length +
this.worklist.dispense_unreported.length +
this.worklist.rejected.length + this.worklist.rejected.length +
this.worklist.errored.length + this.worklist.errored.length +
this.worklist.stuck_pending.length + this.worklist.stuck_pending.length +
@ -1248,50 +1208,12 @@ window.app = Vue.createApp({
openPartialDispense(settlement) { openPartialDispense(settlement) {
this.partialDispenseDialog.settlement = settlement this.partialDispenseDialog.settlement = settlement
this.partialDispenseDialog.mode = 'fraction' this.partialDispenseDialog.mode = 'fraction'
// ADR-005: pre-fill from the machine's report — the hardware's own count this.partialDispenseDialog.dispensed_fraction = null
// 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.dispensed_sats = null
this.partialDispenseDialog.notes = settlement.dispense_error this.partialDispenseDialog.notes = ''
? `Machine reported: ${settlement.dispense_error_code || ''} ${settlement.dispense_raw_code || ''} — ${settlement.dispense_error}`.trim()
: ''
this.partialDispenseDialog.show = true 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() { async submitPartialDispense() {
const d = this.partialDispenseDialog const d = this.partialDispenseDialog
const body = {notes: d.notes || null} const body = {notes: d.notes || null}

View file

@ -169,17 +169,8 @@ async def _handle_payment(payment: Payment) -> None:
if isinstance(nostr_event_id, str) and nostr_event_id: if isinstance(nostr_event_id, str) and nostr_event_id:
data.bitspire_event_id = nostr_event_id data.bitspire_event_id = nostr_event_id
# 3) Insert. ADR-005 §1: payment is authorization, the machine's dispense # 3) Insert + distribute.
# report is capture. A cash_out waits in `awaiting_dispense` for that settlement = await create_settlement_idempotent(data, initial_status="pending")
# 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: if settlement is None:
logger.error( logger.error(
f"spirekeeper: failed to insert settlement for " f"spirekeeper: failed to insert settlement for "
@ -194,10 +185,6 @@ async def _handle_payment(payment: Payment) -> None:
f"(super_fee={data.platform_fee_sats} " f"(super_fee={data.platform_fee_sats} "
f"operator_fee={data.operator_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 # Spawn distribution on a background task so the LNbits invoice queue
# (shared across all extensions) keeps draining while we move sats. # (shared across all extensions) keeps draining while we move sats.
# Concurrency-safe: process_settlement uses claim_settlement_for_processing # Concurrency-safe: process_settlement uses claim_settlement_for_processing
@ -209,23 +196,6 @@ async def _handle_payment(payment: Payment) -> None:
task.add_done_callback(_inflight_distributions.discard) 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: async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) -> None:
"""Insert a minimal `dca_settlements` row with `status='rejected'` and """Insert a minimal `dca_settlements` row with `status='rejected'` and
the exception message for operator forensics. the exception message for operator forensics.
@ -468,38 +438,12 @@ async def _record_counts_uncertainty(
await set_machine_counts_uncertain(machine_id, since) 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( async def _handle_cassette_state_event(
event_message, event_message,
get_machine_by_atm_pubkey_hex, get_machine_by_atm_pubkey_hex,
apply_reported_state, apply_reported_state,
mark_cassette_ops_acked, mark_cassette_ops_acked,
set_machine_counts_uncertain, set_machine_counts_uncertain,
set_machine_cash_out_hold=None,
) -> None: ) -> None:
"""Verify signature, resolve the operator's signer, decrypt via the """Verify signature, resolve the operator's signer, decrypt via the
signer abstraction (bunker round-trip for RemoteBunkerSigner; direct signer abstraction (bunker round-trip for RemoteBunkerSigner; direct
@ -614,8 +558,3 @@ async def _handle_cassette_state_event(
# _record_op_acknowledgements for why. Same for the uncertainty marker. # _record_op_acknowledgements for why. Same for the uncertainty marker.
await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked) await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked)
await _record_counts_uncertainty(machine.id, payload, set_machine_counts_uncertain) 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)

View file

@ -661,12 +661,6 @@
@click="viewMachineFromWorklist(props.row)"> @click="viewMachineFromWorklist(props.row)">
<q-tooltip>Open machine detail</q-tooltip> <q-tooltip>Open machine detail</q-tooltip>
</q-btn> </q-btn>
<q-btn v-if="bucket.key === 'partial_pending'"
flat dense size="sm" icon="call_split"
color="deep-orange"
@click="openPartialDispense(props.row)">
<q-tooltip>Record the resolution (pre-filled from the machine's report)</q-tooltip>
</q-btn>
<q-btn v-if="bucket.key === 'errored'" <q-btn v-if="bucket.key === 'errored'"
flat dense size="sm" icon="restart_alt" flat dense size="sm" icon="restart_alt"
color="primary" color="primary"
@ -1184,25 +1178,6 @@
<span v-text="machineDetail.cassettesError"></span> <span v-text="machineDetail.cassettesError"></span>
</q-banner> </q-banner>
<q-banner v-if="machineDetail.machine
&& machineDetail.machine.cash_out_held_since"
class="bg-red-1 text-grey-9 q-mb-md">
<template v-slot:avatar>
<q-icon name="block" color="negative"></q-icon>
</template>
<b>Cash-out is held.</b>
The machine latched cash-out off after a terminal dispenser fault
(<span v-text="machineDetail.machine.cash_out_held_code || 'fault'"></span>
at <span v-text="formatTime(machineDetail.machine.cash_out_held_since)"></span>):
<span v-text="machineDetail.machine.cash_out_held_reason"></span>.
Clear the transport path, then either record a <b>Recount</b>
(which also fixes the count) or release it here.
<template v-slot:action>
<q-btn flat color="negative" label="Resume cash-out"
@click="confirmResumeCashOut(machineDetail.machine)"></q-btn>
</template>
</q-banner>
<q-banner v-if="machineDetail.machine <q-banner v-if="machineDetail.machine
&& machineDetail.machine.counts_uncertain_since" && machineDetail.machine.counts_uncertain_since"
class="bg-orange-1 text-grey-9 q-mb-md"> class="bg-orange-1 text-grey-9 q-mb-md">

View file

@ -124,8 +124,6 @@ class TestWireShape:
"recount": {"count": 1}, "recount": {"count": 1},
"set_denomination": {"denomination": 1}, "set_denomination": {"denomination": 1},
"empty": {}, "empty": {},
# machine-wide (ADR-005 §5): position 0, no position on the wire
"resume_cash_out": {"position": 0},
}[op_type] }[op_type]
wire = op(op_type=op_type, **kw).to_wire_dict() wire = op(op_type=op_type, **kw).to_wire_dict()
assert None not in wire.values() assert None not in wire.values()

View file

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

View file

@ -19,7 +19,6 @@ from lnbits.core.services.nsec_bunker import (
) )
from lnbits.decorators import check_super_user, check_user_exists from lnbits.decorators import check_super_user, check_user_exists
from lnbits.utils.nostr import normalize_public_key from lnbits.utils.nostr import normalize_public_key
from loguru import logger
from .calculations import MAX_FEE_FRACTION_PER_DIRECTION from .calculations import MAX_FEE_FRACTION_PER_DIRECTION
from .cassette_transport import ( from .cassette_transport import (
@ -29,6 +28,15 @@ from .cassette_transport import (
SignerUnavailable, SignerUnavailable,
publish_ops_to_atm, 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 ( from .crud import (
append_settlement_note, append_settlement_note,
count_completed_legs_for_settlement, count_completed_legs_for_settlement,
@ -77,7 +85,6 @@ from .distribution import (
process_settlement, process_settlement,
settle_lp_balance, settle_lp_balance,
) )
from .fee_transport import publish_fee_config
from .models import ( from .models import (
AppendSettlementNoteData, AppendSettlementNoteData,
CassetteConfig, CassetteConfig,
@ -105,14 +112,6 @@ from .models import (
UpdateMachineData, UpdateMachineData,
UpdateSuperConfigData, UpdateSuperConfigData,
) )
from .pairing import (
PairingError,
PairResult,
RevokeResult,
default_relay_endpoint,
pair_spire,
revoke_spire,
)
spirekeeper_api_router = APIRouter() spirekeeper_api_router = APIRouter()
@ -769,13 +768,7 @@ async def api_list_stuck_settlements(
) -> StuckSettlementsResponse: ) -> StuckSettlementsResponse:
"""Operator worklist of settlements that didn't process cleanly. """Operator worklist of settlements that didn't process cleanly.
Returns seven lists. The first three (ADR-005 §6) mean a customer is owed Returns four lists:
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 - rejected: Nostr attribution cross-check failed — signer didn't
match the machine identity. Investigate; do not retry. match the machine identity. Investigate; do not retry.
- errored: distribution ran and failed; retry endpoint handles these - errored: distribution ran and failed; retry endpoint handles these
@ -790,9 +783,6 @@ async def api_list_stuck_settlements(
buckets = await get_stuck_settlements_for_operator(user.id, threshold_minutes) buckets = await get_stuck_settlements_for_operator(user.id, threshold_minutes)
return StuckSettlementsResponse( return StuckSettlementsResponse(
threshold_minutes=threshold_minutes, threshold_minutes=threshold_minutes,
cash_owed=buckets["cash_owed"],
partial_pending=buckets["partial_pending"],
dispense_unreported=buckets["dispense_unreported"],
rejected=buckets["rejected"], rejected=buckets["rejected"],
errored=buckets["errored"], errored=buckets["errored"],
stuck_pending=buckets["stuck_pending"], stuck_pending=buckets["stuck_pending"],
@ -1239,57 +1229,3 @@ async def api_create_machine_cassette_op(
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
return op 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