diff --git a/crud.py b/crud.py index a622d96..bc5cb0a 100644 --- a/crud.py +++ b/crud.py @@ -12,9 +12,11 @@ from lnbits.helpers import urlsafe_short_hash from .models import ( CassetteConfig, + CassetteOp, ClientBalanceSummary, CommissionSplit, CommissionSplitLeg, + CreateCassetteOpData, CreateDcaClientData, CreateDcaPaymentData, CreateDcaSettlementData, @@ -1639,3 +1641,135 @@ async def apply_reported_state( }, ) 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 diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index 8187fa2..ab6cd88 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -15,6 +15,7 @@ from datetime import datetime, timezone import pytest from pydantic import ValidationError +from ..crud import _should_ack_op from ..models import ( CASSETTE_OP_TYPES, CassetteOp, @@ -140,3 +141,28 @@ class TestWireShape: "schema_version": 2, "ops": [], } + + +class TestShouldAckOp: + """The pure decision behind mark_cassette_ops_acked. + + The machine echoes a WINDOW of applied ids on every state publish, so the + same id arrives repeatedly and from a machine that may not own it. + """ + + def test_acks_an_unacked_op_for_the_reporting_machine(self): + assert _should_ack_op(op(op_type="empty"), "m1") is True + + def test_ignores_an_unknown_id(self): + assert _should_ack_op(None, "m1") is False + + def test_ignores_an_op_belonging_to_another_machine(self): + """Ids are unique, but a report from one machine must never close out + another machine's operation.""" + assert _should_ack_op(op(op_type="empty", machine_id="m2"), "m1") is False + + def test_keeps_the_first_acknowledgement(self): + """Every subsequent window carries the id again. Re-acking would slide + the timestamp forward and lose when the operation actually landed.""" + already = op(op_type="empty", acked_at=AT) + assert _should_ack_op(already, "m1") is False