diff --git a/features/snippets/worker/worker.cs b/features/snippets/worker/worker.cs index 5d0f77a4..b1759025 100644 --- a/features/snippets/worker/worker.cs +++ b/features/snippets/worker/worker.cs @@ -1,5 +1,8 @@ +using Temporalio.Activities; using Temporalio.Client; +using Temporalio.Common; using Temporalio.Worker; +using Temporalio.Workflows; public class WorkerSnippet { @@ -16,4 +19,116 @@ public static async Task Run() }); // @@@SNIPEND } + + public static async Task CreateWorker() + { + var client = await TemporalClient.ConnectAsync(new("localhost:7233")); + + // @@@SNIPSTART dotnet-create-worker + var options = new TemporalWorkerOptions("my-task-queue"); + options.AddWorkflow(); + options.AddAllActivities(typeof(GreetingActivities), null); + + using var worker = new TemporalWorker(client, options); + await worker.ExecuteAsync(CancellationToken.None); + // @@@SNIPEND + } + + public static async Task CreateVersionedWorker() + { + var client = await TemporalClient.ConnectAsync(new("localhost:7233")); + + // @@@SNIPSTART dotnet-versioned-worker + var options = new TemporalWorkerOptions("my-task-queue") + { + DeploymentOptions = new WorkerDeploymentOptions( + new WorkerDeploymentVersion("my-app", "1.0"), + useWorkerVersioning: true), + }; + options.AddWorkflow(); + options.AddAllActivities(typeof(GreetingActivities), null); + + using var worker = new TemporalWorker(client, options); + // @@@SNIPEND + await Task.CompletedTask; + } + + // A Serverless Worker on GCP Cloud Run is a standard long-lived Worker that reads its + // connection settings from the environment and enables Worker Versioning. + public static async Task CloudRunWorker() + { + // @@@SNIPSTART dotnet-cloud-run-worker + var client = await TemporalClient.ConnectAsync( + new(Environment.GetEnvironmentVariable("TEMPORAL_ADDRESS")!) + { + Namespace = Environment.GetEnvironmentVariable("TEMPORAL_NAMESPACE")!, + ApiKey = Environment.GetEnvironmentVariable("TEMPORAL_API_KEY"), + Tls = new(), + }); + + var options = new TemporalWorkerOptions( + Environment.GetEnvironmentVariable("TEMPORAL_TASK_QUEUE")!) + { + DeploymentOptions = new(new("my-app", "build-1"), useWorkerVersioning: true) + { + DefaultVersioningBehavior = VersioningBehavior.Pinned, + }, + }; + options.AddWorkflow(); + options.AddAllActivities(typeof(GreetingActivities), null); + + using var worker = new TemporalWorker(client, options); + await worker.ExecuteAsync(CancellationToken.None); + // @@@SNIPEND + } + + public static async Task ShutdownWorker() + { + var client = await TemporalClient.ConnectAsync(new("localhost:7233")); + + // @@@SNIPSTART dotnet-worker-graceful-shutdown + using var tokenSource = new CancellationTokenSource(); + Console.CancelKeyPress += (_, eventArgs) => + { + tokenSource.Cancel(); + eventArgs.Cancel = true; + }; + + var options = new TemporalWorkerOptions("my-task-queue") + { + GracefulShutdownTimeout = TimeSpan.FromSeconds(30), + }; + options.AddWorkflow(); + + using var worker = new TemporalWorker(client, options); + await worker.ExecuteAsync(tokenSource.Token); + // @@@SNIPEND + } + + public static class GreetingActivities + { + [Activity] + public static string SayHello(string name) => $"Hello, {name}!"; + } + + [Workflow] + public class GreetingWorkflow + { + [WorkflowRun] + public async Task RunAsync(string name) => + await Workflow.ExecuteActivityAsync( + () => GreetingActivities.SayHello(name), + new() { StartToCloseTimeout = TimeSpan.FromSeconds(10) }); + } + + // A versioning behavior is only valid on a Worker that has versioning enabled. + [Workflow(VersioningBehavior = VersioningBehavior.Pinned)] + public class VersionedGreetingWorkflow + { + [WorkflowRun] + public async Task RunAsync(string name) => + await Workflow.ExecuteActivityAsync( + () => GreetingActivities.SayHello(name), + new() { StartToCloseTimeout = TimeSpan.FromSeconds(10) }); + } } diff --git a/features/snippets/worker/worker.java b/features/snippets/worker/worker.java index 3a3ee8c4..1489bf07 100644 --- a/features/snippets/worker/worker.java +++ b/features/snippets/worker/worker.java @@ -1,10 +1,63 @@ +import io.temporal.activity.ActivityInterface; +import io.temporal.activity.ActivityMethod; import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.common.VersioningBehavior; +import io.temporal.common.WorkerDeploymentVersion; import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; import io.temporal.worker.Worker; +import io.temporal.worker.WorkerDeploymentOptions; import io.temporal.worker.WorkerFactory; import io.temporal.worker.WorkerFactoryOptions; +import io.temporal.worker.WorkerOptions; +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; +import io.temporal.workflow.WorkflowVersioningBehavior; +import java.util.concurrent.TimeUnit; class WorkerSnippet { + @ActivityInterface + public interface GreetingActivities { + @ActivityMethod + String sayHello(String name); + } + + public static class GreetingActivitiesImpl implements GreetingActivities { + @Override + public String sayHello(String name) { + return "Hello, " + name + "!"; + } + } + + @WorkflowInterface + public interface GreetingWorkflow { + @WorkflowMethod + String greet(String name); + } + + public static class GreetingWorkflowImpl implements GreetingWorkflow { + @Override + public String greet(String name) { + return "Hello, " + name + "!"; + } + } + + @WorkflowInterface + public interface VersionedGreetingWorkflow { + @WorkflowMethod + String greet(String name); + } + + // A versioning behavior is only valid on a Worker that has versioning enabled. + public static class VersionedGreetingWorkflowImpl implements VersionedGreetingWorkflow { + @Override + @WorkflowVersioningBehavior(VersioningBehavior.PINNED) + public String greet(String name) { + return "Hello, " + name + "!"; + } + } + public static void main(String[] args) { WorkflowServiceStubs service = WorkflowServiceStubs.newLocalServiceStubs(); WorkflowClient client = WorkflowClient.newInstance(service); @@ -18,4 +71,90 @@ public static void main(String[] args) { factory.start(); } + + static void createWorker(WorkflowClient client) { + // @@@SNIPSTART java-create-worker + WorkerFactory factory = WorkerFactory.newInstance(client); + + Worker worker = factory.newWorker("my-task-queue"); + worker.registerWorkflowImplementationTypes(GreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivitiesImpl()); + + factory.start(); + // @@@SNIPEND + } + + static void createVersionedWorker(WorkflowClient client) { + WorkerFactory factory = WorkerFactory.newInstance(client); + + // @@@SNIPSTART java-versioned-worker + WorkerOptions options = + WorkerOptions.newBuilder() + .setDeploymentOptions( + WorkerDeploymentOptions.newBuilder() + .setVersion(new WorkerDeploymentVersion("my-app", "1.0")) + .setUseVersioning(true) + .build()) + .build(); + + Worker worker = factory.newWorker("my-task-queue", options); + worker.registerWorkflowImplementationTypes(VersionedGreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivitiesImpl()); + // @@@SNIPEND + + factory.start(); + } + + // A Serverless Worker on GCP Cloud Run is a standard long-lived Worker that reads its + // connection settings from the environment and enables Worker Versioning. + static void cloudRunWorker() { + // @@@SNIPSTART java-cloud-run-worker + String apiKey = System.getenv("TEMPORAL_API_KEY"); + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs( + WorkflowServiceStubsOptions.newBuilder() + .setTarget(System.getenv("TEMPORAL_ADDRESS")) + .setEnableHttps(true) + .addApiKey(() -> apiKey) + .build()); + + WorkflowClient client = + WorkflowClient.newInstance( + service, + WorkflowClientOptions.newBuilder() + .setNamespace(System.getenv("TEMPORAL_NAMESPACE")) + .build()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + + Worker worker = + factory.newWorker( + System.getenv("TEMPORAL_TASK_QUEUE"), + WorkerOptions.newBuilder() + .setDeploymentOptions( + WorkerDeploymentOptions.newBuilder() + .setUseVersioning(true) + .setVersion(new WorkerDeploymentVersion("my-app", "build-1")) + .setDefaultVersioningBehavior(VersioningBehavior.PINNED) + .build()) + .build()); + + worker.registerWorkflowImplementationTypes(GreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivitiesImpl()); + + factory.start(); + // @@@SNIPEND + } + + static void shutdownWorker(WorkflowClient client) { + WorkerFactory factory = WorkerFactory.newInstance(client); + factory.newWorker("my-task-queue"); + factory.start(); + + // @@@SNIPSTART java-worker-graceful-shutdown + factory.shutdown(); + factory.awaitTermination(30, TimeUnit.SECONDS); + // @@@SNIPEND + } } diff --git a/features/snippets/worker/worker.rb b/features/snippets/worker/worker.rb index f4ab1f5e..31faee7a 100644 --- a/features/snippets/worker/worker.rb +++ b/features/snippets/worker/worker.rb @@ -1,14 +1,39 @@ # frozen_string_literal: true +require 'temporalio/activity' require 'temporalio/client' +require 'temporalio/common_enums' require 'temporalio/worker' +require 'temporalio/worker_deployment_version' +require 'temporalio/workflow' + +class SayHello < Temporalio::Activity::Definition + def execute(name) + "Hello, #{name}!" + end +end + +class GreetingWorkflow < Temporalio::Workflow::Definition + def execute(name) + Temporalio::Workflow.execute_activity(SayHello, name, start_to_close_timeout: 10) + end +end + +# A versioning behavior is only valid on a Worker that has versioning enabled. +class VersionedGreetingWorkflow < Temporalio::Workflow::Definition + workflow_versioning_behavior Temporalio::VersioningBehavior::PINNED + + def execute(name) + Temporalio::Workflow.execute_activity(SayHello, name, start_to_close_timeout: 10) + end +end def run client = Temporalio::Client.connect( 'localhost:7233', 'default' ) - + # @@@SNIPSTART ruby-worker-max-cached-workflows worker = Temporalio::Worker.new( client: client, @@ -16,6 +41,88 @@ def run max_cached_workflows: 0 ) # @@@SNIPEND - + + worker.run +end + +def run_worker + client = Temporalio::Client.connect('localhost:7233', 'default') + + # @@@SNIPSTART ruby-create-worker + worker = Temporalio::Worker.new( + client: client, + task_queue: 'my-task-queue', + workflows: [GreetingWorkflow], + activities: [SayHello] + ) + + worker.run + # @@@SNIPEND +end + +def run_versioned_worker + client = Temporalio::Client.connect('localhost:7233', 'default') + + # @@@SNIPSTART ruby-versioned-worker + worker = Temporalio::Worker.new( + client: client, + task_queue: 'my-task-queue', + workflows: [VersionedGreetingWorkflow], + activities: [SayHello], + deployment_options: Temporalio::Worker::DeploymentOptions.new( + version: Temporalio::WorkerDeploymentVersion.new( + deployment_name: 'my-app', + build_id: '1.0' + ), + use_worker_versioning: true + ) + ) + # @@@SNIPEND + worker.run -end \ No newline at end of file +end + +# A Serverless Worker on GCP Cloud Run is a standard long-lived Worker that reads its +# connection settings from the environment and enables Worker Versioning. +def run_cloud_run_worker + # @@@SNIPSTART ruby-cloud-run-worker + client = Temporalio::Client.connect( + ENV.fetch('TEMPORAL_ADDRESS'), + ENV.fetch('TEMPORAL_NAMESPACE'), + api_key: ENV.fetch('TEMPORAL_API_KEY'), + tls: true + ) + + worker = Temporalio::Worker.new( + client:, + task_queue: ENV.fetch('TEMPORAL_TASK_QUEUE'), + workflows: [GreetingWorkflow], + activities: [SayHello], + deployment_options: Temporalio::Worker::DeploymentOptions.new( + version: Temporalio::WorkerDeploymentVersion.new( + deployment_name: 'my-app', + build_id: 'build-1' + ), + use_worker_versioning: true, + default_versioning_behavior: Temporalio::VersioningBehavior::PINNED + ) + ) + + worker.run + # @@@SNIPEND +end + +def run_worker_until_interrupted + client = Temporalio::Client.connect('localhost:7233', 'default') + + worker = Temporalio::Worker.new( + client: client, + task_queue: 'my-task-queue', + workflows: [GreetingWorkflow], + activities: [SayHello] + ) + + # @@@SNIPSTART ruby-worker-graceful-shutdown + worker.run(shutdown_signals: %w[SIGINT SIGTERM]) + # @@@SNIPEND +end diff --git a/features/snippets/worker/worker.ts b/features/snippets/worker/worker.ts index b294c9b3..fb838d12 100644 --- a/features/snippets/worker/worker.ts +++ b/features/snippets/worker/worker.ts @@ -34,6 +34,32 @@ async function _runVersioned() { // @@@SNIPEND } +// A Serverless Worker on GCP Cloud Run is a standard long-lived Worker that reads its +// connection settings from the environment and enables Worker Versioning. +async function _runCloudRunWorker() { + // @@@SNIPSTART typescript-cloud-run-worker + const connection = await NativeConnection.connect({ + address: process.env.TEMPORAL_ADDRESS, + apiKey: process.env.TEMPORAL_API_KEY, + tls: true, + }); + + const worker = await Worker.create({ + connection, + namespace: process.env.TEMPORAL_NAMESPACE!, + taskQueue: process.env.TEMPORAL_TASK_QUEUE!, + workflowsPath: require.resolve('./workflows'), + workerDeploymentOptions: { + version: { deploymentName: 'my-app', buildId: 'build-1' }, + useWorkerVersioning: true, + defaultVersioningBehavior: 'PINNED', + }, + }); + + await worker.run(); + // @@@SNIPEND +} + async function _runWithGracefulShutdown() { const connection = await NativeConnection.connect({ address: 'localhost:7233',