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.
2029 lines
71 KiB
Python
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
|