Publish cassette operations instead of counts #46
2 changed files with 138 additions and 3 deletions
feat(cassettes): publish operations to the ATM
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 <noreply@anthropic.com>
commit
90bd43d6da
|
|
@ -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
|
The ATM-side consumer (lamassu-next#56) subscribes by the d-tag + its own
|
||||||
npub, decrypts, validates, applies, hot-reloads HAL.
|
npub, decrypts, validates, applies, hot-reloads HAL.
|
||||||
|
|
||||||
Reverse direction (ATM → operator, v1 = one-shot bootstrap on first boot,
|
The operator → ATM direction carries OPERATIONS as of v2 (bitspire ADR-004):
|
||||||
v2 = continuous reverse channel for reconciliation):
|
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
|
kind = 30078
|
||||||
tags = [
|
tags = [
|
||||||
|
|
@ -52,7 +58,12 @@ from lnbits.core.signers.base import (
|
||||||
)
|
)
|
||||||
from lnbits.utils.nostr import normalize_public_key
|
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 .nip44 import Nip44Error
|
||||||
from .nostr_publish import (
|
from .nostr_publish import (
|
||||||
NostrPublishError,
|
NostrPublishError,
|
||||||
|
|
@ -179,6 +190,39 @@ async def publish_to_atm(
|
||||||
return signed
|
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)
|
# Consume — ATM → operator (the bootstrap consumer task)
|
||||||
# =============================================================================
|
# =============================================================================
|
||||||
|
|
|
||||||
|
|
@ -15,11 +15,13 @@ from datetime import datetime, timezone
|
||||||
import pytest
|
import pytest
|
||||||
from pydantic import ValidationError
|
from pydantic import ValidationError
|
||||||
|
|
||||||
|
from .. import cassette_transport
|
||||||
from ..crud import _should_ack_op
|
from ..crud import _should_ack_op
|
||||||
from ..models import (
|
from ..models import (
|
||||||
CASSETTE_OP_TYPES,
|
CASSETTE_OP_TYPES,
|
||||||
CassetteOp,
|
CassetteOp,
|
||||||
CreateCassetteOpData,
|
CreateCassetteOpData,
|
||||||
|
Machine,
|
||||||
PublishCassetteOpsPayload,
|
PublishCassetteOpsPayload,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -166,3 +168,92 @@ class TestShouldAckOp:
|
||||||
the timestamp forward and lose when the operation actually landed."""
|
the timestamp forward and lose when the operation actually landed."""
|
||||||
already = op(op_type="empty", acked_at=AT)
|
already = op(op_type="empty", acked_at=AT)
|
||||||
assert _should_ack_op(already, "m1") is False
|
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}"
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue