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