spirekeeper/tasks.py
Padreug 44c2afa5bb feat(transport): report_dispense RPC — capture a cash-out on the machine's report, not on payment (ADR-005 §1–§2)
The structural fix for bitspire#122. _handle_payment used to spawn
process_settlement the instant a cash_out payment landed — before the
machine had begun to dispense — so a jam two seconds later found the
legs already paid and `processed` was the honest answer. Payment is now
the authorization; the machine's report is the capture.

A cash_out lands as awaiting_dispense and is not distributed. The new
`report_dispense` handler (identity from the VERIFIED transport sender,
same as create_withdraw / get_machine_config) stores every report
append-only and moves the settlement: dispense_confirmed → pending and
distribution runs; some notes out → partial_pending, held whole
(ADR-005 Decision 1, one distribution when the shortfall is resolved);
nothing out → cash_owed, first on the worklist. A report naming
remediates_txid moves the owed settlement it names to pending in full.
Already-captured settlements are recorded but never moved — a report
cannot un-pay legs. A byte-identical resend is acked without a new row.

Both orders of arrival are handled: a report that precedes its payment
(hold invoices settle after the dispense; the invoice listener can lag)
is stored unlinked and adopted when the settlement is inserted, through
the same transition. counts_uncertain on a report mirrors onto the
machine immediately rather than at the next heartbeat. The state-event
consumer mirrors cash_out_held_* onto dca_machines, including clearing it.

Soft-fails like the other RPCs: without register_rpc the settlements sit
in awaiting_dispense and surface as dispense_unreported — the honest state.
2026-10-10 21:51:51 +02:00

621 lines
26 KiB
Python

# Satoshi Machine v2 — invoice listener (P1 + fix bundle 2).
#
# Subscribes to LNbits' invoice dispatcher (register_invoice_listener), then
# for each successful inbound payment:
# 1. Checks if wallet_id belongs to an active dca_machines row. If not, skip.
# 2. Verifies the originating Nostr signer matches the machine identity
# (assert_nostr_attribution; uses Payment.extra.nostr_sender_pubkey
# stamped by lnbits nostr-transport dispatcher).
# 3. Parses Payment.extra for bitSpire's canonical split stamp per
# aiolabs/lamassu-next#44 (`source: "bitspire"`, principal_sats,
# fee_sats, exchange_rate). Raises if the stamp is missing or
# garbage (no more Lamassu-era reverse-derivation fallback).
# 4. Computes the two-stage split (super_fee first, operator remainder).
# 5. Inserts a dca_settlements row idempotently (keyed by payment_hash).
# 6. Spawns the distribution processor on a background task so the
# LNbits invoice queue (which serves ALL extensions on the node)
# keeps draining while we move sats. Concurrency is safe because
# process_settlement now uses an optimistic-lock claim (fix bundle 1).
#
# Rejection paths (settlement still recorded with status='rejected' for
# operator forensics, but distribution is skipped):
# - SettlementAttributionError: signer mismatch (G5).
# - SettlementMetadataError: Payment.extra missing bitSpire stamp.
# - SettlementInvariantError: stamped values violate the canonical
# sat-amount invariants (range/sum).
import asyncio
from lnbits.core.models import Payment
from lnbits.tasks import register_invoice_listener
from loguru import logger
from .bitspire import (
SettlementAttributionError,
SettlementInvariantError,
SettlementMetadataError,
assert_nostr_attribution,
parse_settlement,
)
from .crud import (
create_settlement_idempotent,
get_active_machine_by_wallet_id,
get_super_config,
)
from .distribution import process_settlement
from .models import CreateDcaSettlementData, Machine
LISTENER_NAME = "ext_spirekeeper"
# Holds strong refs to in-flight distribution tasks so Python's GC doesn't
# collect them mid-flight (asyncio.create_task only weakly references its
# task once awaiters drop). Tasks self-clean by removing themselves on
# completion via the done_callback below.
_inflight_distributions: set = set()
async def wait_for_paid_invoices() -> None:
invoice_queue: asyncio.Queue = asyncio.Queue()
register_invoice_listener(invoice_queue, LISTENER_NAME)
logger.info(
"spirekeeper v2: invoice listener registered as "
f"`{LISTENER_NAME}` — waiting for bitSpire settlements."
)
while True:
payment: Payment = await invoice_queue.get()
try:
await _handle_payment(payment)
except Exception as exc: # listener must never die
logger.error(
f"spirekeeper: error handling payment "
f"{payment.payment_hash[:12]}...: {exc}"
)
async def _handle_payment(payment: Payment) -> None:
if not payment.success:
return
machine = await get_active_machine_by_wallet_id(payment.wallet_id)
if machine is None:
return
extra = payment.extra or {}
# Two axes, deliberately named in pairs to avoid the inversion trap
# documented at `~/.claude/projects/.../memory/feedback_naming_business_vs_protocol.md`:
#
# - is_lightning_inbound / is_lightning_outbound: PROTOCOL direction
# at the operator's wallet. `payment.is_in` from LNbits.
# - tx_type ∈ {"cash_out", "cash_in"}: BUSINESS direction at the ATM.
# Sourced from Payment.extra (canonical, stamped by bitSpire).
#
# Canonical mapping:
# cash_out ↔ is_lightning_inbound (customer pays ATM's invoice in BTC,
# operator wallet receives sats)
# cash_in ↔ is_lightning_outbound (customer redeems ATM's LNURL-
# withdraw, operator wallet sends sats)
#
# Process BOTH directions; reject mismatches at the discriminator gate.
is_lightning_inbound = payment.is_in
is_lightning_outbound = not payment.is_in
# Outbound payments from the operator's wallet need an extra
# discriminator before we touch them. An operator may legitimately
# send sats for non-ATM reasons (manual send, different extension,
# etc.). Without `source=bitspire` on Payment.extra we can't tell
# the operator paying their landlord from a cash-in settlement —
# skip silently. (For cash-out / inbound payments we already gate
# on machine-owned wallet via `get_active_machine_by_wallet_id`.)
if is_lightning_outbound and extra.get("source") != "bitspire":
return
# 1) Attribution FIRST — uses only `extra.nostr_sender_pubkey` (no parse
# needed). If this fails, every subsequent field on `extra` is
# attacker-controlled and untrustworthy — record a minimal rejected
# row with placeholder zeros (don't display unverified split numbers
# in the operator dashboard).
try:
assert_nostr_attribution(machine, extra)
except SettlementAttributionError as exc:
await _record_rejected(payment, machine, exc)
return
# 2) Parse + invariants. parse_settlement enforces the canonical
# sat-amount invariants on the bitSpire-stamped numbers (range +
# direction-specific sum). Raises SettlementMetadataError if the
# stamp is missing, SettlementInvariantError on any range/sum
# breach.
super_config = await get_super_config()
assert super_config is not None # m001 inserts the default singleton
try:
data = parse_settlement(
machine=machine,
payment_hash=payment.payment_hash,
# `payment.sat` is signed by protocol direction (negative for an
# outbound cash-in payout, positive for an inbound cash-out
# receipt). The settlement's `wire_sats` is a magnitude — direction
# is carried separately by `tx_type` — so pass the absolute value.
wire_sats=abs(payment.sat),
extra=extra,
super_config=super_config,
)
except (SettlementMetadataError, SettlementInvariantError) as exc:
await _record_rejected(payment, machine, exc)
return
# Cross-axis sanity: protocol direction must agree with business
# direction per the canonical mapping above. A mismatch means
# something upstream is confused — refuse to process. Concrete
# symptom this catches: an attacker (or a buggy extension) stamps
# `source=bitspire, type=cash_out` on an outbound payment from the
# operator's wallet to attempt a fake "we just received sats" row.
expected_inbound = data.tx_type == "cash_out"
if is_lightning_inbound != expected_inbound:
await _record_rejected(
payment,
machine,
SettlementInvariantError(
f"direction mismatch: payment.is_in={is_lightning_inbound} "
f"but tx_type={data.tx_type!r}. Expected cash_out ↔ inbound, "
"cash_in ↔ outbound."
),
)
return
del is_lightning_outbound # only used for the discriminator above
# Stamp the originating Nostr event id (the kind-21000 create_invoice
# RPC) onto the row for post-hoc forensics — an auditor can trace
# settlement → RPC event → signing key without trusting our DB.
nostr_event_id = extra.get("nostr_event_id")
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"
)
if settlement is None:
logger.error(
f"spirekeeper: failed to insert settlement for "
f"payment_hash={payment.payment_hash[:12]}..."
)
return
logger.info(
f"spirekeeper: landed settlement {settlement.id} for "
f"machine={machine.machine_npub[:12]}... "
f"wire={data.wire_sats}sats principal={data.principal_sats}sats "
f"fee={data.fee_sats}sats "
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
# so a listener re-fire can't double-process. Listener latency is now
# bounded by the create_settlement_idempotent insert, not by the N+M
# internal pay_invoice round-trips of a full distribution.
task = asyncio.create_task(process_settlement(settlement.id))
_inflight_distributions.add(task)
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.
Used for every rejection path (attribution / metadata / invariant).
The split fields are zero placeholders — we deliberately do NOT
display attacker-supplied numbers in the operator dashboard. The
wire amount (`payment.sat`) is the only value LNbits authenticated;
everything else from Payment.extra is untrusted in this branch.
"""
data = CreateDcaSettlementData(
machine_id=machine.id,
payment_hash=payment.payment_hash,
# Magnitude, not the signed `payment.sat` (negative for outbound).
wire_sats=abs(payment.sat),
fiat_amount=0.0,
fiat_code=machine.fiat_code,
exchange_rate=0.0,
principal_sats=0,
fee_sats=0,
platform_fee_sats=0,
operator_fee_sats=0,
# The parsed tx_type is unavailable on the rejection path, but the
# authenticated protocol direction is: an outbound payment is a
# cash-in, an inbound one a cash-out. Use that so a rejected row shows
# the right direction instead of always reading "cash-out".
tx_type="cash_in" if not payment.is_in else "cash_out",
)
rejected = await create_settlement_idempotent(
data, initial_status="rejected", error_message=str(exc)
)
if rejected is None:
logger.error(
f"spirekeeper: failed to insert rejected settlement for "
f"payment_hash={payment.payment_hash[:12]}..."
)
return
logger.error(
f"spirekeeper: rejected settlement {rejected.id} "
# An unpaired machine (machine_npub None) reaches here now that
# assert_nostr_attribution rejects it — fall back to the id so the
# log line doesn't crash on None[:12].
f"(machine={(machine.machine_npub or machine.id)[:12]}..., "
f"payment_hash={payment.payment_hash[:12]}...): {exc}"
)
# =============================================================================
# Cassette state consumer (#29)
# =============================================================================
# Subscribes to kind-30078 bitspire-cassettes-state:<atm_pubkey_hex> events
# published by each active machine's ATM. Decrypts the NIP-44 v2 content with
# the operator's privkey + ATM sender pubkey, validates as
# PublishCassettesPayload, and reconciles cassette_configs via
# apply_reported_state.
#
# This is continuous, not one-shot: the ATM publishes on startup, after every
# change to its bays, and on a heartbeat, and each event replaces the last.
# The comments here used to call it a "v1 bootstrap" consumer, which was
# misleading — there has never been a once-per-machine guard, so every event
# an ATM published was already being applied. Ordering is enforced in
# apply_reported_state by created_at; it is not inferred from arrival order.
#
# Implementation: polls nostrclient.router.NostrRouter.received_subscription_
# events keyed by our subscription_id. nostrclient's NostrRouter design is
# per-WebSocket-client; the singleton dict it drains into is the only
# server-side hook to consume events without standing up an in-process
# websocket. The relay manager is the same singleton publish_to_atm uses,
# so add_subscription registers a filter against the same relay pool.
CASSETTE_BOOTSTRAP_SUB_ID = "spirekeeper-cassette-bootstrap"
_CASSETTE_POLL_INTERVAL_S = 2.0
_CASSETTE_BACKOFF_S = 30.0 # when nostrclient isn't installed yet
async def wait_for_cassette_state_events() -> None:
"""Long-running task: subscribe to bitspire-cassettes-state events from
every active machine's ATM and upsert cassette_configs on receipt.
Pattern mirrors wait_for_paid_invoices (try/except wraps each event,
never lets the loop die). Re-derives the subscription filter on each
tick from the current active-machines list — newly-added machines
start receiving bootstrap events without an LNbits restart.
Soft-fail surfaces:
- nostrclient not installed → log + sleep _CASSETTE_BACKOFF_S
between retries (operator may install it later)
- inbound event fails sig-verify / decrypt / parse → log + skip
the event, continue the loop
- apply_reported_state errors → log + skip
"""
logger.info(
"spirekeeper v2: cassette bootstrap consumer starting "
f"(sub_id={CASSETTE_BOOTSTRAP_SUB_ID})"
)
current_filter_key: str | None = None
while True:
try:
current_filter_key = await _cassette_consumer_tick(current_filter_key)
await asyncio.sleep(_CASSETTE_POLL_INTERVAL_S)
except _NostrclientUnavailable:
logger.warning(
"spirekeeper: nostrclient extension not installed; "
f"cassette bootstrap consumer sleeping {_CASSETTE_BACKOFF_S}s "
"before retry. Install + activate nostrclient on this "
"LNbits instance."
)
current_filter_key = None
await asyncio.sleep(_CASSETTE_BACKOFF_S)
except Exception as exc: # listener must never die
logger.error(
f"spirekeeper: cassette consumer loop error (continuing): " f"{exc}"
)
await asyncio.sleep(_CASSETTE_POLL_INTERVAL_S)
class _NostrclientUnavailable(Exception):
"""Internal sentinel — nostrclient extension import failed. Caller
sleeps a backoff then retries; the operator may install nostrclient
at any time."""
async def _cassette_consumer_tick(current_filter_key: str | None) -> str:
"""Single iteration of the bootstrap-consumer loop. Returns the filter
key used this tick so the caller can detect filter-set changes.
Raises _NostrclientUnavailable if nostrclient can't be imported (the
outer loop backs off + retries).
"""
try:
from nostrclient.router import ( # type: ignore[import-not-found]
NostrRouter,
nostr_client,
)
except ImportError as exc:
raise _NostrclientUnavailable() from exc
from .cassette_transport import build_state_d_tags_for_machines
from .crud import (
apply_reported_state,
get_machine_by_atm_pubkey_hex,
list_all_active_machines,
mark_cassette_ops_acked,
set_machine_counts_uncertain,
)
machines = await list_all_active_machines()
d_tags = build_state_d_tags_for_machines(machines)
filter_key = ",".join(sorted(d_tags))
if filter_key != current_filter_key:
if d_tags:
filters = [{"kinds": [30078], "#d": d_tags}]
# nostrclient's add_subscription is typed as list[str] but the
# actual relay protocol accepts list[Filter-dict] — type ignore
# the upstream typing mismatch.
nostr_client.relay_manager.add_subscription(
CASSETTE_BOOTSTRAP_SUB_ID, filters # type: ignore[arg-type]
)
logger.info(
"spirekeeper: (re)registered cassette bootstrap "
f"subscription with {len(d_tags)} d-tag(s)"
)
else:
nostr_client.relay_manager.close_subscription(CASSETTE_BOOTSTRAP_SUB_ID)
logger.info(
"spirekeeper: no active machines; closed cassette "
"bootstrap subscription"
)
inbound = NostrRouter.received_subscription_events.get(CASSETTE_BOOTSTRAP_SUB_ID)
if inbound:
while inbound:
event_message = inbound.pop(0)
try:
await _handle_cassette_state_event(
event_message,
get_machine_by_atm_pubkey_hex,
apply_reported_state,
mark_cassette_ops_acked,
set_machine_counts_uncertain,
)
except Exception as exc:
logger.warning(
f"spirekeeper: cassette state event handler "
f"failed (skipping): {exc}"
)
return filter_key
async def _record_op_acknowledgements(
machine_id: str, payload, mark_cassette_ops_acked
) -> None:
"""Mark the operations a machine reports as applied.
Deliberately not gated on whether the state event advanced the counts. The
machine echoes its applied-op ids on EVERY state publish, so an event
carrying nothing new about the counts can still be the first one to tell us
an operation landed; gating on that would lose the acknowledgement.
This echo is the only acknowledgement the transport can carry. An
addressable event gives its publisher no failure signal at all — the relay
returns OK for an event it then discards — so without it the dashboard
could only ever show an operation as sent, never as delivered.
"""
if not payload.applied_ops:
return
newly_acked = await mark_cassette_ops_acked(machine_id, payload.applied_ops)
if newly_acked:
logger.info(
f"spirekeeper: machine {machine_id} acknowledged "
f"{newly_acked} cassette operation(s)"
)
async def _record_counts_uncertainty(
machine_id: str, payload, set_machine_counts_uncertain
) -> None:
"""Mirror the machine's counts-uncertain marker onto its registry row.
The machine sets this when a dispense ended without a reliable count of
what physically left the bay — a dispenser throw, or a timeout. It cannot
know how many notes moved, so it says so instead of decrementing a number
it would be guessing at.
Written on every state event, including when it is None, because the
machine clearing the marker is exactly as important as setting it: the
operator has recounted, the bay is trustworthy again, and a banner that
never goes away is a banner nobody reads.
"""
from datetime import datetime as _datetime
from datetime import timezone as _timezone
since = None
if payload.counts_uncertain_since is not None:
since = _datetime.fromtimestamp(
int(payload.counts_uncertain_since), tz=_timezone.utc
)
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
prvkey on the LocalSigner transitional fallback inside the transport
helper), parse, upsert.
Each step logs at WARNING (not ERROR) so a noisy attacker can't fill
the logs — this is data on a public relay, garbage is expected.
Two skip outcomes:
- Terminal (CassetteEventDecodeError / SignerUnavailable /
OperatorIdentityMissing / etc.): log + return. `apply_bootstrap_
state` is never called → `state_event_id` is not advanced →
same event would re-process on next poll cycle but the consumer's
WARN log surfaces the underlying issue immediately.
- Transient (CassetteEventTransientError): log at INFO (less noisy)
+ return. Same retry-via-no-advance semantics, just less
alarming in the operator log feed.
"""
import json as _json
from datetime import datetime as _datetime
from datetime import timezone as _timezone
from lnbits.utils.nostr import verify_event
from .cassette_transport import (
CassetteEventDecodeError,
CassetteEventTransientError,
CassetteTransportError,
decrypt_and_parse_state_event,
)
from .nostr_publish import resolve_operator_signer
event_raw = event_message.event
if isinstance(event_raw, str):
event_obj = _json.loads(event_raw)
elif isinstance(event_raw, dict):
event_obj = event_raw
else:
logger.warning(
f"spirekeeper: cassette event of unexpected type "
f"{type(event_raw).__name__}; skipping"
)
return
if not verify_event(event_obj):
logger.warning(
f"spirekeeper: cassette state event sig verify failed "
f"(id={event_obj.get('id', '?')[:12]}...)"
)
return
sender_pubkey = event_obj.get("pubkey", "")
machine = await get_machine_by_atm_pubkey_hex(sender_pubkey)
if machine is None:
# Unknown sender — could be relay noise or an attacker. Don't
# treat as our problem.
logger.warning(
f"spirekeeper: cassette state event from unknown ATM "
f"pubkey {sender_pubkey[:12]}... (not in dca_machines); "
"skipping"
)
return
try:
account, signer = await resolve_operator_signer(machine.operator_user_id)
except CassetteTransportError as exc:
# OperatorIdentityMissing / SignerUnavailable — log + skip.
logger.warning(
f"spirekeeper: can't resolve signer for operator "
f"{machine.operator_user_id[:8]}... (machine {machine.id}): "
f"{exc}"
)
return
try:
payload = await decrypt_and_parse_state_event(event_obj, account, signer)
except CassetteEventTransientError as exc:
logger.info(
f"spirekeeper: cassette state event for machine {machine.id} "
f"hit a transient signer error (will retry next poll): {exc}"
)
return
except CassetteEventDecodeError as exc:
logger.warning(
f"spirekeeper: cassette state event decode failed for "
f"machine {machine.id} (id={event_obj.get('id', '?')[:12]}...): "
f"{exc}"
)
return
event_id = event_obj.get("id", "")
created_at_unix = event_obj.get("created_at", 0)
event_created_at = _datetime.fromtimestamp(int(created_at_unix), tz=_timezone.utc)
applied = await apply_reported_state(
machine.id, event_id, event_created_at, payload
)
if applied:
logger.info(
f"spirekeeper: applied reported state event {event_id[:12]}... "
f"to machine {machine.id} ({len(payload.positions)} cassettes)"
)
else:
# Replay or an older event. Normal on relay reconnect.
logger.debug(
f"spirekeeper: cassette state event {event_id[:12]}... "
f"not newer than stored state for machine {machine.id} (no-op)"
)
# Acknowledgement runs regardless of whether the counts were newer — see
# _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)