Temporalio.Extensions.WorkflowStreams
1.20.0
Prefix Reserved
dotnet add package Temporalio.Extensions.WorkflowStreams --version 1.20.0
NuGet\Install-Package Temporalio.Extensions.WorkflowStreams -Version 1.20.0
<PackageReference Include="Temporalio.Extensions.WorkflowStreams" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.WorkflowStreams" Version="1.20.0" />
<PackageReference Include="Temporalio.Extensions.WorkflowStreams" />
paket add Temporalio.Extensions.WorkflowStreams --version 1.20.0
#r "nuget: Temporalio.Extensions.WorkflowStreams, 1.20.0"
#:package Temporalio.Extensions.WorkflowStreams@1.20.0
#addin nuget:?package=Temporalio.Extensions.WorkflowStreams&version=1.20.0
#tool nuget:?package=Temporalio.Extensions.WorkflowStreams&version=1.20.0
Temporal .NET Workflow Streams
Workflow Streams is experimental. Its API and wire protocol may change before it is declared stable.
Temporalio.Extensions.WorkflowStreams provides a durable, offset-addressed, multi-topic log
hosted by a Temporal Workflow. External publishers append batches with a Signal, subscribers
long-poll with an Update, and a Query reports the current global offset. The implementation includes
publisher deduplication, topic filtering, truncation, response paging, and continue-as-new state
handoff.
This is intended for durable progress, event, and incremental-result streams. Each poll is a Temporal Update round trip, so it is not intended for ultra-low-latency media or token streaming.
Install
dotnet add package Temporalio.Extensions.WorkflowStreams
The package targets netstandard2.0.
Host a stream in a Workflow
Create the stream while the Workflow instance is constructed. This ensures its dynamically registered Signal, Update, and Query handlers exist before the first handler can be dispatched, including on a continue-as-new successor run.
using Temporalio.Extensions.WorkflowStreams;
using Temporalio.Workflows;
public record OrderInput(int CompletedSteps, WorkflowStreamState? StreamState = null);
[Workflow]
public class OrderWorkflow
{
private readonly WorkflowStream stream;
private bool finished;
[WorkflowInit]
public OrderWorkflow(OrderInput input)
{
stream = new(input.StreamState);
}
[WorkflowRun]
public async Task RunAsync(OrderInput input)
{
stream.GetTopic<object>("status").Publish(new { State = "started" });
await Workflow.WaitConditionAsync(() => finished || Workflow.ContinueAsNewSuggested);
if (Workflow.ContinueAsNewSuggested)
{
var state = await stream.CaptureStateForContinueAsNewAsync();
throw Workflow.CreateContinueAsNewException(
(OrderWorkflow workflow) => workflow.RunAsync(
input with { StreamState = state }));
}
}
[WorkflowSignal]
public Task FinishAsync()
{
finished = true;
return Task.CompletedTask;
}
}
CaptureStateForContinueAsNewAsync first detaches admitted pollers, waits until all handlers finish,
and only then captures WorkflowStreamState. Thread that state through the next run's input as shown
above. Once it captures the state, workflow-side publication is disabled so the snapshot describes
the final stream contents. Create and throw the continue-as-new exception immediately.
The log is retained until the Workflow calls Truncate(offset). Offsets are global across every
topic. A subscriber that falls behind truncation automatically resumes at the beginning of the
retained log.
Publish from a client or Activity
Use one asynchronously disposed client per target Workflow ID. Values are converted to Temporal
Payloads when Publish is called, then buffered until the two-second interval, the configured
batch size, a publication with forceFlush: true, a call to FlushAsync, or asynchronous disposal.
await using var streams = new WorkflowStreamClient(temporalClient, workflowId);
var status = streams.GetTopic<object>("status");
status.Publish(new { State = "working" });
status.Publish(new { State = "done" }, forceFlush: true);
await streams.FlushAsync(cancellationToken);
FlushAsync is a barrier for everything buffered before the call. Failed or timed-out Signal RPCs
retain the same publisher ID and sequence for retry, allowing the Workflow to deduplicate ambiguous
delivery. If the retry window expires, FlushTimeoutException is thrown and that ambiguous batch is
dropped locally: it may already be present in the Workflow log, or it may be lost. Later batches use
a new sequence.
Inside an Activity, FromActivity obtains the Temporal client, parent Workflow ID, and payload
converter from the current Activity context:
[Activity]
public async Task ReportAsync(IEnumerable<Progress> values)
{
await using var streams = WorkflowStreamClient.FromActivity();
foreach (var value in values)
{
streams.GetTopic<Progress>("progress").Publish(value);
}
}
Always use await using or call DisposeAsync so the final buffer is drained and publisher resources
are released. Publication after disposal throws ObjectDisposedException.
Subscribe
Strongly typed subscriptions deserialize each payload to the topic handle's value type. The returned
IAsyncEnumerable<WorkflowStreamItem<T>> is reusable: each enumeration starts with its own offset
and polling state. The non-generic overloads remain available for heterogeneous topics and yield raw
Temporal Payloads.
var subscription = streams.SubscribeAsync<MyEvent>(new()
{
Topics = new[] { "status", "progress" },
FromOffset = 0,
});
await foreach (var item in subscription.WithCancellation(cancellationToken))
{
Console.WriteLine($"{item.Offset} {item.Topic}: {item.Value}");
}
For one topic, use streams.GetTopic<StatusUpdate>("status").SubscribeAsync(fromOffset). An empty
topic collection subscribes to every topic; the empty string is the cross-SDK no-topic value.
Consumer cancellation cancels the in-flight RPC and throws OperationCanceledException. Disposing
the owning client ends its active enumerations cleanly. Enumerations also end cleanly when the
Workflow reaches a terminal state and automatically follow continue-as-new chains.
Data conversion and interoperability
The fixed protocol handlers are:
- Signal
__temporal_workflow_stream_publish - Update
__temporal_workflow_stream_poll - Query
__temporal_workflow_stream_offset
The public wire DTOs use the protocol's exact snake-case JSON names. Each item's data is standard
padded base64 containing a serialized temporal.api.common.v1.Payload protobuf. This preserves its
encoding metadata and makes .NET publishers, Workflow hosts, and subscribers interoperable with the
official Workflow Streams implementations in other Temporal SDKs.
Only payload conversion is applied to an individual item. A client's payload codec chain, if any, is applied once to the Signal or Update envelope rather than once per item, avoiding double encoding. The envelope itself must use JSON-compatible conversion for cross-language interoperability.
FromActivity uses the Activity's payload converter for individual stream items. Do not use its
subscriptions with a custom payload converter that requires the serialization context used for
deserialization to match the context used for serialization: Activity publications and Workflow
Stream subscriptions necessarily have different contexts.
Operational limits
- Every waiting subscription uses an admitted Workflow Update. Account for concurrent and total Update limits when choosing subscriber counts and continue-as-new frequency.
- Poll results are paged at an estimated one megabyte and immediately repolled while another page is ready. An individual item must fit in one page.
- The Workflow log has no automatic retention policy. Truncate items once all required consumers have advanced, and carry only the retained state through continue-as-new.
- Publishing uses Signals, so malformed externally supplied wire entries cannot return errors to their sender. The Workflow skips them and emits a replay-safe warning.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net5.0 was computed. net5.0-windows was computed. net6.0 was computed. net6.0-android was computed. net6.0-ios was computed. net6.0-maccatalyst was computed. net6.0-macos was computed. net6.0-tvos was computed. net6.0-windows was computed. net7.0 was computed. net7.0-android was computed. net7.0-ios was computed. net7.0-maccatalyst was computed. net7.0-macos was computed. net7.0-tvos was computed. net7.0-windows was computed. net8.0 was computed. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
| .NET Core | netcoreapp2.0 was computed. netcoreapp2.1 was computed. netcoreapp2.2 was computed. netcoreapp3.0 was computed. netcoreapp3.1 was computed. |
| .NET Standard | netstandard2.0 is compatible. netstandard2.1 was computed. |
| .NET Framework | net461 was computed. net462 was computed. net463 was computed. net47 was computed. net471 was computed. net472 was computed. net48 was computed. net481 was computed. |
| MonoAndroid | monoandroid was computed. |
| MonoMac | monomac was computed. |
| MonoTouch | monotouch was computed. |
| Tizen | tizen40 was computed. tizen60 was computed. |
| Xamarin.iOS | xamarinios was computed. |
| Xamarin.Mac | xamarinmac was computed. |
| Xamarin.TVOS | xamarintvos was computed. |
| Xamarin.WatchOS | xamarinwatchos was computed. |
-
.NETStandard 2.0
- Microsoft.Bcl.AsyncInterfaces (>= 9.0.4)
- System.Text.Json (>= 9.0.4)
- Temporalio (>= 1.20.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 1.20.0 | 177 | 9/28/2026 |