diff --git a/__init__.py b/__init__.py index 08ea43a..5bea74e 100644 --- a/__init__.py +++ b/__init__.py @@ -6,6 +6,7 @@ from loguru import logger from .cashin_transport import register_create_withdraw_rpc from .crud import db +from .dispense_transport import register_dispense_report_rpc from .machine_config_transport import register_machine_config_rpc from .nostr_transport_roster import register_with_lnbits as register_roster_with_lnbits from .tasks import wait_for_cassette_state_events, wait_for_paid_invoices @@ -68,6 +69,11 @@ def spirekeeper_start(): # config over the transport, leaving "awaiting configuration" with no # per-machine env provisioning. Soft-fails if register_rpc isn't exposed. register_machine_config_rpc() + # Dispense outcome capture (bitspire ADR-005 §2 / #122): register the + # report_dispense RPC. A cash-out settlement now waits in awaiting_dispense + # until the machine reports; the success report is what distributes it, + # a failure report puts the customer on the owed-cash worklist. + register_dispense_report_rpc() __all__ = [ diff --git a/dispense_transport.py b/dispense_transport.py new file mode 100644 index 0000000..2b13efa --- /dev/null +++ b/dispense_transport.py @@ -0,0 +1,299 @@ +""" +Dispense outcome capture: the `report_dispense` nostr-transport RPC +(bitspire ADR-005 §1-§2, aiolabs/bitspire#122). + +A cash-out used to be captured the instant its payment landed: `_handle_payment` +spawned distribution in the same breath, so when a dispenser jammed two seconds +later the legs were already paid and `processed` was the honest answer. Payment +is now the authorization and the machine's dispense report is the capture. A +`cash_out` settlement lands as `awaiting_dispense` and this handler moves it: + + dispense_confirmed → pending → distribution runs + some notes out, value short → partial_pending (held whole — ADR-005 + Decision 1: one distribution when the + shortfall is resolved) + nothing out → cash_owed (legs never run; the customer + is owed; first on the worklist) + report carries remediates_txid → the owed/partial settlement it names + goes to pending and distributes in full + +Every report is stored append-only in `dispense_reports` (lamassu-server's +`cash_out_actions` shape), including the success ones — the success report is +what captures. Identity is the VERIFIED transport sender, never the body, same +as `create_withdraw` / `get_machine_config`. Idempotent on (txid, at): the +machine resends until it gets an OK, and a byte-identical resend is acked +without a new row or a second transition. + +A report can also arrive BEFORE its payment lands (hold invoices settle after +the dispense; the invoice listener can lag). It is stored with no settlement; +`_handle_payment` adopts it when the settlement is inserted and applies the +same transition. See `apply_report_to_settlement`. +""" + +from __future__ import annotations + +import asyncio +from datetime import datetime + +from loguru import logger + +from .crud import ( + apply_dispense_outcome, + get_dispense_report, + get_machine_by_atm_pubkey_hex, + get_settlement_by_txid, + insert_dispense_report, + link_dispense_reports_to_settlement, + set_machine_counts_uncertain, +) +from .models import DcaSettlement, DispenseReport, DispenseReportIn, Machine + +_RPC_NAME = "report_dispense" + +# Statuses a first report may move. Anything else (processed, pending, +# processing, errored, rejected, partial, refunded) is recorded but not moved — +# a report cannot un-pay legs, and a late report for an already-captured sale +# is information, not an instruction. +_CAPTURABLE = ("awaiting_dispense",) +# Statuses a remediation report may close out. +_OWED = ("cash_owed", "partial_pending") + +# Strong references to in-flight distribution tasks, same reason as tasks.py. +_inflight: set[asyncio.Task] = set() + + +def _outcome_status(report: DispenseReportIn) -> str: + if report.dispense_confirmed: + return "pending" + if report.total_dispensed_notes > 0: + return "partial_pending" + return "cash_owed" + + +def _spawn_distribution(settlement_id: str) -> None: + # Lazy import: distribution imports crud, and crud is what this module + # already depends on; importing at module load would make a cycle. + from .distribution import process_settlement + + task = asyncio.create_task(process_settlement(settlement_id)) + _inflight.add(task) + task.add_done_callback(_inflight.discard) + + +async def apply_report_to_settlement( + settlement: DcaSettlement, + report: DispenseReportIn, + machine: Machine, +) -> str: + """Move `settlement` according to `report`. Returns the resulting status. + + Shared by the RPC handler (report after payment) and `_handle_payment` + (payment after report), so both orders of arrival take the same path. + """ + reported_at = datetime.fromtimestamp(report.at) + + if report.remediates_txid: + # Handled by the caller against the settlement the remediation names; + # for the remediation's OWN txid there is nothing to capture. + return settlement.status + + if settlement.status not in _CAPTURABLE: + logger.info( + f"spirekeeper: report_dispense for settlement {settlement.id} in status " + f"{settlement.status!r} — recorded, not moved (txid={report.txid})" + ) + return settlement.status + + new_status = _outcome_status(report) + updated = await apply_dispense_outcome( + settlement.id, report, new_status, reported_at + ) + status = updated.status if updated else new_status + + if new_status == "pending": + _spawn_distribution(settlement.id) + logger.info( + f"spirekeeper: dispense CONFIRMED for settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}, " + f"{report.fiat_cents / 100:.2f} {report.currency}) — distributing" + ) + elif new_status == "partial_pending": + logger.warning( + f"spirekeeper: PARTIAL dispense for settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}): " + f"{report.dispensed_fiat_cents / 100:.2f} of " + f"{report.fiat_cents / 100:.2f} " + f"{report.currency} left the machine; {report.error_code or 'no error'} " + f"{report.raw_code or ''}. Held until the operator resolves the shortfall." + ) + else: + logger.error( + f"spirekeeper: CASH OWED — settlement {settlement.id} " + f"(machine={machine.id}, txid={report.txid}): customer paid " + f"{report.fiat_cents / 100:.2f} {report.currency}, nothing dispensed; " + f"{report.error_code or 'no error'} {report.raw_code or ''}: " + f"{report.error or ''}" + ) + return status + + +async def _apply_remediation( + machine: Machine, report: DispenseReportIn +) -> tuple[str | None, str | None]: + """A manual dispense closed out an earlier failed txid. Returns + (settlement_id, status) of the remediated settlement, or (None, None).""" + assert report.remediates_txid + target = await get_settlement_by_txid(machine.id, report.remediates_txid) + if target is None: + logger.warning( + f"spirekeeper: remediation report {report.txid} names txid " + f"{report.remediates_txid} with no settlement on this server" + ) + return None, None + if target.status not in _OWED: + logger.info( + f"spirekeeper: remediation report {report.txid} for settlement " + f"{target.id} in status {target.status!r} — recorded, not moved" + ) + return target.id, target.status + if not report.dispense_confirmed: + logger.warning( + f"spirekeeper: remediation report {report.txid} for settlement " + f"{target.id} did not itself confirm — settlement stays {target.status}" + ) + return target.id, target.status + # The customer is whole: the sale is the full amount. Keep the ORIGINAL + # report's columns on the settlement (that is what happened at the sale); + # the remediation is its own dispense_reports row. + from .crud import mark_settlement_status + + await mark_settlement_status(target.id, "pending", None) + _spawn_distribution(target.id) + logger.info( + f"spirekeeper: settlement {target.id} remediated by manual dispense " + f"{report.txid} — distributing in full" + ) + return target.id, "pending" + + +async def handle_report_dispense(auth, request) -> dict: + """nostr-transport RPC handler. `auth` is the roster-resolved auth context + (unused — the machine is identified from the signature); `request` is a + NostrRpcRequest with `body` and `sender_pubkey` (verified). + + Returns `{txid, received, settlement_status}`; raises ValueError (→ transport + ERROR reply) for an unpaired sender or a malformed body. The machine treats + anything but OK as "resend later", so a malformed report is retried — which + is right: the bug is on one side or the other and the row must not be lost. + """ + sender = (request.sender_pubkey or "").lower() + if not sender: + raise ValueError("missing verified sender_pubkey") + machine = await get_machine_by_atm_pubkey_hex(sender) + if machine is None: + raise ValueError("sender pubkey is not a paired machine") + + try: + report = DispenseReportIn(**(request.body or {})) + except Exception as exc: # pydantic ValidationError, TypeError + raise ValueError(f"invalid report_dispense body: {exc}") from exc + + # Idempotency: the machine resends until acked. + existing = await get_dispense_report(machine.id, report.txid, report.at) + if existing is not None: + settlement = ( + await get_settlement_by_txid(machine.id, report.txid) + if existing.settlement_id + else None + ) + return { + "txid": report.txid, + "received": True, + "settlement_status": settlement.status if settlement else None, + "duplicate": True, + } + + if report.counts_uncertain: + # The state document carries this too; mirroring it here means the + # dashboard learns at report time rather than at the next heartbeat. + await set_machine_counts_uncertain(machine.id, datetime.now()) + + settlement = await get_settlement_by_txid(machine.id, report.txid) + stored: DispenseReport = await insert_dispense_report( + machine.id, settlement.id if settlement else None, report + ) + + status: str | None + if report.remediates_txid: + _, status = await _apply_remediation(machine, report) + elif settlement is None: + # Payment not landed yet (or never will). Kept unlinked; + # _handle_payment adopts it when the settlement is inserted. + logger.warning( + f"spirekeeper: report_dispense {report.txid} from machine {machine.id} " + f"has no settlement yet (confirmed={report.dispense_confirmed}) — " + f"stored unlinked, will attach when the payment lands" + ) + status = None + else: + status = await apply_report_to_settlement(settlement, report, machine) + + logger.info( + f"spirekeeper: report_dispense stored id={stored.id} machine={machine.id} " + f"txid={report.txid} confirmed={report.dispense_confirmed} → {status}" + ) + return {"txid": report.txid, "received": True, "settlement_status": status} + + +async def adopt_unlinked_report( + settlement: DcaSettlement, machine: Machine, report_row: DispenseReport +) -> str: + """`_handle_payment` found a report that arrived before the payment: link + it and apply the same transition the live handler would have.""" + import json as _json + + await link_dispense_reports_to_settlement( + machine.id, report_row.txid, settlement.id + ) + report = DispenseReportIn( + txid=report_row.txid, + payment_hash=report_row.payment_hash, + dispense_confirmed=report_row.dispense_confirmed, + error=report_row.error, + error_code=report_row.error_code, + raw_code=report_row.raw_code, + error_class=report_row.error_class, + fiat_cents=report_row.fiat_cents, + currency=report_row.currency, + bills=_json.loads(report_row.bills_json or "[]"), + cassettes=_json.loads(report_row.cassettes_json or "[]"), + counts_uncertain=report_row.counts_uncertain, + remediates_txid=report_row.remediates_txid, + at=int(report_row.reported_at.timestamp()), + ) + logger.info( + f"spirekeeper: adopting early dispense report {report_row.id} for " + f"settlement {settlement.id} (report preceded the payment)" + ) + return await apply_report_to_settlement(settlement, report, machine) + + +def register_dispense_report_rpc() -> None: + """Register `report_dispense` with the lnbits nostr transport. Soft-fails + if the transport doesn't expose `register_rpc` (older lnbits) — then + cash-out settlements wait in awaiting_dispense and surface on the + worklist as dispense_unreported, which is the honest state.""" + try: + from lnbits.core.services.nostr_transport.dispatcher import ( # type: ignore + AUTH_ACCOUNT, + register_rpc, + ) + except ImportError: + logger.warning( + "spirekeeper: nostr-transport register_rpc unavailable; " + "'report_dispense' not registered (ADR-005 capture disabled — " + "cash-out settlements will sit in awaiting_dispense)" + ) + return + register_rpc(_RPC_NAME, handle_report_dispense, AUTH_ACCOUNT) + logger.info("spirekeeper: registered nostr-transport RPC 'report_dispense'") diff --git a/tasks.py b/tasks.py index 6e5ac55..7001bae 100644 --- a/tasks.py +++ b/tasks.py @@ -169,8 +169,17 @@ async def _handle_payment(payment: Payment) -> None: if isinstance(nostr_event_id, str) and nostr_event_id: data.bitspire_event_id = nostr_event_id - # 3) Insert + distribute. - settlement = await create_settlement_idempotent(data, initial_status="pending") + # 3) Insert. ADR-005 §1: payment is authorization, the machine's dispense + # report is capture. A cash_out waits in `awaiting_dispense` for that + # report (handled in dispense_transport); distribution runs only once it + # says dispense_confirmed. A cash_in has no dispense and proceeds as + # before. Before this gate the legs were paid sub-second, before the + # machine had even begun to dispense — which is how a jam read + # `processed` on 2026-10-09 (bitspire#122). + is_cash_out = data.tx_type == "cash_out" + settlement = await create_settlement_idempotent( + data, initial_status="awaiting_dispense" if is_cash_out else "pending" + ) if settlement is None: logger.error( f"spirekeeper: failed to insert settlement for " @@ -185,6 +194,10 @@ async def _handle_payment(payment: Payment) -> None: f"(super_fee={data.platform_fee_sats} " f"operator_fee={data.operator_fee_sats})" ) + if is_cash_out: + await _await_dispense_or_adopt(settlement, machine, data) + return + # Spawn distribution on a background task so the LNbits invoice queue # (shared across all extensions) keeps draining while we move sats. # Concurrency-safe: process_settlement uses claim_settlement_for_processing @@ -196,6 +209,23 @@ async def _handle_payment(payment: Payment) -> None: task.add_done_callback(_inflight_distributions.discard) +async def _await_dispense_or_adopt( + settlement, machine: Machine, data: CreateDcaSettlementData +) -> None: + """A cash_out waits for the machine's dispense report (ADR-005 §1) — unless + the report is already here. Under hold invoices the payment settles AFTER + the dispense, and the invoice listener can lag the transport, so an orphan + report for this txid is adopted and applied now.""" + if settlement.status != "awaiting_dispense" or not data.bitspire_txid: + return + from .crud import get_latest_unlinked_dispense_report + from .dispense_transport import adopt_unlinked_report + + early = await get_latest_unlinked_dispense_report(machine.id, data.bitspire_txid) + if early is not None: + await adopt_unlinked_report(settlement, machine, early) + + async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) -> None: """Insert a minimal `dca_settlements` row with `status='rejected'` and the exception message for operator forensics. @@ -438,12 +468,38 @@ async def _record_counts_uncertainty( await set_machine_counts_uncertain(machine_id, since) +async def _record_cash_out_hold( + machine_id: str, payload, set_machine_cash_out_hold +) -> None: + """Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row. + + Written on every state event, including when absent, because the machine + clearing the hold — after an operator recount or resume_cash_out — matters + exactly as much as it setting one. + """ + from datetime import datetime as _datetime + from datetime import timezone as _timezone + + since = None + if payload.cash_out_held_since is not None: + since = _datetime.fromtimestamp( + int(payload.cash_out_held_since), tz=_timezone.utc + ) + await set_machine_cash_out_hold( + machine_id, + since, + payload.cash_out_held_reason if since else None, + payload.cash_out_held_code if since else None, + ) + + async def _handle_cassette_state_event( event_message, get_machine_by_atm_pubkey_hex, apply_reported_state, mark_cassette_ops_acked, set_machine_counts_uncertain, + set_machine_cash_out_hold=None, ) -> None: """Verify signature, resolve the operator's signer, decrypt via the signer abstraction (bunker round-trip for RemoteBunkerSigner; direct @@ -558,3 +614,8 @@ async def _handle_cassette_state_event( # _record_op_acknowledgements for why. Same for the uncertainty marker. await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked) await _record_counts_uncertainty(machine.id, payload, set_machine_counts_uncertain) + if set_machine_cash_out_hold is None: + from .crud import set_machine_cash_out_hold as _default_set_hold + + set_machine_cash_out_hold = _default_set_hold + await _record_cash_out_hold(machine.id, payload, set_machine_cash_out_hold)