feat(schema): dispense outcome on settlements, dispense_reports, cash-out hold mirror (ADR-005)

m016: an append-only `dispense_reports` table (lamassu-server's
cash_out_actions shape — one row per report the machine sent, so a
retry, a late report and a remediation report stay distinct);
dispense_confirmed / dispense_error / dispense_error_code /
dispense_raw_code / dispense_error_class / dispense_reported_at /
dispensed_fiat_cents on dca_settlements; cash_out_held_since / _reason /
_code on dca_machines beside counts_uncertain_since.

Settlement lifecycle gains awaiting_dispense (cash_out at insert — paid,
waiting for the machine's report), partial_pending (some notes out,
value short; held whole until the operator records the resolution) and
cash_owed (nothing out; legs never run). dispense_unreported is derived
by the worklist, not stored.

resume_cash_out joins CASSETTE_OP_TYPES as a machine-wide op: position 0,
no position on the wire, no bay fields. It rides the operator channel the
machine already consumes; the machine honours it only if stamped after
the hold began. A recount releases the hold too.

crud: get_settlement_by_txid (the join the machine's extra.txid already
provides), apply_dispense_outcome (copies the report onto the settlement
and finally writes bills_json / cassettes_json with what actually came
out), the dispense_reports accessors incl. adopting a report that
arrived before its payment, set_machine_cash_out_hold, and the three
new worklist buckets.
This commit is contained in:
Padreug 2026-10-10 21:51:51 +02:00
commit b8a5e6352a
3 changed files with 511 additions and 18 deletions

250
crud.py
View file

@ -27,6 +27,8 @@ from .models import (
DcaLpPreferences,
DcaPayment,
DcaSettlement,
DispenseReport,
DispenseReportIn,
Machine,
PublishCassettesPayload,
SuperConfig,
@ -257,6 +259,32 @@ 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.
@ -711,6 +739,171 @@ 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",
@ -752,7 +945,14 @@ async def get_stuck_settlements_for_operator(
) -> dict:
"""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
failed — signer didn't match the machine identity). Distinct
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},
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,
@ -870,9 +1111,10 @@ async def mark_settlement_status(
status: str,
error_message: str | None = None,
) -> DcaSettlement | None:
"""Status: 'pending' | 'processing' | 'processed' | 'partial' |
'refunded' | 'errored'. Clears processing_claim on terminal states so a
fresh claim attempt won't see a stale token."""
"""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."""
await db.execute(
"""
UPDATE spirekeeper.dca_settlements

View file

@ -941,3 +941,82 @@ 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
View file

@ -67,6 +67,13 @@ 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
@ -311,19 +318,39 @@ class DcaSettlement(BaseModel):
fee_mismatch_sats: int | None = None
bills_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)
# '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)
# '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
@ -336,6 +363,110 @@ 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.
# =============================================================================
@ -570,6 +701,13 @@ 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
@ -738,6 +876,11 @@ 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):
@ -803,7 +946,21 @@ class PublishCassettesPayload(BaseModel):
# 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):
@ -837,11 +994,17 @@ class CassetteOp(BaseModel):
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
return v
@validator("position")
def _position_positive(cls, v):
if v <= 0:
raise ValueError(f"position must be > 0, got {v}")
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
def to_wire_dict(self) -> dict:
"""The published form. Drops the fields this op_type does not use, so
@ -850,8 +1013,10 @@ 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":
@ -880,11 +1045,17 @@ class CreateCassetteOpData(BaseModel):
raise ValueError(f"op_type must be one of {CASSETTE_OP_TYPES}, got {v!r}")
return v
@validator("position")
def _position_positive(cls, v):
if v <= 0:
raise ValueError(f"position must be > 0, got {v}")
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("bills")
def _bills_positive(cls, v):
@ -911,6 +1082,7 @@ 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}`")