diff --git a/cassette_transport.py b/cassette_transport.py index f64ddec..bc84523 100644 --- a/cassette_transport.py +++ b/cassette_transport.py @@ -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 `` 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) # ============================================================================= diff --git a/crud.py b/crud.py index bc5cb0a..e83b9af 100644 --- a/crud.py +++ b/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, diff --git a/models.py b/models.py index 7a09093..f2a1f3d 100644 --- a/models.py +++ b/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": {"": {"denomination", "count"}}}`.""" diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index 242aa13..b6adf2d 100644 --- a/tests/test_cassette_configs.py +++ b/tests/test_cassette_configs.py @@ -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 # ============================================================================= diff --git a/views_api.py b/views_api.py index 73b2d0b..33fb2de 100644 --- a/views_api.py +++ b/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: and - p=. +) -> list[CassetteOp]: + """Recent cassette operations for a machine, newest first. - The `` 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