diff --git a/.github/CODEOWNERS b/.github/CODEOWNERS index cf486e1..46d7966 100644 --- a/.github/CODEOWNERS +++ b/.github/CODEOWNERS @@ -8,3 +8,7 @@ # The Nexus team owns any folder whose name starts with "Nexus" /src/Nexus*/ @temporalio/nexus /tests/Nexus*/ @temporalio/nexus + +# SDK & AI SDK teams share ownership of Workflow Streams samples +/src/WorkflowStreams/ @temporalio/sdk @temporalio/ai-sdk +/tests/WorkflowStreams/ @temporalio/sdk @temporalio/ai-sdk diff --git a/Directory.Packages.props b/Directory.Packages.props index c48df34..d478da0 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -17,15 +17,17 @@ + - + + diff --git a/README.md b/README.md index a09c161..afa76f1 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,7 @@ Prerequisites: * [UpdateWithStartLazyInit](src/UpdateWithStartLazyInit) - Use update with start to lazily start a workflow before sending update. * [WorkerSpecificTaskQueues](src/WorkerSpecificTaskQueues) - Use a unique task queue per Worker to have certain Activities only run on that specific Worker. * [WorkerVersioning](src/WorkerVersioning) - How to use the Worker Versioning feature to more easily deploy changes to Workflow & other code. +* [WorkflowStreams](src/WorkflowStreams) - Host a durable publish/subscribe stream inside a Workflow. * [WorkflowUpdate](src/WorkflowUpdate) - How to use the Workflow Update feature while blocking in update method for concurrent updates. ## Development diff --git a/TemporalioSamples.sln b/TemporalioSamples.sln index c945182..85b8bc8 100644 --- a/TemporalioSamples.sln +++ b/TemporalioSamples.sln @@ -141,6 +141,18 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "SearchAttributes", "SearchA EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.SearchAttributes", "src\SearchAttributes\TemporalioSamples.SearchAttributes.csproj", "{97376F57-BA10-464B-AFBD-583E187DE947}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.BasicPublishSubscribe", "src\WorkflowStreams\BasicPublishSubscribe\TemporalioSamples.WorkflowStreams.BasicPublishSubscribe.csproj", "{091DC9D4-A88D-4393-9F84-32627DB1141F}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.BoundedLog", "src\WorkflowStreams\BoundedLog\TemporalioSamples.WorkflowStreams.BoundedLog.csproj", "{A1111111-1111-4111-8111-111111111111}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.ExternalPublisher", "src\WorkflowStreams\ExternalPublisher\TemporalioSamples.WorkflowStreams.ExternalPublisher.csproj", "{B2222222-2222-4222-8222-222222222222}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions", "src\WorkflowStreams\ConcurrentSubscriptions\TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions.csproj", "{C3333333-3333-4333-8333-333333333333}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.LlmTokenStreaming", "src\WorkflowStreams\LlmTokenStreaming\TemporalioSamples.WorkflowStreams.LlmTokenStreaming.csproj", "{D4444444-4444-4444-8444-444444444444}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.WorkflowStreams.ReconnectingSubscriber", "src\WorkflowStreams\ReconnectingSubscriber\TemporalioSamples.WorkflowStreams.ReconnectingSubscriber.csproj", "{E5555555-5555-4555-8555-555555555555}" +EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "NexusStandaloneActivity", "NexusStandaloneActivity", "{5D5AFBCD-B4B7-4A4E-91C8-FD80E4C1E27B}" EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.NexusStandaloneActivity", "src\NexusStandaloneActivity\TemporalioSamples.NexusStandaloneActivity.csproj", "{4D8C9F9B-F8E3-4160-9286-32C3966A0125}" @@ -803,6 +815,78 @@ Global {97376F57-BA10-464B-AFBD-583E187DE947}.Release|x64.Build.0 = Release|Any CPU {97376F57-BA10-464B-AFBD-583E187DE947}.Release|x86.ActiveCfg = Release|Any CPU {97376F57-BA10-464B-AFBD-583E187DE947}.Release|x86.Build.0 = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|Any CPU.Build.0 = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|x64.ActiveCfg = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|x64.Build.0 = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|x86.ActiveCfg = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Debug|x86.Build.0 = Debug|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|Any CPU.ActiveCfg = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|Any CPU.Build.0 = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|x64.ActiveCfg = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|x64.Build.0 = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|x86.ActiveCfg = Release|Any CPU + {091DC9D4-A88D-4393-9F84-32627DB1141F}.Release|x86.Build.0 = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|Any CPU.Build.0 = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|x64.ActiveCfg = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|x64.Build.0 = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|x86.ActiveCfg = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Debug|x86.Build.0 = Debug|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|Any CPU.ActiveCfg = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|Any CPU.Build.0 = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|x64.ActiveCfg = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|x64.Build.0 = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|x86.ActiveCfg = Release|Any CPU + {A1111111-1111-4111-8111-111111111111}.Release|x86.Build.0 = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|Any CPU.Build.0 = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|x64.ActiveCfg = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|x64.Build.0 = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|x86.ActiveCfg = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Debug|x86.Build.0 = Debug|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|Any CPU.ActiveCfg = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|Any CPU.Build.0 = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|x64.ActiveCfg = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|x64.Build.0 = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|x86.ActiveCfg = Release|Any CPU + {B2222222-2222-4222-8222-222222222222}.Release|x86.Build.0 = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|Any CPU.Build.0 = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|x64.ActiveCfg = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|x64.Build.0 = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|x86.ActiveCfg = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Debug|x86.Build.0 = Debug|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|Any CPU.ActiveCfg = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|Any CPU.Build.0 = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|x64.ActiveCfg = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|x64.Build.0 = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|x86.ActiveCfg = Release|Any CPU + {C3333333-3333-4333-8333-333333333333}.Release|x86.Build.0 = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|Any CPU.Build.0 = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|x64.ActiveCfg = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|x64.Build.0 = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|x86.ActiveCfg = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Debug|x86.Build.0 = Debug|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|Any CPU.ActiveCfg = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|Any CPU.Build.0 = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|x64.ActiveCfg = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|x64.Build.0 = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|x86.ActiveCfg = Release|Any CPU + {D4444444-4444-4444-8444-444444444444}.Release|x86.Build.0 = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|Any CPU.Build.0 = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|x64.ActiveCfg = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|x64.Build.0 = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|x86.ActiveCfg = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Debug|x86.Build.0 = Debug|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|Any CPU.ActiveCfg = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|Any CPU.Build.0 = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|x64.ActiveCfg = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|x64.Build.0 = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|x86.ActiveCfg = Release|Any CPU + {E5555555-5555-4555-8555-555555555555}.Release|x86.Build.0 = Release|Any CPU {4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Debug|Any CPU.ActiveCfg = Debug|Any CPU {4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Debug|Any CPU.Build.0 = Debug|Any CPU {4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Debug|x64.ActiveCfg = Debug|Any CPU @@ -885,6 +969,12 @@ Global {6B6622AE-2970-4AAD-B2E4-A7EE1E2C40EA} = {17436B0C-8853-030B-9E2F-BAEB65BE59FB} {34ADAC00-5559-6993-F779-D71CE5A0F93F} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} {97376F57-BA10-464B-AFBD-583E187DE947} = {34ADAC00-5559-6993-F779-D71CE5A0F93F} + {091DC9D4-A88D-4393-9F84-32627DB1141F} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} + {A1111111-1111-4111-8111-111111111111} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} + {B2222222-2222-4222-8222-222222222222} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} + {C3333333-3333-4333-8333-333333333333} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} + {D4444444-4444-4444-8444-444444444444} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} + {E5555555-5555-4555-8555-555555555555} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} {5D5AFBCD-B4B7-4A4E-91C8-FD80E4C1E27B} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC} {4D8C9F9B-F8E3-4160-9286-32C3966A0125} = {5D5AFBCD-B4B7-4A4E-91C8-FD80E4C1E27B} EndGlobalSection diff --git a/src/WorkflowStreams/BasicPublishSubscribe/Constants.cs b/src/WorkflowStreams/BasicPublishSubscribe/Constants.cs new file mode 100644 index 0000000..1d35391 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/Constants.cs @@ -0,0 +1,12 @@ +namespace TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-basic-publish-subscribe"; + + public const string TopicStatus = "status"; + + public const string TopicProgress = "progress"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/BasicPublishSubscribe/Models.cs b/src/WorkflowStreams/BasicPublishSubscribe/Models.cs new file mode 100644 index 0000000..cba780f --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/Models.cs @@ -0,0 +1,9 @@ +namespace TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +using Temporalio.Extensions.WorkflowStreams; + +public record OrderInput(string OrderId, WorkflowStreamState? StreamState = null); + +public record StatusEvent(string Kind, string OrderId); + +public record ProgressEvent(string Message); diff --git a/src/WorkflowStreams/BasicPublishSubscribe/OrderWorkflow.workflow.cs b/src/WorkflowStreams/BasicPublishSubscribe/OrderWorkflow.workflow.cs new file mode 100644 index 0000000..0377df4 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/OrderWorkflow.workflow.cs @@ -0,0 +1,31 @@ +namespace TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[Workflow] +public class OrderWorkflow +{ + private readonly WorkflowStream stream; + + [WorkflowInit] + public OrderWorkflow(OrderInput input) => stream = new(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(OrderInput input) + { + var status = stream.GetTopic(Constants.TopicStatus); + var progress = stream.GetTopic(Constants.TopicProgress); + + status.Publish(new StatusEvent("received", input.OrderId)); + var chargeId = await Workflow.ExecuteActivityAsync( + () => PaymentActivities.ChargeCardAsync(input.OrderId), + new() { StartToCloseTimeout = TimeSpan.FromMinutes(1), }); + status.Publish(new StatusEvent("shipped", input.OrderId)); + progress.Publish(new ProgressEvent($"charge id: {chargeId}")); + status.Publish(new StatusEvent("complete", input.OrderId)); + + await Workflow.DelayAsync(Constants.DrainDelay); + return chargeId; + } +} diff --git a/src/WorkflowStreams/BasicPublishSubscribe/PaymentActivities.cs b/src/WorkflowStreams/BasicPublishSubscribe/PaymentActivities.cs new file mode 100644 index 0000000..4c34d3e --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/PaymentActivities.cs @@ -0,0 +1,26 @@ +namespace TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +using Microsoft.Extensions.Logging; +using Temporalio.Activities; +using Temporalio.Extensions.WorkflowStreams; + +public static class PaymentActivities +{ + [Activity] + public static async Task ChargeCardAsync(string orderId) + { + await using var streamClient = WorkflowStreamClient.FromActivity( + new() { BatchInterval = TimeSpan.FromMilliseconds(200), }); + var progress = streamClient.GetTopic(Constants.TopicProgress); + + progress.Publish(new ProgressEvent("charging card...")); + ActivityExecutionContext.Current.Logger.LogInformation( + "Charging card for order {OrderId}", orderId); + await Task.Delay( + TimeSpan.FromSeconds(1), + ActivityExecutionContext.Current.CancellationToken); + progress.Publish(new ProgressEvent("card charged")); + await streamClient.FlushAsync(); + return $"charge-{orderId}"; + } +} diff --git a/src/WorkflowStreams/BasicPublishSubscribe/Program.cs b/src/WorkflowStreams/BasicPublishSubscribe/Program.cs new file mode 100644 index 0000000..da133e3 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/Program.cs @@ -0,0 +1,33 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue). + AddActivity(PaymentActivities.ChargeCardAsync). + AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunPublisherAsync(client); +} diff --git a/src/WorkflowStreams/BasicPublishSubscribe/README.md b/src/WorkflowStreams/BasicPublishSubscribe/README.md new file mode 100644 index 0000000..6f2c575 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/README.md @@ -0,0 +1,16 @@ +# Basic publish/subscribe + +An order Workflow publishes lifecycle events to `status`, while its payment Activity publishes +progress events to `progress`. The subscriber consumes and decodes both topics. + +Start a Temporal service, then run the worker: + +```bash +dotnet run -- worker +``` + +In another terminal, run the subscriber: + +```bash +dotnet run +``` diff --git a/src/WorkflowStreams/BasicPublishSubscribe/Scenario.cs b/src/WorkflowStreams/BasicPublishSubscribe/Scenario.cs new file mode 100644 index 0000000..6ec3de2 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/Scenario.cs @@ -0,0 +1,49 @@ +namespace TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; + +using Temporalio.Client; +using Temporalio.Converters; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + public static async Task RunPublisherAsync(ITemporalClient client) + { + var workflowId = $"workflow-streams-order-{Guid.NewGuid()}"; + var handle = await client.StartWorkflowAsync( + (OrderWorkflow wf) => wf.RunAsync(new OrderInput("order-42", null)), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine($"Started workflow: {workflowId}"); + + await using var streamClient = new WorkflowStreamClient(client, workflowId); + var options = new WorkflowStreamSubscribeOptions + { + Topics = new List + { + Constants.TopicStatus, + Constants.TopicProgress, + }, + }; + await foreach (var item in streamClient.SubscribeAsync(options)) + { + if (item.Topic == Constants.TopicStatus) + { + var evt = Decode(client, item); + Console.WriteLine($"[status] {evt.Kind}: order={evt.OrderId}"); + if (evt.Kind == "complete") + { + break; + } + } + else if (item.Topic == Constants.TopicProgress) + { + var evt = Decode(client, item); + Console.WriteLine($"[progress] {evt.Message}"); + } + } + + Console.WriteLine($"Workflow result: {await handle.GetResultAsync()}"); + } + + private static T Decode(ITemporalClient client, WorkflowStreamItem item) => + client.Options.DataConverter.PayloadConverter.ToValue(item.Payload); +} diff --git a/src/WorkflowStreams/BasicPublishSubscribe/TemporalioSamples.WorkflowStreams.BasicPublishSubscribe.csproj b/src/WorkflowStreams/BasicPublishSubscribe/TemporalioSamples.WorkflowStreams.BasicPublishSubscribe.csproj new file mode 100644 index 0000000..c1cc2f4 --- /dev/null +++ b/src/WorkflowStreams/BasicPublishSubscribe/TemporalioSamples.WorkflowStreams.BasicPublishSubscribe.csproj @@ -0,0 +1,11 @@ + + + + Exe + + + + + + + diff --git a/src/WorkflowStreams/BoundedLog/Constants.cs b/src/WorkflowStreams/BoundedLog/Constants.cs new file mode 100644 index 0000000..5cc4cac --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/Constants.cs @@ -0,0 +1,10 @@ +namespace TemporalioSamples.WorkflowStreams.BoundedLog; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-bounded-log"; + + public const string TopicTick = "tick"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/BoundedLog/Models.cs b/src/WorkflowStreams/BoundedLog/Models.cs new file mode 100644 index 0000000..42a38c2 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/Models.cs @@ -0,0 +1,12 @@ +namespace TemporalioSamples.WorkflowStreams.BoundedLog; + +using Temporalio.Extensions.WorkflowStreams; + +public record TickerInput( + int Count = 50, + int KeepLast = 10, + int TruncateEvery = 5, + TimeSpan? Interval = null, + WorkflowStreamState? StreamState = null); + +public record TickEvent(int N); diff --git a/src/WorkflowStreams/BoundedLog/Program.cs b/src/WorkflowStreams/BoundedLog/Program.cs new file mode 100644 index 0000000..a872a70 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/Program.cs @@ -0,0 +1,31 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.BoundedLog; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue).AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunTruncatingTickerAsync(client); +} diff --git a/src/WorkflowStreams/BoundedLog/README.md b/src/WorkflowStreams/BoundedLog/README.md new file mode 100644 index 0000000..47d0276 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/README.md @@ -0,0 +1,16 @@ +# Bounded log + +A ticker Workflow periodically truncates its stream. A fast subscriber sees every tick while a +late subscriber is advanced from a stale offset to the retained base offset. + +Start a Temporal service, then run the worker: + +```bash +dotnet run -- worker +``` + +In another terminal, run the subscribers: + +```bash +dotnet run +``` diff --git a/src/WorkflowStreams/BoundedLog/Scenario.cs b/src/WorkflowStreams/BoundedLog/Scenario.cs new file mode 100644 index 0000000..094f214 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/Scenario.cs @@ -0,0 +1,70 @@ +namespace TemporalioSamples.WorkflowStreams.BoundedLog; + +using Temporalio.Client; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + private const int TickCount = 30; + private const int KeepLast = 5; + private const int TruncateEvery = 5; + private const long StaleOffset = 1; + + public static async Task RunTruncatingTickerAsync(ITemporalClient client) + { + var workflowId = $"workflow-streams-ticker-{Guid.NewGuid()}"; + var handle = await client.StartWorkflowAsync( + (TickerWorkflow wf) => wf.RunAsync( + new TickerInput(TickCount, KeepLast, TruncateEvery, null, null)), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine($"Started workflow: {workflowId}"); + + async Task FastSubscriberAsync() + { + await using var streamClient = new WorkflowStreamClient(client, workflowId); + await foreach (var item in streamClient. + GetTopic(Constants.TopicTick).SubscribeAsync()) + { + var evt = item.Value; + Console.WriteLine($"[fast] offset={item.Offset,3} n={evt.N}"); + if (evt.N == TickCount - 1) + { + break; + } + } + } + + async Task LateSubscriberAsync() + { + await using var streamClient = new WorkflowStreamClient(client, workflowId); + var firstTruncate = ((KeepLast / TruncateEvery) + 1) * TruncateEvery; + while (await streamClient.GetOffsetAsync() <= firstTruncate) + { + await Task.Delay(TimeSpan.FromMilliseconds(200)); + } + + var first = true; + await foreach (var item in streamClient. + GetTopic(Constants.TopicTick).SubscribeAsync(StaleOffset)) + { + var evt = item.Value; + if (first && item.Offset > StaleOffset) + { + Console.WriteLine( + $"[late] requested offset {StaleOffset} but it was truncated; " + + $"fast-forwarded to offset {item.Offset} " + + $"(skipped {item.Offset - StaleOffset} tick(s))"); + } + first = false; + Console.WriteLine($"[late] offset={item.Offset,3} n={evt.N}"); + if (evt.N == TickCount - 1) + { + break; + } + } + } + + await Task.WhenAll(FastSubscriberAsync(), LateSubscriberAsync()); + Console.WriteLine($"Workflow result: {await handle.GetResultAsync()}"); + } +} diff --git a/src/WorkflowStreams/BoundedLog/TemporalioSamples.WorkflowStreams.BoundedLog.csproj b/src/WorkflowStreams/BoundedLog/TemporalioSamples.WorkflowStreams.BoundedLog.csproj new file mode 100644 index 0000000..c1cc2f4 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/TemporalioSamples.WorkflowStreams.BoundedLog.csproj @@ -0,0 +1,11 @@ + + + + Exe + + + + + + + diff --git a/src/WorkflowStreams/BoundedLog/TickerWorkflow.workflow.cs b/src/WorkflowStreams/BoundedLog/TickerWorkflow.workflow.cs new file mode 100644 index 0000000..74d0f90 --- /dev/null +++ b/src/WorkflowStreams/BoundedLog/TickerWorkflow.workflow.cs @@ -0,0 +1,38 @@ +namespace TemporalioSamples.WorkflowStreams.BoundedLog; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[Workflow] +public class TickerWorkflow +{ + private readonly WorkflowStream stream; + + [WorkflowInit] + public TickerWorkflow(TickerInput input) => stream = new(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(TickerInput input) + { + var tick = stream.GetTopic(Constants.TopicTick); + var interval = input.Interval ?? TimeSpan.FromMilliseconds(200); + + for (var n = 0; n < input.Count; n++) + { + tick.Publish(new TickEvent(n)); + if (interval > TimeSpan.Zero) + { + await Workflow.DelayAsync(interval); + } + + var published = n + 1; + if (published % input.TruncateEvery == 0 && published > input.KeepLast) + { + stream.Truncate(published - input.KeepLast); + } + } + + await Workflow.DelayAsync(Constants.DrainDelay); + return $"ticker emitted {input.Count} events"; + } +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/Constants.cs b/src/WorkflowStreams/ConcurrentSubscriptions/Constants.cs new file mode 100644 index 0000000..042ac25 --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/Constants.cs @@ -0,0 +1,12 @@ +namespace TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-concurrent-subscriptions"; + + public const string TopicStatus = "status"; + + public const string TopicProgress = "progress"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/Models.cs b/src/WorkflowStreams/ConcurrentSubscriptions/Models.cs new file mode 100644 index 0000000..903b651 --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/Models.cs @@ -0,0 +1,9 @@ +namespace TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +using Temporalio.Extensions.WorkflowStreams; + +public record OrderInput(string OrderId, WorkflowStreamState? StreamState = null); + +public record StatusEvent(string Kind, string OrderId); + +public record ProgressEvent(string Message); diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/OrderWorkflow.workflow.cs b/src/WorkflowStreams/ConcurrentSubscriptions/OrderWorkflow.workflow.cs new file mode 100644 index 0000000..88133f0 --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/OrderWorkflow.workflow.cs @@ -0,0 +1,31 @@ +namespace TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[Workflow] +public class OrderWorkflow +{ + private readonly WorkflowStream stream; + + [WorkflowInit] + public OrderWorkflow(OrderInput input) => stream = new(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(OrderInput input) + { + var status = stream.GetTopic(Constants.TopicStatus); + var progress = stream.GetTopic(Constants.TopicProgress); + + status.Publish(new StatusEvent("received", input.OrderId)); + var chargeId = await Workflow.ExecuteActivityAsync( + () => PaymentActivities.ChargeCardAsync(input.OrderId), + new() { StartToCloseTimeout = TimeSpan.FromMinutes(1), }); + status.Publish(new StatusEvent("shipped", input.OrderId)); + progress.Publish(new ProgressEvent($"charge id: {chargeId}")); + status.Publish(new StatusEvent("complete", input.OrderId)); + + await Workflow.DelayAsync(Constants.DrainDelay); + return chargeId; + } +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/PaymentActivities.cs b/src/WorkflowStreams/ConcurrentSubscriptions/PaymentActivities.cs new file mode 100644 index 0000000..41d9b6d --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/PaymentActivities.cs @@ -0,0 +1,26 @@ +namespace TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +using Microsoft.Extensions.Logging; +using Temporalio.Activities; +using Temporalio.Extensions.WorkflowStreams; + +public static class PaymentActivities +{ + [Activity] + public static async Task ChargeCardAsync(string orderId) + { + await using var streamClient = WorkflowStreamClient.FromActivity( + new() { BatchInterval = TimeSpan.FromMilliseconds(200), }); + var progress = streamClient.GetTopic(Constants.TopicProgress); + + progress.Publish(new ProgressEvent("charging card...")); + ActivityExecutionContext.Current.Logger.LogInformation( + "Charging card for order {OrderId}", orderId); + await Task.Delay( + TimeSpan.FromSeconds(1), + ActivityExecutionContext.Current.CancellationToken); + progress.Publish(new ProgressEvent("card charged")); + await streamClient.FlushAsync(); + return $"charge-{orderId}"; + } +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/Program.cs b/src/WorkflowStreams/ConcurrentSubscriptions/Program.cs new file mode 100644 index 0000000..2f984dc --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/Program.cs @@ -0,0 +1,33 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue). + AddActivity(PaymentActivities.ChargeCardAsync). + AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunSubscriptionsAsync(client); +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/README.md b/src/WorkflowStreams/ConcurrentSubscriptions/README.md new file mode 100644 index 0000000..75a44fe --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/README.md @@ -0,0 +1,16 @@ +# Concurrent subscriptions + +Runs two order Workflows and consumes each reusable `IAsyncEnumerable` +concurrently. Awaiting each item provides natural backpressure within each subscription. + +Start a Temporal service, then run the worker: + +```bash +dotnet run -- worker +``` + +In another terminal, run the subscribers: + +```bash +dotnet run +``` diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/Scenario.cs b/src/WorkflowStreams/ConcurrentSubscriptions/Scenario.cs new file mode 100644 index 0000000..0b56167 --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/Scenario.cs @@ -0,0 +1,81 @@ +namespace TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; + +using Temporalio.Client; +using Temporalio.Converters; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + private static readonly string[] OrderIds = new[] { "order-A", "order-B" }; + + public static async Task RunSubscriptionsAsync(ITemporalClient client) + { + var workflowHandles = new List>(); + var streamClients = new List(); + var subscriptions = new List(); + + foreach (var orderId in OrderIds) + { + var workflowId = $"workflow-streams-concurrent-{orderId}-{Guid.NewGuid()}"; + var workflowHandle = await client.StartWorkflowAsync( + (OrderWorkflow wf) => wf.RunAsync(new OrderInput(orderId, null)), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine($"Started workflow: {workflowId}"); + +#pragma warning disable CA2000 // The client is closed in the method's finally block + var streamClient = new WorkflowStreamClient(client, workflowId); +#pragma warning restore CA2000 + workflowHandles.Add(workflowHandle); + streamClients.Add(streamClient); + subscriptions.Add(RenderSubscriptionAsync(streamClient, client, orderId)); + } + + try + { + await Task.WhenAll(subscriptions); + for (var i = 0; i < workflowHandles.Count; i++) + { + var result = await workflowHandles[i].GetResultAsync(); + Console.WriteLine($"[{OrderIds[i]}] workflow result: {result}"); + } + } + finally + { + foreach (var streamClient in streamClients) + { + await streamClient.DisposeAsync(); + } + } + } + + private static async Task RenderSubscriptionAsync( + WorkflowStreamClient streamClient, + ITemporalClient client, + string orderId) + { + await foreach (var item in streamClient.SubscribeAsync(new() + { + Topics = new List + { + Constants.TopicStatus, + Constants.TopicProgress, + }, + })) + { + if (item.Topic == Constants.TopicStatus) + { + Console.WriteLine( + $"[{orderId}] [status] {Decode(client, item).Kind}"); + } + else if (item.Topic == Constants.TopicProgress) + { + Console.WriteLine( + $"[{orderId}] [progress] {Decode(client, item).Message}"); + } + } + Console.WriteLine($"[{orderId}] stream completed"); + } + + private static T Decode(ITemporalClient client, WorkflowStreamItem item) => + client.Options.DataConverter.PayloadConverter.ToValue(item.Payload); +} diff --git a/src/WorkflowStreams/ConcurrentSubscriptions/TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions.csproj b/src/WorkflowStreams/ConcurrentSubscriptions/TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions.csproj new file mode 100644 index 0000000..c1cc2f4 --- /dev/null +++ b/src/WorkflowStreams/ConcurrentSubscriptions/TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions.csproj @@ -0,0 +1,11 @@ + + + + Exe + + + + + + + diff --git a/src/WorkflowStreams/ExternalPublisher/Constants.cs b/src/WorkflowStreams/ExternalPublisher/Constants.cs new file mode 100644 index 0000000..51249e1 --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/Constants.cs @@ -0,0 +1,12 @@ +namespace TemporalioSamples.WorkflowStreams.ExternalPublisher; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-external-publisher"; + + public const string TopicNews = "news"; + + public const string DoneHeadline = "-- end of feed --"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/ExternalPublisher/HubWorkflow.workflow.cs b/src/WorkflowStreams/ExternalPublisher/HubWorkflow.workflow.cs new file mode 100644 index 0000000..820bdb3 --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/HubWorkflow.workflow.cs @@ -0,0 +1,24 @@ +namespace TemporalioSamples.WorkflowStreams.ExternalPublisher; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[Workflow] +public class HubWorkflow +{ + private bool closed; + + [WorkflowInit] + public HubWorkflow(HubInput input) => _ = new WorkflowStream(input.StreamState); + + [WorkflowRun] + public async Task RunAsync(HubInput input) + { + await Workflow.WaitConditionAsync(() => closed); + await Workflow.DelayAsync(Constants.DrainDelay); + return $"hub {input.HubId} closed"; + } + + [WorkflowSignal] + public async Task CloseAsync() => closed = true; +} diff --git a/src/WorkflowStreams/ExternalPublisher/Models.cs b/src/WorkflowStreams/ExternalPublisher/Models.cs new file mode 100644 index 0000000..1a4ac4f --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/Models.cs @@ -0,0 +1,7 @@ +namespace TemporalioSamples.WorkflowStreams.ExternalPublisher; + +using Temporalio.Extensions.WorkflowStreams; + +public record HubInput(string HubId, WorkflowStreamState? StreamState = null); + +public record NewsEvent(string Headline); diff --git a/src/WorkflowStreams/ExternalPublisher/Program.cs b/src/WorkflowStreams/ExternalPublisher/Program.cs new file mode 100644 index 0000000..fa82847 --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/Program.cs @@ -0,0 +1,31 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.ExternalPublisher; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue).AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunExternalPublisherAsync(client); +} diff --git a/src/WorkflowStreams/ExternalPublisher/README.md b/src/WorkflowStreams/ExternalPublisher/README.md new file mode 100644 index 0000000..9b09f8a --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/README.md @@ -0,0 +1,16 @@ +# External publisher + +A Workflow hosts the stream while a normal client publishes news and another client subscribes. +The client signals the Workflow to close after flushing a sentinel event. + +Start a Temporal service, then run the worker: + +```bash +dotnet run -- worker +``` + +In another terminal, run the publisher and subscriber: + +```bash +dotnet run +``` diff --git a/src/WorkflowStreams/ExternalPublisher/Scenario.cs b/src/WorkflowStreams/ExternalPublisher/Scenario.cs new file mode 100644 index 0000000..03a1ecf --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/Scenario.cs @@ -0,0 +1,57 @@ +namespace TemporalioSamples.WorkflowStreams.ExternalPublisher; + +using Temporalio.Client; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + private static readonly string[] Headlines = + [ + "markets open higher", + "new bridge opens downtown", + "local team wins championship", + ]; + + public static async Task RunExternalPublisherAsync(ITemporalClient client) + { + var workflowId = $"workflow-streams-hub-{Guid.NewGuid()}"; + var handle = await client.StartWorkflowAsync( + (HubWorkflow wf) => wf.RunAsync(new HubInput("newsroom", null)), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine($"Started workflow: {workflowId}"); + + async Task SubscribeAsync() + { + await using var streamClient = new WorkflowStreamClient(client, workflowId); + await foreach (var item in streamClient. + GetTopic(Constants.TopicNews).SubscribeAsync()) + { + var evt = item.Value; + if (evt.Headline == Constants.DoneHeadline) + { + break; + } + Console.WriteLine($"[subscriber] {evt.Headline}"); + } + } + + async Task PublishAsync() + { + await using var streamClient = new WorkflowStreamClient(client, workflowId); + var news = streamClient.GetTopic(Constants.TopicNews); + foreach (var headline in Headlines) + { + news.Publish(new NewsEvent(headline)); + Console.WriteLine($"[publisher] sent: {headline}"); + await Task.Delay(TimeSpan.FromMilliseconds(500)); + } + news.Publish(new NewsEvent(Constants.DoneHeadline), forceFlush: true); + await streamClient.FlushAsync(); + await handle.SignalAsync(wf => wf.CloseAsync()); + Console.WriteLine("[publisher] signaled close"); + } + + await Task.WhenAll(SubscribeAsync(), PublishAsync()); + Console.WriteLine($"Workflow result: {await handle.GetResultAsync()}"); + } +} diff --git a/src/WorkflowStreams/ExternalPublisher/TemporalioSamples.WorkflowStreams.ExternalPublisher.csproj b/src/WorkflowStreams/ExternalPublisher/TemporalioSamples.WorkflowStreams.ExternalPublisher.csproj new file mode 100644 index 0000000..c1cc2f4 --- /dev/null +++ b/src/WorkflowStreams/ExternalPublisher/TemporalioSamples.WorkflowStreams.ExternalPublisher.csproj @@ -0,0 +1,11 @@ + + + + Exe + + + + + + + diff --git a/src/WorkflowStreams/LlmTokenStreaming/Constants.cs b/src/WorkflowStreams/LlmTokenStreaming/Constants.cs new file mode 100644 index 0000000..14adf4c --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/Constants.cs @@ -0,0 +1,14 @@ +namespace TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-llm-token-streaming"; + + public const string TopicDelta = "delta"; + + public const string TopicComplete = "complete"; + + public const string TopicRetry = "retry"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/LlmTokenStreaming/LlmActivities.cs b/src/WorkflowStreams/LlmTokenStreaming/LlmActivities.cs new file mode 100644 index 0000000..1bad542 --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/LlmActivities.cs @@ -0,0 +1,65 @@ +namespace TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +using System.ClientModel; +using System.ClientModel.Primitives; +using System.Text; +using OpenAI; +using OpenAI.Chat; +using Temporalio.Activities; +using Temporalio.Extensions.WorkflowStreams; + +public static class LlmActivities +{ + [Activity] + public static async Task StreamCompletionAsync(LlmInput input) + { + var apiKey = Environment.GetEnvironmentVariable("OPENAI_API_KEY"); + if (string.IsNullOrWhiteSpace(apiKey)) + { + throw new InvalidOperationException("OPENAI_API_KEY must be set for the LLM scenario"); + } + + await using var streamClient = WorkflowStreamClient.FromActivity( + new() { BatchInterval = TimeSpan.FromMilliseconds(200), }); + var deltas = streamClient.GetTopic(Constants.TopicDelta); + var complete = streamClient.GetTopic(Constants.TopicComplete); + var retry = streamClient.GetTopic(Constants.TopicRetry); + + var activityContext = ActivityExecutionContext.Current; + if (activityContext.Info.Attempt > 1) + { + retry.Publish(new RetryEvent(activityContext.Info.Attempt), forceFlush: true); + } + + var clientOptions = new OpenAIClientOptions + { + RetryPolicy = new ClientRetryPolicy(0), + }; + var chatClient = new ChatClient( + input.Model, + new ApiKeyCredential(apiKey), + clientOptions); + + var fullText = new StringBuilder(); + var updates = chatClient.CompleteChatStreamingAsync( + [new UserChatMessage(input.Prompt)], + cancellationToken: activityContext.CancellationToken); + await foreach (var update in updates) + { + foreach (var contentPart in update.ContentUpdate) + { + if (string.IsNullOrEmpty(contentPart.Text)) + { + continue; + } + deltas.Publish(new TextDelta(contentPart.Text)); + fullText.Append(contentPart.Text); + } + } + + var result = fullText.ToString(); + complete.Publish(new TextComplete(result), forceFlush: true); + await streamClient.FlushAsync(); + return result; + } +} diff --git a/src/WorkflowStreams/LlmTokenStreaming/LlmWorkflow.workflow.cs b/src/WorkflowStreams/LlmTokenStreaming/LlmWorkflow.workflow.cs new file mode 100644 index 0000000..50d106f --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/LlmWorkflow.workflow.cs @@ -0,0 +1,21 @@ +namespace TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[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), + new() { StartToCloseTimeout = TimeSpan.FromMinutes(2), }); + await Workflow.DelayAsync(Constants.DrainDelay); + return result; + } +} diff --git a/src/WorkflowStreams/LlmTokenStreaming/Models.cs b/src/WorkflowStreams/LlmTokenStreaming/Models.cs new file mode 100644 index 0000000..0f5d27d --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/Models.cs @@ -0,0 +1,14 @@ +namespace TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +using Temporalio.Extensions.WorkflowStreams; + +public record LlmInput( + string Prompt, + string Model = "gpt-4o-mini", + WorkflowStreamState? StreamState = null); + +public record TextDelta(string Text); + +public record TextComplete(string FullText); + +public record RetryEvent(int Attempt); diff --git a/src/WorkflowStreams/LlmTokenStreaming/Program.cs b/src/WorkflowStreams/LlmTokenStreaming/Program.cs new file mode 100644 index 0000000..f3d71dd --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/Program.cs @@ -0,0 +1,35 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue). + AddActivity(LlmActivities.StreamCompletionAsync). + AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunLlmAsync( + client, + args.Length > 0 ? string.Join(' ', args) : null); +} diff --git a/src/WorkflowStreams/LlmTokenStreaming/README.md b/src/WorkflowStreams/LlmTokenStreaming/README.md new file mode 100644 index 0000000..9f9cac8 --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/README.md @@ -0,0 +1,18 @@ +# LLM token streaming + +An Activity streams token deltas from OpenAI through a Workflow Stream. Retry events tell the +subscriber to clear partial output and render the new attempt from scratch. + +Set `OPENAI_API_KEY`, start a Temporal service, and run the worker: + +```bash +export OPENAI_API_KEY=... +dotnet run -- worker +``` + +In another terminal, run the subscriber with an optional prompt: + +```bash +dotnet run +dotnet run -- "Explain durable execution in one paragraph." +``` diff --git a/src/WorkflowStreams/LlmTokenStreaming/Scenario.cs b/src/WorkflowStreams/LlmTokenStreaming/Scenario.cs new file mode 100644 index 0000000..4de7c2a --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/Scenario.cs @@ -0,0 +1,66 @@ +namespace TemporalioSamples.WorkflowStreams.LlmTokenStreaming; + +using Temporalio.Client; +using Temporalio.Converters; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + public static async Task RunLlmAsync(ITemporalClient client, string? prompt) + { + const string ansiSave = "\u001b[s"; + const string ansiRestoreAndClear = "\u001b[u\u001b[J"; + prompt ??= + "Write a 500-word comparison of Paxos, Raft, and Viewstamped Replication for " + + "a new distributed-systems engineer. Cover the core ideas, leader election, " + + "normal-case operation, reconfiguration, and practical implementation tradeoffs."; + + var input = new LlmInput(prompt); + var workflowId = $"workflow-streams-llm-{Guid.NewGuid()}"; + var handle = await client.StartWorkflowAsync( + (LlmWorkflow wf) => wf.RunAsync(input), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine( + $"[llm {workflowId}] streaming response from {input.Model}, awaiting first token..."); + Console.WriteLine(); + Console.Write(ansiSave); + + await using var streamClient = new WorkflowStreamClient(client, workflowId); + var options = new WorkflowStreamSubscribeOptions + { + Topics = new List + { + Constants.TopicDelta, + Constants.TopicRetry, + Constants.TopicComplete, + }, + }; + await foreach (var item in streamClient.SubscribeAsync(options)) + { + if (item.Topic == Constants.TopicRetry) + { + var evt = Decode(client, item); + Console.Write(ansiRestoreAndClear); + Console.WriteLine($"[retry attempt {evt.Attempt}] resetting output"); + Console.WriteLine(); + Console.Write(ansiSave); + } + else if (item.Topic == Constants.TopicDelta) + { + Console.Write(Decode(client, item).Text); + } + else if (item.Topic == Constants.TopicComplete) + { + _ = Decode(client, item); + Console.WriteLine(); + break; + } + } + + var result = await handle.GetResultAsync(); + Console.WriteLine($"[workflow result: {result.Length} chars]"); + } + + private static T Decode(ITemporalClient client, WorkflowStreamItem item) => + client.Options.DataConverter.PayloadConverter.ToValue(item.Payload); +} diff --git a/src/WorkflowStreams/LlmTokenStreaming/TemporalioSamples.WorkflowStreams.LlmTokenStreaming.csproj b/src/WorkflowStreams/LlmTokenStreaming/TemporalioSamples.WorkflowStreams.LlmTokenStreaming.csproj new file mode 100644 index 0000000..3879625 --- /dev/null +++ b/src/WorkflowStreams/LlmTokenStreaming/TemporalioSamples.WorkflowStreams.LlmTokenStreaming.csproj @@ -0,0 +1,12 @@ + + + + Exe + + + + + + + + diff --git a/src/WorkflowStreams/README.md b/src/WorkflowStreams/README.md new file mode 100644 index 0000000..758e013 --- /dev/null +++ b/src/WorkflowStreams/README.md @@ -0,0 +1,22 @@ +# Workflow Streams + +> **Experimental.** This sample uses `Temporalio.Extensions.WorkflowStreams`, whose API may +> change in future versions. + +A workflow stream is a durable, offset-addressed publish/subscribe log hosted inside a Temporal +Workflow. Workflow code and external clients publish to named topics through Signals, subscribers +long-poll through Updates, and a Query exposes the current global offset. The extension handles +batching, publisher deduplication, topic filtering, Continue-As-New handoff, and truncation. + +Each scenario is a self-contained project with its own worker, client, models, constants, and +instructions: + +- [Basic publish/subscribe](BasicPublishSubscribe) +- [Concurrent subscriptions](ConcurrentSubscriptions) +- [Reconnecting subscriber](ReconnectingSubscriber) +- [External publisher](ExternalPublisher) +- [Bounded log](BoundedLog) +- [LLM token streaming](LlmTokenStreaming) + +See the [repository README](../../README.md) for common prerequisites. Each scenario README has +the commands for running that project. diff --git a/src/WorkflowStreams/ReconnectingSubscriber/Constants.cs b/src/WorkflowStreams/ReconnectingSubscriber/Constants.cs new file mode 100644 index 0000000..d59b2ec --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/Constants.cs @@ -0,0 +1,10 @@ +namespace TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +public static class Constants +{ + public const string TaskQueue = "workflow-streams-reconnecting-subscriber"; + + public const string TopicStatus = "status"; + + public static readonly TimeSpan DrainDelay = TimeSpan.FromMilliseconds(500); +} diff --git a/src/WorkflowStreams/ReconnectingSubscriber/Models.cs b/src/WorkflowStreams/ReconnectingSubscriber/Models.cs new file mode 100644 index 0000000..7cb9043 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/Models.cs @@ -0,0 +1,10 @@ +namespace TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +using Temporalio.Extensions.WorkflowStreams; + +public record PipelineInput( + string PipelineId, + TimeSpan? StageInterval = null, + WorkflowStreamState? StreamState = null); + +public record StageEvent(string Stage); diff --git a/src/WorkflowStreams/ReconnectingSubscriber/PipelineWorkflow.workflow.cs b/src/WorkflowStreams/ReconnectingSubscriber/PipelineWorkflow.workflow.cs new file mode 100644 index 0000000..6c49147 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/PipelineWorkflow.workflow.cs @@ -0,0 +1,40 @@ +namespace TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Workflows; + +[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 status = stream.GetTopic(Constants.TopicStatus); + var stageInterval = input.StageInterval ?? TimeSpan.FromSeconds(2); + var stages = new[] + { + "validating", + "loading data", + "transforming", + "writing output", + "verifying", + "complete", + }; + foreach (var stage in stages) + { + status.Publish(new StageEvent(stage)); + if (stage != "complete") + { + await Workflow.DelayAsync(stageInterval); + } + } + + await Workflow.DelayAsync(Constants.DrainDelay); + return $"pipeline {input.PipelineId} done"; + } +} diff --git a/src/WorkflowStreams/ReconnectingSubscriber/Program.cs b/src/WorkflowStreams/ReconnectingSubscriber/Program.cs new file mode 100644 index 0000000..542b2f3 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/Program.cs @@ -0,0 +1,31 @@ +using Microsoft.Extensions.Logging; +using Temporalio.Client; +using Temporalio.Common.EnvConfig; +using Temporalio.Worker; +using TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +var connectOptions = ClientEnvConfig.LoadClientConnectOptions(); +connectOptions.TargetHost ??= "localhost:7233"; +connectOptions.LoggerFactory = LoggerFactory.Create(builder => + builder. + AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] "). + SetMinimumLevel(LogLevel.Information)); +var client = await TemporalClient.ConnectAsync(connectOptions); + +if (args.ElementAtOrDefault(0) == "worker") +{ + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(Constants.TaskQueue).AddWorkflow()); + await worker.ExecuteAsync(tokenSource.Token); +} +else +{ + await Scenario.RunReconnectingSubscriberAsync(client); +} diff --git a/src/WorkflowStreams/ReconnectingSubscriber/README.md b/src/WorkflowStreams/ReconnectingSubscriber/README.md new file mode 100644 index 0000000..e01e049 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/README.md @@ -0,0 +1,16 @@ +# Reconnecting subscriber + +A subscriber disconnects after two pipeline stages, saves the next offset, and resumes without +gaps or duplicates using a fresh client. + +Start a Temporal service, then run the worker: + +```bash +dotnet run -- worker +``` + +In another terminal, run the subscriber: + +```bash +dotnet run +``` diff --git a/src/WorkflowStreams/ReconnectingSubscriber/Scenario.cs b/src/WorkflowStreams/ReconnectingSubscriber/Scenario.cs new file mode 100644 index 0000000..2822bc3 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/Scenario.cs @@ -0,0 +1,52 @@ +namespace TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +using Temporalio.Client; +using Temporalio.Extensions.WorkflowStreams; + +public static class Scenario +{ + public static async Task RunReconnectingSubscriberAsync(ITemporalClient client) + { + var workflowId = $"workflow-streams-pipeline-{Guid.NewGuid()}"; + var handle = await client.StartWorkflowAsync( + (PipelineWorkflow wf) => wf.RunAsync(new PipelineInput("pipeline-7", null, null)), + new(workflowId, Constants.TaskQueue)); + Console.WriteLine($"Started workflow: {workflowId}"); + + long nextOffset = 0; + Console.WriteLine("--- phase 1: initial subscriber ---"); + await using (var streamClient = new WorkflowStreamClient(client, workflowId)) + { + var seen = 0; + await foreach (var item in streamClient. + GetTopic(Constants.TopicStatus).SubscribeAsync()) + { + var evt = item.Value; + nextOffset = item.Offset + 1; + Console.WriteLine($"offset={item.Offset} stage={evt.Stage}"); + if (++seen == 2) + { + break; + } + } + } + + Console.WriteLine($"--- disconnected; will resume from offset {nextOffset} ---"); + Console.WriteLine("--- phase 2: reconnected subscriber ---"); + await using (var streamClient = new WorkflowStreamClient(client, workflowId)) + { + await foreach (var item in streamClient. + GetTopic(Constants.TopicStatus).SubscribeAsync(nextOffset)) + { + var evt = item.Value; + Console.WriteLine($"offset={item.Offset} stage={evt.Stage}"); + if (evt.Stage == "complete") + { + break; + } + } + } + + Console.WriteLine($"Workflow result: {await handle.GetResultAsync()}"); + } +} diff --git a/src/WorkflowStreams/ReconnectingSubscriber/TemporalioSamples.WorkflowStreams.ReconnectingSubscriber.csproj b/src/WorkflowStreams/ReconnectingSubscriber/TemporalioSamples.WorkflowStreams.ReconnectingSubscriber.csproj new file mode 100644 index 0000000..c1cc2f4 --- /dev/null +++ b/src/WorkflowStreams/ReconnectingSubscriber/TemporalioSamples.WorkflowStreams.ReconnectingSubscriber.csproj @@ -0,0 +1,11 @@ + + + + Exe + + + + + + + diff --git a/tests/TemporalioSamples.Tests.csproj b/tests/TemporalioSamples.Tests.csproj index 66b6b0a..6496a55 100644 --- a/tests/TemporalioSamples.Tests.csproj +++ b/tests/TemporalioSamples.Tests.csproj @@ -40,6 +40,12 @@ + + + + + + diff --git a/tests/WorkflowStreams/WorkflowStreamsTests.cs b/tests/WorkflowStreams/WorkflowStreamsTests.cs new file mode 100644 index 0000000..85ecccd --- /dev/null +++ b/tests/WorkflowStreams/WorkflowStreamsTests.cs @@ -0,0 +1,259 @@ +namespace TemporalioSamples.Tests.WorkflowStreams; + +using Temporalio.Activities; +using Temporalio.Client; +using Temporalio.Converters; +using Temporalio.Extensions.WorkflowStreams; +using Temporalio.Worker; +using Xunit; +using Xunit.Abstractions; +using Basic = TemporalioSamples.WorkflowStreams.BasicPublishSubscribe; +using Bounded = TemporalioSamples.WorkflowStreams.BoundedLog; +using Concurrent = TemporalioSamples.WorkflowStreams.ConcurrentSubscriptions; +using External = TemporalioSamples.WorkflowStreams.ExternalPublisher; +using Llm = TemporalioSamples.WorkflowStreams.LlmTokenStreaming; +using Reconnecting = TemporalioSamples.WorkflowStreams.ReconnectingSubscriber; + +public class WorkflowStreamsTests : WorkflowEnvironmentTestBase +{ + private static readonly string[] ExpectedOrderStatuses = + ["received", "shipped", "complete"]; + + public WorkflowStreamsTests(ITestOutputHelper output, WorkflowEnvironment env) + : base(output, env) + { + } + + [Fact] + public async Task OrderWorkflow_PublishesWorkflowAndActivityEvents() + { + using var worker = new TemporalWorker( + Client, + NewWorker(). + AddActivity(Basic.PaymentActivities.ChargeCardAsync). + AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var workflowId = $"workflow-streams-order-{Guid.NewGuid()}"; + var handle = await Client.StartWorkflowAsync( + (Basic.OrderWorkflow wf) => wf.RunAsync(new Basic.OrderInput("order-42", null)), + new(workflowId, worker.Options.TaskQueue!)); + await using var streamClient = new WorkflowStreamClient(Client, workflowId); + + var statuses = new List(); + var progressCount = 0; + await foreach (var item in streamClient.SubscribeAsync(new() + { + Topics = new List + { + Basic.Constants.TopicStatus, + Basic.Constants.TopicProgress, + }, + })) + { + if (item.Topic == Basic.Constants.TopicStatus) + { + var status = Decode(item); + statuses.Add(status.Kind); + if (status.Kind == "complete") + { + break; + } + } + else if (item.Topic == Basic.Constants.TopicProgress) + { + progressCount++; + } + } + + Assert.Equal("charge-order-42", await handle.GetResultAsync()); + Assert.Equal(ExpectedOrderStatuses, statuses); + Assert.True(progressCount >= 2); + }); + } + + [Fact] + public async Task SubscriptionAsyncEnumerable_DeliversItemsInOrder() + { + using var worker = new TemporalWorker( + Client, + NewWorker(). + AddActivity(Concurrent.PaymentActivities.ChargeCardAsync). + AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var workflowId = $"workflow-streams-concurrent-{Guid.NewGuid()}"; + var handle = await Client.StartWorkflowAsync( + (Concurrent.OrderWorkflow wf) => + wf.RunAsync(new Concurrent.OrderInput("order-concurrent", null)), + new(workflowId, worker.Options.TaskQueue!)); + await using var streamClient = new WorkflowStreamClient(Client, workflowId); + var statuses = new List(); + await foreach (var item in streamClient.SubscribeAsync( + new WorkflowStreamSubscribeOptions + { + Topics = new List + { + Concurrent.Constants.TopicStatus, + Concurrent.Constants.TopicProgress, + }, + })) + { + await Task.Yield(); + if (item.Topic == Concurrent.Constants.TopicStatus) + { + statuses.Add(Decode(item).Kind); + } + } + + Assert.Equal("charge-order-concurrent", await handle.GetResultAsync()); + Assert.Equal(ExpectedOrderStatuses, statuses); + }); + } + + [Fact] + public async Task ReconnectingSubscriber_ResumesAtNextOffset() + { + using var worker = new TemporalWorker( + Client, + NewWorker().AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var workflowId = $"workflow-streams-pipeline-{Guid.NewGuid()}"; + var handle = await Client.StartWorkflowAsync( + (Reconnecting.PipelineWorkflow wf) => wf.RunAsync( + new Reconnecting.PipelineInput( + "pipeline-test", + TimeSpan.FromMilliseconds(50), + null)), + new(workflowId, worker.Options.TaskQueue!)); + + var offsets = new List(); + long nextOffset = 0; + await using (var firstClient = new WorkflowStreamClient(Client, workflowId)) + { + await foreach (var item in firstClient. + GetTopic(Reconnecting.Constants.TopicStatus). + SubscribeAsync()) + { + offsets.Add(item.Offset); + nextOffset = item.Offset + 1; + if (offsets.Count == 2) + { + break; + } + } + } + + var remainingStages = new List(); + await using (var secondClient = new WorkflowStreamClient(Client, workflowId)) + { + await foreach (var item in secondClient. + GetTopic(Reconnecting.Constants.TopicStatus). + SubscribeAsync(nextOffset)) + { + offsets.Add(item.Offset); + var stage = item.Value.Stage; + remainingStages.Add(stage); + if (stage == "complete") + { + break; + } + } + } + + Assert.Equal("pipeline pipeline-test done", await handle.GetResultAsync()); + Assert.Equal(offsets.Distinct().Count(), offsets.Count); + Assert.Equal(nextOffset, offsets[2]); + Assert.Equal("complete", remainingStages[^1]); + }); + } + + [Fact] + public async Task ExternalPublisher_PublishesAndClosesHub() + { + using var worker = new TemporalWorker( + Client, + NewWorker().AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var workflowId = $"workflow-streams-hub-{Guid.NewGuid()}"; + var handle = await Client.StartWorkflowAsync( + (External.HubWorkflow wf) => + wf.RunAsync(new External.HubInput("test-hub", null)), + new(workflowId, worker.Options.TaskQueue!)); + await using var subscriber = new WorkflowStreamClient(Client, workflowId); + await using var publisher = new WorkflowStreamClient(Client, workflowId); + publisher.GetTopic(External.Constants.TopicNews). + Publish(new External.NewsEvent("test headline"), forceFlush: true); + await publisher.FlushAsync(); + + await foreach (var item in subscriber. + GetTopic(External.Constants.TopicNews). + SubscribeAsync()) + { + Assert.Equal("test headline", item.Value.Headline); + break; + } + await handle.SignalAsync(wf => wf.CloseAsync()); + Assert.Equal("hub test-hub closed", await handle.GetResultAsync()); + }); + } + + [Fact] + public async Task TruncatingTicker_FastForwardsStaleOffset() + { + using var worker = new TemporalWorker( + Client, + NewWorker().AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var workflowId = $"workflow-streams-ticker-{Guid.NewGuid()}"; + var handle = await Client.StartWorkflowAsync( + (Bounded.TickerWorkflow wf) => wf.RunAsync( + new Bounded.TickerInput(20, 5, 5, TimeSpan.Zero, null)), + new(workflowId, worker.Options.TaskQueue!)); + await using var streamClient = new WorkflowStreamClient(Client, workflowId); + + await AssertMore.EventuallyAsync(async () => + Assert.True(await streamClient.GetOffsetAsync() >= 10)); + + await foreach (var item in streamClient. + GetTopic(Bounded.Constants.TopicTick). + SubscribeAsync(1)) + { + Assert.True(item.Offset >= 5); + break; + } + Assert.Equal("ticker emitted 20 events", await handle.GetResultAsync()); + }); + } + + [Fact] + public async Task LlmWorkflow_ReturnsMockedStreamingResult() + { + [Activity("StreamCompletion")] + static Task StreamCompletionAsync(Llm.LlmInput input) => + Task.FromResult("a streamed answer"); + + using var worker = new TemporalWorker( + Client, + NewWorker(). + AddActivity(StreamCompletionAsync). + AddWorkflow()); + await worker.ExecuteAsync(async () => + { + var result = await Client.ExecuteWorkflowAsync( + (Llm.LlmWorkflow wf) => + wf.RunAsync(new Llm.LlmInput("hello", "gpt-4o-mini", null)), + new($"workflow-streams-llm-{Guid.NewGuid()}", worker.Options.TaskQueue!)); + Assert.Equal("a streamed answer", result); + }); + } + + private TemporalWorkerOptions NewWorker() => + new($"workflow-streams-{Guid.NewGuid()}"); + + private T Decode(WorkflowStreamItem item) => + Client.Options.DataConverter.PayloadConverter.ToValue(item.Payload); +}