diff --git a/__init__.py b/__init__.py index 5bea74e..08ea43a 100644 --- a/__init__.py +++ b/__init__.py @@ -6,7 +6,6 @@ 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 @@ -69,11 +68,6 @@ 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 1858cb6..c4bffd5 100644 --- a/crud.py +++ b/crud.py @@ -27,8 +27,6 @@ from .models import ( DcaLpPreferences, DcaPayment, DcaSettlement, - DispenseReport, - DispenseReportIn, Machine, PublishCassettesPayload, SuperConfig, @@ -259,32 +257,6 @@ 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. @@ -739,171 +711,6 @@ 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", @@ -945,14 +752,7 @@ async def get_stuck_settlements_for_operator( ) -> dict: """Operator worklist of settlements that didn't process cleanly. - 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: + Returns a dict with four keyed lists: - '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, @@ -1017,48 +817,7 @@ 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, @@ -1111,10 +870,9 @@ async def mark_settlement_status( status: str, error_message: str | None = None, ) -> DcaSettlement | None: - """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.""" + """Status: 'pending' | 'processing' | 'processed' | '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 deleted file mode 100644 index 2b13efa..0000000 --- a/dispense_transport.py +++ /dev/null @@ -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'") diff --git a/migrations.py b/migrations.py index ed3d6f1..f96fe53 100644 --- a/migrations.py +++ b/migrations.py @@ -941,82 +941,3 @@ 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 f1422fc..dc6f533 100644 --- a/models.py +++ b/models.py @@ -67,13 +67,6 @@ 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 @@ -318,39 +311,19 @@ class DcaSettlement(BaseModel): fee_mismatch_sats: int | None = None bills_json: str | None cassettes_json: str | None - # 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) + # 'pending' (default at insert) # 'processing' (claim taken by distribution processor) # 'processed' (all legs paid) - # '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) + # 'partial' (operator marked partial-dispense after the fact) # '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 @@ -363,110 +336,6 @@ 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. # ============================================================================= @@ -701,13 +570,6 @@ 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 @@ -876,11 +738,6 @@ 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): @@ -946,21 +803,7 @@ class PublishCassettesPayload(BaseModel): # landed after the same problem. -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",) +CASSETTE_OP_TYPES = ("refill", "empty", "recount", "set_denomination") 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}") 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("position") + def _position_positive(cls, v): + if v <= 0: + raise ValueError(f"position must be > 0, got {v}") + return v def to_wire_dict(self) -> dict: """The published form. Drops the fields this op_type does not use, so @@ -1013,10 +850,8 @@ 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": @@ -1045,17 +880,11 @@ class CreateCassetteOpData(BaseModel): raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}") 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("position") + def _position_positive(cls, v): + if v <= 0: + raise ValueError(f"position must be > 0, got {v}") + return v @validator("bills") def _bills_positive(cls, v): @@ -1082,7 +911,6 @@ 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 54b6014..e8a31e2 100644 --- a/static/js/index.js +++ b/static/js/index.js @@ -70,10 +70,6 @@ window.app = Vue.createApp({ // Worklist (P9g) worklist: { - // ADR-005 §6 — owed-cash buckets first - cash_owed: [], - partial_pending: [], - dispense_unreported: [], rejected: [], errored: [], stuck_pending: [], @@ -376,33 +372,6 @@ 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', @@ -591,9 +560,6 @@ 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) + @@ -609,17 +575,11 @@ 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 + @@ -1248,50 +1208,12 @@ window.app = Vue.createApp({ openPartialDispense(settlement) { this.partialDispenseDialog.settlement = settlement this.partialDispenseDialog.mode = 'fraction' - // 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_fraction = null this.partialDispenseDialog.dispensed_sats = null - this.partialDispenseDialog.notes = settlement.dispense_error - ? `Machine reported: ${settlement.dispense_error_code || ''} ${settlement.dispense_raw_code || ''} — ${settlement.dispense_error}`.trim() - : '' + this.partialDispenseDialog.notes = '' 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 7001bae..6e5ac55 100644 --- a/tasks.py +++ b/tasks.py @@ -169,17 +169,8 @@ 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. 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" - ) + # 3) Insert + distribute. + settlement = await create_settlement_idempotent(data, initial_status="pending") if settlement is None: logger.error( 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"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 @@ -209,23 +196,6 @@ 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. @@ -468,38 +438,12 @@ 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 @@ -614,8 +558,3 @@ 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 77d9612..ad01503 100644 --- a/templates/spirekeeper/index.html +++ b/templates/spirekeeper/index.html @@ -661,12 +661,6 @@ @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 35ca48c..150d493 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -124,8 +124,6 @@ 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 deleted file mode 100644 index 363dc5d..0000000 --- a/tests/test_dispense_outcome.py +++ /dev/null @@ -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) diff --git a/views_api.py b/views_api.py index 7d6bf9f..33fb2de 100644 --- a/views_api.py +++ b/views_api.py @@ -19,7 +19,6 @@ 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 ( @@ -29,6 +28,15 @@ 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, @@ -77,7 +85,6 @@ from .distribution import ( process_settlement, settle_lp_balance, ) -from .fee_transport import publish_fee_config from .models import ( AppendSettlementNoteData, CassetteConfig, @@ -105,14 +112,6 @@ from .models import ( UpdateMachineData, UpdateSuperConfigData, ) -from .pairing import ( - PairingError, - PairResult, - RevokeResult, - default_relay_endpoint, - pair_spire, - revoke_spire, -) spirekeeper_api_router = APIRouter() @@ -769,13 +768,7 @@ async def api_list_stuck_settlements( ) -> StuckSettlementsResponse: """Operator worklist of settlements that didn't process cleanly. - 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: + Returns four lists: - 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 @@ -790,9 +783,6 @@ 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"], @@ -1239,57 +1229,3 @@ 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 - -