Publish cassette operations instead of counts #46
2 changed files with 160 additions and 0 deletions
feat(cassettes): record and read operator operations
Append-only, and deliberately nothing here writes cassette_configs. That table now holds only what the machine has reported; letting an operation write it would put back the second writer this whole design exists to remove. get_cassette_ops_window returns oldest-first because order is meaning: a recount followed by a refill is not the same as the reverse. It takes the most recent N and reverses, so the window slides without the machine ever seeing them out of sequence. The window is what makes the channel self-healing, so it has to cover a plausible outage rather than just the newest change — a machine that missed one event still sees the operation in the next. _should_ack_op is extracted pure, the same way the state-event gate is, because three of its rules are easy to get wrong and none need a database to test: an unknown id closes out nothing, one machine must never be able to ack another machine's operation, and the FIRST acknowledgement is the one worth keeping. That last one matters because the machine echoes a window, so every id comes back many times over; overwriting would keep sliding the timestamp forward and lose when the operation actually landed. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
commit
b2157db219
134
crud.py
134
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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue