Some checks failed
ci.yml / fix(cassettes): break the same-second tie with the machine's counter (pull_request) Failing after 0s
The ordering gate compares created_at, which NIP-01 defines at one-second granularity. A dispense and the publish that follows it land inside one second routinely, so the report was dropped and the operator kept the pre-dispense count until the next heartbeat five minutes later. The machine bumps a counter on every local change to a bay count and carries it in its state document. m015 stores it per row, and the gate consults it only when the stamps are equal, where created_at carries no information at all. Only on equality, deliberately. A machine whose state.db was replaced restarts its counter at zero while its wall clock keeps moving forward; gating on the counter across different stamps would lock that machine out for good. Equal stamps with no counter on either side stay closed, which costs one heartbeat and risks nothing.
1787 lines
62 KiB
Python
1787 lines
62 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,
|
|
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_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(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 four keyed lists:
|
|
- 'rejected': any status='rejected' (Nostr attribution cross-check
|
|
failed — signer didn't match the machine identity). Distinct
|
|
from 'errored' because retry is wrong: the row was misrouted,
|
|
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,
|
|
)
|
|
return {
|
|
"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: 'pending' | 'processing' | 'processed' | 'partial' |
|
|
'refunded' | 'errored'. Clears processing_claim on terminal states so a
|
|
fresh claim attempt won't see a stale token."""
|
|
await db.execute(
|
|
"""
|
|
UPDATE spirekeeper.dca_settlements
|
|
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
|