ADR-005 rollout step 2 (slice 1): capture cash-out settlements on the machine's dispense report #49
3 changed files with 368 additions and 2 deletions
feat(transport): report_dispense RPC — capture a cash-out on the machine's report, not on payment (ADR-005 §1–§2)
The structural fix for bitspire#122. _handle_payment used to spawn process_settlement the instant a cash_out payment landed — before the machine had begun to dispense — so a jam two seconds later found the legs already paid and `processed` was the honest answer. Payment is now the authorization; the machine's report is the capture. A cash_out lands as awaiting_dispense and is not distributed. The new `report_dispense` handler (identity from the VERIFIED transport sender, same as create_withdraw / get_machine_config) stores every report append-only and moves the settlement: dispense_confirmed → pending and distribution runs; some notes out → partial_pending, held whole (ADR-005 Decision 1, one distribution when the shortfall is resolved); nothing out → cash_owed, first on the worklist. A report naming remediates_txid moves the owed settlement it names to pending in full. Already-captured settlements are recorded but never moved — a report cannot un-pay legs. A byte-identical resend is acked without a new row. Both orders of arrival are handled: a report that precedes its payment (hold invoices settle after the dispense; the invoice listener can lag) is stored unlinked and adopted when the settlement is inserted, through the same transition. counts_uncertain on a report mirrors onto the machine immediately rather than at the next heartbeat. The state-event consumer mirrors cash_out_held_* onto dca_machines, including clearing it. Soft-fails like the other RPCs: without register_rpc the settlements sit in awaiting_dispense and surface as dispense_unreported — the honest state.
commit
44c2afa5bb
|
|
@ -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__ = [
|
||||
|
|
|
|||
299
dispense_transport.py
Normal file
299
dispense_transport.py
Normal file
|
|
@ -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'")
|
||||
65
tasks.py
65
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)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue