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