spirekeeper/crud.py
Padreug b8a5e6352a 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.
2026-10-10 21:51:51 +02:00

2029 lines
71 KiB
Python

# Satoshi Machine v2 — CRUD layer over the m005 schema.
#
# All operator-scoped queries take an operator_user_id and enforce isolation
# at the SQL boundary. Cross-operator LP queries (for satmachineclient) join
# through dca_machines.operator_user_id. See plan section "Identity & multi-
# machine model".
from datetime import datetime, timezone
from lnbits.db import Database
from lnbits.helpers import urlsafe_short_hash
from .models import (
CassetteConfig,
CassetteOp,
ClientBalanceSummary,
CommissionSplit,
CommissionSplitLeg,
CreateCassetteOpData,
CreateDcaClientData,
CreateDcaPaymentData,
CreateDcaSettlementData,
CreateDepositData,
CreateMachineData,
DcaClient,
DcaDeposit,
DcaLpPreferences,
DcaPayment,
DcaSettlement,
DispenseReport,
DispenseReportIn,
Machine,
PublishCassettesPayload,
SuperConfig,
TelemetrySnapshot,
UpdateDcaClientData,
UpdateDepositData,
UpdateDepositStatusData,
UpdateMachineData,
UpdateSuperConfigData,
UpsertDcaLpData,
)
db = Database("ext_spirekeeper")
# =============================================================================
# Super config
# =============================================================================
async def get_super_config() -> SuperConfig | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.super_config WHERE id = :id",
{"id": "default"},
SuperConfig,
)
async def update_super_config(data: UpdateSuperConfigData) -> SuperConfig | None:
update_data = {k: v for k, v in data.dict().items() if v is not None}
if not update_data:
return await get_super_config()
update_data["updated_at"] = datetime.now()
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
update_data["id"] = "default"
await db.execute(
f"UPDATE spirekeeper.super_config SET {set_clause} WHERE id = :id",
update_data,
)
return await get_super_config()
# =============================================================================
# Machines
# =============================================================================
async def create_machine(operator_user_id: str, data: CreateMachineData) -> Machine:
machine_id = urlsafe_short_hash()
now = datetime.now()
await db.execute(
"""
INSERT INTO spirekeeper.dca_machines
(id, operator_user_id, machine_npub, wallet_id, name, location,
fiat_code, is_active,
operator_cash_in_fee_fraction, operator_cash_out_fee_fraction,
created_at, updated_at)
VALUES (:id, :operator_user_id, :machine_npub, :wallet_id, :name,
:location, :fiat_code, :is_active,
:operator_cash_in_fee_fraction, :operator_cash_out_fee_fraction,
:created_at, :updated_at)
""",
{
"id": machine_id,
"operator_user_id": operator_user_id,
"machine_npub": data.machine_npub,
"wallet_id": data.wallet_id,
"name": data.name,
"location": data.location,
"fiat_code": data.fiat_code,
"is_active": True,
"operator_cash_in_fee_fraction": data.operator_cash_in_fee_fraction,
"operator_cash_out_fee_fraction": data.operator_cash_out_fee_fraction,
"created_at": now,
"updated_at": now,
},
)
machine = await get_machine(machine_id)
assert machine is not None
return machine
async def get_machine(machine_id: str) -> Machine | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_machines WHERE id = :id",
{"id": machine_id},
Machine,
)
async def get_machine_by_npub(machine_npub: str) -> Machine | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_machines WHERE machine_npub = :npub",
{"npub": machine_npub},
Machine,
)
async def get_active_machine_by_wallet_id(wallet_id: str) -> Machine | None:
"""Used by the invoice listener to route an incoming payment to a machine."""
return await db.fetchone(
"""
SELECT * FROM spirekeeper.dca_machines
WHERE wallet_id = :wid AND is_active = true
LIMIT 1
""",
{"wid": wallet_id},
Machine,
)
async def get_machines_for_operator(operator_user_id: str) -> list[Machine]:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_machines
WHERE operator_user_id = :uid
ORDER BY created_at DESC
""",
{"uid": operator_user_id},
Machine,
)
async def list_all_active_machines() -> list[Machine]:
"""Used by the cassette bootstrap consumer task to build a single
cross-operator subscription filter. Each event's pubkey routes to
the right operator via get_machine_by_atm_pubkey_hex + the machine's
operator_user_id.
"""
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_machines
WHERE is_active = true
ORDER BY created_at DESC
""",
{},
Machine,
)
async def get_machine_by_atm_pubkey_hex(atm_pubkey_hex: str) -> Machine | None:
"""Look up an active machine by its ATM pubkey, accepting hex or bech32
in machine_npub. Used by the cassette bootstrap consumer to route an
incoming state event to the right machine row (and therefore operator
privkey for decryption).
O(N) over active machines — fine for small fleets. If fleet sizes
grow, normalise machine_npub-at-write to hex and add an index.
"""
from lnbits.utils.nostr import normalize_public_key
target = atm_pubkey_hex.lower()
machines = await list_all_active_machines()
for m in machines:
# Unpaired machines (machine_npub is None — nullable since #29/m011)
# have no identity to match and would raise AttributeError in
# normalize_public_key (not caught below); skip them.
if not m.machine_npub:
continue
try:
if normalize_public_key(m.machine_npub).lower() == target:
return m
except (ValueError, AssertionError):
continue
return None
async def update_machine(machine_id: str, data: UpdateMachineData) -> Machine | None:
update_data = {k: v for k, v in data.dict().items() if v is not None}
if not update_data:
return await get_machine(machine_id)
update_data["updated_at"] = datetime.now()
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
update_data["id"] = machine_id
await db.execute(
f"UPDATE spirekeeper.dca_machines SET {set_clause} WHERE id = :id",
update_data,
)
return await get_machine(machine_id)
async def set_machine_pairing(
machine_id: str,
*,
machine_npub: str,
bunker_spire_key_name: str,
paired_at: datetime,
) -> Machine | None:
"""Persist the result of a (re-)pair: the bunker-minted spire identity
becomes the machine's npub (so lnbits' path-B roster routes it), and we
record the bunker key name + pair time. Stored as lowercase hex — the
roster + collision guard normalise either form, hex is canonical."""
await db.execute(
"""
UPDATE spirekeeper.dca_machines
SET machine_npub = :npub,
bunker_spire_key_name = :key_name,
paired_at = :paired_at,
updated_at = :updated_at
WHERE id = :id
""",
{
"npub": machine_npub.lower(),
"key_name": bunker_spire_key_name,
"paired_at": paired_at,
"updated_at": datetime.now(),
"id": machine_id,
},
)
return await get_machine(machine_id)
async def set_machine_unpaired(machine_id: str) -> Machine | None:
"""Mark a machine unpaired after revoking its spire's bunker access
(POST /revoke). Clears `paired_at`; keeps `machine_npub` +
`bunker_spire_key_name` for audit / re-pair. The bunker-side
`KeyUser.revokedAt` (set by `revoke_spire`) is what actually stops the
spire signing — this just records the operator-visible state."""
await db.execute(
"""
UPDATE spirekeeper.dca_machines
SET paired_at = NULL,
updated_at = :updated_at
WHERE id = :id
""",
{"updated_at": datetime.now(), "id": machine_id},
)
return await get_machine(machine_id)
async def set_machine_cash_out_hold(
machine_id: str,
since: datetime | None,
reason: str | None,
code: str | None,
) -> None:
"""Mirror the machine's cash-out hold (ADR-005 §5) onto its registry row.
Written on every state event, including when it is None: the machine
clearing the hold — after an operator recount or resume_cash_out — is as
important as it setting one. `updated_at` is left alone for the same
reason as counts_uncertain_since: this is the machine reporting about
itself on a heartbeat, not an operator editing the machine.
"""
await db.execute(
"""
UPDATE spirekeeper.dca_machines
SET cash_out_held_since = :since,
cash_out_held_reason = :reason,
cash_out_held_code = :code
WHERE id = :id
""",
{"id": machine_id, "since": since, "reason": reason, "code": code},
)
async def set_machine_counts_uncertain(machine_id: str, since: datetime | None) -> None:
"""Record (or clear) the machine's own "I can't vouch for these counts"
marker, straight from its state document.
`updated_at` is deliberately left alone. This is the machine reporting
about itself on a five-minute heartbeat, not an operator editing the
machine, and touching the timestamp on every heartbeat would make the
registry look perpetually just-modified.
"""
await db.execute(
"""
UPDATE spirekeeper.dca_machines
SET counts_uncertain_since = :since
WHERE id = :id
""",
{"since": since, "id": machine_id},
)
async def delete_machine(machine_id: str) -> None:
await db.execute(
"DELETE FROM spirekeeper.dca_machines WHERE id = :id",
{"id": machine_id},
)
# =============================================================================
# DCA Clients (LPs)
# =============================================================================
async def create_dca_client(data: CreateDcaClientData) -> DcaClient:
"""Operator enrols an LP at one of their machines.
Pure (machine, LP) record. Wallet / mode / autoforward live on
dca_lp (per-user) — populated by the LP via satmachineclient.
Enrolment doesn't require the LP to be onboarded yet, but deposits
do (see `create_deposit`).
"""
client_id = urlsafe_short_hash()
now = datetime.now()
await db.execute(
"""
INSERT INTO spirekeeper.dca_clients
(id, machine_id, user_id, username, status, created_at, updated_at)
VALUES (:id, :machine_id, :user_id, :username, :status,
:created_at, :updated_at)
""",
{
"id": client_id,
"machine_id": data.machine_id,
"user_id": data.user_id,
"username": data.username,
"status": "active",
"created_at": now,
"updated_at": now,
},
)
client = await get_dca_client(client_id)
assert client is not None
return client
# Shared SELECT fragment: client columns plus the LP-onboarded flag
# computed via LEFT JOIN on dca_lp. Returned as `lp_onboarded` (boolean
# 0/1 in SQLite, which Pydantic coerces to bool on the DcaClient model).
_CLIENT_SELECT = """
c.id, c.machine_id, c.user_id, c.username, c.status,
c.created_at, c.updated_at,
(lp.user_id IS NOT NULL) AS lp_onboarded
"""
_CLIENT_FROM = (
"spirekeeper.dca_clients c "
"LEFT JOIN spirekeeper.dca_lp lp ON lp.user_id = c.user_id"
)
async def get_dca_client(client_id: str) -> DcaClient | None:
return await db.fetchone(
f"SELECT {_CLIENT_SELECT} FROM {_CLIENT_FROM} WHERE c.id = :id",
{"id": client_id},
DcaClient,
)
async def get_dca_client_for_machine_user(
machine_id: str, user_id: str
) -> DcaClient | None:
return await db.fetchone(
f"""
SELECT {_CLIENT_SELECT} FROM {_CLIENT_FROM}
WHERE c.machine_id = :machine_id AND c.user_id = :user_id
""",
{"machine_id": machine_id, "user_id": user_id},
DcaClient,
)
async def get_dca_clients_for_machine(machine_id: str) -> list[DcaClient]:
return await db.fetchall(
f"""
SELECT {_CLIENT_SELECT} FROM {_CLIENT_FROM}
WHERE c.machine_id = :machine_id
ORDER BY c.created_at DESC
""",
{"machine_id": machine_id},
DcaClient,
)
async def get_dca_clients_for_operator(operator_user_id: str) -> list[DcaClient]:
"""All clients across every machine this operator owns."""
return await db.fetchall(
f"""
SELECT {_CLIENT_SELECT}
FROM {_CLIENT_FROM}
JOIN spirekeeper.dca_machines m ON m.id = c.machine_id
WHERE m.operator_user_id = :uid
ORDER BY c.created_at DESC
""",
{"uid": operator_user_id},
DcaClient,
)
async def get_dca_clients_for_user(user_id: str) -> list[DcaClient]:
"""LP cross-operator view — every machine this LP is registered at."""
return await db.fetchall(
f"""
SELECT {_CLIENT_SELECT} FROM {_CLIENT_FROM}
WHERE c.user_id = :user_id
ORDER BY c.created_at DESC
""",
{"user_id": user_id},
DcaClient,
)
async def get_flow_mode_clients_for_machine(machine_id: str) -> list[DcaClient]:
"""Active LPs enrolled at this machine whose per-user `dca_lp` row
has `default_dca_mode = 'flow'`. Used by the distribution algorithm.
An LP enrolment without a matching `dca_lp` row (i.e., the LP hasn't
onboarded via satmachineclient yet) is filtered out by the INNER
JOIN — there's no destination wallet to pay to.
"""
return await db.fetchall(
"""
SELECT c.*
FROM spirekeeper.dca_clients c
JOIN spirekeeper.dca_lp lp ON lp.user_id = c.user_id
WHERE c.machine_id = :machine_id
AND lp.default_dca_mode = 'flow'
AND c.status = 'active'
ORDER BY c.created_at ASC
""",
{"machine_id": machine_id},
DcaClient,
)
# =============================================================================
# DCA LP preferences (per-user) — wallet + mode + autoforward
# =============================================================================
async def get_dca_lp(user_id: str) -> DcaLpPreferences | None:
"""Return the LP's preferences row, or None if they haven't onboarded
via satmachineclient yet."""
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_lp WHERE user_id = :uid",
{"uid": user_id},
DcaLpPreferences,
)
async def lp_is_onboarded(user_id: str) -> bool:
"""Cheap existence check used by the deposit-creation gate."""
row = await db.fetchone(
"SELECT user_id FROM spirekeeper.dca_lp WHERE user_id = :uid",
{"uid": user_id},
)
return row is not None
async def upsert_dca_lp(
user_id: str,
data: UpsertDcaLpData,
*,
fallback_wallet_id: str | None = None,
) -> DcaLpPreferences:
"""Create or update the LP's preferences row.
First call (no row yet): `data.dca_wallet_id` must be set OR
`fallback_wallet_id` must be provided (satmachineclient passes the
LP's default LNbits wallet here when auto-seeding on first dashboard
visit). Subsequent calls update only the fields in `data` that are
non-None.
"""
existing = await get_dca_lp(user_id)
now = datetime.now()
if existing is None:
wallet_id = data.dca_wallet_id or fallback_wallet_id
if not wallet_id:
raise ValueError(
"first upsert requires dca_wallet_id (or fallback_wallet_id)"
)
await db.execute(
"""
INSERT INTO spirekeeper.dca_lp
(user_id, dca_wallet_id, default_dca_mode, fixed_mode_daily_limit,
autoforward_ln_address, autoforward_enabled,
created_at, updated_at)
VALUES (:uid, :wallet, :mode, :limit, :ln_addr, :auto,
:now, :now)
""",
{
"uid": user_id,
"wallet": wallet_id,
"mode": data.default_dca_mode or "flow",
"limit": data.fixed_mode_daily_limit,
"ln_addr": data.autoforward_ln_address,
"auto": data.autoforward_enabled or False,
"now": now,
},
)
else:
update_data: dict = {k: v for k, v in data.dict().items() if v is not None}
if not update_data:
return existing
update_data["updated_at"] = now
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
update_data["uid"] = user_id
await db.execute(
f"UPDATE spirekeeper.dca_lp SET {set_clause} WHERE user_id = :uid",
update_data,
)
refreshed = await get_dca_lp(user_id)
assert refreshed is not None
return refreshed
async def update_dca_client(
client_id: str, data: UpdateDcaClientData
) -> DcaClient | None:
update_data = {k: v for k, v in data.dict().items() if v is not None}
if not update_data:
return await get_dca_client(client_id)
update_data["updated_at"] = datetime.now()
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
update_data["id"] = client_id
await db.execute(
f"UPDATE spirekeeper.dca_clients SET {set_clause} WHERE id = :id",
update_data,
)
return await get_dca_client(client_id)
async def delete_dca_client(client_id: str) -> None:
await db.execute(
"DELETE FROM spirekeeper.dca_clients WHERE id = :id",
{"id": client_id},
)
# =============================================================================
# Deposits
# =============================================================================
async def create_deposit(
creator_user_id: str, data: CreateDepositData, *, currency: str
) -> DcaDeposit:
"""Insert a deposit row.
`currency` is passed explicitly by the caller (the API endpoint
resolves it from the target machine's `fiat_code`) rather than
coming off the request body — the operator doesn't get to choose
it (`aiolabs/satmachineadmin#26`).
"""
deposit_id = urlsafe_short_hash()
await db.execute(
"""
INSERT INTO spirekeeper.dca_deposits
(id, client_id, machine_id, creator_user_id, amount, currency,
status, notes, created_at)
VALUES (:id, :client_id, :machine_id, :creator_user_id, :amount,
:currency, :status, :notes, :created_at)
""",
{
"id": deposit_id,
"client_id": data.client_id,
"machine_id": data.machine_id,
"creator_user_id": creator_user_id,
"amount": data.amount,
"currency": currency,
"status": "pending",
"notes": data.notes,
"created_at": datetime.now(),
},
)
deposit = await get_deposit(deposit_id)
assert deposit is not None
return deposit
async def get_deposit(deposit_id: str) -> DcaDeposit | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_deposits WHERE id = :id",
{"id": deposit_id},
DcaDeposit,
)
async def get_deposits_for_client(client_id: str) -> list[DcaDeposit]:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_deposits
WHERE client_id = :client_id
ORDER BY created_at DESC
""",
{"client_id": client_id},
DcaDeposit,
)
async def get_deposits_for_operator(operator_user_id: str) -> list[DcaDeposit]:
return await db.fetchall(
"""
SELECT d.*
FROM spirekeeper.dca_deposits d
JOIN spirekeeper.dca_machines m ON m.id = d.machine_id
WHERE m.operator_user_id = :uid
ORDER BY d.created_at DESC
""",
{"uid": operator_user_id},
DcaDeposit,
)
async def update_deposit(
deposit_id: str, data: UpdateDepositData
) -> DcaDeposit | None:
update_data = {k: v for k, v in data.dict().items() if v is not None}
if not update_data:
return await get_deposit(deposit_id)
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
update_data["id"] = deposit_id
await db.execute(
f"UPDATE spirekeeper.dca_deposits SET {set_clause} WHERE id = :id",
update_data,
)
return await get_deposit(deposit_id)
async def update_deposit_status(
deposit_id: str, data: UpdateDepositStatusData
) -> DcaDeposit | None:
payload = {
"id": deposit_id,
"status": data.status,
"notes": data.notes,
"confirmed_at": datetime.now() if data.status == "confirmed" else None,
}
await db.execute(
"""
UPDATE spirekeeper.dca_deposits
SET status = :status,
notes = COALESCE(:notes, notes),
confirmed_at = COALESCE(:confirmed_at, confirmed_at)
WHERE id = :id
""",
payload,
)
return await get_deposit(deposit_id)
async def delete_deposit(deposit_id: str) -> None:
await db.execute(
"DELETE FROM spirekeeper.dca_deposits WHERE id = :id",
{"id": deposit_id},
)
# =============================================================================
# Settlements (bitSpire kind-21000 events)
# =============================================================================
async def create_settlement_idempotent(
data: CreateDcaSettlementData,
initial_status: str,
error_message: str | None = None,
) -> DcaSettlement | None:
"""Insert a settlement keyed by payment_hash.
Returns the inserted row on first sight; returns the existing row
if the payment_hash was already seen (subscription replay,
dispatcher double-fire). The UNIQUE constraint on payment_hash is
the source of truth.
`initial_status` is the row's status at insert time. Normal
settlements arrive as 'pending' and the distribution processor
transitions them through 'processing' → 'processed' / 'errored'.
A row that fails the Nostr attribution cross-check (bitspire.
assert_nostr_attribution) is inserted directly as 'rejected' with
the failure reason in `error_message` — never goes near the
distribution path.
"""
existing = await get_settlement_by_payment_hash(data.payment_hash)
if existing is not None:
return existing
settlement_id = urlsafe_short_hash()
await db.execute(
"""
INSERT INTO spirekeeper.dca_settlements
(id, machine_id, payment_hash, bitspire_event_id, bitspire_txid,
wire_sats, fiat_amount, fiat_code, exchange_rate, principal_sats,
fee_sats, platform_fee_sats, operator_fee_sats, fee_mismatch_sats,
tx_type, bills_json, cassettes_json,
status, error_message, created_at)
VALUES (:id, :machine_id, :payment_hash, :bitspire_event_id,
:bitspire_txid, :wire_sats, :fiat_amount, :fiat_code,
:exchange_rate, :principal_sats, :fee_sats,
:platform_fee_sats, :operator_fee_sats, :fee_mismatch_sats,
:tx_type, :bills_json, :cassettes_json, :status,
:error_message, :created_at)
""",
{
"id": settlement_id,
"machine_id": data.machine_id,
"payment_hash": data.payment_hash,
"bitspire_event_id": data.bitspire_event_id,
"bitspire_txid": data.bitspire_txid,
"wire_sats": data.wire_sats,
"fiat_amount": data.fiat_amount,
"fiat_code": data.fiat_code,
"exchange_rate": data.exchange_rate,
"principal_sats": data.principal_sats,
"fee_sats": data.fee_sats,
"platform_fee_sats": data.platform_fee_sats,
"operator_fee_sats": data.operator_fee_sats,
"fee_mismatch_sats": data.fee_mismatch_sats,
"tx_type": data.tx_type,
"bills_json": data.bills_json,
"cassettes_json": data.cassettes_json,
"status": initial_status,
"error_message": error_message,
"created_at": datetime.now(),
},
)
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",
{"id": settlement_id},
DcaSettlement,
)
async def get_settlement_by_payment_hash(
payment_hash: str,
) -> DcaSettlement | None:
return await db.fetchone(
"""
SELECT * FROM spirekeeper.dca_settlements
WHERE payment_hash = :hash
""",
{"hash": payment_hash},
DcaSettlement,
)
async def get_settlements_for_machine(
machine_id: str, limit: int = 100
) -> list[DcaSettlement]:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_settlements
WHERE machine_id = :machine_id
ORDER BY created_at DESC
LIMIT :lim
""",
{"machine_id": machine_id, "lim": limit},
DcaSettlement,
)
async def get_stuck_settlements_for_operator(
operator_user_id: str, threshold_minutes: int = 30
) -> 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:
- '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,
not operationally failed. Operator must investigate the machine.
- 'errored': any status='errored' (distribution failed for an
operational reason — wallet error, network, downstream payment).
Operator retries from this bucket.
- 'stuck_pending': status='pending' AND older than threshold
(listener crashed before invoking process_settlement).
- 'stuck_processing': status='processing' AND older than threshold
(processor crashed mid-flight; processing_claim is set but no
completion landed).
"""
from datetime import timedelta
threshold_at = datetime.now() - timedelta(minutes=threshold_minutes)
rejected = 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 = 'rejected'
ORDER BY s.created_at DESC
""",
{"uid": operator_user_id},
DcaSettlement,
)
errored = 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 = 'errored'
ORDER BY s.created_at DESC
""",
{"uid": operator_user_id},
DcaSettlement,
)
stuck_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 = 'pending'
AND s.created_at < :threshold
ORDER BY s.created_at ASC
""",
{"uid": operator_user_id, "threshold": threshold_at},
DcaSettlement,
)
stuck_processing = 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 = 'processing'
AND s.created_at < :threshold
ORDER BY s.created_at ASC
""",
{"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,
"stuck_processing": stuck_processing,
}
async def force_reset_stuck_settlement(
settlement_id: str,
) -> DcaSettlement | None:
"""Operator escape hatch for genuinely stuck settlements (processor
crashed mid-flight, etc.). Flips 'pending'/'processing' → 'errored' so
the existing retry endpoint can take over. Clears processing_claim.
Caller is responsible for verifying the settlement is *actually* stuck
(e.g., via threshold check on created_at). This function trusts the
decision."""
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET status = 'errored',
processing_claim = NULL,
error_message = 'force-reset by operator (was stuck)'
WHERE id = :id AND status IN ('pending', 'processing')
""",
{"id": settlement_id},
)
return await get_settlement(settlement_id)
async def get_settlements_for_operator(
operator_user_id: str, limit: int = 200
) -> list[DcaSettlement]:
return 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
ORDER BY s.created_at DESC
LIMIT :lim
""",
{"uid": operator_user_id, "lim": limit},
DcaSettlement,
)
async def mark_settlement_status(
settlement_id: str,
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."""
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET status = :status,
error_message = :err,
processed_at = CASE
WHEN :status IN ('processed', 'partial', 'refunded')
THEN :now ELSE processed_at
END,
processing_claim = CASE
WHEN :status = 'processing' THEN processing_claim
ELSE NULL
END
WHERE id = :id
""",
{
"id": settlement_id,
"status": status,
"err": error_message,
"now": datetime.now(),
},
)
return await get_settlement(settlement_id)
async def claim_settlement_for_processing(
settlement_id: str,
) -> DcaSettlement | None:
"""Optimistic-lock claim: atomically flip a settlement to 'processing'
and tag it with a per-invocation token. Returns the claimed row on
success; None if another caller already won the claim or the settlement
is not in a claimable state ('pending').
Pattern is portable across SQLite + PostgreSQL (doesn't rely on
UPDATE ... RETURNING). Two concurrent invocations may both run the
UPDATE, but only one row matches the WHERE clause; the loser's UPDATE
is a no-op against status='processing'. The read-back check on the
token disambiguates."""
token = urlsafe_short_hash()
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET status = 'processing', processing_claim = :token
WHERE id = :id AND status = 'pending'
""",
{"id": settlement_id, "token": token},
)
after = await get_settlement(settlement_id)
if after is None:
return None
if after.processing_claim != token:
return None
return after
async def reset_settlement_for_retry(
settlement_id: str,
) -> DcaSettlement | None:
"""Operator retry path. Flips 'errored' → 'pending' and voids any
'failed' legs so process_settlement re-runs them fresh. Completed legs
are left in place — we never re-pay sats that already moved."""
await db.execute(
"""
UPDATE spirekeeper.dca_payments
SET status = 'voided'
WHERE settlement_id = :sid AND status = 'failed'
""",
{"sid": settlement_id},
)
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET status = 'pending',
error_message = NULL,
processing_claim = NULL,
processed_at = NULL
WHERE id = :id AND status = 'errored'
""",
{"id": settlement_id},
)
return await get_settlement(settlement_id)
async def apply_partial_dispense(
settlement_id: str,
*,
new_wire_sats: int,
new_principal_sats: int,
new_fee_sats: int,
new_platform_fee_sats: int,
new_operator_fee_sats: int,
new_fiat_amount: float,
appended_note: str,
) -> DcaSettlement | None:
"""Overwrite the monetary fields on a settlement (partial-dispense
recompute) and prepend `appended_note` to the notes column.
Notes are append-only: new lines go at the top (newest first) so the
settlement detail view shows the most recent adjustment first without
needing to scroll. Resets status to 'pending' so process_settlement
can re-distribute via the existing idempotent path."""
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET wire_sats = :gross,
principal_sats = :principal,
fee_sats = :commission,
platform_fee_sats = :platform,
operator_fee_sats = :operator,
fiat_amount = :fiat,
status = 'pending',
error_message = NULL,
processed_at = NULL,
notes = CASE
WHEN notes IS NULL OR notes = '' THEN :note
ELSE :note || char(10) || char(10) || notes
END
WHERE id = :id
""",
{
"id": settlement_id,
"gross": new_wire_sats,
"principal": new_principal_sats,
"commission": new_fee_sats,
"platform": new_platform_fee_sats,
"operator": new_operator_fee_sats,
"fiat": new_fiat_amount,
"note": appended_note,
},
)
return await get_settlement(settlement_id)
async def count_completed_legs_for_settlement(settlement_id: str) -> int:
"""Used by partial-dispense to refuse adjustments after any leg has
successfully moved sats (Lightning payments can't be clawed back)."""
row = await db.fetchone(
"""
SELECT COUNT(*) AS n FROM spirekeeper.dca_payments
WHERE settlement_id = :sid AND status = 'completed'
""",
{"sid": settlement_id},
)
return int(row["n"]) if row else 0
async def append_settlement_note(
settlement_id: str, note: str, author_user_id: str
) -> DcaSettlement | None:
"""Prepend an operator-authored note to settlement.notes. Each entry is
timestamped (UTC) and tagged with the author's user id so the trail
is accountable. Append-only: existing entries are never edited."""
from datetime import timezone
ts = datetime.now(timezone.utc).isoformat(timespec="seconds")
formatted = f"[{ts} by {author_user_id}] {note}"
await db.execute(
"""
UPDATE spirekeeper.dca_settlements
SET notes = CASE
WHEN notes IS NULL OR notes = '' THEN :note
ELSE :note || char(10) || char(10) || notes
END
WHERE id = :id
""",
{"id": settlement_id, "note": formatted},
)
return await get_settlement(settlement_id)
async def void_open_legs_for_settlement(settlement_id: str) -> None:
"""Marks open legs as 'voided' before re-running distribution on a
partial-dispense recompute. Preserves the rows for audit but stops
them from being interpreted as live. Includes 'skipped' so that audit
rows from a prior attempt don't double-count once the new attempt
writes its own (possibly different) skipped reasons."""
await db.execute(
"""
UPDATE spirekeeper.dca_payments
SET status = 'voided'
WHERE settlement_id = :sid
AND status IN ('pending', 'failed', 'skipped')
""",
{"sid": settlement_id},
)
# =============================================================================
# Commission splits — operator's remainder-distribution rules.
# =============================================================================
async def get_commission_splits(
operator_user_id: str, machine_id: str | None = None
) -> list[CommissionSplit]:
"""Returns the rule set for the given scope.
Precedence (caller's responsibility): try per-machine override first;
if empty, fall back to operator default (machine_id IS NULL).
"""
if machine_id is None:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_commission_splits
WHERE operator_user_id = :uid AND machine_id IS NULL
ORDER BY sort_order ASC
""",
{"uid": operator_user_id},
CommissionSplit,
)
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_commission_splits
WHERE operator_user_id = :uid AND machine_id = :mid
ORDER BY sort_order ASC
""",
{"uid": operator_user_id, "mid": machine_id},
CommissionSplit,
)
async def get_effective_commission_splits(
operator_user_id: str, machine_id: str
) -> list[CommissionSplit]:
"""Per-machine override if set, otherwise operator's default ruleset."""
overrides = await get_commission_splits(operator_user_id, machine_id)
if overrides:
return overrides
return await get_commission_splits(operator_user_id, None)
async def replace_commission_splits(
operator_user_id: str,
machine_id: str | None,
legs: list[CommissionSplitLeg],
) -> list[CommissionSplit]:
"""Atomic replace for the (operator, machine) scope. Caller should have
already validated legs sum to 1.0 via the Pydantic model."""
if machine_id is None:
await db.execute(
"""
DELETE FROM spirekeeper.dca_commission_splits
WHERE operator_user_id = :uid AND machine_id IS NULL
""",
{"uid": operator_user_id},
)
else:
await db.execute(
"""
DELETE FROM spirekeeper.dca_commission_splits
WHERE operator_user_id = :uid AND machine_id = :mid
""",
{"uid": operator_user_id, "mid": machine_id},
)
now = datetime.now()
for leg in legs:
await db.execute(
"""
INSERT INTO spirekeeper.dca_commission_splits
(id, machine_id, operator_user_id, target, label, fraction,
sort_order, created_at)
VALUES (:id, :machine_id, :uid, :target, :label, :fraction,
:sort_order, :created_at)
""",
{
"id": urlsafe_short_hash(),
"machine_id": machine_id,
"uid": operator_user_id,
"target": leg.target,
"label": leg.label,
"fraction": leg.fraction,
"sort_order": leg.sort_order,
"created_at": now,
},
)
return await get_commission_splits(operator_user_id, machine_id)
# =============================================================================
# Payments — distribution legs.
# =============================================================================
async def create_dca_payment(data: CreateDcaPaymentData) -> DcaPayment:
payment_id = urlsafe_short_hash()
await db.execute(
"""
INSERT INTO spirekeeper.dca_payments
(id, settlement_id, client_id, machine_id, operator_user_id,
leg_type, destination_wallet_id, destination_ln_address,
amount_sats, amount_fiat, exchange_rate, transaction_time,
external_payment_hash, status, created_at)
VALUES (:id, :settlement_id, :client_id, :machine_id,
:operator_user_id, :leg_type, :destination_wallet_id,
:destination_ln_address, :amount_sats, :amount_fiat,
:exchange_rate, :transaction_time, :external_payment_hash,
:status, :created_at)
""",
{
"id": payment_id,
"settlement_id": data.settlement_id,
"client_id": data.client_id,
"machine_id": data.machine_id,
"operator_user_id": data.operator_user_id,
"leg_type": data.leg_type,
"destination_wallet_id": data.destination_wallet_id,
"destination_ln_address": data.destination_ln_address,
"amount_sats": data.amount_sats,
"amount_fiat": data.amount_fiat,
"exchange_rate": data.exchange_rate,
"transaction_time": data.transaction_time,
"external_payment_hash": data.external_payment_hash,
"status": "pending",
"created_at": datetime.now(),
},
)
payment = await get_dca_payment(payment_id)
assert payment is not None
return payment
async def get_dca_payment(payment_id: str) -> DcaPayment | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_payments WHERE id = :id",
{"id": payment_id},
DcaPayment,
)
async def get_payments_for_settlement(settlement_id: str) -> list[DcaPayment]:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_payments
WHERE settlement_id = :sid
ORDER BY created_at ASC
""",
{"sid": settlement_id},
DcaPayment,
)
async def get_payments_for_client(client_id: str) -> list[DcaPayment]:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_payments
WHERE client_id = :cid
ORDER BY created_at DESC
""",
{"cid": client_id},
DcaPayment,
)
async def get_payments_for_operator(
operator_user_id: str, leg_type: str | None = None, limit: int = 200
) -> list[DcaPayment]:
if leg_type is None:
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_payments
WHERE operator_user_id = :uid
ORDER BY created_at DESC
LIMIT :lim
""",
{"uid": operator_user_id, "lim": limit},
DcaPayment,
)
return await db.fetchall(
"""
SELECT * FROM spirekeeper.dca_payments
WHERE operator_user_id = :uid AND leg_type = :leg
ORDER BY created_at DESC
LIMIT :lim
""",
{"uid": operator_user_id, "leg": leg_type, "lim": limit},
DcaPayment,
)
async def update_payment_status(
payment_id: str,
status: str,
external_payment_hash: str | None = None,
error_message: str | None = None,
) -> DcaPayment | None:
await db.execute(
"""
UPDATE spirekeeper.dca_payments
SET status = :status,
external_payment_hash = COALESCE(:hash, external_payment_hash),
error_message = :err
WHERE id = :id
""",
{
"id": payment_id,
"status": status,
"hash": external_payment_hash,
"err": error_message,
},
)
return await get_dca_payment(payment_id)
# =============================================================================
# Balance summaries
# =============================================================================
async def get_client_balance_summary(
client_id: str,
) -> ClientBalanceSummary | None:
"""Per-client (and per-machine, since clients are per-machine in v2) summary.
DCA legs only — settlement/autoforward/super_fee/operator_split legs are
not credited against an LP's balance.
"""
client = await get_dca_client(client_id)
if client is None:
return None
deposits_row = await db.fetchone(
"""
SELECT COALESCE(SUM(amount), 0) AS total
FROM spirekeeper.dca_deposits
WHERE client_id = :cid AND status = 'confirmed'
""",
{"cid": client_id},
)
# Both DCA legs (auto, from bitSpire settlements) and balance-settle legs
# (operator-initiated under #4) reduce the LP's remaining fiat balance.
payments_row = await db.fetchone(
"""
SELECT COALESCE(SUM(amount_fiat), 0) AS total
FROM spirekeeper.dca_payments
WHERE client_id = :cid
AND leg_type IN ('dca', 'settlement')
AND status = 'completed'
""",
{"cid": client_id},
)
total_deposits = float(deposits_row["total"]) if deposits_row else 0.0
total_payments = float(payments_row["total"]) if payments_row else 0.0
# fiat code: take it from the machine (clients inherit their machine's fiat)
machine = await get_machine(client.machine_id)
currency = machine.fiat_code if machine else "GTQ"
return ClientBalanceSummary(
client_id=client_id,
machine_id=client.machine_id,
total_deposits=round(total_deposits, 2),
total_payments=round(total_payments, 2),
remaining_balance=round(total_deposits - total_payments, 2),
currency=currency,
)
# =============================================================================
# Telemetry — sparse beacon (kind-30078) and fleet snapshot (kind-30079) state.
# =============================================================================
async def get_telemetry(machine_id: str) -> TelemetrySnapshot | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.dca_telemetry WHERE machine_id = :mid",
{"mid": machine_id},
TelemetrySnapshot,
)
async def upsert_beacon_snapshot(
machine_id: str,
*,
cash_in: bool | None = None,
cash_out: bool | None = None,
cash_level: str | None = None,
fiat: str | None = None,
model: str | None = None,
name: str | None = None,
location: str | None = None,
geo: str | None = None,
fees_json: str | None = None,
limits_json: str | None = None,
denominations_json: str | None = None,
version: str | None = None,
) -> TelemetrySnapshot | None:
"""Upsert kind-30078 beacon fields. All fields are nullable because today's
upstream payload only carries cash_in/cash_out/cash_level/fiat/model (see
lamassu-next#43 — the enrichment is not yet shipped)."""
existing = await get_telemetry(machine_id)
now = datetime.now()
if existing is None:
await db.execute(
"""
INSERT INTO spirekeeper.dca_telemetry
(machine_id, beacon_cash_in, beacon_cash_out, beacon_cash_level,
beacon_fiat, beacon_model, beacon_name, beacon_location,
beacon_geo, beacon_fees_json, beacon_limits_json,
beacon_denominations_json, beacon_version, beacon_received_at)
VALUES (:mid, :cash_in, :cash_out, :cash_level, :fiat, :model,
:name, :location, :geo, :fees, :limits, :denoms,
:version, :now)
""",
{
"mid": machine_id,
"cash_in": cash_in,
"cash_out": cash_out,
"cash_level": cash_level,
"fiat": fiat,
"model": model,
"name": name,
"location": location,
"geo": geo,
"fees": fees_json,
"limits": limits_json,
"denoms": denominations_json,
"version": version,
"now": now,
},
)
else:
await db.execute(
"""
UPDATE spirekeeper.dca_telemetry SET
beacon_cash_in = COALESCE(:cash_in, beacon_cash_in),
beacon_cash_out = COALESCE(:cash_out, beacon_cash_out),
beacon_cash_level = COALESCE(:cash_level, beacon_cash_level),
beacon_fiat = COALESCE(:fiat, beacon_fiat),
beacon_model = COALESCE(:model, beacon_model),
beacon_name = COALESCE(:name, beacon_name),
beacon_location = COALESCE(:location, beacon_location),
beacon_geo = COALESCE(:geo, beacon_geo),
beacon_fees_json = COALESCE(:fees, beacon_fees_json),
beacon_limits_json = COALESCE(:limits, beacon_limits_json),
beacon_denominations_json =
COALESCE(:denoms, beacon_denominations_json),
beacon_version = COALESCE(:version, beacon_version),
beacon_received_at = :now
WHERE machine_id = :mid
""",
{
"mid": machine_id,
"cash_in": cash_in,
"cash_out": cash_out,
"cash_level": cash_level,
"fiat": fiat,
"model": model,
"name": name,
"location": location,
"geo": geo,
"fees": fees_json,
"limits": limits_json,
"denoms": denominations_json,
"version": version,
"now": now,
},
)
return await get_telemetry(machine_id)
async def upsert_fleet_snapshot(
machine_id: str, telemetry_json: str
) -> TelemetrySnapshot | None:
"""Upsert kind-30079 operator-only telemetry. Awaits lamassu-next#42 to
produce a real schema; we store the raw JSON blob until then."""
existing = await get_telemetry(machine_id)
now = datetime.now()
if existing is None:
await db.execute(
"""
INSERT INTO spirekeeper.dca_telemetry
(machine_id, telemetry_json, telemetry_received_at)
VALUES (:mid, :json, :now)
""",
{"mid": machine_id, "json": telemetry_json, "now": now},
)
else:
await db.execute(
"""
UPDATE spirekeeper.dca_telemetry
SET telemetry_json = :json, telemetry_received_at = :now
WHERE machine_id = :mid
""",
{"mid": machine_id, "json": telemetry_json, "now": now},
)
return await get_telemetry(machine_id)
# =============================================================================
# Cassette configs — operator-driven ATM cassette inventory (#29 v1.1).
# =============================================================================
# Row lifecycle per #29:
# - First population for a (machine_id, position) pair → apply_reported_state
# (consumer reading the ATM's one-shot bitspire-cassettes-state event)
# - The operator does NOT write these rows. It records operations
# (cassette_ops) and the machine keeps the running count; these columns
# hold what the machine last reported.
# - Rows appear and disappear only as the machine reports its bay set —
# the slot count is hardware-determined.
def _as_unix(value) -> float | None:
"""Normalise whatever the driver hands back for state_at to a unix float.
SQLite stores these as integers and Postgres as timestamps, and an event's
created_at arrives tz-aware, so comparing the raw values risks either a
TypeError (aware vs naive) or a silently wrong answer. Everything is
compared as seconds since the epoch instead.
"""
if value is None:
return None
if isinstance(value, datetime):
if value.tzinfo is None:
value = value.replace(tzinfo=timezone.utc)
return value.timestamp()
try:
return float(value)
except (TypeError, ValueError):
return None
def _row_field(row, name):
"""Read one column from a row the data layer may hand back as either a
mapping or an object, depending on the driver."""
if isinstance(row, dict):
return row.get(name)
return getattr(row, name, None)
def _should_apply_state_event(
oldest_state_at, incoming_created_at, oldest_seq=None, incoming_seq=None
) -> bool:
"""Ordering gate for apply_reported_state.
Applies only when the incoming event is strictly newer than the OLDEST
state stamp on file for the machine.
Three deliberate choices:
- Compare created_at, not event ids. The old gate asked only whether the
incoming id differed from one stored row, which is a one-event memory:
a re-delivered A, B, A applied three times, and an event that arrived
late overwrote newer state because nothing ever looked at the clock.
Strict `>` also subsumes exact-replay dedup, since a replay carries the
same stamp.
- Oldest, not newest. Every execute in this data layer commits on its own,
so a multi-row apply cannot be made atomic here; a crash mid-loop leaves
some rows advanced and some not. Gating on the oldest means a partial
apply is re-applied on the next event rather than being mistaken for a
complete one, and the ATM republishes on a heartbeat, so it converges.
- `seq` breaks a same-second tie, and only that. NIP-01 stamps at
one-second granularity, so a dispense and the publish that follows it
can share a stamp and the later report would be dropped. The machine
bumps seq on every local count change, so a higher seq at an equal stamp
is strictly newer. It is consulted ONLY on equality: a machine whose
state.db was replaced restarts its counter at zero, and gating on seq
across different stamps would lock that machine out permanently while
its wall clock kept moving forward.
"""
oldest = _as_unix(oldest_state_at)
if oldest is None:
return True
incoming = _as_unix(incoming_created_at)
if incoming is None:
return False
if incoming != oldest:
return incoming > oldest
if incoming_seq is None or oldest_seq is None:
return False
return incoming_seq > oldest_seq
async def get_cassette_config(machine_id: str, position: int) -> CassetteConfig | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid AND position = :pos",
{"mid": machine_id, "pos": position},
CassetteConfig,
)
async def list_cassette_configs_for_machine(
machine_id: str,
) -> list[CassetteConfig]:
return await db.fetchall(
"SELECT * FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid ORDER BY position",
{"mid": machine_id},
CassetteConfig,
)
async def apply_reported_state(
machine_id: str,
event_id: str,
event_created_at: datetime,
payload: PublishCassettesPayload,
) -> bool:
"""Consume an ATM-published kind-30078 bitspire-cassettes-state:<m> event
and reconcile cassette_configs for the machine against it.
Returns True if the state was applied, False if the event was not newer
than what is already on file (see _should_apply_state_event).
The payload is the machine's full bay set, and the machine owns that
layout — the bay count is hardware-determined. So positions absent from
the payload are DELETED here. Leaving them behind was half the problem a
shrinking bay count caused: the stale row stayed, the dashboard showed a
mix of old and new, and every later publish was rejected for a position
mismatch with no way out but hand-written SQL.
Populates both the operator-believed columns (denomination, count,
updated_at, updated_by) and the reported columns (state_denomination,
state_count, state_at, state_event_id), so the UI can show reported
against believed.
"""
oldest: dict | None = await db.fetchone(
"SELECT state_at, state_seq FROM spirekeeper.cassette_configs "
"WHERE machine_id = :mid AND state_at IS NOT NULL "
"ORDER BY state_at ASC, state_seq ASC LIMIT 1",
{"mid": machine_id},
)
oldest_state_at = None
oldest_seq = None
if oldest is not None:
oldest_state_at = _row_field(oldest, "state_at")
oldest_seq = _row_field(oldest, "state_seq")
if not _should_apply_state_event(
oldest_state_at, event_created_at, oldest_seq, payload.seq
):
return False
# Drop bays the machine no longer reports, before writing the rest. A crash
# between the two leaves rows missing rather than stale, and the next
# heartbeat re-inserts them — the safe direction.
keep = sorted(payload.positions.keys())
placeholders = ", ".join(f":p{i}" for i in range(len(keep)))
delete_values: dict = {"mid": machine_id}
for i, pos in enumerate(keep):
delete_values[f"p{i}"] = pos
await db.execute(
"DELETE FROM spirekeeper.cassette_configs "
f"WHERE machine_id = :mid AND position NOT IN ({placeholders})",
delete_values,
)
now = datetime.now()
for pos, row in payload.positions.items():
await db.execute(
"""
INSERT INTO spirekeeper.cassette_configs
(machine_id, position, denomination, count, updated_at,
updated_by, state_denomination, state_count, state_at,
state_event_id, state_seq)
VALUES (:mid, :pos, :denom, :count, :now, :by,
:state_denom, :state_count, :state_at, :event_id,
:state_seq)
ON CONFLICT (machine_id, position) DO UPDATE SET
denomination = excluded.denomination,
count = excluded.count,
updated_at = excluded.updated_at,
updated_by = excluded.updated_by,
state_denomination = excluded.state_denomination,
state_count = excluded.state_count,
state_at = excluded.state_at,
state_event_id = excluded.state_event_id,
state_seq = excluded.state_seq
""",
{
"mid": machine_id,
"pos": pos,
"denom": row.denomination,
"count": row.count,
"now": now,
"by": "atm-report",
"state_denom": row.denomination,
"state_count": row.count,
"state_at": event_created_at,
"event_id": event_id,
"state_seq": payload.seq,
},
)
return True
# ---------------------------------------------------------------------------
# Cassette operations — the v2 operator → ATM wire (bitspire ADR-004)
# ---------------------------------------------------------------------------
# Append-only. The operator records what it DID to a bay; the machine keeps the
# running count. Nothing here ever updates cassette_configs: that table now
# holds only what the machine has reported, and letting an operation write it
# would reintroduce the second writer this design exists to remove.
# How many operations ride along in each published event. The window is what
# makes the channel self-healing — a machine that missed one event still sees
# the operation in the next — so it needs to cover a plausible outage, not just
# the newest change. Twenty is roughly a fortnight of refills on a busy fleet
# and still a small payload.
CASSETTE_OPS_WINDOW = 20
async def create_cassette_op(
machine_id: str, data: CreateCassetteOpData, created_by: str | None
) -> CassetteOp:
"""Record one operator intent against one bay.
The id minted here is the idempotency key: it travels on the wire, the
machine records the ones it has applied, and a re-delivered event is
therefore free rather than double-counted.
"""
op_id = urlsafe_short_hash()
await db.execute(
"""
INSERT INTO spirekeeper.cassette_ops
(id, machine_id, position, op_type, bills, count, denomination,
created_at, created_by)
VALUES (:id, :machine_id, :position, :op_type, :bills, :count,
:denomination, :created_at, :created_by)
""",
{
"id": op_id,
"machine_id": machine_id,
"position": data.position,
"op_type": data.op_type,
"bills": data.bills,
"count": data.count,
"denomination": data.denomination,
"created_at": datetime.now(),
"created_by": created_by,
},
)
op = await get_cassette_op(op_id)
assert op is not None, "Newly recorded cassette op couldn't be retrieved"
return op
async def get_cassette_op(op_id: str) -> CassetteOp | None:
return await db.fetchone(
"SELECT * FROM spirekeeper.cassette_ops WHERE id = :id",
{"id": op_id},
CassetteOp,
)
async def get_cassette_ops_window(
machine_id: str, limit: int = CASSETTE_OPS_WINDOW
) -> list[CassetteOp]:
"""The operations a publish carries, oldest first.
Selected newest-first to take the most recent `limit`, then reversed so the
machine applies them in the order they happened. That ordering matters:
a recount followed by a refill is not the same as the reverse.
"""
rows = await db.fetchall(
"SELECT * FROM spirekeeper.cassette_ops WHERE machine_id = :mid "
"ORDER BY created_at DESC, id DESC LIMIT :limit",
{"mid": machine_id, "limit": limit},
CassetteOp,
)
return list(reversed(rows))
async def list_cassette_ops(
machine_id: str, limit: int = CASSETTE_OPS_WINDOW
) -> list[CassetteOp]:
"""Newest-first, for the dashboard's history panel."""
return await db.fetchall(
"SELECT * FROM spirekeeper.cassette_ops WHERE machine_id = :mid "
"ORDER BY created_at DESC, id DESC LIMIT :limit",
{"mid": machine_id, "limit": limit},
CassetteOp,
)
def _should_ack_op(existing: CassetteOp | None, machine_id: str) -> bool:
"""Pure decision behind mark_cassette_ops_acked, extracted so it is
testable without a database — same approach as the state-event gate.
Three reasons not to ack, and each matters:
- the id is unknown, so there is nothing to close out;
- it belongs to another machine, and one machine reporting an id must
never be able to close out another machine's operation;
- it is already acked, and the first acknowledgement is the interesting
one. The machine echoes a WINDOW, so every id arrives many times over;
overwriting would keep sliding the timestamp forward and lose when the
operation actually landed.
"""
if existing is None:
return False
if existing.machine_id != machine_id:
return False
return existing.acked_at is None
async def mark_cassette_ops_acked(machine_id: str, op_ids: list[str]) -> int:
"""Record that the machine reported these operation ids as applied.
Returns how many rows moved from unacked to acked. See _should_ack_op for
which reports are ignored and why.
"""
if not op_ids:
return 0
acked = 0
now = datetime.now()
for op_id in op_ids:
existing = await get_cassette_op(op_id)
if not _should_ack_op(existing, machine_id):
continue
await db.execute(
"UPDATE spirekeeper.cassette_ops SET acked_at = :now "
"WHERE id = :id AND machine_id = :mid",
{"now": now, "id": op_id, "mid": machine_id},
)
acked += 1
return acked