Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 15 additions & 2 deletions temporalio/streams/_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@

from __future__ import annotations

from collections.abc import Callable
from typing import Any

from temporalio.streams._record import (
Expand All @@ -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
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion temporalio/streams/_wire.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
21 changes: 21 additions & 0 deletions tests/streams/test_streams_internals.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading