fix(cassettes): break the same-second tie with the machine's counter
Some checks failed
ci.yml / fix(cassettes): break the same-second tie with the machine's counter (pull_request) Failing after 0s
Some checks failed
ci.yml / fix(cassettes): break the same-second tie with the machine's counter (pull_request) Failing after 0s
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.
This commit is contained in:
parent
79a4f83293
commit
5f60b3fe31
4 changed files with 89 additions and 14 deletions
53
crud.py
53
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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue