Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,7 @@ class AnnotationBudgetExceeded(temporalio.exceptions.ApplicationError):
terminal could not fit an empty annotation, at the point the Workflow makes
it.

An :py:class:`~temporalio.exceptions.ApplicationError` marked
An :py:class:`temporalio.exceptions.ApplicationError` marked
non-retryable, and that is the substance of this class rather than a detail.
A plain exception here fails the *Workflow Task*, and the server retries
Workflow Task failures forever: the encoding that overflowed overflows again
Expand Down
6 changes: 3 additions & 3 deletions temporalio/contrib/external_workflow_streams/_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@
"""

#: Where per-Run subscription state hangs off the Workflow instance. Reserved,
#: and deliberately not in the `__temporal_workflow_stream*` namespace the
#: and deliberately not in the ``__temporal_workflow_stream*`` namespace the
#: shipped contrib feature already owns.
_RUN_STATE_ATTR = "__temporal_external_stream_state"

Expand Down Expand Up @@ -182,7 +182,7 @@ class _RunState:

runtime: ExternalStreamRuntime | None = None
next_wait_id: int = 1
#: `wait_id -> Future`, resolved by the readiness activation handler.
#: ``wait_id -> Future``, resolved by the readiness activation handler.
pending: dict[int, Any] = field(default_factory=dict)


Expand Down Expand Up @@ -538,7 +538,7 @@ class ExternalStreamSubscription(Generic[AnyType]):
**One consumer.** The cursor, the readiness future, and the blocked flag are
all the subscription's rather than an iterator's, so two coroutines waiting on
it at once is refused with
:class:`~temporalio.contrib.external_workflow_streams._errors.ConcurrentStreamConsumerError`
:class:`temporalio.contrib.external_workflow_streams.ConcurrentStreamConsumerError`
rather than served -- see :meth:`_refuse_a_second_waiter` for what sharing
them would do. Two consumers of the same *stream* is a supported shape and
the way to ask for it is a second ``subscribe()``: delivery is a broadcast
Expand Down
32 changes: 17 additions & 15 deletions temporalio/contrib/external_workflow_streams/_errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,20 +11,22 @@
key on the metrics here rather than on Workflow Task failure counts, which a
transient backend outage also increments.

============================ ========================== ====================
Condition Error type Operator response
============================ ========================== ====================
Backend unreachable/erroring :class:`StreamStorageError` None -- clears when
the backend recovers
Recorded offset missing, :class:`StreamIntegrityError` Repair or restore
expired, reordered, or the backend, or
miscounted terminate the Run
Bytes intact but undecodable :class:`StreamDecodeError` Align the consumer's
converter with the
producer's
Annotation does not match ordinary nondeterminism Fix or version the
the subscriptions made Workflow code
============================ ========================== ====================
::

============================ ========================== ====================
Condition Error type Operator response
============================ ========================== ====================
Backend unreachable/erroring ``StreamStorageError`` None -- clears when
the backend recovers
Recorded offset missing, ``StreamIntegrityError`` Repair or restore
expired, reordered, or the backend, or
miscounted terminate the Run
Bytes intact but undecodable ``StreamDecodeError`` Align the consumer's
converter with the
producer's
Annotation does not match ordinary nondeterminism Fix or version the
the subscriptions made Workflow code
============================ ========================== ====================

There is deliberately **no** ``workflow_failure_exception_types`` registration
here: integrity loss *blocks* a Workflow rather than terminating it (ADR-014).
Expand Down Expand Up @@ -83,7 +85,7 @@ class StreamIntegrityError(StreamError):
class StreamDecodeError(StreamError):
"""A record was present and intact but could not be decoded.

Separated from :class:`StreamIntegrityError` because it is a configuration
Separated from ``StreamIntegrityError`` because it is a configuration
error on the *consumer* -- a DataConverter or codec that does not match the
producer's -- and reporting it as integrity loss sends an operator to
restore a backend that was never damaged (ADR-015).
Expand Down
50 changes: 25 additions & 25 deletions temporalio/contrib/external_workflow_streams/_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,13 +28,13 @@
Three cursors, and conflating them is how a speculative read becomes a durable
claim:

=================== =================== =========== ==================
Cursor Advances on Owner Survives eviction
=================== =================== =========== ==================
``committed`` marker commit only the marker yes
``delivery`` hand-off to workflow the instance no
``prefetch`` buffering the manager no
=================== =================== =========== ==================
============= ==================== ============ =================
Cursor Advances on Owner Survives eviction
============= ==================== ============ =================
``committed`` marker commit only the marker yes
``delivery`` hand-off to workflow the instance no
``prefetch`` buffering the manager no
============= ==================== ============ =================

``prefetch`` is speculative: reading a record is not consuming it, and consuming
it is not committing it. The manager may only move it *backwards* to
Expand Down Expand Up @@ -227,7 +227,7 @@ class RunStatus:
out is only the first thing the activation then has to do.

Holds an activation takes *itself* are deliberately not bounded. Those are the
same backend exposure the park handshake already has -- `install_park_intent`
same backend exposure the park handshake already has -- ``install_park_intent``
can hang exactly as a removal can -- so a bound there would move the wait rather
than remove it.

Expand Down Expand Up @@ -316,7 +316,7 @@ class Subscription:
_has_room: asyncio.Event = field(default_factory=asyncio.Event, repr=False)
_watcher: asyncio.Task[None] | None = field(default=None, repr=False)
_cancelled: bool = field(default=False, repr=False)
#: Bumped whenever every speculative read is discarded. A `read_after`
#: Bumped whenever every speculative read is discarded. A ``read_after``
#: already in flight started from the cursor being discarded, so its records
#: are speculative too; the watcher compares this across the await and drops
#: them rather than appending them on top of the reset -- which would put the
Expand Down Expand Up @@ -747,7 +747,7 @@ def __init__(
self._pending_wake_runs: dict[str, _PendingWake] = {}
self._activation_sequences: dict[str, int] = {}
#: Wakes the shutdown sweep could not get acknowledged. Reported through
#: `external_stream_shutdown_wake_failed`; kept here so a test can tell
#: ``external_stream_shutdown_wake_failed``; kept here so a test can tell
#: "no wake was needed" from "a wake was needed and lost".
self.shutdown_wake_failures = 0
self._buffer_size = buffer_size
Expand Down Expand Up @@ -780,7 +780,7 @@ def __init__(
#: Strong references to the in-flight registration-time reconciliations.
self._reconciliations: set[asyncio.Task[None]] = set()
#: Removals this manager decided on and did not get confirmed, per Run
#: and then per `(stream key, wait_id)`. The durable half of the intent
#: and then per ``(stream key, wait_id)``. The durable half of the intent
#: invariant: a removal that failed is a claim about the *backend*, and
#: recording it anywhere that a close or an eviction takes away is the
#: same as not recording it at all. Drained under `_park_lock`.
Expand Down Expand Up @@ -858,7 +858,7 @@ def _cancel_watcher(self, subscription: Subscription) -> None:
def _start_watcher(self, subscription: Subscription) -> None:
"""Starts a watcher, and reconciles the park state it inherited.

Both on the manager's own loop, because both are `create_task` calls and
Both on the manager's own loop, because both are ``create_task`` calls and
`register` runs on the Workflow executor thread.
"""
if subscription._cancelled or subscription._watcher is not None:
Expand Down Expand Up @@ -886,10 +886,10 @@ async def _reconcile_inherited_park(self, subscription: Subscription) -> None:
Registration is where a Worker learns such an intent exists, and it is
also the moment its status is unambiguous. A subscription is registered
by user Workflow code running, and no user code runs inside a park
(`wft-lifecycle.md`), so an intent found here belongs to a park that is
(``wft-lifecycle.md``), so an intent found here belongs to a park that is
over: the Core that confirmed it has either moved on or gone with the
Worker that held it. What leaving it costs is the invariant's whole
point -- `current_park_generation` keeps answering a generation Core has
point -- ``current_park_generation`` keeps answering a generation Core has
discarded, every producer wake names that generation and Core discards
it as stale, and because a parked wake's request ID ignores sender
identity the second such wake is byte-identical to the first and the
Expand Down Expand Up @@ -1379,7 +1379,7 @@ def _reannounce_after_cleanup(

A stale intent does not only leak. While it was installed every wake for
this stream named the generation behind it -- the producer reads
`current_park_generation`, and so does the Worker's own sender -- and
``current_park_generation``, and so does the Worker's own sender -- and
Core discards a non-zero generation that is not the park it is holding.
A record that arrived in that window was therefore announced to nobody,
and *removing the intent does not announce it*: `_report_ready` counted
Expand All @@ -1390,12 +1390,12 @@ def _reannounce_after_cleanup(
create a Workflow Task.

Both outcomes that clear the key announce, and that is why the provider
contract does not answer with a Boolean. `ABSENT` is not "someone else's
contract does not answer with a Boolean. ``ABSENT`` is not "someone else's
intent is in the way", it is "the intent this entry named is gone" --
which is what a retry sees after the reply to a delete that in fact
succeeded was lost, and what one cleanup owner sees after another
finished the job. Reading that as a mismatch would leave the record the
intent silenced silent for good. `MISMATCH` announces nothing: an intent
intent silenced silent for good. ``MISMATCH`` announces nothing: an intent
is still installed there, the suppression it causes has not ended, and it
belongs to a park this entry knows nothing about.

Expand Down Expand Up @@ -1428,7 +1428,7 @@ def _park_lock(self, run_id: str) -> asyncio.Lock:
"""Serializes one Run's park-intent work on the manager's loop.

The install/recheck handshake, the resolve, and the reconciliation above
all read-then-write the same `(stream key, wait_id)` objects, and the
all read-then-write the same ``(stream key, wait_id)`` objects, and the
reconciliation is scheduled from another thread, so their interleaving
is not otherwise constrained. Without this, a reconciliation that
overlapped a confirming park could remove the intent that park had just
Expand Down Expand Up @@ -1493,7 +1493,7 @@ def rearm_ready(self, run_id: str) -> None:
Workflow Task whose data had already arrived.

Called from the Workflow thread at activation completion, so the work is
hopped onto the manager's loop: `create_task` is not thread-safe, and a
hopped onto the manager's loop: ``create_task`` is not thread-safe, and a
task created from the Workflow executor thread is silently never
scheduled -- indistinguishable from a stream that never delivers.
"""
Expand Down Expand Up @@ -1803,7 +1803,7 @@ async def _report_ready(self, subscription: Subscription) -> None:
leave this method. The watcher calls it in its loop, so an exception
escaping ends that watcher for good -- the subscription stays registered,
its buffer keeps its records, and nothing ever announces them again.
`_rearm_ready` also launches it with `create_task`, where an exception
`_rearm_ready` also launches it with ``create_task``, where an exception
becomes a task result nobody retrieves and the failure is not even
logged.
"""
Expand Down Expand Up @@ -1917,7 +1917,7 @@ async def _retry_stale(self, subscription: Subscription) -> str:
Returns the result the retries ended on, so the caller can act on it. A
Boolean would answer only "was it announced", and the four
non-``Accepted`` answers are not interchangeable: they differ in what
happens to the watcher, and `RunNotFound` in particular requires the
happens to the watcher, and ``RunNotFound`` in particular requires the
watcher to be torn down. Discarding them left that teardown unreachable
from here.

Expand All @@ -1927,7 +1927,7 @@ async def _retry_stale(self, subscription: Subscription) -> str:
Workflow consuming happily, and a wake owed at the end of that costs one
empty Workflow Task rather than a silent stall.

`RunNotFound` ends the retries rather than using them up. It is the one
``RunNotFound`` ends the retries rather than using them up. It is the one
answer that cannot change back: the Run is gone from this Worker, so a
further report can only be answered the same way, and each attempt costs a
delay before the wake this record still needs.
Expand Down Expand Up @@ -1967,7 +1967,7 @@ async def prepare_park(
is waiting on and aborts a legitimate park, which then runs again on the
next idle timeout and aborts again; and an intent installed for a wait
outside the set is an intent with no park behind it, which is exactly
what `backend-contract.md` forbids leaving in a backend.
what ``backend-contract.md`` forbids leaving in a backend.

The order is what closes the append/park race: a producer appends its
record *before* it observes the park generation, so an append is either
Expand Down Expand Up @@ -2328,7 +2328,7 @@ async def cancel(self, run_id: str, wait_id: int) -> None:
behind for good.

And it is not confined to the closed wait. A stale intent keeps
`parked_wait_ids` non-empty, which suppresses the unparked-wake fallback
``parked_wait_ids`` non-empty, which suppresses the unparked-wake fallback
for the **whole stream**: with no live wait parked, the producer sends
only the dead generation, Core discards it as stale, and dedup silences
every later publish -- so live waits across the Continue-As-New chain
Expand Down Expand Up @@ -2684,7 +2684,7 @@ async def _sweep_wake(self, subscription: Subscription) -> None:
Counting the failure is deliberately **not** done here for the
cancellation case. The grace period expiring cancels this coroutine
wherever it is, `_send_owed_wake` re-raises `CancelledError` by design,
and no `except` here could both record the failure and leave the
and no ``except`` here could both record the failure and leave the
cancellation intact for the subscriptions after this one -- which are not
reached either. The subscription stays in the unaccounted set instead and
:meth:`_account_unswept` counts it, which covers being cancelled and
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -286,7 +286,7 @@ async def append_output(
input record's ``(session_id, sequence)`` and byte identity exactly as
:meth:`StreamBackend.append`: an identical retry returns the original
output record and offset; different bytes under the same key raise
:class:`~temporalio.contrib.external_workflow_streams._backend.AppendConflictError`.
:class:`temporalio.contrib.external_workflow_streams.AppendConflictError`.
It may not pass an already-positioned pending stage in client reads.
"""

Expand Down
12 changes: 6 additions & 6 deletions temporalio/contrib/external_workflow_streams/_producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -118,17 +118,17 @@ def __init__(
whatever ended the attempt (ADR-036).
"""
self.restart = restart
"""Whether the caller must call :meth:`.wake` again instead of retrying.
"""Whether the caller must call ``wake()`` again instead of retrying.

The wake is three steps -- observe the parked set, claim the generation,
Signal -- and only the third produces the requests
:meth:`ProducerTopicHandle.retry_wake` re-sends. A failure in the first
``ProducerTopicHandle.retry_wake()`` re-sends. A failure in the first
two leaves nothing to re-send, so ``pending`` is empty and retrying it
would silently do nothing at all: the record would stay durable and
unannounced while the caller believed it had recovered.

``True`` therefore says "no wake was composed; compose one". Calling
:meth:`ProducerTopicHandle.wake` again is safe and is the whole recovery
``ProducerTopicHandle.wake()`` again is safe and is the whole recovery
-- it re-observes the parked set, and a parked wake's request ID is
derived from the generation rather than from the sender, so a wake some
other producer already sent deduplicates against it.
Expand Down Expand Up @@ -198,7 +198,7 @@ def __init__(
"""The stream the append was for. Where it must be settled.

Carried because a record does not name its own stream and the backend's
idempotency scope does: `(session_id, sequence)` is unused on every
idempotency scope does: ``(session_id, sequence)`` is unused on every
*other* stream, so the same record handed to another topic's
``resolve_append`` would append a second copy there rather than
deduplicate. The recovery refuses that, and this is what a caller
Expand Down Expand Up @@ -455,7 +455,7 @@ def __init__(
#: Per stream, the appends that never reported an outcome.
#:
#: Kept because the operation *is* the recovery: only these exact bytes
#: under these exact `(session_id, sequence)` pairs re-append as a no-op
#: under these exact ``(session_id, sequence)`` pairs re-append as a no-op
#: if the first attempt landed, and only what the interrupted call owed
#: says whether settling it still has a wake to send. A list rather than
#: a single slot because concurrent publishes to one stream are supported
Expand Down Expand Up @@ -872,7 +872,7 @@ def _outstanding(self, record: StreamRecord) -> _UnresolvedAppend:

The lookup **is** the safety check, and it is three checks at once. It
binds the recovery to the stream, because a record does not name its own
stream and `(session_id, sequence)` is unused on every other one -- so a
stream and ``(session_id, sequence)`` is unused on every other one -- so a
record settled against the wrong topic appends a second copy of the value
rather than deduplicating, and leaves the real stream still blocked. It
binds the recovery to the exact bytes, because idempotency is on identity
Expand Down
Loading
Loading