ADR-005 rollout step 2 (slice 1): capture cash-out settlements on the machine's dispense report #49
11 changed files with 1602 additions and 32 deletions
|
|
@ -6,6 +6,7 @@ from loguru import logger
|
||||||
|
|
||||||
from .cashin_transport import register_create_withdraw_rpc
|
from .cashin_transport import register_create_withdraw_rpc
|
||||||
from .crud import db
|
from .crud import db
|
||||||
|
from .dispense_transport import register_dispense_report_rpc
|
||||||
from .machine_config_transport import register_machine_config_rpc
|
from .machine_config_transport import register_machine_config_rpc
|
||||||
from .nostr_transport_roster import register_with_lnbits as register_roster_with_lnbits
|
from .nostr_transport_roster import register_with_lnbits as register_roster_with_lnbits
|
||||||
from .tasks import wait_for_cassette_state_events, wait_for_paid_invoices
|
from .tasks import wait_for_cassette_state_events, wait_for_paid_invoices
|
||||||
|
|
@ -68,6 +69,11 @@ def spirekeeper_start():
|
||||||
# config over the transport, leaving "awaiting configuration" with no
|
# config over the transport, leaving "awaiting configuration" with no
|
||||||
# per-machine env provisioning. Soft-fails if register_rpc isn't exposed.
|
# per-machine env provisioning. Soft-fails if register_rpc isn't exposed.
|
||||||
register_machine_config_rpc()
|
register_machine_config_rpc()
|
||||||
|
# Dispense outcome capture (bitspire ADR-005 §2 / #122): register the
|
||||||
|
# report_dispense RPC. A cash-out settlement now waits in awaiting_dispense
|
||||||
|
# until the machine reports; the success report is what distributes it,
|
||||||
|
# a failure report puts the customer on the owed-cash worklist.
|
||||||
|
register_dispense_report_rpc()
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
|
|
|
||||||
250
crud.py
250
crud.py
|
|
@ -27,6 +27,8 @@ from .models import (
|
||||||
DcaLpPreferences,
|
DcaLpPreferences,
|
||||||
DcaPayment,
|
DcaPayment,
|
||||||
DcaSettlement,
|
DcaSettlement,
|
||||||
|
DispenseReport,
|
||||||
|
DispenseReportIn,
|
||||||
Machine,
|
Machine,
|
||||||
PublishCassettesPayload,
|
PublishCassettesPayload,
|
||||||
SuperConfig,
|
SuperConfig,
|
||||||
|
|
@ -257,6 +259,32 @@ async def set_machine_unpaired(machine_id: str) -> Machine | None:
|
||||||
return await get_machine(machine_id)
|
return await get_machine(machine_id)
|
||||||
|
|
||||||
|
|
||||||
|
async def set_machine_cash_out_hold(
|
||||||
|
machine_id: str,
|
||||||
|
since: datetime | None,
|
||||||
|
reason: str | None,
|
||||||
|
code: str | None,
|
||||||
|
) -> None:
|
||||||
|
"""Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row.
|
||||||
|
|
||||||
|
Written on every state event, including when it is None: the machine
|
||||||
|
clearing the hold — after an operator recount or resume_cash_out — is as
|
||||||
|
important as it setting one. `updated_at` is left alone for the same
|
||||||
|
reason as counts_uncertain_since: this is the machine reporting about
|
||||||
|
itself on a heartbeat, not an operator editing the machine.
|
||||||
|
"""
|
||||||
|
await db.execute(
|
||||||
|
"""
|
||||||
|
UPDATE spirekeeper.dca_machines
|
||||||
|
SET cash_out_held_since = :since,
|
||||||
|
cash_out_held_reason = :reason,
|
||||||
|
cash_out_held_code = :code
|
||||||
|
WHERE id = :id
|
||||||
|
""",
|
||||||
|
{"id": machine_id, "since": since, "reason": reason, "code": code},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
async def set_machine_counts_uncertain(machine_id: str, since: datetime | None) -> None:
|
async def set_machine_counts_uncertain(machine_id: str, since: datetime | None) -> None:
|
||||||
"""Record (or clear) the machine's own "I can't vouch for these counts"
|
"""Record (or clear) the machine's own "I can't vouch for these counts"
|
||||||
marker, straight from its state document.
|
marker, straight from its state document.
|
||||||
|
|
@ -711,6 +739,171 @@ async def create_settlement_idempotent(
|
||||||
return await get_settlement(settlement_id)
|
return await get_settlement(settlement_id)
|
||||||
|
|
||||||
|
|
||||||
|
async def get_settlement_by_txid(
|
||||||
|
machine_id: str, bitspire_txid: str
|
||||||
|
) -> DcaSettlement | None:
|
||||||
|
"""The settlement a machine report refers to. `bitspire_txid` comes from
|
||||||
|
the invoice's extra.txid, stamped by the machine at create_invoice time,
|
||||||
|
so it is the natural join for a report that names the same txid."""
|
||||||
|
return await db.fetchone(
|
||||||
|
"SELECT * FROM spirekeeper.dca_settlements "
|
||||||
|
"WHERE machine_id = :mid AND bitspire_txid = :txid",
|
||||||
|
{"mid": machine_id, "txid": bitspire_txid},
|
||||||
|
DcaSettlement,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def apply_dispense_outcome(
|
||||||
|
settlement_id: str,
|
||||||
|
report: DispenseReportIn,
|
||||||
|
new_status: str,
|
||||||
|
reported_at: datetime,
|
||||||
|
) -> DcaSettlement | None:
|
||||||
|
"""Copy the machine's report onto the settlement and move it (ADR-005 §1).
|
||||||
|
|
||||||
|
Fills the never-before-written bills_json / cassettes_json with what
|
||||||
|
actually came out, not what was provisioned. `error_message` is left to
|
||||||
|
the distribution path; the dispense error lives in its own columns.
|
||||||
|
"""
|
||||||
|
import json as _json
|
||||||
|
|
||||||
|
await db.execute(
|
||||||
|
"""
|
||||||
|
UPDATE spirekeeper.dca_settlements
|
||||||
|
SET status = :status,
|
||||||
|
dispense_confirmed = :confirmed,
|
||||||
|
dispense_error = :error,
|
||||||
|
dispense_error_code = :error_code,
|
||||||
|
dispense_raw_code = :raw_code,
|
||||||
|
dispense_error_class = :error_class,
|
||||||
|
dispense_reported_at = :reported_at,
|
||||||
|
dispensed_fiat_cents = :dispensed_fiat_cents,
|
||||||
|
bills_json = :bills_json,
|
||||||
|
cassettes_json = :cassettes_json,
|
||||||
|
processing_claim = NULL
|
||||||
|
WHERE id = :id
|
||||||
|
""",
|
||||||
|
{
|
||||||
|
"id": settlement_id,
|
||||||
|
"status": new_status,
|
||||||
|
"confirmed": report.dispense_confirmed,
|
||||||
|
"error": report.error,
|
||||||
|
"error_code": report.error_code,
|
||||||
|
"raw_code": report.raw_code,
|
||||||
|
"error_class": report.error_class,
|
||||||
|
"reported_at": reported_at,
|
||||||
|
"dispensed_fiat_cents": report.dispensed_fiat_cents,
|
||||||
|
"bills_json": _json.dumps([b.dict() for b in report.bills]),
|
||||||
|
"cassettes_json": _json.dumps([c.dict() for c in report.cassettes]),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
return await get_settlement(settlement_id)
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Dispense reports (ADR-005 §2) — append-only
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def get_dispense_report(
|
||||||
|
machine_id: str, txid: str, reported_at: int
|
||||||
|
) -> DispenseReport | None:
|
||||||
|
"""A report is identified by (machine, txid, at): the machine resends the
|
||||||
|
same report until acked, and a byte-identical resend must not grow the
|
||||||
|
log. A remediation report for the same txid carries a later `at`."""
|
||||||
|
return await db.fetchone(
|
||||||
|
"SELECT * FROM spirekeeper.dispense_reports "
|
||||||
|
"WHERE machine_id = :mid AND txid = :txid AND reported_at = :at",
|
||||||
|
{"mid": machine_id, "txid": txid, "at": datetime.fromtimestamp(reported_at)},
|
||||||
|
DispenseReport,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def insert_dispense_report(
|
||||||
|
machine_id: str, settlement_id: str | None, report: DispenseReportIn
|
||||||
|
) -> DispenseReport:
|
||||||
|
import json as _json
|
||||||
|
|
||||||
|
report_id = urlsafe_short_hash()
|
||||||
|
await db.execute(
|
||||||
|
"""
|
||||||
|
INSERT INTO spirekeeper.dispense_reports
|
||||||
|
(id, machine_id, settlement_id, txid, payment_hash, dispense_confirmed,
|
||||||
|
error, error_code, raw_code, error_class, fiat_cents, currency,
|
||||||
|
bills_json, cassettes_json, counts_uncertain, remediates_txid,
|
||||||
|
reported_at, received_at)
|
||||||
|
VALUES (:id, :machine_id, :settlement_id, :txid, :payment_hash,
|
||||||
|
:confirmed, :error, :error_code, :raw_code, :error_class,
|
||||||
|
:fiat_cents, :currency, :bills_json, :cassettes_json,
|
||||||
|
:counts_uncertain, :remediates_txid, :reported_at, :received_at)
|
||||||
|
""",
|
||||||
|
{
|
||||||
|
"id": report_id,
|
||||||
|
"machine_id": machine_id,
|
||||||
|
"settlement_id": settlement_id,
|
||||||
|
"txid": report.txid,
|
||||||
|
"payment_hash": report.payment_hash,
|
||||||
|
"confirmed": report.dispense_confirmed,
|
||||||
|
"error": report.error,
|
||||||
|
"error_code": report.error_code,
|
||||||
|
"raw_code": report.raw_code,
|
||||||
|
"error_class": report.error_class,
|
||||||
|
"fiat_cents": report.fiat_cents,
|
||||||
|
"currency": report.currency,
|
||||||
|
"bills_json": _json.dumps([b.dict() for b in report.bills]),
|
||||||
|
"cassettes_json": _json.dumps([c.dict() for c in report.cassettes]),
|
||||||
|
"counts_uncertain": report.counts_uncertain,
|
||||||
|
"remediates_txid": report.remediates_txid,
|
||||||
|
"reported_at": datetime.fromtimestamp(report.at),
|
||||||
|
"received_at": datetime.now(),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
row = await db.fetchone(
|
||||||
|
"SELECT * FROM spirekeeper.dispense_reports WHERE id = :id",
|
||||||
|
{"id": report_id},
|
||||||
|
DispenseReport,
|
||||||
|
)
|
||||||
|
assert row is not None, "Newly inserted dispense report couldn't be retrieved"
|
||||||
|
return row
|
||||||
|
|
||||||
|
|
||||||
|
async def link_dispense_reports_to_settlement(
|
||||||
|
machine_id: str, txid: str, settlement_id: str
|
||||||
|
) -> int:
|
||||||
|
"""A report can arrive before its payment lands (hold invoices settle
|
||||||
|
after the dispense; the invoice listener can lag). When the settlement is
|
||||||
|
finally inserted, adopt the orphan rows."""
|
||||||
|
result = await db.execute(
|
||||||
|
"UPDATE spirekeeper.dispense_reports SET settlement_id = :sid "
|
||||||
|
"WHERE machine_id = :mid AND txid = :txid AND settlement_id IS NULL",
|
||||||
|
{"sid": settlement_id, "mid": machine_id, "txid": txid},
|
||||||
|
)
|
||||||
|
return getattr(result, "rowcount", 0) or 0
|
||||||
|
|
||||||
|
|
||||||
|
async def get_latest_unlinked_dispense_report(
|
||||||
|
machine_id: str, txid: str
|
||||||
|
) -> DispenseReport | None:
|
||||||
|
return await db.fetchone(
|
||||||
|
"SELECT * FROM spirekeeper.dispense_reports "
|
||||||
|
"WHERE machine_id = :mid AND txid = :txid AND settlement_id IS NULL "
|
||||||
|
"ORDER BY reported_at DESC LIMIT 1",
|
||||||
|
{"mid": machine_id, "txid": txid},
|
||||||
|
DispenseReport,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def get_dispense_reports_for_settlement(
|
||||||
|
settlement_id: str,
|
||||||
|
) -> list[DispenseReport]:
|
||||||
|
return await db.fetchall(
|
||||||
|
"SELECT * FROM spirekeeper.dispense_reports WHERE settlement_id = :sid "
|
||||||
|
"ORDER BY reported_at ASC",
|
||||||
|
{"sid": settlement_id},
|
||||||
|
DispenseReport,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
async def get_settlement(settlement_id: str) -> DcaSettlement | None:
|
async def get_settlement(settlement_id: str) -> DcaSettlement | None:
|
||||||
return await db.fetchone(
|
return await db.fetchone(
|
||||||
"SELECT * FROM spirekeeper.dca_settlements WHERE id = :id",
|
"SELECT * FROM spirekeeper.dca_settlements WHERE id = :id",
|
||||||
|
|
@ -752,7 +945,14 @@ async def get_stuck_settlements_for_operator(
|
||||||
) -> dict:
|
) -> dict:
|
||||||
"""Operator worklist of settlements that didn't process cleanly.
|
"""Operator worklist of settlements that didn't process cleanly.
|
||||||
|
|
||||||
Returns a dict with four keyed lists:
|
Returns a dict with seven keyed lists. The first three are ADR-005 §6 —
|
||||||
|
the only ones whose meaning is "a customer is owed money":
|
||||||
|
- 'cash_owed': the machine reported nothing dispensed; legs never ran.
|
||||||
|
- 'partial_pending': some notes out, value short; held whole until the
|
||||||
|
operator records the resolution.
|
||||||
|
- 'dispense_unreported': awaiting_dispense older than the threshold —
|
||||||
|
the machine never reported (crashed, offline, or an old build).
|
||||||
|
Then the original four:
|
||||||
- 'rejected': any status='rejected' (Nostr attribution cross-check
|
- 'rejected': any status='rejected' (Nostr attribution cross-check
|
||||||
failed — signer didn't match the machine identity). Distinct
|
failed — signer didn't match the machine identity). Distinct
|
||||||
from 'errored' because retry is wrong: the row was misrouted,
|
from 'errored' because retry is wrong: the row was misrouted,
|
||||||
|
|
@ -817,7 +1017,48 @@ async def get_stuck_settlements_for_operator(
|
||||||
{"uid": operator_user_id, "threshold": threshold_at},
|
{"uid": operator_user_id, "threshold": threshold_at},
|
||||||
DcaSettlement,
|
DcaSettlement,
|
||||||
)
|
)
|
||||||
|
# ADR-005 §6 — the owed-cash buckets. cash_owed / partial_pending are
|
||||||
|
# stored statuses; dispense_unreported is derived: a cash-out that landed
|
||||||
|
# and never heard from its machine within the threshold.
|
||||||
|
cash_owed = await db.fetchall(
|
||||||
|
"""
|
||||||
|
SELECT s.*
|
||||||
|
FROM spirekeeper.dca_settlements s
|
||||||
|
JOIN spirekeeper.dca_machines m ON m.id = s.machine_id
|
||||||
|
WHERE m.operator_user_id = :uid AND s.status = 'cash_owed'
|
||||||
|
ORDER BY s.created_at DESC
|
||||||
|
""",
|
||||||
|
{"uid": operator_user_id},
|
||||||
|
DcaSettlement,
|
||||||
|
)
|
||||||
|
partial_pending = await db.fetchall(
|
||||||
|
"""
|
||||||
|
SELECT s.*
|
||||||
|
FROM spirekeeper.dca_settlements s
|
||||||
|
JOIN spirekeeper.dca_machines m ON m.id = s.machine_id
|
||||||
|
WHERE m.operator_user_id = :uid AND s.status = 'partial_pending'
|
||||||
|
ORDER BY s.created_at DESC
|
||||||
|
""",
|
||||||
|
{"uid": operator_user_id},
|
||||||
|
DcaSettlement,
|
||||||
|
)
|
||||||
|
dispense_unreported = await db.fetchall(
|
||||||
|
"""
|
||||||
|
SELECT s.*
|
||||||
|
FROM spirekeeper.dca_settlements s
|
||||||
|
JOIN spirekeeper.dca_machines m ON m.id = s.machine_id
|
||||||
|
WHERE m.operator_user_id = :uid
|
||||||
|
AND s.status = 'awaiting_dispense'
|
||||||
|
AND s.created_at < :threshold
|
||||||
|
ORDER BY s.created_at ASC
|
||||||
|
""",
|
||||||
|
{"uid": operator_user_id, "threshold": threshold_at},
|
||||||
|
DcaSettlement,
|
||||||
|
)
|
||||||
return {
|
return {
|
||||||
|
"cash_owed": cash_owed,
|
||||||
|
"partial_pending": partial_pending,
|
||||||
|
"dispense_unreported": dispense_unreported,
|
||||||
"rejected": rejected,
|
"rejected": rejected,
|
||||||
"errored": errored,
|
"errored": errored,
|
||||||
"stuck_pending": stuck_pending,
|
"stuck_pending": stuck_pending,
|
||||||
|
|
@ -870,9 +1111,10 @@ async def mark_settlement_status(
|
||||||
status: str,
|
status: str,
|
||||||
error_message: str | None = None,
|
error_message: str | None = None,
|
||||||
) -> DcaSettlement | None:
|
) -> DcaSettlement | None:
|
||||||
"""Status: 'pending' | 'processing' | 'processed' | 'partial' |
|
"""Status: 'awaiting_dispense' | 'pending' | 'processing' | 'processed' |
|
||||||
'refunded' | 'errored'. Clears processing_claim on terminal states so a
|
'partial_pending' | 'cash_owed' | 'partial' | 'refunded' | 'errored'.
|
||||||
fresh claim attempt won't see a stale token."""
|
Clears processing_claim on terminal states so a fresh claim attempt won't
|
||||||
|
see a stale token."""
|
||||||
await db.execute(
|
await db.execute(
|
||||||
"""
|
"""
|
||||||
UPDATE spirekeeper.dca_settlements
|
UPDATE spirekeeper.dca_settlements
|
||||||
|
|
|
||||||
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'")
|
||||||
|
|
@ -941,3 +941,82 @@ async def m015_add_cassette_state_seq(db):
|
||||||
await db.execute(
|
await db.execute(
|
||||||
"ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER"
|
"ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def m016_dispense_outcome(db):
|
||||||
|
"""The dispense outcome becomes a first-class fact (bitspire ADR-005 §2).
|
||||||
|
|
||||||
|
Until now a cash-out settlement was captured the instant the payment
|
||||||
|
landed: `_handle_payment` spawned distribution in the same breath, so by
|
||||||
|
the time a dispenser jammed two seconds later the legs were already paid
|
||||||
|
and the dashboard honestly reported `processed`. The machine now reports
|
||||||
|
every cash-out's outcome over a `report_dispense` RPC and the settlement
|
||||||
|
waits for it (`awaiting_dispense`) before anything moves.
|
||||||
|
|
||||||
|
`dispense_reports` is append-only, one row per report the machine sent —
|
||||||
|
lamassu-server's `cash_out_actions` shape — so a retry, a late report and
|
||||||
|
a remediation report are all visible as distinct rows. `settlement_id` is
|
||||||
|
NULL for a report whose payment this server never saw.
|
||||||
|
|
||||||
|
The settlement carries the three lamassu fields (dispense_confirmed,
|
||||||
|
error, error_code) plus raw_code / error_class and the fiat value that
|
||||||
|
actually left the machine, so the partial-dispense dialog can be
|
||||||
|
pre-filled with the hardware's own number instead of a typed one.
|
||||||
|
|
||||||
|
The machine's cash-out hold is mirrored onto its registry row beside
|
||||||
|
counts_uncertain_since: a latched machine refuses cash-out until an
|
||||||
|
operator recounts or publishes `resume_cash_out`, and the dashboard needs
|
||||||
|
to show that and offer the button.
|
||||||
|
"""
|
||||||
|
await db.execute(
|
||||||
|
f"""
|
||||||
|
CREATE TABLE IF NOT EXISTS spirekeeper.dispense_reports (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
machine_id TEXT NOT NULL,
|
||||||
|
settlement_id TEXT,
|
||||||
|
txid TEXT NOT NULL,
|
||||||
|
payment_hash TEXT,
|
||||||
|
dispense_confirmed BOOLEAN NOT NULL,
|
||||||
|
error TEXT,
|
||||||
|
error_code TEXT,
|
||||||
|
raw_code TEXT,
|
||||||
|
error_class TEXT,
|
||||||
|
fiat_cents INTEGER NOT NULL,
|
||||||
|
currency TEXT NOT NULL,
|
||||||
|
bills_json TEXT NOT NULL,
|
||||||
|
cassettes_json TEXT NOT NULL,
|
||||||
|
counts_uncertain BOOLEAN NOT NULL DEFAULT false,
|
||||||
|
remediates_txid TEXT,
|
||||||
|
reported_at TIMESTAMP NOT NULL,
|
||||||
|
received_at TIMESTAMP NOT NULL DEFAULT {db.timestamp_now}
|
||||||
|
);
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
await db.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS dispense_reports_txid_idx "
|
||||||
|
"ON dispense_reports (machine_id, txid)"
|
||||||
|
)
|
||||||
|
await db.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS dispense_reports_settlement_idx "
|
||||||
|
"ON dispense_reports (settlement_id)"
|
||||||
|
)
|
||||||
|
for col, typ in (
|
||||||
|
("dispense_confirmed", "BOOLEAN"),
|
||||||
|
("dispense_error", "TEXT"),
|
||||||
|
("dispense_error_code", "TEXT"),
|
||||||
|
("dispense_raw_code", "TEXT"),
|
||||||
|
("dispense_error_class", "TEXT"),
|
||||||
|
("dispense_reported_at", "TIMESTAMP"),
|
||||||
|
("dispensed_fiat_cents", "INTEGER"),
|
||||||
|
):
|
||||||
|
await db.execute(
|
||||||
|
f"ALTER TABLE spirekeeper.dca_settlements ADD COLUMN {col} {typ}"
|
||||||
|
)
|
||||||
|
for col, typ in (
|
||||||
|
("cash_out_held_since", "TIMESTAMP"),
|
||||||
|
("cash_out_held_reason", "TEXT"),
|
||||||
|
("cash_out_held_code", "TEXT"),
|
||||||
|
):
|
||||||
|
await db.execute(
|
||||||
|
f"ALTER TABLE spirekeeper.dca_machines ADD COLUMN {col} {typ}"
|
||||||
|
)
|
||||||
|
|
|
||||||
200
models.py
200
models.py
|
|
@ -67,6 +67,13 @@ class Machine(BaseModel):
|
||||||
# count of what left the bay; cleared by the machine's own report. The
|
# count of what left the bay; cleared by the machine's own report. The
|
||||||
# dashboard turns this into a prompt to open the bay and recount.
|
# dashboard turns this into a prompt to open the bay and recount.
|
||||||
counts_uncertain_since: datetime | None = None
|
counts_uncertain_since: datetime | None = None
|
||||||
|
# ADR-005 §5: the machine has latched cash-out off after a terminal
|
||||||
|
# dispenser fault. Mirrored from its state document. Cleared when the
|
||||||
|
# machine reports the hold released (an operator recount or a
|
||||||
|
# resume_cash_out op). The dashboard shows it and offers the button.
|
||||||
|
cash_out_held_since: datetime | None = None
|
||||||
|
cash_out_held_reason: str | None = None
|
||||||
|
cash_out_held_code: str | None = None
|
||||||
created_at: datetime
|
created_at: datetime
|
||||||
updated_at: datetime
|
updated_at: datetime
|
||||||
|
|
||||||
|
|
@ -311,19 +318,39 @@ class DcaSettlement(BaseModel):
|
||||||
fee_mismatch_sats: int | None = None
|
fee_mismatch_sats: int | None = None
|
||||||
bills_json: str | None
|
bills_json: str | None
|
||||||
cassettes_json: str | None
|
cassettes_json: str | None
|
||||||
# 'pending' (default at insert)
|
# Lifecycle (bitspire ADR-005 §1 — payment is authorization, the
|
||||||
|
# machine's dispense report is capture; distribution waits for capture):
|
||||||
|
# 'awaiting_dispense' (cash_out at insert: paid, waiting for the report)
|
||||||
|
# 'pending' (cash_in at insert; cash_out once dispense_confirmed)
|
||||||
# 'processing' (claim taken by distribution processor)
|
# 'processing' (claim taken by distribution processor)
|
||||||
# 'processed' (all legs paid)
|
# 'processed' (all legs paid)
|
||||||
# 'partial' (operator marked partial-dispense after the fact)
|
# 'partial_pending' (report says some notes out, value short — holds
|
||||||
|
# EVERYTHING until the operator records how the shortfall
|
||||||
|
# was resolved; then one distribution at the final amount)
|
||||||
|
# 'cash_owed' (report says nothing out — legs never run, funds stay in
|
||||||
|
# the machine wallet, the customer is owed; worklist first)
|
||||||
|
# 'partial' (operator confirmed a partial amount; distributed scaled)
|
||||||
# 'refunded' (operator-initiated refund)
|
# 'refunded' (operator-initiated refund)
|
||||||
# 'errored' (operational distribution failure — retry path applies)
|
# 'errored' (operational distribution failure — retry path applies)
|
||||||
# 'rejected' (Nostr attribution cross-check failed at land time;
|
# 'rejected' (Nostr attribution cross-check failed at land time;
|
||||||
# never went near distribution. error_message holds the
|
# never went near distribution. error_message holds the
|
||||||
# reason. Retry is wrong — investigate the machine.)
|
# reason. Retry is wrong — investigate the machine.)
|
||||||
|
# 'dispense_unreported' is NOT stored: the worklist derives it from
|
||||||
|
# awaiting_dispense rows older than its threshold.
|
||||||
status: str
|
status: str
|
||||||
error_message: str | None
|
error_message: str | None
|
||||||
processed_at: datetime | None
|
processed_at: datetime | None
|
||||||
created_at: datetime
|
created_at: datetime
|
||||||
|
# ADR-005 §2 — copied from the machine's report. The three lamassu
|
||||||
|
# fields plus raw_code / error_class; dispensed_fiat_cents is what the
|
||||||
|
# hardware says physically left, which pre-fills partial-dispense.
|
||||||
|
dispense_confirmed: bool | None = None
|
||||||
|
dispense_error: str | None = None
|
||||||
|
dispense_error_code: str | None = None
|
||||||
|
dispense_raw_code: str | None = None
|
||||||
|
dispense_error_class: str | None = None
|
||||||
|
dispense_reported_at: datetime | None = None
|
||||||
|
dispensed_fiat_cents: int | None = None
|
||||||
# Append-only audit memo. Populated when an operator triggers an in-place
|
# Append-only audit memo. Populated when an operator triggers an in-place
|
||||||
# adjustment (partial-dispense, manual reconciliation override). Each
|
# adjustment (partial-dispense, manual reconciliation override). Each
|
||||||
# entry timestamped + records original values so the overwrite is
|
# entry timestamped + records original values so the overwrite is
|
||||||
|
|
@ -336,6 +363,110 @@ class DcaSettlement(BaseModel):
|
||||||
processing_claim: str | None = None
|
processing_claim: str | None = None
|
||||||
|
|
||||||
|
|
||||||
|
# =============================================================================
|
||||||
|
# Dispense outcome (bitspire ADR-005 §2) — machine → spirekeeper `report_dispense`
|
||||||
|
# =============================================================================
|
||||||
|
|
||||||
|
|
||||||
|
class DispenseReportBill(BaseModel):
|
||||||
|
denomination: int
|
||||||
|
requested: int
|
||||||
|
dispensed: int
|
||||||
|
rejected: int
|
||||||
|
|
||||||
|
|
||||||
|
class DispenseReportCassette(BaseModel):
|
||||||
|
position: int
|
||||||
|
denomination: int
|
||||||
|
provisioned: int
|
||||||
|
dispensed: int
|
||||||
|
rejected: int
|
||||||
|
|
||||||
|
|
||||||
|
class DispenseReportIn(BaseModel):
|
||||||
|
"""The RPC body as the machine sends it. Mirrors @bitSpire/lnbits
|
||||||
|
DispenseReportBody. Field names follow lamassu-server's cash_out_txs /
|
||||||
|
cash_out_actions (dispense_confirmed, error, error_code). Idempotent on
|
||||||
|
(txid, at): the machine resends until acked; a byte-identical resend is
|
||||||
|
acknowledged without a new row."""
|
||||||
|
|
||||||
|
txid: str
|
||||||
|
payment_hash: str | None = None
|
||||||
|
tx_type: str = "cash_out"
|
||||||
|
dispense_confirmed: bool
|
||||||
|
error: str | None = None
|
||||||
|
error_code: str | None = None
|
||||||
|
raw_code: str | None = None
|
||||||
|
error_class: str | None = None # 'terminal' | 'recoverable' | 'inventory'
|
||||||
|
fiat_cents: int
|
||||||
|
currency: str
|
||||||
|
bills: list[DispenseReportBill] = []
|
||||||
|
cassettes: list[DispenseReportCassette] = []
|
||||||
|
counts_uncertain: bool = False
|
||||||
|
remediates_txid: str | None = None
|
||||||
|
at: int
|
||||||
|
|
||||||
|
@validator("txid")
|
||||||
|
def _txid_present(cls, v):
|
||||||
|
if not v or not v.strip():
|
||||||
|
raise ValueError("txid is required")
|
||||||
|
return v.strip()
|
||||||
|
|
||||||
|
@validator("tx_type")
|
||||||
|
def _cash_out_only(cls, v):
|
||||||
|
if v != "cash_out":
|
||||||
|
raise ValueError("report_dispense is for cash_out transactions only")
|
||||||
|
return v
|
||||||
|
|
||||||
|
@validator("error_class")
|
||||||
|
def _known_class(cls, v):
|
||||||
|
if v is not None and v not in ("terminal", "recoverable", "inventory"):
|
||||||
|
raise ValueError(
|
||||||
|
f"error_class must be terminal|recoverable|inventory, got {v!r}"
|
||||||
|
)
|
||||||
|
return v
|
||||||
|
|
||||||
|
@validator("fiat_cents")
|
||||||
|
def _fiat_non_negative(cls, v):
|
||||||
|
if v < 0:
|
||||||
|
raise ValueError("fiat_cents must be >= 0")
|
||||||
|
return v
|
||||||
|
|
||||||
|
@property
|
||||||
|
def dispensed_fiat_cents(self) -> int:
|
||||||
|
"""What the hardware says physically left, in cents."""
|
||||||
|
return sum(b.denomination * b.dispensed for b in self.bills) * 100
|
||||||
|
|
||||||
|
@property
|
||||||
|
def total_dispensed_notes(self) -> int:
|
||||||
|
return sum(b.dispensed for b in self.bills)
|
||||||
|
|
||||||
|
|
||||||
|
class DispenseReport(BaseModel):
|
||||||
|
"""One stored report (append-only — a retry, a late report and a
|
||||||
|
remediation report are distinct rows). `settlement_id` is NULL when the
|
||||||
|
payment this report refers to was never seen by this server."""
|
||||||
|
|
||||||
|
id: str
|
||||||
|
machine_id: str
|
||||||
|
settlement_id: str | None
|
||||||
|
txid: str
|
||||||
|
payment_hash: str | None
|
||||||
|
dispense_confirmed: bool
|
||||||
|
error: str | None
|
||||||
|
error_code: str | None
|
||||||
|
raw_code: str | None
|
||||||
|
error_class: str | None
|
||||||
|
fiat_cents: int
|
||||||
|
currency: str
|
||||||
|
bills_json: str
|
||||||
|
cassettes_json: str
|
||||||
|
counts_uncertain: bool = False
|
||||||
|
remediates_txid: str | None = None
|
||||||
|
reported_at: datetime
|
||||||
|
received_at: datetime
|
||||||
|
|
||||||
|
|
||||||
# =============================================================================
|
# =============================================================================
|
||||||
# Commission splits — operator-defined remainder allocation per machine.
|
# Commission splits — operator-defined remainder allocation per machine.
|
||||||
# =============================================================================
|
# =============================================================================
|
||||||
|
|
@ -570,6 +701,13 @@ class StuckSettlementsResponse(BaseModel):
|
||||||
"""
|
"""
|
||||||
|
|
||||||
threshold_minutes: int
|
threshold_minutes: int
|
||||||
|
# ADR-005 §6 — the only buckets whose meaning is "a customer is owed
|
||||||
|
# money"; rendered first.
|
||||||
|
cash_owed: list = [] # list[DcaSettlement]
|
||||||
|
partial_pending: list = []
|
||||||
|
# awaiting_dispense older than the threshold: the machine never reported.
|
||||||
|
# A machine on an old build lands here too — that is the upgrade path.
|
||||||
|
dispense_unreported: list = []
|
||||||
rejected: list # list[DcaSettlement]
|
rejected: list # list[DcaSettlement]
|
||||||
errored: list
|
errored: list
|
||||||
stuck_pending: list
|
stuck_pending: list
|
||||||
|
|
@ -738,6 +876,11 @@ class PublishCassettesPayload(BaseModel):
|
||||||
applied_ops: list[str] = []
|
applied_ops: list[str] = []
|
||||||
seq: int | None = None
|
seq: int | None = None
|
||||||
counts_uncertain_since: int | None = None
|
counts_uncertain_since: int | None = None
|
||||||
|
# ADR-005 §5 — additive like counts_uncertain_since. Present while the
|
||||||
|
# machine refuses cash-out after a terminal dispenser fault.
|
||||||
|
cash_out_held_since: int | None = None
|
||||||
|
cash_out_held_reason: str | None = None
|
||||||
|
cash_out_held_code: str | None = None
|
||||||
|
|
||||||
@validator("positions", pre=True)
|
@validator("positions", pre=True)
|
||||||
def coerce_string_keys_to_int(cls, v):
|
def coerce_string_keys_to_int(cls, v):
|
||||||
|
|
@ -803,7 +946,21 @@ class PublishCassettesPayload(BaseModel):
|
||||||
# landed after the same problem.
|
# landed after the same problem.
|
||||||
|
|
||||||
|
|
||||||
CASSETTE_OP_TYPES = ("refill", "empty", "recount", "set_denomination")
|
CASSETTE_OP_TYPES = (
|
||||||
|
"refill",
|
||||||
|
"empty",
|
||||||
|
"recount",
|
||||||
|
"set_denomination",
|
||||||
|
"resume_cash_out",
|
||||||
|
)
|
||||||
|
|
||||||
|
# ADR-005 §5: `resume_cash_out` is not a cassette operation — it releases the
|
||||||
|
# machine's cash-out hold without touching a bay — but it rides the same
|
||||||
|
# operator event, with the same id/at dedup shape, because that channel is the
|
||||||
|
# one the machine already consumes. It is machine-wide, so its position is 0
|
||||||
|
# and its wire form carries no position at all. A recount also releases the
|
||||||
|
# hold (same "operator at the open machine" gesture).
|
||||||
|
MACHINE_WIDE_OP_TYPES = ("resume_cash_out",)
|
||||||
|
|
||||||
|
|
||||||
class CassetteOp(BaseModel):
|
class CassetteOp(BaseModel):
|
||||||
|
|
@ -837,11 +994,17 @@ class CassetteOp(BaseModel):
|
||||||
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
|
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
|
||||||
return v
|
return v
|
||||||
|
|
||||||
@validator("position")
|
@root_validator(skip_on_failure=True)
|
||||||
def _position_positive(cls, v):
|
def _position_matches_scope(cls, values):
|
||||||
if v <= 0:
|
pos, typ = values.get("position"), values.get("op_type")
|
||||||
raise ValueError(f"position must be > 0, got {v}")
|
if typ in MACHINE_WIDE_OP_TYPES:
|
||||||
return v
|
if pos != 0:
|
||||||
|
raise ValueError(
|
||||||
|
f"{typ} is machine-wide; position must be 0, got {pos}"
|
||||||
|
)
|
||||||
|
elif pos is None or pos <= 0:
|
||||||
|
raise ValueError(f"position must be > 0, got {pos}")
|
||||||
|
return values
|
||||||
|
|
||||||
def to_wire_dict(self) -> dict:
|
def to_wire_dict(self) -> dict:
|
||||||
"""The published form. Drops the fields this op_type does not use, so
|
"""The published form. Drops the fields this op_type does not use, so
|
||||||
|
|
@ -850,8 +1013,10 @@ class CassetteOp(BaseModel):
|
||||||
"id": self.id,
|
"id": self.id,
|
||||||
"at": int(self.created_at.timestamp()),
|
"at": int(self.created_at.timestamp()),
|
||||||
"type": self.op_type,
|
"type": self.op_type,
|
||||||
"position": self.position,
|
|
||||||
}
|
}
|
||||||
|
if self.op_type in MACHINE_WIDE_OP_TYPES:
|
||||||
|
return out
|
||||||
|
out["position"] = self.position
|
||||||
if self.op_type == "refill":
|
if self.op_type == "refill":
|
||||||
out["bills"] = self.bills
|
out["bills"] = self.bills
|
||||||
elif self.op_type == "recount":
|
elif self.op_type == "recount":
|
||||||
|
|
@ -880,11 +1045,17 @@ class CreateCassetteOpData(BaseModel):
|
||||||
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
|
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
|
||||||
return v
|
return v
|
||||||
|
|
||||||
@validator("position")
|
@root_validator(skip_on_failure=True)
|
||||||
def _position_positive(cls, v):
|
def _position_matches_scope(cls, values):
|
||||||
if v <= 0:
|
pos, typ = values.get("position"), values.get("op_type")
|
||||||
raise ValueError(f"position must be > 0, got {v}")
|
if typ in MACHINE_WIDE_OP_TYPES:
|
||||||
return v
|
if pos != 0:
|
||||||
|
raise ValueError(
|
||||||
|
f"{typ} is machine-wide; position must be 0, got {pos}"
|
||||||
|
)
|
||||||
|
elif pos is None or pos <= 0:
|
||||||
|
raise ValueError(f"position must be > 0, got {pos}")
|
||||||
|
return values
|
||||||
|
|
||||||
@validator("bills")
|
@validator("bills")
|
||||||
def _bills_positive(cls, v):
|
def _bills_positive(cls, v):
|
||||||
|
|
@ -911,6 +1082,7 @@ class CreateCassetteOpData(BaseModel):
|
||||||
"recount": "count",
|
"recount": "count",
|
||||||
"set_denomination": "denomination",
|
"set_denomination": "denomination",
|
||||||
"empty": None,
|
"empty": None,
|
||||||
|
"resume_cash_out": None,
|
||||||
}[values.get("op_type")]
|
}[values.get("op_type")]
|
||||||
if required is not None and values.get(required) is None:
|
if required is not None and values.get(required) is None:
|
||||||
raise ValueError(f"{values['op_type']} requires `{required}`")
|
raise ValueError(f"{values['op_type']} requires `{required}`")
|
||||||
|
|
|
||||||
|
|
@ -70,6 +70,10 @@ window.app = Vue.createApp({
|
||||||
|
|
||||||
// Worklist (P9g)
|
// Worklist (P9g)
|
||||||
worklist: {
|
worklist: {
|
||||||
|
// ADR-005 §6 — owed-cash buckets first
|
||||||
|
cash_owed: [],
|
||||||
|
partial_pending: [],
|
||||||
|
dispense_unreported: [],
|
||||||
rejected: [],
|
rejected: [],
|
||||||
errored: [],
|
errored: [],
|
||||||
stuck_pending: [],
|
stuck_pending: [],
|
||||||
|
|
@ -372,6 +376,33 @@ window.app = Vue.createApp({
|
||||||
},
|
},
|
||||||
worklistBuckets() {
|
worklistBuckets() {
|
||||||
return [
|
return [
|
||||||
|
{
|
||||||
|
key: 'cash_owed',
|
||||||
|
label:
|
||||||
|
'Cash owed — customer paid, machine dispensed nothing. ' +
|
||||||
|
'Legs never ran; funds are in the machine wallet.',
|
||||||
|
icon: 'money_off',
|
||||||
|
color: 'negative',
|
||||||
|
rows: this.worklist.cash_owed
|
||||||
|
},
|
||||||
|
{
|
||||||
|
key: 'partial_pending',
|
||||||
|
label:
|
||||||
|
'Partial dispense — some notes out, value short. Held whole ' +
|
||||||
|
'until you record how the shortfall was resolved.',
|
||||||
|
icon: 'call_split',
|
||||||
|
color: 'deep-orange',
|
||||||
|
rows: this.worklist.partial_pending
|
||||||
|
},
|
||||||
|
{
|
||||||
|
key: 'dispense_unreported',
|
||||||
|
label:
|
||||||
|
'Unreported — cash-out paid, machine never reported the dispense. ' +
|
||||||
|
'Check the machine; an old build lands here too.',
|
||||||
|
icon: 'help_outline',
|
||||||
|
color: 'amber',
|
||||||
|
rows: this.worklist.dispense_unreported
|
||||||
|
},
|
||||||
{
|
{
|
||||||
key: 'rejected',
|
key: 'rejected',
|
||||||
label: 'Rejected — Nostr attribution failed; investigate machine',
|
label: 'Rejected — Nostr attribution failed; investigate machine',
|
||||||
|
|
@ -560,6 +591,9 @@ window.app = Vue.createApp({
|
||||||
try {
|
try {
|
||||||
const {data} = await LNbits.api.request('GET', STUCK_PATH)
|
const {data} = await LNbits.api.request('GET', STUCK_PATH)
|
||||||
this.worklistCount =
|
this.worklistCount =
|
||||||
|
(data?.cash_owed?.length || 0) +
|
||||||
|
(data?.partial_pending?.length || 0) +
|
||||||
|
(data?.dispense_unreported?.length || 0) +
|
||||||
(data?.rejected?.length || 0) +
|
(data?.rejected?.length || 0) +
|
||||||
(data?.errored?.length || 0) +
|
(data?.errored?.length || 0) +
|
||||||
(data?.stuck_pending?.length || 0) +
|
(data?.stuck_pending?.length || 0) +
|
||||||
|
|
@ -575,11 +609,17 @@ window.app = Vue.createApp({
|
||||||
const {data} = await LNbits.api.request(
|
const {data} = await LNbits.api.request(
|
||||||
'GET', `${STUCK_PATH}?threshold_minutes=${this.worklistThreshold}`
|
'GET', `${STUCK_PATH}?threshold_minutes=${this.worklistThreshold}`
|
||||||
)
|
)
|
||||||
|
this.worklist.cash_owed = data?.cash_owed || []
|
||||||
|
this.worklist.partial_pending = data?.partial_pending || []
|
||||||
|
this.worklist.dispense_unreported = data?.dispense_unreported || []
|
||||||
this.worklist.rejected = data?.rejected || []
|
this.worklist.rejected = data?.rejected || []
|
||||||
this.worklist.errored = data?.errored || []
|
this.worklist.errored = data?.errored || []
|
||||||
this.worklist.stuck_pending = data?.stuck_pending || []
|
this.worklist.stuck_pending = data?.stuck_pending || []
|
||||||
this.worklist.stuck_processing = data?.stuck_processing || []
|
this.worklist.stuck_processing = data?.stuck_processing || []
|
||||||
this.worklist.totalCount =
|
this.worklist.totalCount =
|
||||||
|
this.worklist.cash_owed.length +
|
||||||
|
this.worklist.partial_pending.length +
|
||||||
|
this.worklist.dispense_unreported.length +
|
||||||
this.worklist.rejected.length +
|
this.worklist.rejected.length +
|
||||||
this.worklist.errored.length +
|
this.worklist.errored.length +
|
||||||
this.worklist.stuck_pending.length +
|
this.worklist.stuck_pending.length +
|
||||||
|
|
@ -1208,12 +1248,50 @@ window.app = Vue.createApp({
|
||||||
openPartialDispense(settlement) {
|
openPartialDispense(settlement) {
|
||||||
this.partialDispenseDialog.settlement = settlement
|
this.partialDispenseDialog.settlement = settlement
|
||||||
this.partialDispenseDialog.mode = 'fraction'
|
this.partialDispenseDialog.mode = 'fraction'
|
||||||
this.partialDispenseDialog.dispensed_fraction = null
|
// ADR-005: pre-fill from the machine's report — the hardware's own count
|
||||||
|
// of what left — so the operator confirms a number rather than typing one.
|
||||||
|
const dispensedCents = settlement.dispensed_fiat_cents
|
||||||
|
const fiat = Number(settlement.fiat_amount)
|
||||||
|
this.partialDispenseDialog.dispensed_fraction =
|
||||||
|
dispensedCents != null && fiat > 0
|
||||||
|
? Math.round((dispensedCents / 100 / fiat) * 10000) / 10000
|
||||||
|
: null
|
||||||
this.partialDispenseDialog.dispensed_sats = null
|
this.partialDispenseDialog.dispensed_sats = null
|
||||||
this.partialDispenseDialog.notes = ''
|
this.partialDispenseDialog.notes = settlement.dispense_error
|
||||||
|
? `Machine reported: ${settlement.dispense_error_code || ''} ${settlement.dispense_raw_code || ''} — ${settlement.dispense_error}`.trim()
|
||||||
|
: ''
|
||||||
this.partialDispenseDialog.show = true
|
this.partialDispenseDialog.show = true
|
||||||
},
|
},
|
||||||
|
|
||||||
|
// ADR-005 §5 — release a machine's cash-out hold after a terminal
|
||||||
|
// dispenser fault, when the jam was cleared without a recount.
|
||||||
|
confirmResumeCashOut(machine) {
|
||||||
|
Quasar.Dialog.create({
|
||||||
|
title: 'Resume cash-out?',
|
||||||
|
message:
|
||||||
|
'The machine latched cash-out off after a dispenser fault' +
|
||||||
|
(machine.cash_out_held_code ? ` (${machine.cash_out_held_code})` : '') +
|
||||||
|
'. Only do this after the transport path has been physically cleared. ' +
|
||||||
|
'A recount releases the hold too, and also fixes the bay count.',
|
||||||
|
cancel: true,
|
||||||
|
persistent: true
|
||||||
|
}).onOk(async () => {
|
||||||
|
try {
|
||||||
|
await LNbits.api.request(
|
||||||
|
'POST',
|
||||||
|
`/spirekeeper/api/v1/dca/machines/${machine.id}/resume-cash-out`
|
||||||
|
)
|
||||||
|
Quasar.Notify.create({
|
||||||
|
type: 'positive',
|
||||||
|
message: 'Resume published — the machine clears the hold on receipt'
|
||||||
|
})
|
||||||
|
if (this.machineDetail && this.machineDetail.machine) await this.reloadMachineDetail()
|
||||||
|
} catch (e) {
|
||||||
|
this._notifyError(e, 'Resume cash-out failed')
|
||||||
|
}
|
||||||
|
})
|
||||||
|
},
|
||||||
|
|
||||||
async submitPartialDispense() {
|
async submitPartialDispense() {
|
||||||
const d = this.partialDispenseDialog
|
const d = this.partialDispenseDialog
|
||||||
const body = {notes: d.notes || null}
|
const body = {notes: d.notes || null}
|
||||||
|
|
|
||||||
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:
|
if isinstance(nostr_event_id, str) and nostr_event_id:
|
||||||
data.bitspire_event_id = nostr_event_id
|
data.bitspire_event_id = nostr_event_id
|
||||||
|
|
||||||
# 3) Insert + distribute.
|
# 3) Insert. ADR-005 §1: payment is authorization, the machine's dispense
|
||||||
settlement = await create_settlement_idempotent(data, initial_status="pending")
|
# 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:
|
if settlement is None:
|
||||||
logger.error(
|
logger.error(
|
||||||
f"spirekeeper: failed to insert settlement for "
|
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"(super_fee={data.platform_fee_sats} "
|
||||||
f"operator_fee={data.operator_fee_sats})"
|
f"operator_fee={data.operator_fee_sats})"
|
||||||
)
|
)
|
||||||
|
if is_cash_out:
|
||||||
|
await _await_dispense_or_adopt(settlement, machine, data)
|
||||||
|
return
|
||||||
|
|
||||||
# Spawn distribution on a background task so the LNbits invoice queue
|
# Spawn distribution on a background task so the LNbits invoice queue
|
||||||
# (shared across all extensions) keeps draining while we move sats.
|
# (shared across all extensions) keeps draining while we move sats.
|
||||||
# Concurrency-safe: process_settlement uses claim_settlement_for_processing
|
# Concurrency-safe: process_settlement uses claim_settlement_for_processing
|
||||||
|
|
@ -196,6 +209,23 @@ async def _handle_payment(payment: Payment) -> None:
|
||||||
task.add_done_callback(_inflight_distributions.discard)
|
task.add_done_callback(_inflight_distributions.discard)
|
||||||
|
|
||||||
|
|
||||||
|
async def _await_dispense_or_adopt(
|
||||||
|
settlement, machine: Machine, data: CreateDcaSettlementData
|
||||||
|
) -> None:
|
||||||
|
"""A cash_out waits for the machine's dispense report (ADR-005 §1) — unless
|
||||||
|
the report is already here. Under hold invoices the payment settles AFTER
|
||||||
|
the dispense, and the invoice listener can lag the transport, so an orphan
|
||||||
|
report for this txid is adopted and applied now."""
|
||||||
|
if settlement.status != "awaiting_dispense" or not data.bitspire_txid:
|
||||||
|
return
|
||||||
|
from .crud import get_latest_unlinked_dispense_report
|
||||||
|
from .dispense_transport import adopt_unlinked_report
|
||||||
|
|
||||||
|
early = await get_latest_unlinked_dispense_report(machine.id, data.bitspire_txid)
|
||||||
|
if early is not None:
|
||||||
|
await adopt_unlinked_report(settlement, machine, early)
|
||||||
|
|
||||||
|
|
||||||
async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) -> None:
|
async def _record_rejected(payment: Payment, machine: Machine, exc: Exception) -> None:
|
||||||
"""Insert a minimal `dca_settlements` row with `status='rejected'` and
|
"""Insert a minimal `dca_settlements` row with `status='rejected'` and
|
||||||
the exception message for operator forensics.
|
the exception message for operator forensics.
|
||||||
|
|
@ -438,12 +468,38 @@ async def _record_counts_uncertainty(
|
||||||
await set_machine_counts_uncertain(machine_id, since)
|
await set_machine_counts_uncertain(machine_id, since)
|
||||||
|
|
||||||
|
|
||||||
|
async def _record_cash_out_hold(
|
||||||
|
machine_id: str, payload, set_machine_cash_out_hold
|
||||||
|
) -> None:
|
||||||
|
"""Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row.
|
||||||
|
|
||||||
|
Written on every state event, including when absent, because the machine
|
||||||
|
clearing the hold — after an operator recount or resume_cash_out — matters
|
||||||
|
exactly as much as it setting one.
|
||||||
|
"""
|
||||||
|
from datetime import datetime as _datetime
|
||||||
|
from datetime import timezone as _timezone
|
||||||
|
|
||||||
|
since = None
|
||||||
|
if payload.cash_out_held_since is not None:
|
||||||
|
since = _datetime.fromtimestamp(
|
||||||
|
int(payload.cash_out_held_since), tz=_timezone.utc
|
||||||
|
)
|
||||||
|
await set_machine_cash_out_hold(
|
||||||
|
machine_id,
|
||||||
|
since,
|
||||||
|
payload.cash_out_held_reason if since else None,
|
||||||
|
payload.cash_out_held_code if since else None,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
async def _handle_cassette_state_event(
|
async def _handle_cassette_state_event(
|
||||||
event_message,
|
event_message,
|
||||||
get_machine_by_atm_pubkey_hex,
|
get_machine_by_atm_pubkey_hex,
|
||||||
apply_reported_state,
|
apply_reported_state,
|
||||||
mark_cassette_ops_acked,
|
mark_cassette_ops_acked,
|
||||||
set_machine_counts_uncertain,
|
set_machine_counts_uncertain,
|
||||||
|
set_machine_cash_out_hold=None,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Verify signature, resolve the operator's signer, decrypt via the
|
"""Verify signature, resolve the operator's signer, decrypt via the
|
||||||
signer abstraction (bunker round-trip for RemoteBunkerSigner; direct
|
signer abstraction (bunker round-trip for RemoteBunkerSigner; direct
|
||||||
|
|
@ -558,3 +614,8 @@ async def _handle_cassette_state_event(
|
||||||
# _record_op_acknowledgements for why. Same for the uncertainty marker.
|
# _record_op_acknowledgements for why. Same for the uncertainty marker.
|
||||||
await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked)
|
await _record_op_acknowledgements(machine.id, payload, mark_cassette_ops_acked)
|
||||||
await _record_counts_uncertainty(machine.id, payload, set_machine_counts_uncertain)
|
await _record_counts_uncertainty(machine.id, payload, set_machine_counts_uncertain)
|
||||||
|
if set_machine_cash_out_hold is None:
|
||||||
|
from .crud import set_machine_cash_out_hold as _default_set_hold
|
||||||
|
|
||||||
|
set_machine_cash_out_hold = _default_set_hold
|
||||||
|
await _record_cash_out_hold(machine.id, payload, set_machine_cash_out_hold)
|
||||||
|
|
|
||||||
|
|
@ -661,6 +661,12 @@
|
||||||
@click="viewMachineFromWorklist(props.row)">
|
@click="viewMachineFromWorklist(props.row)">
|
||||||
<q-tooltip>Open machine detail</q-tooltip>
|
<q-tooltip>Open machine detail</q-tooltip>
|
||||||
</q-btn>
|
</q-btn>
|
||||||
|
<q-btn v-if="bucket.key === 'partial_pending'"
|
||||||
|
flat dense size="sm" icon="call_split"
|
||||||
|
color="deep-orange"
|
||||||
|
@click="openPartialDispense(props.row)">
|
||||||
|
<q-tooltip>Record the resolution (pre-filled from the machine's report)</q-tooltip>
|
||||||
|
</q-btn>
|
||||||
<q-btn v-if="bucket.key === 'errored'"
|
<q-btn v-if="bucket.key === 'errored'"
|
||||||
flat dense size="sm" icon="restart_alt"
|
flat dense size="sm" icon="restart_alt"
|
||||||
color="primary"
|
color="primary"
|
||||||
|
|
@ -1178,6 +1184,25 @@
|
||||||
<span v-text="machineDetail.cassettesError"></span>
|
<span v-text="machineDetail.cassettesError"></span>
|
||||||
</q-banner>
|
</q-banner>
|
||||||
|
|
||||||
|
<q-banner v-if="machineDetail.machine
|
||||||
|
&& machineDetail.machine.cash_out_held_since"
|
||||||
|
class="bg-red-1 text-grey-9 q-mb-md">
|
||||||
|
<template v-slot:avatar>
|
||||||
|
<q-icon name="block" color="negative"></q-icon>
|
||||||
|
</template>
|
||||||
|
<b>Cash-out is held.</b>
|
||||||
|
The machine latched cash-out off after a terminal dispenser fault
|
||||||
|
(<span v-text="machineDetail.machine.cash_out_held_code || 'fault'"></span>
|
||||||
|
at <span v-text="formatTime(machineDetail.machine.cash_out_held_since)"></span>):
|
||||||
|
<span v-text="machineDetail.machine.cash_out_held_reason"></span>.
|
||||||
|
Clear the transport path, then either record a <b>Recount</b>
|
||||||
|
(which also fixes the count) or release it here.
|
||||||
|
<template v-slot:action>
|
||||||
|
<q-btn flat color="negative" label="Resume cash-out"
|
||||||
|
@click="confirmResumeCashOut(machineDetail.machine)"></q-btn>
|
||||||
|
</template>
|
||||||
|
</q-banner>
|
||||||
|
|
||||||
<q-banner v-if="machineDetail.machine
|
<q-banner v-if="machineDetail.machine
|
||||||
&& machineDetail.machine.counts_uncertain_since"
|
&& machineDetail.machine.counts_uncertain_since"
|
||||||
class="bg-orange-1 text-grey-9 q-mb-md">
|
class="bg-orange-1 text-grey-9 q-mb-md">
|
||||||
|
|
|
||||||
|
|
@ -124,6 +124,8 @@ class TestWireShape:
|
||||||
"recount": {"count": 1},
|
"recount": {"count": 1},
|
||||||
"set_denomination": {"denomination": 1},
|
"set_denomination": {"denomination": 1},
|
||||||
"empty": {},
|
"empty": {},
|
||||||
|
# machine-wide (ADR-005 §5): position 0, no position on the wire
|
||||||
|
"resume_cash_out": {"position": 0},
|
||||||
}[op_type]
|
}[op_type]
|
||||||
wire = op(op_type=op_type, **kw).to_wire_dict()
|
wire = op(op_type=op_type, **kw).to_wire_dict()
|
||||||
assert None not in wire.values()
|
assert None not in wire.values()
|
||||||
|
|
|
||||||
542
tests/test_dispense_outcome.py
Normal file
542
tests/test_dispense_outcome.py
Normal file
|
|
@ -0,0 +1,542 @@
|
||||||
|
"""
|
||||||
|
Tests for dispense-outcome capture (bitspire ADR-005 §1-§2, #122).
|
||||||
|
|
||||||
|
Covers the pure pieces and the handler's transitions with the crud layer
|
||||||
|
monkeypatched (no DB), in the project's established style: asyncio.run inside
|
||||||
|
the test body, SimpleNamespace for request/payment shapes.
|
||||||
|
|
||||||
|
- models: DispenseReportIn validation + derived numbers; resume_cash_out as a
|
||||||
|
machine-wide op (position 0, no position on the wire); the three new
|
||||||
|
worklist buckets default empty.
|
||||||
|
- handler: confirmed → pending + distribution spawned; nothing out →
|
||||||
|
cash_owed; some out → partial_pending; already-captured settlements are
|
||||||
|
recorded but not moved; a byte-identical resend is acked without a new
|
||||||
|
row; an orphan report (payment not landed) is stored unlinked; a
|
||||||
|
remediation report moves the owed settlement to pending; unpaired sender
|
||||||
|
and bad bodies are refused.
|
||||||
|
- gate: _handle_payment inserts cash_out as awaiting_dispense and does NOT
|
||||||
|
spawn distribution; cash_in is unchanged; an early report is adopted.
|
||||||
|
- consumer: the state document's cash_out_held_* mirrors onto the machine,
|
||||||
|
including clearing it.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from pydantic import ValidationError
|
||||||
|
|
||||||
|
from .. import crud as crud_mod
|
||||||
|
from .. import dispense_transport, tasks
|
||||||
|
from ..dispense_transport import (
|
||||||
|
_outcome_status,
|
||||||
|
handle_report_dispense,
|
||||||
|
)
|
||||||
|
from ..models import (
|
||||||
|
CASSETTE_OP_TYPES,
|
||||||
|
CassetteOp,
|
||||||
|
CreateCassetteOpData,
|
||||||
|
CreateDcaSettlementData,
|
||||||
|
DcaSettlement,
|
||||||
|
DispenseReportIn,
|
||||||
|
Machine,
|
||||||
|
PublishCassettesPayload,
|
||||||
|
StuckSettlementsResponse,
|
||||||
|
)
|
||||||
|
|
||||||
|
_NOW = datetime(2026, 10, 9, 7, 2, 33)
|
||||||
|
_ATM_HEX = "df2003343784b69cb813b2a4fd231f83ae81133279251c735414f9909baa7ac6"
|
||||||
|
_TXID = "tx_mv0madw6_wdhtea1v"
|
||||||
|
_HASH = "6f216df32c36" + "0" * 52
|
||||||
|
|
||||||
|
|
||||||
|
def _machine(**over) -> Machine:
|
||||||
|
base: dict[str, Any] = {
|
||||||
|
"id": "m1",
|
||||||
|
"operator_user_id": "op1",
|
||||||
|
"machine_npub": _ATM_HEX,
|
||||||
|
"wallet_id": "w1",
|
||||||
|
"name": "sintra",
|
||||||
|
"location": None,
|
||||||
|
"fiat_code": "EUR",
|
||||||
|
"is_active": True,
|
||||||
|
"created_at": _NOW,
|
||||||
|
"updated_at": _NOW,
|
||||||
|
}
|
||||||
|
base.update(over)
|
||||||
|
return Machine(**base)
|
||||||
|
|
||||||
|
|
||||||
|
def _settlement(status="awaiting_dispense", **over) -> DcaSettlement:
|
||||||
|
base: dict[str, Any] = {
|
||||||
|
"id": "s1",
|
||||||
|
"machine_id": "m1",
|
||||||
|
"payment_hash": _HASH,
|
||||||
|
"bitspire_event_id": None,
|
||||||
|
"bitspire_txid": _TXID,
|
||||||
|
"wire_sats": 54440,
|
||||||
|
"fiat_amount": 40.0,
|
||||||
|
"fiat_code": "EUR",
|
||||||
|
"exchange_rate": 1361.0,
|
||||||
|
"principal_sats": 54440,
|
||||||
|
"fee_sats": 0,
|
||||||
|
"platform_fee_sats": 0,
|
||||||
|
"operator_fee_sats": 0,
|
||||||
|
"tx_type": "cash_out",
|
||||||
|
"bills_json": None,
|
||||||
|
"cassettes_json": None,
|
||||||
|
"status": status,
|
||||||
|
"error_message": None,
|
||||||
|
"processed_at": None,
|
||||||
|
"created_at": _NOW,
|
||||||
|
}
|
||||||
|
base.update(over)
|
||||||
|
return DcaSettlement(**base)
|
||||||
|
|
||||||
|
|
||||||
|
def _report(**over) -> dict:
|
||||||
|
"""The wire body for sintra's 2026-10-09 jam, as a dict."""
|
||||||
|
body: dict[str, Any] = {
|
||||||
|
"txid": _TXID,
|
||||||
|
"payment_hash": _HASH,
|
||||||
|
"tx_type": "cash_out",
|
||||||
|
"dispense_confirmed": False,
|
||||||
|
"error": "Note stopped at the cassette exit",
|
||||||
|
"error_code": "F56DispenseError",
|
||||||
|
"raw_code": "78 42",
|
||||||
|
"error_class": "terminal",
|
||||||
|
"fiat_cents": 4000,
|
||||||
|
"currency": "EUR",
|
||||||
|
"bills": [{"denomination": 20, "requested": 2, "dispensed": 0, "rejected": 0}],
|
||||||
|
"cassettes": [
|
||||||
|
{
|
||||||
|
"position": 1,
|
||||||
|
"denomination": 50,
|
||||||
|
"provisioned": 0,
|
||||||
|
"dispensed": 0,
|
||||||
|
"rejected": 0,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"position": 2,
|
||||||
|
"denomination": 20,
|
||||||
|
"provisioned": 2,
|
||||||
|
"dispensed": 0,
|
||||||
|
"rejected": 0,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"counts_uncertain": True,
|
||||||
|
"at": 1791529353,
|
||||||
|
}
|
||||||
|
body.update(over)
|
||||||
|
return body
|
||||||
|
|
||||||
|
|
||||||
|
def _confirmed() -> dict:
|
||||||
|
return _report(
|
||||||
|
dispense_confirmed=True,
|
||||||
|
error=None,
|
||||||
|
error_code=None,
|
||||||
|
raw_code=None,
|
||||||
|
error_class=None,
|
||||||
|
bills=[{"denomination": 20, "requested": 2, "dispensed": 2, "rejected": 0}],
|
||||||
|
counts_uncertain=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _partial() -> dict:
|
||||||
|
return _report(
|
||||||
|
bills=[{"denomination": 20, "requested": 2, "dispensed": 1, "rejected": 1}],
|
||||||
|
counts_uncertain=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Models
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
class TestDispenseReportIn:
|
||||||
|
def test_derived_numbers(self):
|
||||||
|
r = DispenseReportIn(**_partial())
|
||||||
|
assert r.dispensed_fiat_cents == 2000
|
||||||
|
assert r.total_dispensed_notes == 1
|
||||||
|
assert DispenseReportIn(**_confirmed()).dispensed_fiat_cents == 4000
|
||||||
|
assert DispenseReportIn(**_report()).total_dispensed_notes == 0
|
||||||
|
|
||||||
|
def test_outcome_routing(self):
|
||||||
|
assert _outcome_status(DispenseReportIn(**_confirmed())) == "pending"
|
||||||
|
assert _outcome_status(DispenseReportIn(**_partial())) == "partial_pending"
|
||||||
|
assert _outcome_status(DispenseReportIn(**_report())) == "cash_owed"
|
||||||
|
|
||||||
|
def test_rejects_cash_in_and_unknown_class(self):
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
DispenseReportIn(**_report(tx_type="cash_in"))
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
DispenseReportIn(**_report(error_class="weird"))
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
DispenseReportIn(**_report(txid=" "))
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
DispenseReportIn(**_report(fiat_cents=-1))
|
||||||
|
|
||||||
|
|
||||||
|
class TestResumeCashOutOp:
|
||||||
|
def test_is_a_known_type_and_machine_wide(self):
|
||||||
|
assert "resume_cash_out" in CASSETTE_OP_TYPES
|
||||||
|
op = CassetteOp(
|
||||||
|
id="r1",
|
||||||
|
machine_id="m1",
|
||||||
|
position=0,
|
||||||
|
op_type="resume_cash_out",
|
||||||
|
created_at=_NOW,
|
||||||
|
)
|
||||||
|
wire = op.to_wire_dict()
|
||||||
|
assert wire == {
|
||||||
|
"id": "r1",
|
||||||
|
"at": int(_NOW.timestamp()),
|
||||||
|
"type": "resume_cash_out",
|
||||||
|
}
|
||||||
|
assert "position" not in wire
|
||||||
|
|
||||||
|
def test_position_must_be_zero_for_resume_and_positive_otherwise(self):
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
CassetteOp(
|
||||||
|
id="r1",
|
||||||
|
machine_id="m1",
|
||||||
|
position=2,
|
||||||
|
op_type="resume_cash_out",
|
||||||
|
created_at=_NOW,
|
||||||
|
)
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
CreateCassetteOpData(position=1, op_type="resume_cash_out")
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
CreateCassetteOpData(position=0, op_type="refill", bills=5)
|
||||||
|
CreateCassetteOpData(position=0, op_type="resume_cash_out") # ok
|
||||||
|
|
||||||
|
def test_resume_carries_no_bay_fields(self):
|
||||||
|
with pytest.raises(ValidationError):
|
||||||
|
CreateCassetteOpData(position=0, op_type="resume_cash_out", count=3)
|
||||||
|
|
||||||
|
|
||||||
|
class TestWorklistModel:
|
||||||
|
def test_new_buckets_default_empty(self):
|
||||||
|
r = StuckSettlementsResponse(
|
||||||
|
threshold_minutes=30,
|
||||||
|
rejected=[],
|
||||||
|
errored=[],
|
||||||
|
stuck_pending=[],
|
||||||
|
stuck_processing=[],
|
||||||
|
)
|
||||||
|
assert (
|
||||||
|
r.cash_owed == []
|
||||||
|
and r.partial_pending == []
|
||||||
|
and r.dispense_unreported == []
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Handler
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
class _Wired:
|
||||||
|
"""Monkeypatched crud layer for handle_report_dispense."""
|
||||||
|
|
||||||
|
def __init__(self, monkeypatch, *, machine, settlement, existing_report=None):
|
||||||
|
self.inserted = []
|
||||||
|
self.applied = []
|
||||||
|
self.spawned = []
|
||||||
|
self.uncertain = []
|
||||||
|
self.statuses = []
|
||||||
|
self.settlement = settlement
|
||||||
|
|
||||||
|
async def get_machine(_hex):
|
||||||
|
return machine
|
||||||
|
|
||||||
|
async def get_report(_mid, _txid, _at):
|
||||||
|
return existing_report
|
||||||
|
|
||||||
|
async def get_settlement(_mid, txid):
|
||||||
|
if self.settlement is not None and self.settlement.bitspire_txid == txid:
|
||||||
|
return self.settlement
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def insert(mid, sid, report):
|
||||||
|
self.inserted.append((mid, sid, report))
|
||||||
|
return SimpleNamespace(id="rep1", settlement_id=sid, txid=report.txid)
|
||||||
|
|
||||||
|
async def apply(sid, report, new_status, reported_at):
|
||||||
|
self.applied.append((sid, new_status, reported_at))
|
||||||
|
return (
|
||||||
|
self.settlement.copy(update={"status": new_status})
|
||||||
|
if self.settlement
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
|
||||||
|
async def set_uncertain(mid, since):
|
||||||
|
self.uncertain.append((mid, since))
|
||||||
|
|
||||||
|
async def mark_status(sid, status, _msg):
|
||||||
|
self.statuses.append((sid, status))
|
||||||
|
return None
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
dispense_transport, "get_machine_by_atm_pubkey_hex", get_machine
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(dispense_transport, "get_dispense_report", get_report)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
dispense_transport, "get_settlement_by_txid", get_settlement
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(dispense_transport, "insert_dispense_report", insert)
|
||||||
|
monkeypatch.setattr(dispense_transport, "apply_dispense_outcome", apply)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
dispense_transport, "set_machine_counts_uncertain", set_uncertain
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
dispense_transport, "_spawn_distribution", self.spawned.append
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(crud_mod, "mark_settlement_status", mark_status)
|
||||||
|
|
||||||
|
|
||||||
|
def _req(body, sender=_ATM_HEX):
|
||||||
|
return SimpleNamespace(body=body, sender_pubkey=sender, event_id="ev1")
|
||||||
|
|
||||||
|
|
||||||
|
class TestHandleReportDispense:
|
||||||
|
def test_confirmed_captures_and_distributes(self, monkeypatch):
|
||||||
|
w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement())
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_confirmed())))
|
||||||
|
assert out["received"] is True
|
||||||
|
assert out["settlement_status"] == "pending"
|
||||||
|
assert w.applied == [("s1", "pending", datetime.fromtimestamp(1791529353))]
|
||||||
|
assert w.spawned == ["s1"]
|
||||||
|
assert w.inserted[0][1] == "s1" # linked to the settlement
|
||||||
|
assert w.uncertain == []
|
||||||
|
|
||||||
|
def test_nothing_out_is_cash_owed_and_nothing_moves(self, monkeypatch):
|
||||||
|
w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement())
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_report())))
|
||||||
|
assert out["settlement_status"] == "cash_owed"
|
||||||
|
assert w.applied[0][1] == "cash_owed"
|
||||||
|
assert w.spawned == []
|
||||||
|
# counts_uncertain on the report mirrors onto the machine immediately
|
||||||
|
assert len(w.uncertain) == 1 and w.uncertain[0][0] == "m1"
|
||||||
|
|
||||||
|
def test_some_out_is_partial_pending_held_whole(self, monkeypatch):
|
||||||
|
w = _Wired(monkeypatch, machine=_machine(), settlement=_settlement())
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_partial())))
|
||||||
|
assert out["settlement_status"] == "partial_pending"
|
||||||
|
assert w.spawned == [] # ADR-005 Decision 1: one distribution, when final
|
||||||
|
|
||||||
|
def test_already_captured_settlement_is_recorded_not_moved(self, monkeypatch):
|
||||||
|
w = _Wired(
|
||||||
|
monkeypatch, machine=_machine(), settlement=_settlement(status="processed")
|
||||||
|
)
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_report())))
|
||||||
|
assert out["settlement_status"] == "processed"
|
||||||
|
assert w.applied == [] and w.spawned == []
|
||||||
|
assert len(w.inserted) == 1 # the row still lands — it is information
|
||||||
|
|
||||||
|
def test_identical_resend_is_acked_without_a_new_row(self, monkeypatch):
|
||||||
|
existing = SimpleNamespace(id="rep0", settlement_id="s1", txid=_TXID)
|
||||||
|
w = _Wired(
|
||||||
|
monkeypatch,
|
||||||
|
machine=_machine(),
|
||||||
|
settlement=_settlement(status="cash_owed"),
|
||||||
|
existing_report=existing,
|
||||||
|
)
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_report())))
|
||||||
|
assert out["received"] is True and out.get("duplicate") is True
|
||||||
|
assert out["settlement_status"] == "cash_owed"
|
||||||
|
assert w.inserted == [] and w.applied == []
|
||||||
|
|
||||||
|
def test_report_before_payment_is_stored_unlinked(self, monkeypatch):
|
||||||
|
w = _Wired(monkeypatch, machine=_machine(), settlement=None)
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(_confirmed())))
|
||||||
|
assert out["settlement_status"] is None
|
||||||
|
assert w.inserted[0][1] is None
|
||||||
|
assert w.spawned == []
|
||||||
|
|
||||||
|
def test_remediation_moves_the_owed_settlement_to_pending(self, monkeypatch):
|
||||||
|
owed = _settlement(status="cash_owed")
|
||||||
|
w = _Wired(monkeypatch, machine=_machine(), settlement=owed)
|
||||||
|
body = _confirmed()
|
||||||
|
body.update(txid="manual-1", remediates_txid=_TXID, at=1791530000)
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(body)))
|
||||||
|
assert out["settlement_status"] == "pending"
|
||||||
|
assert w.statuses == [("s1", "pending")]
|
||||||
|
assert w.spawned == ["s1"]
|
||||||
|
assert w.applied == [] # the ORIGINAL report's columns stay on the settlement
|
||||||
|
|
||||||
|
def test_remediation_that_did_not_confirm_leaves_it_owed(self, monkeypatch):
|
||||||
|
w = _Wired(
|
||||||
|
monkeypatch, machine=_machine(), settlement=_settlement(status="cash_owed")
|
||||||
|
)
|
||||||
|
body = _report()
|
||||||
|
body.update(txid="manual-2", remediates_txid=_TXID, at=1791530001)
|
||||||
|
out = asyncio.run(handle_report_dispense(None, _req(body)))
|
||||||
|
assert out["settlement_status"] == "cash_owed"
|
||||||
|
assert w.statuses == [] and w.spawned == []
|
||||||
|
|
||||||
|
def test_unpaired_sender_and_bad_body_are_refused(self, monkeypatch):
|
||||||
|
_Wired(monkeypatch, machine=None, settlement=None)
|
||||||
|
with pytest.raises(ValueError, match="not a paired machine"):
|
||||||
|
asyncio.run(handle_report_dispense(None, _req(_report())))
|
||||||
|
_Wired(monkeypatch, machine=_machine(), settlement=_settlement())
|
||||||
|
with pytest.raises(ValueError, match="invalid report_dispense body"):
|
||||||
|
asyncio.run(handle_report_dispense(None, _req({"txid": _TXID})))
|
||||||
|
with pytest.raises(ValueError, match="sender_pubkey"):
|
||||||
|
asyncio.run(handle_report_dispense(None, _req(_report(), sender="")))
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Gate in _handle_payment
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
def _payment(is_in=True):
|
||||||
|
return SimpleNamespace(
|
||||||
|
success=True,
|
||||||
|
wallet_id="w1",
|
||||||
|
extra={
|
||||||
|
"source": "bitspire",
|
||||||
|
"type": "cash_out" if is_in else "cash_in",
|
||||||
|
"txid": _TXID,
|
||||||
|
},
|
||||||
|
is_in=is_in,
|
||||||
|
sat=54440 if is_in else -54440,
|
||||||
|
payment_hash=_HASH,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _data(tx_type) -> CreateDcaSettlementData:
|
||||||
|
return CreateDcaSettlementData(
|
||||||
|
machine_id="m1",
|
||||||
|
payment_hash=_HASH,
|
||||||
|
bitspire_txid=_TXID,
|
||||||
|
wire_sats=54440,
|
||||||
|
fiat_amount=40.0,
|
||||||
|
fiat_code="EUR",
|
||||||
|
exchange_rate=1361.0,
|
||||||
|
principal_sats=54440,
|
||||||
|
fee_sats=0,
|
||||||
|
platform_fee_sats=0,
|
||||||
|
operator_fee_sats=0,
|
||||||
|
tx_type=tx_type,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class _GateWired:
|
||||||
|
def __init__(self, monkeypatch, *, tx_type, early_report=None):
|
||||||
|
self.created = []
|
||||||
|
self.spawned = []
|
||||||
|
self.adopted = []
|
||||||
|
|
||||||
|
async def get_machine(_wid):
|
||||||
|
return _machine()
|
||||||
|
|
||||||
|
def attribution(_machine, _extra):
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def get_super():
|
||||||
|
return SimpleNamespace(id="default")
|
||||||
|
|
||||||
|
def parse(**_kw):
|
||||||
|
return _data(tx_type)
|
||||||
|
|
||||||
|
async def create(data, initial_status, error_message=None):
|
||||||
|
self.created.append((data.tx_type, initial_status))
|
||||||
|
return _settlement(status=initial_status, tx_type=data.tx_type)
|
||||||
|
|
||||||
|
async def process(sid):
|
||||||
|
self.spawned.append(sid)
|
||||||
|
|
||||||
|
async def early(_mid, _txid):
|
||||||
|
return early_report
|
||||||
|
|
||||||
|
async def adopt(settlement, machine, row):
|
||||||
|
self.adopted.append((settlement.id, row.id))
|
||||||
|
return "pending"
|
||||||
|
|
||||||
|
monkeypatch.setattr(tasks, "get_active_machine_by_wallet_id", get_machine)
|
||||||
|
monkeypatch.setattr(tasks, "assert_nostr_attribution", attribution)
|
||||||
|
monkeypatch.setattr(tasks, "get_super_config", get_super)
|
||||||
|
monkeypatch.setattr(tasks, "parse_settlement", parse)
|
||||||
|
monkeypatch.setattr(tasks, "create_settlement_idempotent", create)
|
||||||
|
monkeypatch.setattr(tasks, "process_settlement", process)
|
||||||
|
monkeypatch.setattr(crud_mod, "get_latest_unlinked_dispense_report", early)
|
||||||
|
monkeypatch.setattr(dispense_transport, "adopt_unlinked_report", adopt)
|
||||||
|
|
||||||
|
|
||||||
|
async def _drain():
|
||||||
|
# let any create_task'd distribution run
|
||||||
|
await asyncio.sleep(0)
|
||||||
|
|
||||||
|
|
||||||
|
class TestPaymentGate:
|
||||||
|
def test_cash_out_lands_awaiting_dispense_and_does_not_distribute(
|
||||||
|
self, monkeypatch
|
||||||
|
):
|
||||||
|
w = _GateWired(monkeypatch, tx_type="cash_out")
|
||||||
|
|
||||||
|
async def run():
|
||||||
|
await tasks._handle_payment(_payment(is_in=True))
|
||||||
|
await _drain()
|
||||||
|
|
||||||
|
asyncio.run(run())
|
||||||
|
assert w.created == [("cash_out", "awaiting_dispense")]
|
||||||
|
assert w.spawned == []
|
||||||
|
assert w.adopted == []
|
||||||
|
|
||||||
|
def test_cash_in_is_unchanged(self, monkeypatch):
|
||||||
|
w = _GateWired(monkeypatch, tx_type="cash_in")
|
||||||
|
|
||||||
|
async def run():
|
||||||
|
await tasks._handle_payment(_payment(is_in=False))
|
||||||
|
await _drain()
|
||||||
|
|
||||||
|
asyncio.run(run())
|
||||||
|
assert w.created == [("cash_in", "pending")]
|
||||||
|
assert w.spawned == ["s1"]
|
||||||
|
|
||||||
|
def test_early_report_is_adopted_when_the_payment_lands(self, monkeypatch):
|
||||||
|
early = SimpleNamespace(id="rep-early", txid=_TXID)
|
||||||
|
w = _GateWired(monkeypatch, tx_type="cash_out", early_report=early)
|
||||||
|
|
||||||
|
async def run():
|
||||||
|
await tasks._handle_payment(_payment(is_in=True))
|
||||||
|
await _drain()
|
||||||
|
|
||||||
|
asyncio.run(run())
|
||||||
|
assert w.adopted == [("s1", "rep-early")]
|
||||||
|
assert w.spawned == [] # adoption decides; the fake adopt did not spawn
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Consumer mirror
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
class TestCashOutHoldMirror:
|
||||||
|
def test_sets_and_clears_from_the_state_document(self):
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
async def setter(mid, since, reason, code):
|
||||||
|
calls.append((mid, since, reason, code))
|
||||||
|
|
||||||
|
held = PublishCassettesPayload(
|
||||||
|
positions={"2": {"denomination": 20, "count": 54}},
|
||||||
|
cash_out_held_since=1791529353,
|
||||||
|
cash_out_held_reason="Note stopped at the cassette exit",
|
||||||
|
cash_out_held_code="78 42",
|
||||||
|
)
|
||||||
|
clear = PublishCassettesPayload(
|
||||||
|
positions={"2": {"denomination": 20, "count": 54}}
|
||||||
|
)
|
||||||
|
asyncio.run(tasks._record_cash_out_hold("m1", held, setter))
|
||||||
|
asyncio.run(tasks._record_cash_out_hold("m1", clear, setter))
|
||||||
|
assert calls[0][0] == "m1"
|
||||||
|
assert calls[0][1] == datetime.fromtimestamp(1791529353, tz=timezone.utc)
|
||||||
|
assert calls[0][2:] == ("Note stopped at the cassette exit", "78 42")
|
||||||
|
assert calls[1] == ("m1", None, None, None)
|
||||||
84
views_api.py
84
views_api.py
|
|
@ -19,6 +19,7 @@ from lnbits.core.services.nsec_bunker import (
|
||||||
)
|
)
|
||||||
from lnbits.decorators import check_super_user, check_user_exists
|
from lnbits.decorators import check_super_user, check_user_exists
|
||||||
from lnbits.utils.nostr import normalize_public_key
|
from lnbits.utils.nostr import normalize_public_key
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from .calculations import MAX_FEE_FRACTION_PER_DIRECTION
|
from .calculations import MAX_FEE_FRACTION_PER_DIRECTION
|
||||||
from .cassette_transport import (
|
from .cassette_transport import (
|
||||||
|
|
@ -28,15 +29,6 @@ from .cassette_transport import (
|
||||||
SignerUnavailable,
|
SignerUnavailable,
|
||||||
publish_ops_to_atm,
|
publish_ops_to_atm,
|
||||||
)
|
)
|
||||||
from .fee_transport import publish_fee_config
|
|
||||||
from .pairing import (
|
|
||||||
PairResult,
|
|
||||||
PairingError,
|
|
||||||
RevokeResult,
|
|
||||||
default_relay_endpoint,
|
|
||||||
pair_spire,
|
|
||||||
revoke_spire,
|
|
||||||
)
|
|
||||||
from .crud import (
|
from .crud import (
|
||||||
append_settlement_note,
|
append_settlement_note,
|
||||||
count_completed_legs_for_settlement,
|
count_completed_legs_for_settlement,
|
||||||
|
|
@ -85,6 +77,7 @@ from .distribution import (
|
||||||
process_settlement,
|
process_settlement,
|
||||||
settle_lp_balance,
|
settle_lp_balance,
|
||||||
)
|
)
|
||||||
|
from .fee_transport import publish_fee_config
|
||||||
from .models import (
|
from .models import (
|
||||||
AppendSettlementNoteData,
|
AppendSettlementNoteData,
|
||||||
CassetteConfig,
|
CassetteConfig,
|
||||||
|
|
@ -112,6 +105,14 @@ from .models import (
|
||||||
UpdateMachineData,
|
UpdateMachineData,
|
||||||
UpdateSuperConfigData,
|
UpdateSuperConfigData,
|
||||||
)
|
)
|
||||||
|
from .pairing import (
|
||||||
|
PairingError,
|
||||||
|
PairResult,
|
||||||
|
RevokeResult,
|
||||||
|
default_relay_endpoint,
|
||||||
|
pair_spire,
|
||||||
|
revoke_spire,
|
||||||
|
)
|
||||||
|
|
||||||
spirekeeper_api_router = APIRouter()
|
spirekeeper_api_router = APIRouter()
|
||||||
|
|
||||||
|
|
@ -768,7 +769,13 @@ async def api_list_stuck_settlements(
|
||||||
) -> StuckSettlementsResponse:
|
) -> StuckSettlementsResponse:
|
||||||
"""Operator worklist of settlements that didn't process cleanly.
|
"""Operator worklist of settlements that didn't process cleanly.
|
||||||
|
|
||||||
Returns four lists:
|
Returns seven lists. The first three (ADR-005 §6) mean a customer is owed
|
||||||
|
money and render first:
|
||||||
|
- cash_owed: the machine reported nothing dispensed; nothing moved
|
||||||
|
- partial_pending: some notes out, value short; held until resolved
|
||||||
|
- dispense_unreported: cash-out landed, machine never reported within
|
||||||
|
the threshold
|
||||||
|
Then:
|
||||||
- rejected: Nostr attribution cross-check failed — signer didn't
|
- rejected: Nostr attribution cross-check failed — signer didn't
|
||||||
match the machine identity. Investigate; do not retry.
|
match the machine identity. Investigate; do not retry.
|
||||||
- errored: distribution ran and failed; retry endpoint handles these
|
- errored: distribution ran and failed; retry endpoint handles these
|
||||||
|
|
@ -783,6 +790,9 @@ async def api_list_stuck_settlements(
|
||||||
buckets = await get_stuck_settlements_for_operator(user.id, threshold_minutes)
|
buckets = await get_stuck_settlements_for_operator(user.id, threshold_minutes)
|
||||||
return StuckSettlementsResponse(
|
return StuckSettlementsResponse(
|
||||||
threshold_minutes=threshold_minutes,
|
threshold_minutes=threshold_minutes,
|
||||||
|
cash_owed=buckets["cash_owed"],
|
||||||
|
partial_pending=buckets["partial_pending"],
|
||||||
|
dispense_unreported=buckets["dispense_unreported"],
|
||||||
rejected=buckets["rejected"],
|
rejected=buckets["rejected"],
|
||||||
errored=buckets["errored"],
|
errored=buckets["errored"],
|
||||||
stuck_pending=buckets["stuck_pending"],
|
stuck_pending=buckets["stuck_pending"],
|
||||||
|
|
@ -1229,3 +1239,57 @@ async def api_create_machine_cassette_op(
|
||||||
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
|
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
|
||||||
|
|
||||||
return op
|
return op
|
||||||
|
@spirekeeper_api_router.post(
|
||||||
|
"/api/v1/dca/machines/{machine_id}/resume-cash-out",
|
||||||
|
response_model=CassetteOp,
|
||||||
|
)
|
||||||
|
async def api_resume_cash_out(
|
||||||
|
machine_id: str,
|
||||||
|
user: User = Depends(check_user_exists),
|
||||||
|
) -> CassetteOp:
|
||||||
|
"""Release a machine's cash-out hold (bitspire ADR-005 §5).
|
||||||
|
|
||||||
|
After a terminal dispenser fault the machine refuses cash-out until an
|
||||||
|
operator has been to it. A `recount` releases the hold as a side effect;
|
||||||
|
this is for the case where the jam was cleared without touching a bay
|
||||||
|
count. Recorded as a machine-wide op (position 0) and published on the
|
||||||
|
same operator channel as the cassette ops — the machine honours it only if
|
||||||
|
it is stamped after the hold began, so a re-delivered old resume cannot
|
||||||
|
clear a newer fault.
|
||||||
|
|
||||||
|
Errors mirror the cassette-op endpoint: 400 unpaired, 503 signer/relay
|
||||||
|
unavailable (the op is recorded and rides out with the next publish).
|
||||||
|
"""
|
||||||
|
machine = await _machine_owned_by(machine_id, user.id)
|
||||||
|
if not machine.machine_npub:
|
||||||
|
raise HTTPException(
|
||||||
|
HTTPStatus.BAD_REQUEST,
|
||||||
|
"machine is not paired — there is no ATM identity to publish to",
|
||||||
|
)
|
||||||
|
if machine.cash_out_held_since is None:
|
||||||
|
logger.info(
|
||||||
|
f"spirekeeper: resume_cash_out for machine {machine_id} with no hold "
|
||||||
|
"on file — publishing anyway (the machine is the authority)"
|
||||||
|
)
|
||||||
|
|
||||||
|
op = await create_cassette_op(
|
||||||
|
machine_id,
|
||||||
|
CreateCassetteOpData(position=0, op_type="resume_cash_out"),
|
||||||
|
created_by=user.id,
|
||||||
|
)
|
||||||
|
window = await get_cassette_ops_window(machine_id)
|
||||||
|
try:
|
||||||
|
await publish_ops_to_atm(machine, window, user.id)
|
||||||
|
except OperatorIdentityMissing as exc:
|
||||||
|
raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc
|
||||||
|
except (SignerUnavailable, RelayUnavailable) as exc:
|
||||||
|
raise HTTPException(
|
||||||
|
HTTPStatus.SERVICE_UNAVAILABLE,
|
||||||
|
f"{exc} — the resume was recorded and will be delivered with the "
|
||||||
|
"next publish",
|
||||||
|
) from exc
|
||||||
|
except CassetteTransportError as exc:
|
||||||
|
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
|
||||||
|
return op
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue