spirekeeper/tests/test_cassette_ops.py
Padreug 786857f568 feat(models): settle_transaction op, alerts_pubkey, operator_notified_at (m017)
ADR-005 §6 groundwork. settle_transaction is a machine-wide op like
resume_cash_out; it carries the machine txid it closes and the operator's
provenance note. super_config.alerts_pubkey overrides where owed-cash
alerts go; dca_settlements.operator_notified_at stops a report resend
from re-alerting.
2026-10-10 22:26:45 +02:00

360 lines
14 KiB
Python

"""
Tests for the v2 cassette-operations models (bitspire ADR-004).
The operator no longer publishes counts; it publishes operations and the
machine keeps the running total. These cover the pure pieces: per-type field
validation, and the wire shape the publisher ships.
A CreateCassetteOpData instance is meant to be publishable by construction —
same contract as FeeConfigPayload — so the type/field agreement is enforced in
the model rather than at the endpoint.
"""
from datetime import datetime, timezone
import pytest
from pydantic import ValidationError
from .. import cassette_transport, tasks
from ..crud import _should_ack_op
from ..models import (
CASSETTE_OP_TYPES,
CassetteOp,
CreateCassetteOpData,
Machine,
PublishCassetteOpsPayload,
PublishCassettesPayload,
)
AT = datetime.fromtimestamp(1790106060, timezone.utc)
def op(**kw) -> CassetteOp:
base = {"id": "op-1", "machine_id": "m1", "position": 2, "created_at": AT}
return CassetteOp(**{**base, **kw})
class TestCreateCassetteOpData:
def test_accepts_one_of_each_type(self):
CreateCassetteOpData(position=2, op_type="refill", bills=100)
CreateCassetteOpData(position=3, op_type="empty")
CreateCassetteOpData(position=1, op_type="recount", count=37)
CreateCassetteOpData(position=1, op_type="set_denomination", denomination=50)
@pytest.mark.parametrize(
"kwargs",
[
{"position": 2, "op_type": "refill"},
{"position": 1, "op_type": "recount"},
{"position": 1, "op_type": "set_denomination"},
],
)
def test_rejects_a_type_missing_its_field(self, kwargs):
with pytest.raises(ValidationError):
CreateCassetteOpData(**kwargs)
@pytest.mark.parametrize(
"kwargs",
[
{"position": 2, "op_type": "refill", "bills": 1, "count": 5},
{"position": 3, "op_type": "empty", "bills": 1},
{"position": 1, "op_type": "recount", "count": 1, "denomination": 50},
],
)
def test_rejects_a_type_carrying_a_foreign_field(self, kwargs):
"""An op that carries two meanings is ambiguous on the wire, and the
machine would have to guess which one to apply."""
with pytest.raises(ValidationError):
CreateCassetteOpData(**kwargs)
def test_rejects_a_refill_of_zero_or_fewer_notes(self):
"""A refill is a delta that adds notes. Zero is a no-op an operator
did not mean, and negative is a withdrawal wearing a refill's name."""
for bills in (0, -5):
with pytest.raises(ValidationError):
CreateCassetteOpData(position=2, op_type="refill", bills=bills)
def test_allows_a_recount_to_zero(self):
"""Distinct from refill: counting a bay and finding it empty is a real
and important observation."""
assert CreateCassetteOpData(position=2, op_type="recount", count=0).count == 0
def test_rejects_a_negative_recount_and_a_non_positive_denomination(self):
with pytest.raises(ValidationError):
CreateCassetteOpData(position=2, op_type="recount", count=-1)
with pytest.raises(ValidationError):
CreateCassetteOpData(position=2, op_type="set_denomination", denomination=0)
def test_rejects_an_unknown_type_and_a_non_positive_position(self):
with pytest.raises(ValidationError):
CreateCassetteOpData(position=1, op_type="drain")
with pytest.raises(ValidationError):
CreateCassetteOpData(position=0, op_type="empty")
class TestWireShape:
def test_each_type_ships_only_its_own_field(self):
assert op(op_type="refill", bills=100).to_wire_dict() == {
"id": "op-1",
"at": 1790106060,
"type": "refill",
"position": 2,
"bills": 100,
}
assert op(op_type="empty").to_wire_dict() == {
"id": "op-1",
"at": 1790106060,
"type": "empty",
"position": 2,
}
assert op(op_type="recount", count=37).to_wire_dict()["count"] == 37
assert (
op(op_type="set_denomination", denomination=50).to_wire_dict()[
"denomination"
]
== 50
)
def test_nulls_never_reach_the_wire(self):
"""The row has three nullable columns and one op only ever means one
of them. Shipping the other two as null would make the machine guess."""
for op_type in CASSETTE_OP_TYPES:
kw = {
"refill": {"bills": 1},
"recount": {"count": 1},
"set_denomination": {"denomination": 1},
"empty": {},
# machine-wide (ADR-005 §5): position 0, no position on the wire
"resume_cash_out": {"position": 0},
# machine-wide too (ADR-005 §6): carries the txid it closes
"settle_transaction": {"position": 0, "txid": "tx_1"},
}[op_type]
wire = op(op_type=op_type, **kw).to_wire_dict()
assert None not in wire.values()
def test_payload_declares_v2_and_preserves_order(self):
ops = [
op(id="a", op_type="refill", bills=1),
op(id="b", op_type="empty"),
]
wire = PublishCassetteOpsPayload(ops=ops).to_wire_dict()
assert wire["schema_version"] == 2
assert [o["id"] for o in wire["ops"]] == ["a", "b"]
def test_an_empty_window_is_representable(self):
"""A machine with no operator history still gets a well-formed
payload rather than the publisher having to special-case it."""
assert PublishCassetteOpsPayload(ops=[]).to_wire_dict() == {
"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
# =============================================================================
# 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}"
# =============================================================================
# _record_op_acknowledgements — the only ack this transport can carry
# =============================================================================
def state_payload(**kw) -> PublishCassettesPayload:
base = {"positions": {"1": {"denomination": 50, "count": 24}}}
return PublishCassettesPayload(**{**base, **kw})
class TestRecordOpAcknowledgements:
@pytest.mark.asyncio
async def test_marks_the_reported_ids(self):
calls = []
async def mark(machine_id, op_ids):
calls.append((machine_id, op_ids))
return len(op_ids)
await tasks._record_op_acknowledgements(
"m1", state_payload(applied_ops=["a", "b"]), mark
)
assert calls == [("m1", ["a", "b"])]
@pytest.mark.asyncio
async def test_does_nothing_when_the_machine_reports_none(self):
"""An older machine sends no applied_ops at all, and a new one with
nothing applied sends an empty list. Neither should write."""
calls = []
async def mark(machine_id, op_ids):
calls.append((machine_id, op_ids))
return 0
await tasks._record_op_acknowledgements("m1", state_payload(), mark)
await tasks._record_op_acknowledgements(
"m1", state_payload(applied_ops=[]), mark
)
assert calls == []
@pytest.mark.asyncio
async def test_acks_even_when_the_counts_were_not_newer(self):
"""The machine echoes its applied ids on every publish, including
heartbeats that carry nothing new about the counts. One of those can
still be the first event to tell us an operation landed, so the ack
must not depend on the state having advanced."""
seen = []
async def mark(machine_id, op_ids):
seen.extend(op_ids)
return len(op_ids)
# Same positions as already on file — a pure heartbeat.
await tasks._record_op_acknowledgements(
"m1", state_payload(applied_ops=["late-ack"]), mark
)
assert seen == ["late-ack"]
# =============================================================================
# _record_counts_uncertainty — the machine saying "don't trust these counts"
# =============================================================================
class TestRecordCountsUncertainty:
@pytest.mark.asyncio
async def test_stores_the_reported_moment_as_utc(self):
calls = []
async def setter(machine_id, since):
calls.append((machine_id, since))
await tasks._record_counts_uncertainty(
"m1", state_payload(counts_uncertain_since=1790110546), setter
)
assert len(calls) == 1
machine_id, since = calls[0]
assert machine_id == "m1"
assert since is not None
assert since.tzinfo is not None
assert int(since.timestamp()) == 1790110546
@pytest.mark.asyncio
async def test_clears_the_marker_when_the_machine_is_confident_again(self):
"""A banner that never goes away is a banner nobody reads. The
machine dropping the field is how the operator learns the recount
took, so None must be written through rather than skipped."""
calls = []
async def setter(machine_id, since):
calls.append((machine_id, since))
await tasks._record_counts_uncertainty("m1", state_payload(), setter)
assert calls == [("m1", None)]