Publish cassette operations instead of counts #46
5 changed files with 89 additions and 214 deletions
feat(cassettes): swap the count-publish endpoint for operation endpoints
The operator can no longer write a count. POST .../cassettes/ops records one operation — refill, empty, recount, set_denomination — and publishes the machine's recent window; GET .../cassettes/ops lists them newest first with acked_at, so the dashboard can tell a delivered operation from one merely sent. POST .../cassettes/publish is gone, along with update_cassette_config and UpsertCassetteConfigData. Nothing in the operator can now set a count, which is the point: a value with one writer cannot be clobbered. Under the old endpoint a dashboard form loaded before a dispense silently discarded that dispense on publish, and neither side could detect it — addressable events order by created_at at second granularity and a relay returns OK for an event it then drops, so the losing writer is never told. The op is recorded before the publish and is deliberately not rolled back when the publish fails. It records something that physically happened; notes went into a bay whether or not a relay was reachable. The window carries recent operations rather than just the newest, so an op that missed its own publish rides out with the next one. Validation rejects an unpaired machine and a position the machine has not reported. Bay count stays hardware-determined.
commit
2ac3e2064e
|
|
@ -36,7 +36,7 @@ startup, after every change to its bays, and on a heartbeat):
|
|||
|
||||
This module owns the wire-format side of both directions. The consumer
|
||||
task (tasks.py) calls `decrypt_and_parse_state_event` per incoming event;
|
||||
the API endpoint (views_api.py) calls `publish_to_atm` per operator submit.
|
||||
the API endpoint (views_api.py) calls `publish_ops_to_atm` per operation.
|
||||
|
||||
The `<m>` placeholder semantics (load-bearing per the 2026-05-30T11:50Z
|
||||
coord-log entry): always the ATM's hex pubkey, NEVER spirekeeper's
|
||||
|
|
@ -86,7 +86,7 @@ __all__ = [
|
|||
"RelayUnavailable",
|
||||
"build_state_d_tags_for_machines",
|
||||
"decrypt_and_parse_state_event",
|
||||
"publish_to_atm",
|
||||
"publish_ops_to_atm",
|
||||
]
|
||||
|
||||
_D_TAG_CONFIG_PREFIX = "bitspire-cassettes:" # operator → ATM
|
||||
|
|
@ -163,33 +163,6 @@ def build_state_d_tags_for_machines(machines: list[Machine]) -> list[str]:
|
|||
# =============================================================================
|
||||
|
||||
|
||||
async def publish_to_atm(
|
||||
machine: Machine,
|
||||
payload: PublishCassettesPayload,
|
||||
operator_user_id: str,
|
||||
) -> dict:
|
||||
"""Build, encrypt, sign, and publish a kind-30078 cassette config event
|
||||
from the operator to the target ATM.
|
||||
|
||||
Returns the signed event dict on success (caller may log event.id for
|
||||
audit). Raises NostrPublishError subclasses (re-exported here as
|
||||
CassetteTransportError, OperatorIdentityMissing, SignerUnavailable,
|
||||
RelayUnavailable) on hard failures.
|
||||
"""
|
||||
atm_pubkey_hex = _atm_hex_pubkey(machine)
|
||||
signed = await publish_encrypted_kind_30078(
|
||||
operator_user_id=operator_user_id,
|
||||
recipient_pubkey_hex=atm_pubkey_hex,
|
||||
d_tag=_config_d_tag(atm_pubkey_hex),
|
||||
payload=payload.to_wire_dict(),
|
||||
log_context=(
|
||||
f"cassette config (machine={machine.id}, "
|
||||
f"positions={sorted(payload.positions.keys())})"
|
||||
),
|
||||
)
|
||||
return signed
|
||||
|
||||
|
||||
async def publish_ops_to_atm(
|
||||
machine: Machine,
|
||||
ops: list[CassetteOp],
|
||||
|
|
@ -224,7 +197,7 @@ async def publish_ops_to_atm(
|
|||
|
||||
|
||||
# =============================================================================
|
||||
# Consume — ATM → operator (the bootstrap consumer task)
|
||||
# Consume — ATM → operator (the machine's state reports)
|
||||
# =============================================================================
|
||||
|
||||
|
||||
|
|
|
|||
42
crud.py
42
crud.py
|
|
@ -36,7 +36,6 @@ from .models import (
|
|||
UpdateDepositStatusData,
|
||||
UpdateMachineData,
|
||||
UpdateSuperConfigData,
|
||||
UpsertCassetteConfigData,
|
||||
UpsertDcaLpData,
|
||||
)
|
||||
|
||||
|
|
@ -1446,10 +1445,11 @@ async def upsert_fleet_snapshot(
|
|||
# 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)
|
||||
# - Operator edit of denomination or count → update_cassette_config
|
||||
# (refuses to create new rows; the slot count is hardware-determined)
|
||||
# - Row creation/deletion for a new position → admin only, via ATM
|
||||
# re-provisioning + new bootstrap event (not exposed in v1 here)
|
||||
# - 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:
|
||||
|
|
@ -1521,38 +1521,6 @@ async def list_cassette_configs_for_machine(
|
|||
)
|
||||
|
||||
|
||||
async def update_cassette_config(
|
||||
machine_id: str,
|
||||
position: int,
|
||||
data: UpsertCassetteConfigData,
|
||||
*,
|
||||
updated_by: str | None = None,
|
||||
) -> CassetteConfig | None:
|
||||
"""Operator-driven row update: change denomination and/or count for a
|
||||
single cassette slot. Refuses to create new rows — those only land via
|
||||
apply_reported_state() consuming an ATM bootstrap event (per #29 row
|
||||
lifecycle: hardware-determined slot count, not operator-creatable).
|
||||
Returns None if the (machine_id, position) row doesn't exist.
|
||||
"""
|
||||
existing = await get_cassette_config(machine_id, position)
|
||||
if existing is None:
|
||||
return None
|
||||
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"] = datetime.now()
|
||||
update_data["updated_by"] = updated_by
|
||||
set_clause = ", ".join(f"{k} = :{k}" for k in update_data)
|
||||
update_data["mid"] = machine_id
|
||||
update_data["pos"] = position
|
||||
await db.execute(
|
||||
f"UPDATE spirekeeper.cassette_configs SET {set_clause} "
|
||||
"WHERE machine_id = :mid AND position = :pos",
|
||||
update_data,
|
||||
)
|
||||
return await get_cassette_config(machine_id, position)
|
||||
|
||||
|
||||
async def apply_reported_state(
|
||||
machine_id: str,
|
||||
event_id: str,
|
||||
|
|
|
|||
25
models.py
25
models.py
|
|
@ -675,31 +675,6 @@ class CassetteConfig(BaseModel):
|
|||
state_event_id: str | None
|
||||
|
||||
|
||||
class UpsertCassetteConfigData(BaseModel):
|
||||
"""Operator edits a single cassette row's denomination or count from
|
||||
the dashboard. Both fields optional; pass only those changed.
|
||||
Position is not edited — it's the row's identity (hardware bay)."""
|
||||
|
||||
denomination: int | None = None
|
||||
count: int | None = None
|
||||
|
||||
@validator("denomination")
|
||||
def denomination_positive(cls, v):
|
||||
if v is None:
|
||||
return v
|
||||
if v <= 0:
|
||||
raise ValueError("denomination must be > 0")
|
||||
return v
|
||||
|
||||
@validator("count")
|
||||
def count_non_negative(cls, v):
|
||||
if v is None:
|
||||
return v
|
||||
if v < 0:
|
||||
raise ValueError("count must be >= 0")
|
||||
return v
|
||||
|
||||
|
||||
class CassettePayloadRow(BaseModel):
|
||||
"""One position's payload values in the wire-format
|
||||
`{"positions": {"<pos>": {"denomination", "count"}}}`."""
|
||||
|
|
|
|||
|
|
@ -2,9 +2,9 @@
|
|||
Tests for the v1.1 cassette-config layer (aiolabs/satmachineadmin#29).
|
||||
|
||||
Covers the pure pieces that don't need a live DB:
|
||||
- Pydantic validator behaviour on PublishCassettesPayload + the row /
|
||||
upsert models (position key coercion, integer ranges, multiple-same-
|
||||
denomination payloads, wire-format round-trip)
|
||||
- Pydantic validator behaviour on PublishCassettesPayload + the row model
|
||||
(position key coercion, integer ranges, multiple-same-denomination
|
||||
payloads, wire-format round-trip)
|
||||
- _should_apply_state_event ordering gate (extracted from
|
||||
apply_reported_state so the decision is testable without a database
|
||||
round-trip)
|
||||
|
|
@ -29,7 +29,6 @@ from ..crud import _as_unix, _should_apply_state_event
|
|||
from ..models import (
|
||||
CassettePayloadRow,
|
||||
PublishCassettesPayload,
|
||||
UpsertCassetteConfigData,
|
||||
)
|
||||
|
||||
# =============================================================================
|
||||
|
|
@ -149,44 +148,6 @@ class TestCassettePayloadRow:
|
|||
CassettePayloadRow(denomination=20, count=-1)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# UpsertCassetteConfigData — operator-edit form
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestUpsertCassetteConfigData:
|
||||
"""Operator-driven row edit. Both fields optional; same int constraints
|
||||
as the wire-format row but applied independently per-edit. Position is
|
||||
NOT editable — it's the row's identity (the hardware bay number)."""
|
||||
|
||||
def test_partial_update_count_only(self):
|
||||
d = UpsertCassetteConfigData(count=80)
|
||||
assert d.count == 80
|
||||
assert d.denomination is None
|
||||
|
||||
def test_partial_update_denomination_only(self):
|
||||
"""v1.1 operational case: operator records a cartridge swap at
|
||||
refill — slot 1 was $20, dispatcher replaced with $50."""
|
||||
d = UpsertCassetteConfigData(denomination=50)
|
||||
assert d.denomination == 50
|
||||
assert d.count is None
|
||||
|
||||
def test_empty_update_is_legal(self):
|
||||
"""An empty UpsertCassetteConfigData parses fine; the CRUD short-
|
||||
circuits a no-op on empty payload (no SQL emitted)."""
|
||||
d = UpsertCassetteConfigData()
|
||||
assert d.count is None
|
||||
assert d.denomination is None
|
||||
|
||||
def test_rejects_negative_count(self):
|
||||
with pytest.raises(ValueError):
|
||||
UpsertCassetteConfigData(count=-1)
|
||||
|
||||
def test_rejects_non_positive_denomination(self):
|
||||
with pytest.raises(ValueError):
|
||||
UpsertCassetteConfigData(denomination=0)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# _should_apply_state_event — ordering gate
|
||||
# =============================================================================
|
||||
|
|
|
|||
158
views_api.py
158
views_api.py
|
|
@ -26,7 +26,7 @@ from .cassette_transport import (
|
|||
OperatorIdentityMissing,
|
||||
RelayUnavailable,
|
||||
SignerUnavailable,
|
||||
publish_to_atm,
|
||||
publish_ops_to_atm,
|
||||
)
|
||||
from .fee_transport import publish_fee_config
|
||||
from .pairing import (
|
||||
|
|
@ -40,6 +40,7 @@ from .pairing import (
|
|||
from .crud import (
|
||||
append_settlement_note,
|
||||
count_completed_legs_for_settlement,
|
||||
create_cassette_op,
|
||||
create_dca_client,
|
||||
create_deposit,
|
||||
create_machine,
|
||||
|
|
@ -47,6 +48,7 @@ from .crud import (
|
|||
delete_deposit,
|
||||
delete_machine,
|
||||
force_reset_stuck_settlement,
|
||||
get_cassette_ops_window,
|
||||
get_client_balance_summary,
|
||||
get_commission_splits,
|
||||
get_dca_client,
|
||||
|
|
@ -66,12 +68,12 @@ from .crud import (
|
|||
get_super_config,
|
||||
list_all_active_machines,
|
||||
list_cassette_configs_for_machine,
|
||||
list_cassette_ops,
|
||||
lp_is_onboarded,
|
||||
replace_commission_splits,
|
||||
reset_settlement_for_retry,
|
||||
set_machine_pairing,
|
||||
set_machine_unpaired,
|
||||
update_cassette_config,
|
||||
update_dca_client,
|
||||
update_deposit,
|
||||
update_deposit_status,
|
||||
|
|
@ -86,8 +88,10 @@ from .distribution import (
|
|||
from .models import (
|
||||
AppendSettlementNoteData,
|
||||
CassetteConfig,
|
||||
CassetteOp,
|
||||
ClientBalanceSummary,
|
||||
CommissionSplit,
|
||||
CreateCassetteOpData,
|
||||
CreateDcaClientData,
|
||||
CreateDepositData,
|
||||
CreateMachineData,
|
||||
|
|
@ -98,7 +102,6 @@ from .models import (
|
|||
Machine,
|
||||
PairMachineData,
|
||||
PartialDispenseData,
|
||||
PublishCassettesPayload,
|
||||
SetCommissionSplitsData,
|
||||
SettleBalanceData,
|
||||
StuckSettlementsResponse,
|
||||
|
|
@ -108,7 +111,6 @@ from .models import (
|
|||
UpdateDepositStatusData,
|
||||
UpdateMachineData,
|
||||
UpdateSuperConfigData,
|
||||
UpsertCassetteConfigData,
|
||||
)
|
||||
|
||||
spirekeeper_api_router = APIRouter()
|
||||
|
|
@ -1102,16 +1104,21 @@ async def api_update_super_config(
|
|||
|
||||
|
||||
# =============================================================================
|
||||
# Cassette configs (#29 v1.1) — per-machine ATM cassette inventory
|
||||
# Cassettes — per-machine ATM inventory (bitspire ADR-004)
|
||||
# =============================================================================
|
||||
# v1.1 surface, paired with aiolabs/lamassu-next#56 ATM-side. Two endpoints:
|
||||
# GET /machines/{id}/cassettes — list rows for the operator UI
|
||||
# POST /machines/{id}/cassettes/publish — apply edits + publish kind-30078
|
||||
# GET /machines/{id}/cassettes — bays as the machine last reported them
|
||||
# GET /machines/{id}/cassettes/ops — recent operations, with their acks
|
||||
# POST /machines/{id}/cassettes/ops — record one operation and publish
|
||||
#
|
||||
# Row creation (new (machine_id, position) pairs) is admin-only via the
|
||||
# bootstrap consumer task — slot count is hardware-determined. Operator-
|
||||
# side flow is edit-and-publish over the existing rows only; the editable
|
||||
# fields per row are denomination and count.
|
||||
# The operator does not write counts. It records what it DID to a bay and the
|
||||
# machine, which holds the notes, keeps the running total. The rows behind the
|
||||
# first endpoint are the machine's report, not an operator draft.
|
||||
#
|
||||
# This replaced an edit-and-publish form over absolute counts. Both sides wrote
|
||||
# the same value across a transport that never tells a writer it lost, so a
|
||||
# form loaded before a dispense discarded that dispense when published — and
|
||||
# nothing could detect it afterwards. Bay count stays hardware-determined:
|
||||
# rows appear and disappear only as the machine reports them.
|
||||
|
||||
|
||||
@spirekeeper_api_router.get(
|
||||
|
|
@ -1129,105 +1136,96 @@ async def api_list_machine_cassettes(
|
|||
return await list_cassette_configs_for_machine(machine_id)
|
||||
|
||||
|
||||
@spirekeeper_api_router.post(
|
||||
"/api/v1/dca/machines/{machine_id}/cassettes/publish",
|
||||
response_model=list[CassetteConfig],
|
||||
@spirekeeper_api_router.get(
|
||||
"/api/v1/dca/machines/{machine_id}/cassettes/ops",
|
||||
response_model=list[CassetteOp],
|
||||
)
|
||||
async def api_publish_machine_cassettes(
|
||||
async def api_list_machine_cassette_ops(
|
||||
machine_id: str,
|
||||
payload: PublishCassettesPayload,
|
||||
user: User = Depends(check_user_exists),
|
||||
) -> list[CassetteConfig]:
|
||||
"""Operator submits the full per-machine cassette state for publish to
|
||||
the ATM. Validates the position set matches what's currently in
|
||||
cassette_configs for the machine (slot count is hardware-fixed),
|
||||
upserts each row, then encrypts + signs + publishes a kind-30078
|
||||
event tagged with d=bitspire-cassettes:<atm_pubkey_hex> and
|
||||
p=<atm_pubkey_hex>.
|
||||
) -> list[CassetteOp]:
|
||||
"""Recent cassette operations for a machine, newest first.
|
||||
|
||||
The `<m>` placeholder in the published d-tag is the ATM's hex pubkey
|
||||
from machine.machine_npub (canonicalised via normalize_public_key),
|
||||
NOT the internal dca_machines.id UUID — see #29 'machine_id semantics'
|
||||
section and coord-log 2026-05-30T11:50Z load-bearing nudge.
|
||||
Each carries acked_at: null until the machine has reported that id back in
|
||||
its state document, which is how the dashboard distinguishes an operation
|
||||
that has been delivered from one that has merely been sent.
|
||||
"""
|
||||
await _machine_owned_by(machine_id, user.id)
|
||||
return await list_cassette_ops(machine_id)
|
||||
|
||||
Returns the fresh cassette_configs rows after the upserts so the UI
|
||||
can refresh its table from one round-trip.
|
||||
|
||||
@spirekeeper_api_router.post(
|
||||
"/api/v1/dca/machines/{machine_id}/cassettes/ops",
|
||||
response_model=CassetteOp,
|
||||
)
|
||||
async def api_create_machine_cassette_op(
|
||||
machine_id: str,
|
||||
data: CreateCassetteOpData,
|
||||
user: User = Depends(check_user_exists),
|
||||
) -> CassetteOp:
|
||||
"""Record one cassette operation and publish the machine's recent window.
|
||||
|
||||
This replaces publishing absolute counts. The operator now records what it
|
||||
DID — a refill, an empty, a recount, a denomination change — and the
|
||||
machine keeps the running total. Both sides used to write the same value
|
||||
over a transport that never tells a writer it lost, so a dashboard form
|
||||
loaded before a dispense silently discarded that dispense when published.
|
||||
|
||||
The op is recorded BEFORE the publish and is deliberately not rolled back
|
||||
if the publish fails. It represents something that physically happened —
|
||||
notes went into a bay — and that stays true whether or not a relay was
|
||||
reachable. Because each publish carries a window of recent operations
|
||||
rather than just the newest, an op that missed its own publish rides out
|
||||
with the next one.
|
||||
|
||||
Errors:
|
||||
400 — payload position set doesn't match the machine's stored set
|
||||
(operator publishing for a slot that doesn't exist on the
|
||||
ATM; or the bootstrap hasn't landed yet so no rows exist)
|
||||
400 — operator hasn't onboarded a Nostr identity
|
||||
503 — signer offline / client-side-only, or nostrclient extension
|
||||
not installed on this LNbits instance
|
||||
500 — anything else from the publish path
|
||||
400 — machine not paired, so there is no ATM identity to publish to
|
||||
400 — position is not a bay this machine has reported
|
||||
503 — signer offline, or relay/nostrclient unreachable. The operation is
|
||||
still recorded and will be delivered with the next publish.
|
||||
"""
|
||||
machine = await _machine_owned_by(machine_id, user.id)
|
||||
if not machine.machine_npub:
|
||||
# Unpaired machine (machine_npub None — nullable since #29/m011) has no
|
||||
# ATM identity to publish a cassette config to. Fail fast with a clean
|
||||
# 400 instead of crashing publish_to_atm's normalize_public_key(None).
|
||||
raise HTTPException(
|
||||
HTTPStatus.BAD_REQUEST,
|
||||
"machine is not paired — pair it before publishing cassette config",
|
||||
"machine is not paired — pair it before recording cassette operations",
|
||||
)
|
||||
|
||||
existing = await list_cassette_configs_for_machine(machine_id)
|
||||
existing_positions = {row.position for row in existing}
|
||||
incoming_positions = set(payload.positions.keys())
|
||||
|
||||
if not existing:
|
||||
raise HTTPException(
|
||||
HTTPStatus.BAD_REQUEST,
|
||||
(
|
||||
"No cassette_configs rows exist for this machine yet — "
|
||||
"waiting for the ATM's bootstrap state event. Power on the "
|
||||
"ATM and confirm it has reached the configured relay; "
|
||||
"spirekeeper will auto-populate cassette_configs on "
|
||||
"receipt."
|
||||
"No cassette rows for this machine yet — waiting for its state "
|
||||
"event. Power on the ATM and confirm it has reached the "
|
||||
"configured relay; spirekeeper populates the bays on receipt."
|
||||
),
|
||||
)
|
||||
if existing_positions != incoming_positions:
|
||||
missing = existing_positions - incoming_positions
|
||||
extra = incoming_positions - existing_positions
|
||||
known_positions = {row.position for row in existing}
|
||||
if data.position not in known_positions:
|
||||
raise HTTPException(
|
||||
HTTPStatus.BAD_REQUEST,
|
||||
(
|
||||
"Payload position set doesn't match the machine's stored "
|
||||
f"set. Missing from payload: {sorted(missing)}; extra in "
|
||||
f"payload: {sorted(extra)}. Slot count is hardware-fixed "
|
||||
"— re-provision the ATM via atm-tui to add/remove physical "
|
||||
"bays, then re-publish."
|
||||
f"position {data.position} is not a bay this machine has "
|
||||
f"reported (has: {sorted(known_positions)}). Bay count is "
|
||||
"hardware-determined; re-provision via atm-tui to change it."
|
||||
),
|
||||
)
|
||||
|
||||
# Apply each per-row edit so the operator-believed state on
|
||||
# spirekeeper reflects the published payload, even if the ATM
|
||||
# ack lands later (v2). updated_by audit-stamps the operator user id.
|
||||
for pos, row in payload.positions.items():
|
||||
updated = await update_cassette_config(
|
||||
machine_id,
|
||||
pos,
|
||||
UpsertCassetteConfigData(denomination=row.denomination, count=row.count),
|
||||
updated_by=user.id,
|
||||
)
|
||||
if updated is None:
|
||||
# Defensive — we just validated the row exists, but a
|
||||
# concurrent delete could land between. Surface as 500.
|
||||
raise HTTPException(
|
||||
HTTPStatus.INTERNAL_SERVER_ERROR,
|
||||
f"cassette row for position {pos} disappeared mid-publish",
|
||||
)
|
||||
op = await create_cassette_op(machine_id, data, created_by=user.id)
|
||||
|
||||
window = await get_cassette_ops_window(machine_id)
|
||||
try:
|
||||
await publish_to_atm(machine, payload, user.id)
|
||||
await publish_ops_to_atm(machine, window, user.id)
|
||||
except OperatorIdentityMissing as exc:
|
||||
raise HTTPException(HTTPStatus.BAD_REQUEST, str(exc)) from exc
|
||||
except SignerUnavailable as exc:
|
||||
raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc
|
||||
except RelayUnavailable as exc:
|
||||
raise HTTPException(HTTPStatus.SERVICE_UNAVAILABLE, str(exc)) from exc
|
||||
except (SignerUnavailable, RelayUnavailable) as exc:
|
||||
raise HTTPException(
|
||||
HTTPStatus.SERVICE_UNAVAILABLE,
|
||||
f"{exc} — the operation was recorded and will be delivered with "
|
||||
"the next publish",
|
||||
) from exc
|
||||
except CassetteTransportError as exc:
|
||||
raise HTTPException(HTTPStatus.INTERNAL_SERVER_ERROR, str(exc)) from exc
|
||||
|
||||
return await list_cassette_configs_for_machine(machine_id)
|
||||
return op
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue