Skip to content

Keep Workflow Streams publishing after transient errors - #1811

Open
1fanwang wants to merge 5 commits into
temporalio:mainfrom
1fanwang:1fannnw/retry-background-stream-flush
Open

1fanwang wants to merge 5 commits into
temporalio:mainfrom
1fanwang:1fannnw/retry-background-stream-flush

Conversation

@1fanwang

@1fanwang 1fanwang commented Sep 3, 2026 •

Copy link
Copy Markdown
Contributor

What was changed

A Workflow Streams publisher now keeps delivering after one transient signal error. Before, the background flusher stopped, and items arrived only on an explicit flush or context exit, which raised the old error.

The flusher retries statuses SDK Core treats as retryable; other errors and an expired retry window still propagate. Pending items keep their order.

Why?

Long-lived publishers went silent after one UNAVAILABLE.

Checklist

  1. Closes: none.

  2. How was this tested:

The probe fails the first signal with UNAVAILABLE and publishes two items a second apart.

Reproducer source: probe.py
import asyncio
from datetime import timedelta
from unittest.mock import patch

from temporalio.contrib.workflow_streams import WorkflowStreamClient
from temporalio.service import RPCError, RPCStatusCode
from temporalio.testing import WorkflowEnvironment

from tests.contrib.workflow_streams.test_workflow_streams import BasicWorkflowStreamWorkflow
from tests.helpers import new_worker


async def main(wait: float) -> None:
    async with await WorkflowEnvironment.start_local() as env:
        async with new_worker(env.client, BasicWorkflowStreamWorkflow) as worker:
            handle = await env.client.start_workflow(
                BasicWorkflowStreamWorkflow.run, id="probe", task_queue=worker.task_queue
            )
            real_signal, sent = handle.signal, []

            async def flaky_signal(*args, **kwargs):
                if not sent:
                    sent.append("failed")
                    raise RPCError("UNAVAILABLE", RPCStatusCode.UNAVAILABLE, b"")
                await real_signal(*args, **kwargs)
                sent.append("delivered")

            stream = WorkflowStreamClient(handle, batch_interval=timedelta(milliseconds=10))
            with patch.object(handle, "signal", side_effect=flaky_signal):
                async with stream:
                    for item in (b"first", b"second"):
                        stream.topic("events", type=bytes).publish(item)
                        await asyncio.sleep(wait)
                    print("signals while open:", sent)


asyncio.run(main(wait=1))

On main:

$ uv run python probe.py
signals while open: ['failed']
temporalio.service.RPCError: UNAVAILABLE

On this branch:

$ uv run python probe.py
signals while open: ['failed', 'delivered', 'delivered']
  1. Any docs updates needed?

No.

@1fanwang
1fanwang requested review from a team as code owners September 3, 2026 06:13
@tconley1428 tconley1428 added the ai-sdk Related to AI integrations label Sep 3, 2026
@brianstrauch

Copy link
Copy Markdown
Member

Two findings from my review:

  1. Please avoid classifying every signal-path exception as retryable.

    _pending is set before WorkflowHandle.signal(), but signal() performs the envelope DataConverter.encode step—including payload codecs and external storage—before issuing the RPC. Consequently, a codec, configuration, external-storage, or interceptor error occurs while _pending is non-None, and the new catch suppresses it and retries until max_retry_duration. With the default settings, this hides the original diagnostic for ten minutes and ultimately replaces it with TimeoutError.

    I reproduced this with a PayloadCodec.encode() implementation that raises: the background flusher remained pending on this PR's head, while the same probe propagated the error with the prior implementation. Could the suppression be restricted to genuinely retryable delivery/RPC failures, with a regression test confirming that codec failures still propagate?

  2. Please add an Unreleased / Fixed changelog entry.

    This fixes user-visible Workflow Streams delivery behavior, and the repository guidance requires a high-level changelog entry for user-facing changes.

Validation on the PR head: the focused regression passed, and the complete Workflow Streams test file passed (43 passed).

@1fanwang 1fanwang changed the title Retry transient Workflow Stream background flushes Keep Workflow Streams publishing after transient errors Sep 3, 2026
@1fanwang

1fanwang commented Sep 3, 2026

Copy link
Copy Markdown
Contributor Author

Done in 7213ce2. Retry scope now matches SDK Core, codec and oversized-message failures propagate, and the changelog is updated.

1fanwang and others added 3 commits September 28, 2026 13:26
@1fanwang
1fanwang force-pushed the 1fannnw/retry-background-stream-flush branch from 12b7db8 to 606a8d3 Compare September 28, 2026 20:28
@brianstrauch

Copy link
Copy Markdown
Member

@1fanwang From CI:

Error: CHANGELOG.md additions must be under the Unreleased section.

@1fanwang

1fanwang commented Oct 1, 2026

Copy link
Copy Markdown
Contributor Author

CHANGELOG.md additions must be under the Unreleased section.

@brianstrauch Done in 453e722.

@brianstrauch
brianstrauch enabled auto-merge (squash) October 1, 2026 22:03

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

ai-sdk Related to AI integrations

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants