From 90bd43d6dae200f671e13c5f92bc748820d2b27f Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:35:42 +0200 Subject: [PATCH] feat(cassettes): publish operations to the ATM MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The v2 operator to ATM wire. Same kind-30078 document and the same d-tag the counts wire used, because the machine subscribes by that tag and the document is addressable, so v2 replaces v1 in place. Sends a WINDOW of recent operations, oldest-first, not just the newest change. Each publish replaces the last, so a machine that was offline for one of them would otherwise never see that operation again; carrying the recent history means the channel heals itself without anyone noticing it broke. Re-delivery costs nothing because every op carries an id the machine dedups on. Tests pin the contract rather than the implementation: the d-tag, that the payload declares v2 and carries no positions key, that window order survives the publisher untouched, that an empty window still ships a well-formed document so a machine can tell "no operations" from "operator still on v1", and that an npub entered in the UI is normalised to hex — get that last one wrong and the machine's subscription filter silently never matches. Additive. The endpoints still publish counts until the next commit, so the tree is not left half-switched. Co-Authored-By: Claude Fable 5.1 --- cassette_transport.py | 50 +++++++++++++++++++-- tests/test_cassette_ops.py | 91 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 138 insertions(+), 3 deletions(-) diff --git a/cassette_transport.py b/cassette_transport.py index e6517a1..f64ddec 100644 --- a/cassette_transport.py +++ b/cassette_transport.py @@ -17,8 +17,14 @@ publishes position-keyed cassette config to a target ATM via: The ATM-side consumer (lamassu-next#56) subscribes by the d-tag + its own npub, decrypts, validates, applies, hot-reloads HAL. -Reverse direction (ATM → operator, v1 = one-shot bootstrap on first boot, -v2 = continuous reverse channel for reconciliation): +The operator → ATM direction carries OPERATIONS as of v2 (bitspire ADR-004): +refill, empty, recount, set_denomination, each with an id the machine dedups +on. It used to carry absolute counts, which meant the operator and the machine +both wrote the same value over a transport that never tells a writer it lost — +so a form loaded before a dispense discarded that dispense when published. + +Reverse direction (ATM → operator, continuous: the machine publishes on +startup, after every change to its bays, and on a heartbeat): kind = 30078 tags = [ @@ -52,7 +58,12 @@ from lnbits.core.signers.base import ( ) from lnbits.utils.nostr import normalize_public_key -from .models import Machine, PublishCassettesPayload +from .models import ( + CassetteOp, + Machine, + PublishCassetteOpsPayload, + PublishCassettesPayload, +) from .nip44 import Nip44Error from .nostr_publish import ( NostrPublishError, @@ -179,6 +190,39 @@ async def publish_to_atm( return signed +async def publish_ops_to_atm( + machine: Machine, + ops: list[CassetteOp], + operator_user_id: str, +) -> dict: + """Publish the operator's recent cassette OPERATIONS to the target ATM. + + The v2 wire (bitspire ADR-004). Replaces sending absolute counts, which + let a dashboard form loaded before a dispense silently discard that + dispense — the operator and the machine were both writing the same value + over a transport that never tells a writer it lost. + + `ops` is a WINDOW, oldest-first, not just the newest change. The event is + addressable, so each publish replaces the last, and a machine that was + offline for one of them would otherwise never see that operation again. + Carrying the recent history means the channel heals itself without anyone + noticing it broke. Re-delivery is harmless because each op carries an id + the machine dedups on. + """ + atm_pubkey_hex = _atm_hex_pubkey(machine) + payload = PublishCassetteOpsPayload(ops=ops) + 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 ops (machine={machine.id}, ops={[o.op_type for o in ops]})" + ), + ) + return signed + + # ============================================================================= # Consume — ATM → operator (the bootstrap consumer task) # ============================================================================= diff --git a/tests/test_cassette_ops.py b/tests/test_cassette_ops.py index ab6cd88..7014a94 100644 --- a/tests/test_cassette_ops.py +++ b/tests/test_cassette_ops.py @@ -15,11 +15,13 @@ from datetime import datetime, timezone import pytest from pydantic import ValidationError +from .. import cassette_transport from ..crud import _should_ack_op from ..models import ( CASSETTE_OP_TYPES, CassetteOp, CreateCassetteOpData, + Machine, PublishCassetteOpsPayload, ) @@ -166,3 +168,92 @@ class TestShouldAckOp: 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 + + +# ============================================================================= +# publish_ops_to_atm — the v2 wire contract +# ============================================================================= + +ATM_HEX = "df2003343784b69cb813b2a4fd231f83ae81133279251c735414f9909baa7ac6" + + +def machine() -> Machine: + return Machine( + id="m1", + operator_user_id="op1", + machine_npub=ATM_HEX, + wallet_id="w1", + name="Cinderella", + location=None, + fiat_code="EUR", + is_active=True, + created_at=AT, + updated_at=AT, + ) + + +@pytest.fixture +def captured(monkeypatch): + """Capture what the transport would publish, without a relay or signer.""" + seen: dict = {} + + async def fake_publish(**kwargs): + seen.update(kwargs) + return {"id": "event-id"} + + monkeypatch.setattr( + cassette_transport, "publish_encrypted_kind_30078", fake_publish + ) + return seen + + +class TestPublishOpsToAtm: + @pytest.mark.asyncio + async def test_publishes_v2_ops_to_the_config_d_tag(self, captured): + ops = [ + op(id="a", op_type="refill", bills=100), + op(id="b", op_type="empty", position=3), + ] + await cassette_transport.publish_ops_to_atm(machine(), ops, "op1") + + # Same d-tag as the counts wire it replaces: the machine subscribes by + # this tag, and the document is addressable, so v2 replaces v1 in place. + assert captured["d_tag"] == f"bitspire-cassettes:{ATM_HEX}" + assert captured["recipient_pubkey_hex"] == ATM_HEX + assert captured["operator_user_id"] == "op1" + + payload = captured["payload"] + assert payload["schema_version"] == 2 + assert [o["id"] for o in payload["ops"]] == ["a", "b"] + assert payload["ops"][0]["bills"] == 100 + assert "positions" not in payload + + @pytest.mark.asyncio + async def test_preserves_window_order(self, captured): + """Order is meaning: a recount then a refill is not the same as the + reverse, so the publisher must not re-sort what crud handed it.""" + ops = [ + op(id="first", op_type="recount", count=10), + op(id="second", op_type="refill", bills=5), + ] + await cassette_transport.publish_ops_to_atm(machine(), ops, "op1") + assert [o["id"] for o in captured["payload"]["ops"]] == ["first", "second"] + + @pytest.mark.asyncio + async def test_publishes_an_empty_window_rather_than_skipping(self, captured): + """A machine with no operator history still gets a well-formed v2 + document, so it can tell 'no operations' from 'operator still on v1'.""" + await cassette_transport.publish_ops_to_atm(machine(), [], "op1") + assert captured["payload"] == {"schema_version": 2, "ops": []} + + @pytest.mark.asyncio + async def test_accepts_an_npub_and_publishes_hex(self, captured): + """Operators enter either form in the UI; the d-tag is always hex, or + the machine's subscription filter silently never matches.""" + import bech32 + + data = bech32.convertbits(bytes.fromhex(ATM_HEX), 8, 5) + npub = bech32.bech32_encode("npub", data) + m = machine().copy(update={"machine_npub": npub}) + await cassette_transport.publish_ops_to_atm(m, [], "op1") + assert captured["d_tag"] == f"bitspire-cassettes:{ATM_HEX}"