diff --git a/docs/develop/dotnet/index.mdx b/docs/develop/dotnet/index.mdx index a9c62b4d61..1b6587bd99 100644 --- a/docs/develop/dotnet/index.mdx +++ b/docs/develop/dotnet/index.mdx @@ -46,6 +46,7 @@ From there, you can dive deeper into any of the Temporal primitives to start bui - [Timers](/develop/dotnet/workflows/timers) - [Dynamic Workflow](/develop/dotnet/workflows/dynamic-workflow) - [Versioning](/develop/dotnet/workflows/versioning) +- [Workflow Streams](/develop/dotnet/workflows/workflow-streams) ## [Activities](/develop/dotnet/activities) diff --git a/docs/develop/dotnet/workflows/index.mdx b/docs/develop/dotnet/workflows/index.mdx index 2073a3cf0c..2c2b155307 100644 --- a/docs/develop/dotnet/workflows/index.mdx +++ b/docs/develop/dotnet/workflows/index.mdx @@ -30,3 +30,4 @@ import * as Components from '@site/src/components'; - [Timers](/develop/dotnet/workflows/timers) - [Dynamic Workflow](/develop/dotnet/workflows/dynamic-workflow) - [Versioning](/develop/dotnet/workflows/versioning) +- [Workflow Streams](/develop/dotnet/workflows/workflow-streams) diff --git a/docs/develop/dotnet/workflows/workflow-streams.mdx b/docs/develop/dotnet/workflows/workflow-streams.mdx new file mode 100644 index 0000000000..3818f4b754 --- /dev/null +++ b/docs/develop/dotnet/workflows/workflow-streams.mdx @@ -0,0 +1,539 @@ +--- +id: workflow-streams +title: Workflow Streams - .NET SDK +sidebar_label: Workflow Streams +toc_max_heading_level: 4 +tags: + - Workflows + - Messages + - Streaming + - .NET SDK + - Temporal SDKs +description: + Workflow Streams in the Temporal .NET SDK sends events from Workflows and Activities to external subscribers through + durable topics. +--- + +import { ReleaseNoteHeader } from '@site/src/components'; + + + +[Workflow Streams](/workflow-streams) adds a durable event channel to a Workflow, letting outside observers follow its +progress in real time. This page explains how to enable a stream, publish events from Workflows and Activities, subscribe +to a stream, and keep a stream running across long-lived Workflows. + +Workflow Streams is available in the .NET SDK v1.20.0 and later. Install the extension package: + +```bash +dotnet add package Temporalio.Extensions.WorkflowStreams --version 1.20.0 +``` + +## Enable streaming on a Workflow + +Enable streaming by constructing one `WorkflowStream` in a [`[WorkflowInit]`](https://dotnet.temporal.io/api/Temporalio.Workflows.WorkflowInitAttribute.html) +constructor. The constructor runs before message handlers can be dispatched, so the stream can register its publish +Signal, poll Update, and offset Query handlers before the first request arrives. + +```csharp +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +public record OrderInput( + string OrderId, + WorkflowStreamState? StreamState = null); + +[Workflow] +public class OrderWorkflow +{ + private readonly WorkflowStream stream; + + [WorkflowInit] + public OrderWorkflow(OrderInput input) => stream = new(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(OrderInput input) + { + // ... rest of the Workflow + } +} +``` + +Pass `null` on the first run. After Continue-As-New, pass the `WorkflowStreamState` captured from the previous run. See +[Stream from long-running Workflows](#continue-as-new) for the complete rollover sequence. + +Construct exactly one `WorkflowStream` per Workflow. Each construction registers the same protocol handler names, so a +second instance replaces the first instance's handlers. + +## Publish from a Workflow + +Bind a [topic name](/workflow-streams#topics) to an event type with `GetTopic()`, then call `Publish()` on the returned +handle. The generic type gives publish call sites compile-time type checking and tells typed subscribers how to decode +the topic. + +```csharp +public record StatusEvent(string State, int Progress = 0, string? Detail = null); + +[Workflow] +public class OrderWorkflow +{ + private readonly WorkflowStream stream; + private readonly WorkflowStreamTopicHandle status; + + [WorkflowInit] + public OrderWorkflow(OrderInput input) + { + stream = new(input.StreamState); + status = stream.GetTopic("status"); + } + + [WorkflowRun] + public async Task RunAsync(OrderInput input) + { + status.Publish(new("validating", Detail: "checking inventory")); + await ValidateOrderAsync(input.OrderId); + + status.Publish(new("charging", 33, "authorizing payment")); + await ChargePaymentAsync(input.OrderId); + + status.Publish(new("shipping", 66, "dispatching to warehouse")); + await DispatchOrderAsync(input.OrderId); + + status.Publish(new("completed", 100)); + } +} +``` + +`Publish()` runs the Workflow's payload converter and appends the resulting `Payload` directly to the in-memory log. A +Workflow-side publication isn't buffered and doesn't need to be flushed. + +The type binding belongs to the handle, not to the stream protocol. If separate publishers use different types for the +same topic name, publication succeeds and a typed subscriber can fail while decoding the mismatched value. + +## Publish from a client + +Any process with a Temporal Client and the target Workflow ID can publish by constructing a `WorkflowStreamClient`. +This includes HTTP services, Workflow starters, scripts, Activities from other Workflows, and standalone Activities. + +Use `await using` so the client flushes buffered events and releases its publishing resources when the scope exits: + +```csharp +using Temporalio.Client; +using Temporalio.Extensions.WorkflowStreams; + +public static async Task PublishStatusAsync( + ITemporalClient temporalClient, + string workflowId, + CancellationToken cancellationToken) +{ + await using var streamClient = new WorkflowStreamClient( + temporalClient, + workflowId, + new() + { + BatchInterval = TimeSpan.FromMilliseconds(200), + }); + + var status = streamClient.GetTopic("status"); + status.Publish(new("started")); + status.Publish(new("completed", 100), forceFlush: true); + await streamClient.FlushAsync(cancellationToken); +} +``` + +By default, the client sends buffered publications every two seconds. Set `MaxBatchSize` to flush after the buffer +reaches a specific item count. A value of `0`, the default, disables size-based flushing. + +Pass `forceFlush: true` when an event should wake the background flusher without waiting for the next interval. The +`Publish()` call still returns after buffering the event; it doesn't wait for the Signal to reach the Temporal Service. +Use `FlushAsync()` when subsequent work must wait until all earlier publications have been acknowledged. + +When an event originates in an Activity, publish it from the Activity instead of returning it for the Workflow to +forward. `FromActivity()` gets the Temporal Client, parent Workflow ID, and payload converter from the current Activity +context: + +```csharp +using Temporalio.Activities; +using Temporalio.Extensions.WorkflowStreams; + +public record ProgressEvent(string Message); + +public static class OrderActivities +{ + [Activity] + public static async Task ProcessOrderAsync(IEnumerable steps) + { + await using var streamClient = WorkflowStreamClient.FromActivity(); + var progress = streamClient.GetTopic("progress"); + + foreach (var step in steps) + { + progress.Publish(new(step)); + ActivityExecutionContext.Current.Heartbeat(step); + } + } +} +``` + +`FromActivity()` requires an Activity started by a Workflow. For a standalone Activity, pass the target Workflow ID in +the Activity input and construct `WorkflowStreamClient` with `ActivityExecutionContext.Current.TemporalClient`. + +Publishing from the Activity lets subscribers see partial output from failed and retried attempts without adding that +output to the Workflow's durable application state. Each Activity attempt creates a new publisher, so subscribers can +receive partial output from an unsuccessful attempt before receiving the retried attempt's output. See +[How events are delivered](/workflow-streams#how-events-are-delivered). + +`Publish()` applies no backpressure. Client-side publication appends to an in-memory buffer, while subscribers read the +Workflow log on their own schedule. If a publisher produces events faster than batches reach the Temporal Service, its +buffer grows. Apply an application-specific policy before `Publish()` if you need to block, drop, sample, or reject +events when the publisher outruns the network. + +## Subscribe + +Subscribe from any process that has a Temporal Client and the Workflow ID. A typed topic handle returns a reusable +`IAsyncEnumerable>`. Each enumeration keeps its own offset and polling state. + +```csharp +using Temporalio.Client; +using Temporalio.Extensions.WorkflowStreams; + +public static async Task WatchOrderAsync( + ITemporalClient temporalClient, + string workflowId, + CancellationToken cancellationToken) +{ + await using var streamClient = new WorkflowStreamClient(temporalClient, workflowId); + var status = streamClient.GetTopic("status"); + + await foreach (var item in status.SubscribeAsync().WithCancellation(cancellationToken)) + { + Console.WriteLine( + $"[{item.Value.Progress,3}%] {item.Value.State}: {item.Value.Detail}"); + if (item.Value.State == "completed") + { + break; + } + } +} +``` + +Pass the next global offset to `SubscribeAsync(fromOffset)` when reconnecting. The iterator follows Continue-As-New +chains, repolls immediately when a response reaches the approximate 1 MB page limit, and advances to the retained +beginning of the log if its requested offset was truncated. + +Canceling the enumeration cancels its in-flight RPC and throws `OperationCanceledException`. Disposing the owning +`WorkflowStreamClient` ends its active enumerations without an exception. An enumeration also ends cleanly when the +Workflow reaches a terminal state. + +### Subscribe to heterogeneous topics + +A typed subscription decodes every item as one type. To consume topics with different payload types, call the untyped +`SubscribeAsync()` overload with a topic filter. It yields `WorkflowStreamItem` values containing raw Temporal +`Payload`s. Dispatch on `Topic` and decode each payload with the matching type: + +```csharp +using Temporalio.Converters; + +var options = new WorkflowStreamSubscribeOptions +{ + Topics = new[] { "status", "progress" }, +}; +var converter = temporalClient.Options.DataConverter.PayloadConverter; + +await foreach (var item in streamClient.SubscribeAsync(options)) +{ + if (item.Topic == "status") + { + var evt = converter.ToValue(item.Payload); + Console.WriteLine($"[status] {evt.State}: {evt.Detail}"); + } + else if (item.Topic == "progress") + { + var evt = converter.ToValue(item.Payload); + Console.WriteLine($"[progress] {evt.Message}"); + } +} +``` + +An empty topic collection subscribes to every topic. The empty string represents the cross-SDK no-topic value. + +### Close the stream {/* #closing-the-stream */} + +A subscription doesn't know which event your application considers final. Publish an in-band terminal event and have +the subscriber break when it receives that event. In the preceding example, the `completed` status is the terminal +event. + +Keep the Workflow open briefly after publishing the terminal event. If the Workflow completes immediately, an in-flight +poll might close before the subscriber receives the event. You can provide this overlap in either of these ways: + +- **Fixed delay.** Call `Workflow.DelayAsync()` after publishing the terminal event and before returning. + + ```csharp + status.Publish(new("completed", 100)); + await Workflow.DelayAsync(TimeSpan.FromSeconds(30)); + return result; + ``` + +- **Acknowledgment handshake.** Have the subscriber Signal the Workflow after it receives the terminal event. Wait for + that Signal with a timeout so the Workflow still completes when no subscriber is attached. + + ```csharp + private bool subscriberAcknowledged; + + [WorkflowSignal] + public Task AcknowledgeStreamAsync() + { + subscriberAcknowledged = true; + return Task.CompletedTask; + } + + // After publishing the terminal event: + await Workflow.WaitConditionAsync( + () => subscriberAcknowledged, + TimeSpan.FromSeconds(30)); + return result; + ``` + +The subscription ends cleanly for every terminal Workflow status. If the application must distinguish completion from +failure, cancellation, termination, or timeout, describe the Workflow after the enumeration ends and inspect its +status. + +## Stream from long-running Workflows {/* #continue-as-new */} + +Use [Continue-As-New](/workflow-execution/continue-as-new) to keep Event History bounded for a Workflow that runs for +hours or accumulates many events. Subscribers automatically follow the new run, but you must carry the stream state +through the Workflow input so they don't see a gap. + +```csharp +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +public record PipelineInput( + int ItemsProcessed = 0, + WorkflowStreamState? StreamState = null); + +[Workflow] +public class PipelineWorkflow +{ + private readonly WorkflowStream stream; + + [WorkflowInit] + public PipelineWorkflow(PipelineInput input) => stream = new(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(PipelineInput input) + { + var itemsProcessed = input.ItemsProcessed; + var progress = stream.GetTopic("progress"); + + while (true) + { + await ProcessNextItemAsync(); + itemsProcessed++; + progress.Publish(new($"processed {itemsProcessed} items")); + + if (Workflow.ContinueAsNewSuggested) + { + var streamState = await stream.CaptureStateForContinueAsNewAsync(); + throw Workflow.CreateContinueAsNewException( + (PipelineWorkflow workflow) => workflow.RunAsync( + new(itemsProcessed, streamState))); + } + } + } +} +``` + +`CaptureStateForContinueAsNewAsync()` detaches waiting pollers, waits for message handlers to finish, and captures the +retained log and publisher deduplication state. It also disables Workflow-side publication. Create and throw the +Continue-As-New exception immediately after the snapshot. + +The snapshot contains the entire retained log. Use `stream.Truncate(offset)` after all required consumers have advanced +to keep it small. If individual events contain large data, store that data in [External Storage](/external-storage) and +publish a reference instead. + +## Configure the deduplication window + +See [How events are delivered](/workflow-streams#how-events-are-delivered) for publisher and subscriber guarantees. Two +settings bound the deduplication window: + +- `WorkflowStreamOptions.PublisherTtl` controls how long publisher sequence state remains in a Continue-As-New + snapshot. Its default is 15 minutes. A publisher that returns after its entry expires can produce a duplicate. +- `WorkflowStreamClientOptions.MaxRetryDuration` controls how long a client retries an ambiguously delivered batch. Its + default is 10 minutes. Keep this value lower than the publisher TTL. + +```csharp +// Workflow initialization +stream = new( + input.StreamState, + new() { PublisherTtl = TimeSpan.FromMinutes(30) }); + +// External client or Activity +await using var streamClient = new WorkflowStreamClient( + temporalClient, + workflowId, + new() { MaxRetryDuration = TimeSpan.FromMinutes(20) }); +``` + +When the retry duration expires, the client drops the ambiguous batch and reports `FlushTimeoutException` from +`FlushAsync()` or asynchronous disposal. The batch might already be in the Workflow log, or it might be lost. Later +batches continue with a new sequence. + +## Follow best practices + +- **Construct one stream during Workflow initialization.** This registers protocol handlers before the first request + and prevents another stream from replacing them. +- **Dispose publishing clients asynchronously.** Use `await using` or call `DisposeAsync()` so the last buffer is + drained. Publishing after disposal throws `ObjectDisposedException`. +- **Use consistent topic types.** Generic topic handles provide local compile-time checking, but publishers don't share + their bindings. A subscriber reports a conversion error if a publisher sends a different type. +- **Keep payload conversion consistent.** Stream items use the payload converter. Payload codecs run once on the Signal + or Update envelope, not once per item. `FromActivity()` subscriptions aren't compatible with custom converters that + require the serialization context used for deserialization to match publication's Activity context. +- **Subscribe outside the hosting Workflow.** The Workflow owns the durable log but doesn't consume it. This keeps + partial output from retried Activities separate from the Workflow's successful Activity result. + +## Example: Stream LLM output {/* #stream-llm-output */} + +An Activity can publish model output as it arrives while the Workflow waits for the Activity result. The subscriber +renders deltas, clears output when an Activity retry begins, and stops on a terminal event. This pattern works for a +terminal client, desktop UI, or Server-Sent Events endpoint. + + + + +```csharp +using System.Text; +using Temporalio.Activities; +using Temporalio.Extensions.WorkflowStreams; + +public record TextDelta(string Text); +public record RetryEvent(int Attempt); +public record TextComplete(bool Done); + +public static class LlmActivities +{ + [Activity] + public static async Task StreamCompletionAsync(string prompt) + { + var context = ActivityExecutionContext.Current; + await using var streamClient = WorkflowStreamClient.FromActivity( + new() { BatchInterval = TimeSpan.FromMilliseconds(200) }); + + var deltas = streamClient.GetTopic("delta"); + var retries = streamClient.GetTopic("retry"); + var complete = streamClient.GetTopic("complete"); + + if (context.Info.Attempt > 1) + { + retries.Publish(new(context.Info.Attempt), forceFlush: true); + } + + var fullText = new StringBuilder(); + var first = true; + // GenerateDeltasAsync wraps the model call and yields text as it arrives. + // Disable provider retries and let the Activity Retry Policy own retries. + await foreach (var text in GenerateDeltasAsync( + prompt, + context.CancellationToken)) + { + deltas.Publish(new(text), forceFlush: first); + first = false; + fullText.Append(text); + } + + complete.Publish(new(true), forceFlush: true); + return fullText.ToString(); + } +} +``` + + + + +```csharp +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +public record LlmInput( + string Prompt, + WorkflowStreamState? StreamState = null); + +[Workflow] +public class LlmWorkflow +{ + [WorkflowInit] + public LlmWorkflow(LlmInput input) => _ = new WorkflowStream(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(LlmInput input) + { + var result = await Workflow.ExecuteActivityAsync( + () => LlmActivities.StreamCompletionAsync(input.Prompt), + new() { StartToCloseTimeout = TimeSpan.FromMinutes(2) }); + + // Give the subscriber's in-flight poll time to receive the terminal event. + await Workflow.DelayAsync(TimeSpan.FromSeconds(30)); + return result; + } +} +``` + + + + +```csharp +using System.Text; +using Temporalio.Client; +using Temporalio.Converters; +using Temporalio.Extensions.WorkflowStreams; + +public static async Task StreamResponseAsync( + ITemporalClient temporalClient, + string workflowId, + CancellationToken cancellationToken) +{ + await using var streamClient = new WorkflowStreamClient(temporalClient, workflowId); + var converter = temporalClient.Options.DataConverter.PayloadConverter; + var output = new StringBuilder(); + var options = new WorkflowStreamSubscribeOptions + { + Topics = new[] { "delta", "retry", "complete" }, + }; + + await foreach (var item in streamClient.SubscribeAsync(options). + WithCancellation(cancellationToken)) + { + if (item.Topic == "retry") + { + output.Clear(); + Render(output.ToString()); + } + else if (item.Topic == "delta") + { + output.Append(converter.ToValue(item.Payload).Text); + Render(output.ToString()); + } + else if (item.Topic == "complete") + { + _ = converter.ToValue(item.Payload); + break; + } + } + + return output.ToString(); +} +``` + + + + +The Activity publishes a `RetryEvent` when its attempt number is greater than one. This tells the subscriber to discard +partial output from the previous attempt. Only the first delta, retries, and the terminal event force a flush; later +deltas use the 200 millisecond batch interval to reduce Signal volume. The Workflow uses the fixed-delay closing pattern; +use an [acknowledgment handshake](#closing-the-stream) when it should return as soon as a subscriber confirms receipt. + +## See also + +- [Workflow Streams samples (samples-dotnet)](https://github.com/temporalio/samples-dotnet/tree/main/src/WorkflowStreams): six runnable scenarios covering basic publish/subscribe, concurrent subscriptions, reconnecting subscribers, external publishers, bounded logs, and LLM streaming. +- [`Temporalio.Extensions.WorkflowStreams` API reference](https://dotnet.temporal.io/api/Temporalio.Extensions.WorkflowStreams.html). +- [Workflow message passing](/develop/dotnet/workflows/message-passing): Signals, Updates, and Queries that Workflow Streams is built on. diff --git a/docs/encyclopedia/workflow-message-passing/workflow-streams.mdx b/docs/encyclopedia/workflow-message-passing/workflow-streams.mdx index e7923a1738..96e3243dcf 100644 --- a/docs/encyclopedia/workflow-message-passing/workflow-streams.mdx +++ b/docs/encyclopedia/workflow-message-passing/workflow-streams.mdx @@ -61,6 +61,11 @@ publishers and subscribers per Workflow, not thousands. text="Workflow Streams with the Java SDK" archetype="feature-guide" /> + ## How Workflow Streams works @@ -138,6 +143,16 @@ Signal retried by the SDK or the network deduplicates to a single landing in the text="Publish from a client with the Java SDK" archetype="feature-guide" /> + + ### Topics @@ -191,6 +206,11 @@ The Workflow is the single source of truth for stream state, so any process brid text="Subscribe with the Java SDK" archetype="feature-guide" /> + ### Log management @@ -406,4 +426,9 @@ storage so each item is a small reference, and use truncation to drop entries th text="Long-running streams with the Java SDK" archetype="feature-guide" /> + diff --git a/sidebars.js b/sidebars.js index 95c9896ed6..6c00a756fc 100644 --- a/sidebars.js +++ b/sidebars.js @@ -29,6 +29,7 @@ const developDotnetCategory = { 'develop/dotnet/workflows/timers', 'develop/dotnet/workflows/dynamic-workflow', 'develop/dotnet/workflows/versioning', + 'develop/dotnet/workflows/workflow-streams', ], }, {