diff --git a/temporalio/contrib/external_workflow_streams/_annotation.py b/temporalio/contrib/external_workflow_streams/_annotation.py index d5c55bb53..a55cfb0f8 100644 --- a/temporalio/contrib/external_workflow_streams/_annotation.py +++ b/temporalio/contrib/external_workflow_streams/_annotation.py @@ -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 diff --git a/temporalio/contrib/external_workflow_streams/_api.py b/temporalio/contrib/external_workflow_streams/_api.py index 2ae02d391..c89986954 100644 --- a/temporalio/contrib/external_workflow_streams/_api.py +++ b/temporalio/contrib/external_workflow_streams/_api.py @@ -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" @@ -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) @@ -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 diff --git a/temporalio/contrib/external_workflow_streams/_errors.py b/temporalio/contrib/external_workflow_streams/_errors.py index c9450d235..6025ca3f7 100644 --- a/temporalio/contrib/external_workflow_streams/_errors.py +++ b/temporalio/contrib/external_workflow_streams/_errors.py @@ -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). @@ -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). diff --git a/temporalio/contrib/external_workflow_streams/_manager.py b/temporalio/contrib/external_workflow_streams/_manager.py index 86b521e8d..2323b42e1 100644 --- a/temporalio/contrib/external_workflow_streams/_manager.py +++ b/temporalio/contrib/external_workflow_streams/_manager.py @@ -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 @@ -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. @@ -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 @@ -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 @@ -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`. @@ -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: @@ -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 @@ -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 @@ -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. @@ -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 @@ -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. """ @@ -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. """ @@ -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. @@ -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. @@ -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 @@ -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 @@ -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 diff --git a/temporalio/contrib/external_workflow_streams/_output_backend.py b/temporalio/contrib/external_workflow_streams/_output_backend.py index ec4b1f58e..bcdb154d3 100644 --- a/temporalio/contrib/external_workflow_streams/_output_backend.py +++ b/temporalio/contrib/external_workflow_streams/_output_backend.py @@ -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. """ diff --git a/temporalio/contrib/external_workflow_streams/_producer.py b/temporalio/contrib/external_workflow_streams/_producer.py index 49ce2876f..18fb0c29b 100644 --- a/temporalio/contrib/external_workflow_streams/_producer.py +++ b/temporalio/contrib/external_workflow_streams/_producer.py @@ -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. @@ -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 @@ -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 @@ -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 diff --git a/temporalio/contrib/external_workflow_streams/_redis.py b/temporalio/contrib/external_workflow_streams/_redis.py index dd7cfceb1..c4d69b1cf 100644 --- a/temporalio/contrib/external_workflow_streams/_redis.py +++ b/temporalio/contrib/external_workflow_streams/_redis.py @@ -71,7 +71,7 @@ #: Redis' beginning-of-stream sentinel for `XREAD`. Not an offset: no record #: ever has this id, which is why `BEGINNING` is a distinct cursor form rather -#: than `AFTER(Offset("0-0"))`. +#: than ``AFTER(Offset("0-0"))``. _BEGINNING_SENTINEL: Final = "0-0" _OUTPUT_STAGE_FIELD: Final = "__tes_output_stage" @@ -99,7 +99,7 @@ def _content_hash(record: StreamRecord) -> str: #: Append-if-new-or-identical, atomically. #: -#: Split across two commands without a script, a crash between `XADD` and the +#: Split across two commands without a script, a crash between ``XADD`` and the #: idempotency write would leave a record no retry could recognise as its own, #: so the retry would append a duplicate. _APPEND_LUA: Final = """ @@ -201,7 +201,7 @@ class RedisStreamBackend(StreamBackend, OutputStreamBackend): """Redis Streams as an external workflow stream provider.""" guarantees_immutability: ClassVar[bool | None] = True - """`XADD` entries cannot be rewritten in place -- only deleted or trimmed.""" + """``XADD`` entries cannot be rewritten in place -- only deleted or trimmed.""" provider_id = "redis-streams" provider_format_version = 1 @@ -738,14 +738,14 @@ async def delete_for_test(self, key: StreamKey, offset: Offset) -> None: await self._client.xdel(self.stream_key(key), offset.serialize()) -#: What Redis' glob matcher treats as more than itself. `]` and `^` are special +#: What Redis' glob matcher treats as more than itself. ``]`` and ``^`` are special #: only inside a class, but escaping them too costs nothing and keeps the rule #: one line long. _GLOB_METACHARACTERS: Final = frozenset("*?[]\\") def _as_glob_literal(text: str) -> str: - """A literal string, made safe to embed in a `SCAN MATCH` pattern.""" + """A literal string, made safe to embed in a ``SCAN MATCH`` pattern.""" return "".join( f"\\{character}" if character in _GLOB_METACHARACTERS else character for character in text @@ -756,25 +756,25 @@ def _escaped(key: StreamKey) -> str: r"""One stream identity as a single, unambiguous key component. Percent-encoded per field, then joined -- **not** joined raw. A Workflow ID - and a stream name are user-chosen strings in which `:` is an ordinary + and a stream name are user-chosen strings in which ``:`` is an ordinary character, so joining the raw fields is not injective: ("ns", "wf", r1, f"{r2}:tokens") and ("ns", f"wf:{r1}", r2, "tokens") - both render as `ns:wf:r1:r2:tokens`. Two unrelated Workflows would then share + both render as ``ns:wf:r1:r2:tokens``. Two unrelated Workflows would then share one stream, one idempotency hash, one park intent and one claim -- delivering each other's records, and each concluding the other's claim had already taken - its wake. Same reasoning as `_wake.py`'s length-prefixed request-ID material, + its wake. Same reasoning as ``_wake.py``'s length-prefixed request-ID material, applied to a key rather than to a digest. Percent-encoding rather than length prefixes because a key is read by humans: - an ordinary identity still renders verbatim in `redis-cli`, and only a field + an ordinary identity still renders verbatim in ``redis-cli``, and only a field that actually contains a delimiter pays for it. It buys one property the - length prefix does not -- the encoded form contains no `:`, `*`, `?`, `[` or - `\`, so the derived `:idem`, `:park:` and `:claim:` suffixes stay - unambiguous and `parked_wait_ids`' pattern cannot be widened by a stream name. + length prefix does not -- the encoded form contains no ``:``, ``*``, ``?``, ``[`` or + ``\``, so the derived ``:idem``, ``:park:`` and ``:claim:`` suffixes stay + unambiguous and ``parked_wait_ids``' pattern cannot be widened by a stream name. - Not reversible in practice, and not meant to be: `key_prefix` is + Not reversible in practice, and not meant to be: ``key_prefix`` is operator-supplied and unescaped, so only the identity half round-trips. """ components: tuple[str, ...] = ( diff --git a/temporalio/contrib/external_workflow_streams/_runtime.py b/temporalio/contrib/external_workflow_streams/_runtime.py index d3e091445..739e8111b 100644 --- a/temporalio/contrib/external_workflow_streams/_runtime.py +++ b/temporalio/contrib/external_workflow_streams/_runtime.py @@ -276,7 +276,7 @@ def __init__( #: the Worker that built this runtime, so `codec_for` hands Workflow #: code a converter carrying the same context every other payload in the #: activation was converted with. Bound out there rather than in here - #: because `with_context` runs user code, and this object lives on the + #: because ``with_context`` runs user code, and this object lives on the #: far side of the sandbox boundary. self._data_converter = data_converter self._default_idle_timeout = default_idle_timeout @@ -330,7 +330,7 @@ def __init__( #: Set when a subscription is registered or a record delivered, so an #: activation that changed nothing at all emits nothing. self._observed_this_activation = False - #: `wait_id -> Future`, awaited by Workflow code and resolved by the + #: ``wait_id -> Future``, awaited by Workflow code and resolved by the #: readiness activation. It lives here rather than on either side alone #: because the two halves are in different modules and a second map #: would mean the side that resolves is never the side that registered. @@ -341,7 +341,7 @@ def __init__( #: only for the length of one replay job, which is what makes a #: registration made during it checkable against what was recorded. self._replay_bindings: dict[int, StreamBinding] | None = None - #: `wait_id -> the converter a *recorded* wait's records convert with`. + #: ``wait_id -> the converter a *recorded* wait's records convert with``. #: Installed by the Worker before a replay job reaches this thread; see #: :meth:`install_replay_converters` for why it exists and why it is not #: torn down when the replay ends. @@ -409,7 +409,7 @@ def begin_activation(self, history_floor_event_id: int | None = None) -> None: drain a full batch, consume one record and block elsewhere on every activation in turn, so n subscriptions arrive at an activation holding roughly n times the cap between them and hand all of it over in one - `activate()` call. Starting the count at the carry-over makes what an + ``activate()`` call. Starting the count at the carry-over makes what an activation may hand over -- carried-over plus newly delivered -- exactly the cap, whatever the schedule. """ @@ -1308,7 +1308,7 @@ def _verify_binding( Only the **stream name** is compared, not the whole key. The other three components -- namespace, Workflow id, first execution Run id -- are the Run's identity rather than anything the code chose, and a replay harness - legitimately supplies its own: `Replayer` runs under `ReplayNamespace`, + legitimately supplies its own: `Replayer` runs under ``ReplayNamespace``, so comparing the full key would report every replayed history as nondeterministic. The key is still *recorded* whole, because replay has to read the ranges it names. diff --git a/temporalio/worker/_workflow.py b/temporalio/worker/_workflow.py index b22b779f3..42fb07ff8 100644 --- a/temporalio/worker/_workflow.py +++ b/temporalio/worker/_workflow.py @@ -249,7 +249,7 @@ def __init__( #: A sandboxed Workflow's `instance` is a proxy that exposes only the #: `WorkflowInstance` protocol, so reaching through it for the runtime #: silently found nothing -- and the jobs that must never reach - #: `activate()` quietly went to `_apply` instead. + #: `activate()` quietly went to ``_apply`` instead. self._external_stream_runtimes: dict[str, Any] = {} # Stages survive activations within the cached Run because Core may # defer the server completion behind a local activity. A marker can be @@ -1252,7 +1252,7 @@ def _replay_stream_converters( to have re-created the subscription that would otherwise carry it. Bound from the Worker's own converter rather than from the runtime's, - which is already bound to this Run: a second `with_context` over the + which is already bound to this Run: a second ``with_context`` over the first would ask a user's component converter to rebind itself, and nothing in the protocol promises that composes. diff --git a/tests/contrib/external_workflow_streams/conftest.py b/tests/contrib/external_workflow_streams/conftest.py index 52dec2f67..a58504e1e 100644 --- a/tests/contrib/external_workflow_streams/conftest.py +++ b/tests/contrib/external_workflow_streams/conftest.py @@ -17,6 +17,10 @@ import pytest import pytest_asyncio +#: Asking for either of these is what gives a case a running server. A case +#: that asks for neither never observes a clock, so no environment can fail it. +_SERVER_FIXTURES = frozenset({"client", "env"}) + #: The environments whose server the suite starts for itself. None of them #: accepts the subscribe-notification-channel command. _ENVIRONMENTS_WITHOUT_CHANNELS = ("local", "time-skipping", "envconfig") @@ -107,6 +111,24 @@ async def server_channel_support(client: Any) -> Any: return ChannelKind.INDEPENDENT +@pytest.fixture(autouse=True) +def skip_under_time_skipping(request: pytest.FixtureRequest) -> None: + """Hold the server-backed cases to a clock the tests can reason about. + + Those cases measure real-server timing: how long a Workflow Task was held + open, the interval a wake sweep runs on, the deadline a shutdown waits out. + The time-skipping server advances the clock whenever workers go idle, which + removes exactly the quantities being measured, so the failures it produces + say nothing about the feature. Which of them fail drifts run to run, so the + server-backed cases are held as a group rather than by name. Everything + else here is offline and keeps running on both environments. + """ + if _SERVER_FIXTURES.isdisjoint(request.fixturenames): + return + if request.getfixturevalue("env").supports_time_skipping: + pytest.skip("this case measures real-server timing; see conftest") + + DEFAULT_REDIS_URL = "redis://127.0.0.1:6379" #: Every key this suite creates starts with this, so a leaked key is