From 5f60b3fe31ca14b3e9b0cb014ddc08617d303eb6 Mon Sep 17 00:00:00 2001 From: Padreug Date: Wed, 23 Sep 2026 12:58:04 +0200 Subject: [PATCH] fix(cassettes): break the same-second tie with the machine's counter The ordering gate compares created_at, which NIP-01 defines at one-second granularity. A dispense and the publish that follows it land inside one second routinely, so the report was dropped and the operator kept the pre-dispense count until the next heartbeat five minutes later. The machine bumps a counter on every local change to a bay count and carries it in its state document. m015 stores it per row, and the gate consults it only when the stamps are equal, where created_at carries no information at all. Only on equality, deliberately. A machine whose state.db was replaced restarts its counter at zero while its wall clock keeps moving forward; gating on the counter across different stamps would lock that machine out for good. Equal stamps with no counter on either side stay closed, which costs one heartbeat and risks nothing. --- crud.py | 53 +++++++++++++++++++++++++--------- migrations.py | 19 ++++++++++++ models.py | 3 ++ tests/test_cassette_configs.py | 28 ++++++++++++++++++ 4 files changed, 89 insertions(+), 14 deletions(-) diff --git a/crud.py b/crud.py index 8508e9a..c4bffd5 100644 --- a/crud.py +++ b/crud.py @@ -1491,13 +1491,23 @@ def _as_unix(value) -> float | None: return None -def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: +def _row_field(row, name): + """Read one column from a row the data layer may hand back as either a + mapping or an object, depending on the driver.""" + if isinstance(row, dict): + return row.get(name) + return getattr(row, name, None) + + +def _should_apply_state_event( + oldest_state_at, incoming_created_at, oldest_seq=None, incoming_seq=None +) -> bool: """Ordering gate for apply_reported_state. Applies only when the incoming event is strictly newer than the OLDEST state stamp on file for the machine. - Two deliberate choices: + Three deliberate choices: - Compare created_at, not event ids. The old gate asked only whether the incoming id differed from one stored row, which is a one-event memory: @@ -1510,6 +1520,14 @@ def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: some rows advanced and some not. Gating on the oldest means a partial apply is re-applied on the next event rather than being mistaken for a complete one, and the ATM republishes on a heartbeat, so it converges. + - `seq` breaks a same-second tie, and only that. NIP-01 stamps at + one-second granularity, so a dispense and the publish that follows it + can share a stamp and the later report would be dropped. The machine + bumps seq on every local count change, so a higher seq at an equal stamp + is strictly newer. It is consulted ONLY on equality: a machine whose + state.db was replaced restarts its counter at zero, and gating on seq + across different stamps would lock that machine out permanently while + its wall clock kept moving forward. """ oldest = _as_unix(oldest_state_at) if oldest is None: @@ -1517,7 +1535,11 @@ def _should_apply_state_event(oldest_state_at, incoming_created_at) -> bool: incoming = _as_unix(incoming_created_at) if incoming is None: return False - return incoming > oldest + if incoming != oldest: + return incoming > oldest + if incoming_seq is None or oldest_seq is None: + return False + return incoming_seq > oldest_seq async def get_cassette_config(machine_id: str, position: int) -> CassetteConfig | None: @@ -1565,19 +1587,19 @@ async def apply_reported_state( against believed. """ oldest: dict | None = await db.fetchone( - "SELECT state_at FROM spirekeeper.cassette_configs " + "SELECT state_at, state_seq FROM spirekeeper.cassette_configs " "WHERE machine_id = :mid AND state_at IS NOT NULL " - "ORDER BY state_at ASC LIMIT 1", + "ORDER BY state_at ASC, state_seq ASC LIMIT 1", {"mid": machine_id}, ) oldest_state_at = None + oldest_seq = None if oldest is not None: - oldest_state_at = ( - oldest.get("state_at") - if isinstance(oldest, dict) - else getattr(oldest, "state_at", None) - ) - if not _should_apply_state_event(oldest_state_at, event_created_at): + oldest_state_at = _row_field(oldest, "state_at") + oldest_seq = _row_field(oldest, "state_seq") + if not _should_apply_state_event( + oldest_state_at, event_created_at, oldest_seq, payload.seq + ): return False # Drop bays the machine no longer reports, before writing the rest. A crash @@ -1601,9 +1623,10 @@ async def apply_reported_state( INSERT INTO spirekeeper.cassette_configs (machine_id, position, denomination, count, updated_at, updated_by, state_denomination, state_count, state_at, - state_event_id) + state_event_id, state_seq) VALUES (:mid, :pos, :denom, :count, :now, :by, - :state_denom, :state_count, :state_at, :event_id) + :state_denom, :state_count, :state_at, :event_id, + :state_seq) ON CONFLICT (machine_id, position) DO UPDATE SET denomination = excluded.denomination, count = excluded.count, @@ -1612,7 +1635,8 @@ async def apply_reported_state( state_denomination = excluded.state_denomination, state_count = excluded.state_count, state_at = excluded.state_at, - state_event_id = excluded.state_event_id + state_event_id = excluded.state_event_id, + state_seq = excluded.state_seq """, { "mid": machine_id, @@ -1625,6 +1649,7 @@ async def apply_reported_state( "state_count": row.count, "state_at": event_created_at, "event_id": event_id, + "state_seq": payload.seq, }, ) return True diff --git a/migrations.py b/migrations.py index 8f5ddcf..f96fe53 100644 --- a/migrations.py +++ b/migrations.py @@ -922,3 +922,22 @@ async def m014_add_counts_uncertain_since(db): "ALTER TABLE spirekeeper.dca_machines " "ADD COLUMN counts_uncertain_since TIMESTAMP" ) + + +async def m015_add_cassette_state_seq(db): + """Break the same-second tie in the state-event ordering gate. + + The gate compares `created_at`, which NIP-01 defines at one-second + granularity — so two reports from the same second are indistinguishable to + it, and the later one is dropped. A dispense and the publish that follows + it land inside one second routinely. + + The machine bumps `seq` on every local change to a bay count, whatever + caused it, and carries it in the state document. Stored per row alongside + `state_at` and consulted ONLY when the stamps are equal, so a machine whose + state.db was replaced — seq back to zero, wall clock still moving forward — + is not locked out by its own counter. + """ + await db.execute( + "ALTER TABLE spirekeeper.cassette_configs ADD COLUMN state_seq INTEGER" + ) diff --git a/models.py b/models.py index 1a0658c..dc6f533 100644 --- a/models.py +++ b/models.py @@ -677,6 +677,9 @@ class CassetteConfig(BaseModel): state_count: int | None state_at: datetime | None state_event_id: str | None + # The machine's own counter, bumped on every local count change. Breaks a + # same-second tie in the ordering gate, where created_at cannot. + state_seq: int | None = None class CassettePayloadRow(BaseModel): diff --git a/tests/test_cassette_configs.py b/tests/test_cassette_configs.py index b6adf2d..bf597d9 100644 --- a/tests/test_cassette_configs.py +++ b/tests/test_cassette_configs.py @@ -208,3 +208,31 @@ class TestShouldApplyStateEvent: """Fail closed: an unparseable incoming stamp must not overwrite state that is known-good.""" assert _should_apply_state_event(NOW, "nonsense") is False + + def test_seq_breaks_a_same_second_tie(self): + """NIP-01 stamps at one-second granularity, so a dispense and the + publish that follows it share a stamp routinely. Without the counter + the later report is dropped and the operator keeps the pre-dispense + count until the next heartbeat.""" + assert _should_apply_state_event(NOW, NOW, 4, 5) is True + assert _should_apply_state_event(NOW, NOW, 5, 5) is False + assert _should_apply_state_event(NOW, NOW, 5, 4) is False + + def test_seq_is_ignored_when_the_stamps_differ(self): + """A machine whose state.db was replaced restarts its counter at zero + while its wall clock keeps moving forward. Gating on the counter + across different stamps would lock that machine out for good.""" + assert ( + _should_apply_state_event(NOW, NOW + timedelta(seconds=1), 900, 0) is True + ) + assert ( + _should_apply_state_event(NOW, NOW - timedelta(seconds=1), 0, 900) is False + ) + + def test_same_second_without_a_counter_stays_closed(self): + """An older machine sends no seq at all. Equal stamps then carry no + evidence the report is newer, and fail-closed is the safe read: the + machine republishes on its heartbeat a second later.""" + assert _should_apply_state_event(NOW, NOW, None, None) is False + assert _should_apply_state_event(NOW, NOW, None, 5) is False + assert _should_apply_state_event(NOW, NOW, 5, None) is False