diff --git a/temporalio/streams/_policy.py b/temporalio/streams/_policy.py index 81949e4fc..b27c839e5 100644 --- a/temporalio/streams/_policy.py +++ b/temporalio/streams/_policy.py @@ -12,6 +12,7 @@ from __future__ import annotations +from collections.abc import Callable from typing import Any from temporalio.streams._record import ( @@ -27,9 +28,10 @@ class AttemptTracker: """Watches producer attempts on one subscription.""" - def __init__(self) -> None: - """Start with no producer seen.""" + def __init__(self, warn: Callable[[str], None] | None = None) -> None: + """Start with no producer seen, saying anything odd through ``warn``.""" self._attempts: dict[str, int] = {} + self._warn = warn def note( self, producer_id: str, attempt: int, *, topic: str, previous: Cursor @@ -45,10 +47,21 @@ def note( A producer that declares no attempt supersedes nothing, because there is no generation to compare. That is the same answer as an unnumbered record: the interface reports what it was told and invents nothing. + + An attempt that goes backwards supersedes nothing either, and is said + rather than passed off as ordinary data: attempts only ever rise on one + producer, so a lower one means a store reordered two generations, and a + consumer reading it as the current answer would render a stale one. """ if not producer_id or attempt <= 0: return None seen = self._attempts.get(producer_id, 0) + if attempt < seen and self._warn is not None: + self._warn( + f"stream record on {topic!r} at {previous} is from attempt {attempt} of " + f"producer {producer_id!r}, behind attempt {seen}, which this reader has " + "already delivered: the store handed back two generations out of order" + ) if attempt <= seen: return None self._attempts[producer_id] = attempt diff --git a/temporalio/streams/_wire.py b/temporalio/streams/_wire.py index 4913cc308..7d05dbe5d 100644 --- a/temporalio/streams/_wire.py +++ b/temporalio/streams/_wire.py @@ -118,7 +118,7 @@ def __init__( self._result_type = result_type self._previous = after self._warn = warn - self._attempts = AttemptTracker() + self._attempts = AttemptTracker(warn) def decode(self, cursor: Cursor, wire: WireRecord) -> list[StreamRecord[Any]]: """The records to yield for one stored record, in order.""" diff --git a/tests/streams/test_streams_internals.py b/tests/streams/test_streams_internals.py index 8e3d09feb..52020fb41 100644 --- a/tests/streams/test_streams_internals.py +++ b/tests/streams/test_streams_internals.py @@ -103,6 +103,27 @@ def test_supersession_is_synthesized_from_observations(): assert attempts.note("model", 2, topic="t", previous=Cursor("memory:1")) is None +def test_an_attempt_that_goes_backwards_is_said_rather_than_passed_off(): + # Attempts only rise on one producer, so a lower one means the store handed + # two generations back out of order. Yielded as data with no signal, a + # consumer renders the stale generation as the current answer. + said: list[str] = [] + attempts = AttemptTracker(said.append) + assert attempts.note("model", 2, topic="t", previous=BEGINNING) is None + assert attempts.note("model", 1, topic="t", previous=Cursor("memory:3")) is None + assert len(said) == 1 + assert "attempt 1" in said[0] and "behind attempt 2" in said[0] + assert "model" in said[0] + + +def test_a_repeat_of_the_current_attempt_is_not_worth_saying(): + said: list[str] = [] + attempts = AttemptTracker(said.append) + attempts.note("model", 1, topic="t", previous=BEGINNING) + assert attempts.note("model", 1, topic="t", previous=Cursor("memory:1")) is None + assert said == [] + + def test_topic_keys_cannot_collide(): # A colon in a workflow id must not make two addresses one key. assert _ids.topic_key("a:b", "c") != _ids.topic_key("a", "b:c")