diff --git a/README.md b/README.md index 284e4c320..7a19472a7 100644 --- a/README.md +++ b/README.md @@ -185,6 +185,9 @@ Load client configuration from TOML files with programmatic overrides. - [**Nexus Messaging**](/core/src/main/java/io/temporal/samples/nexusmessaging): Demonstrates how to send signal, update and query messages through Nexus. This contains two samples, one sending messages to an existing workflow and a second that creates a workflow through Nexus and sends messages to it. + +- [**Nexus Messaging Temporal Operation**](/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation): Demonstrates how to send signal, update and query messages through Nexus. + This version uses the `@TemporalOperation` annotation to declare Temporal-backed Nexus operations. ### Running SpringBoot Samples diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/README.md b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/README.md index e52576f89..b369a95fc 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/README.md +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/README.md @@ -1,11 +1,11 @@ ## Caller pattern The handler worker starts a `GreetingWorkflow` for a User ID. -`NexusGreetingServiceImpl` derives the Workflow ID and routes every Nexus operation to it. +`NexusGreetingServiceImpl` holds that ID and routes every Nexus operation to it. The caller's input does not have that Workflow ID as the caller doesn't know it - but the caller sends in the User ID, -and `NexusGreetingServiceImpl` knows how to get the desired Workflow ID from that User ID (see the `getWorkflowId` call). +and `NexusGreetingServiceImpl` knows how to get the desired Workflow ID from that User ID (see the getWorkflowId call). -`HandlerWorker` is using the same `getWorkflowId` call to generate a Workflow ID from a User ID when it launches the Workflow. +HandlerWorker is using the same getWorkflowId call to generate a Workflow ID from a User ID when it launches the Workflow. The caller Workflow: 1. Queries for supported languages (`getLanguages` — backed by a `@QueryMethod`) @@ -15,23 +15,19 @@ The caller Workflow: ### Running -This sample requires a Temporal dev server build that supports Workflow Update callbacks. Download the compatible -binary from the [Temporal CLI pre-release instructions](https://docs.temporal.io/standalone-nexus-operation#temporal-cli-support). - -Start the Temporal dev server with the required namespaces pre-created and Workflow Update callbacks enabled: +Start a Temporal server: ```bash -./temporal server start-dev \ - --dynamic-config-value history.enableUpdateCallbacks=true \ - --dynamic-config-value history.enableCHASMSignalBacklinks=true \ - --namespace nexus-messaging-handler-namespace \ - --namespace nexus-messaging-caller-namespace +temporal server start-dev ``` -Create the Nexus endpoint: +Create the namespaces and Nexus endpoint: ```bash -./temporal operator nexus endpoint create \ +temporal operator namespace create --namespace nexus-messaging-handler-namespace +temporal operator namespace create --namespace nexus-messaging-caller-namespace + +temporal operator nexus endpoint create \ --name nexus-messaging-nexus-endpoint \ --target-namespace nexus-messaging-handler-namespace \ --target-task-queue nexus-messaging-handler-task-queue diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/GreetingWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/GreetingWorkflow.java index 0f266b2b2..37651ed57 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/GreetingWorkflow.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/GreetingWorkflow.java @@ -17,10 +17,6 @@ @WorkflowInterface public interface GreetingWorkflow { - // The wire name of the setLanguageUsingActivity Update, needed by the Nexus handler when it - // starts the Update through TemporalNexusClient. - String SET_LANGUAGE_USING_ACTIVITY_UPDATE = "setLanguageUsingActivity"; - @WorkflowMethod String run(); diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/NexusGreetingServiceImpl.java b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/NexusGreetingServiceImpl.java index 09dd3fc1b..2a6d04c65 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/NexusGreetingServiceImpl.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/callerpattern/handler/NexusGreetingServiceImpl.java @@ -1,13 +1,9 @@ package io.temporal.samples.nexusmessaging.callerpattern.handler; -import io.nexusrpc.OperationException; +import io.nexusrpc.handler.OperationHandler; +import io.nexusrpc.handler.OperationImpl; import io.nexusrpc.handler.ServiceImpl; -import io.temporal.client.UpdateOptions; -import io.temporal.client.WorkflowUpdateStage; -import io.temporal.nexus.TemporalNexusClient; -import io.temporal.nexus.TemporalOperation; -import io.temporal.nexus.TemporalOperationResult; -import io.temporal.nexus.TemporalOperationStartContext; +import io.temporal.nexus.Nexus; import io.temporal.samples.nexusmessaging.callerpattern.service.Language; import io.temporal.samples.nexusmessaging.callerpattern.service.NexusGreetingService; import org.slf4j.Logger; @@ -33,64 +29,51 @@ public static String getWorkflowId(String userId) { return WORKFLOW_ID_PREFIX + userId; } - private GreetingWorkflow getWorkflowStub(TemporalNexusClient client, String userId) { - return client + private GreetingWorkflow getWorkflowStub(String userId) { + return Nexus.getOperationContext() .getWorkflowClient() .newWorkflowStub(GreetingWorkflow.class, getWorkflowId(userId)); } - @TemporalOperation - public TemporalOperationResult getLanguages( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusGreetingService.GetLanguagesInput input) { - logger.info("Query for GetLanguages was received for user {}", input.getUserId()); - return TemporalOperationResult.sync( - getWorkflowStub(client, input.getUserId()).getLanguages(input)); + @OperationImpl + public OperationHandler< + NexusGreetingService.GetLanguagesInput, NexusGreetingService.GetLanguagesOutput> + getLanguages() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Query for GetLanguages was received for user {}", input.getUserId()); + return getWorkflowStub(input.getUserId()).getLanguages(input); + }); } - @TemporalOperation - public TemporalOperationResult getLanguage( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusGreetingService.GetLanguageInput input) { - logger.info("Query for GetLanguage was received for user {}", input.getUserId()); - return TemporalOperationResult.sync(getWorkflowStub(client, input.getUserId()).getLanguage()); + @OperationImpl + public OperationHandler getLanguage() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Query for GetLanguage was received for user {}", input.getUserId()); + return getWorkflowStub(input.getUserId()).getLanguage(); + }); } // Routes to setLanguageUsingActivity (not setLanguage) so that new languages not already in the // greetings map can be fetched via an activity. - @TemporalOperation - public TemporalOperationResult setLanguage( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusGreetingService.SetLanguageInput input) - throws OperationException { - logger.info("Update for SetLanguage was received for user {}", input.getUserId()); - return client.startWorkflowUpdate( - GreetingWorkflow.class, - getWorkflowId(input.getUserId()), - GreetingWorkflow::setLanguageUsingActivity, - input, - UpdateOptions.newBuilder() - .setResultClass(Language.class) - // The Update to invoke has to be named explicitly; the method reference above - // supplies the argument and result types but not the wire name. - .setUpdateName(GreetingWorkflow.SET_LANGUAGE_USING_ACTIVITY_UPDATE) - // An Update-backed Operation must wait for the ACCEPTED stage. Any other stage - // is rejected with "nexus op workflow updates only support - // WorkflowUpdateStageAccepted for async updates". - .setWaitForStage(WorkflowUpdateStage.ACCEPTED) - .build()); + @OperationImpl + public OperationHandler setLanguage() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Update for SetLanguage was received for user {}", input.getUserId()); + return getWorkflowStub(input.getUserId()).setLanguageUsingActivity(input); + }); } - @TemporalOperation - public TemporalOperationResult approve( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusGreetingService.ApproveInput input) { - logger.info("Signal for Approve was received for user {}", input.getUserId()); - getWorkflowStub(client, input.getUserId()).approve(input); - return TemporalOperationResult.sync(new NexusGreetingService.ApproveOutput()); + @OperationImpl + public OperationHandler + approve() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Signal for Approve was received for user {}", input.getUserId()); + getWorkflowStub(input.getUserId()).approve(input); + return new NexusGreetingService.ApproveOutput(); + }); } } diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/README.md b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/README.md index e05873e90..c987e85f8 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/README.md +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/README.md @@ -6,37 +6,28 @@ operations. `NexusRemoteGreetingService` adds a `runFromRemote` operation that s instance to target. The caller Workflow: -1. Attaches approval context for the first user via `attachApprovalContext`, before anything has - started that user's Workflow -2. Starts or attaches to two remote `GreetingWorkflow` instances via `runFromRemote` (backed by a Workflow started - through `TemporalNexusClient.startWorkflow`) -3. Attaches approval context for the second user, whose Workflow now already exists -4. Queries each for supported languages -5. Changes the language on each (Arabic and Hindi) -6. Confirms the changes via queries -7. Approves both Workflows -8. Waits for each to complete and returns their results +1. Starts two remote `GreetingWorkflow` instances via `runFromRemote` (backed by `WorkflowRunOperation`) +2. Queries each for supported languages +3. Changes the language on each (Arabic and Hindi) +4. Confirms the changes via queries +5. Approves both Workflows +6. Waits for each to complete and returns their results ### Running -This sample requires a Temporal dev server build that supports Workflow Update callbacks. Download the compatible -binary from the [Temporal CLI pre-release instructions](https://docs.temporal.io/standalone-nexus-operation#temporal-cli-support). - -Start the Temporal dev server with the required namespaces pre-created and Workflow Update callbacks enabled: +Start a Temporal server: ```bash -./temporal server start-dev \ - --dynamic-config-value history.enableUpdateCallbacks=true \ - --dynamic-config-value history.enableCHASMSignalBacklinks=true \ - --dynamic-config-value history.enableSignalWithStartFromWorkflow=true \ - --namespace nexus-messaging-handler-namespace \ - --namespace nexus-messaging-caller-namespace +temporal server start-dev ``` -Create the Nexus endpoint: +Create the namespaces and Nexus endpoint: ```bash -./temporal operator nexus endpoint create \ +temporal operator namespace create --namespace nexus-messaging-handler-namespace +temporal operator namespace create --namespace nexus-messaging-caller-namespace + +temporal operator nexus endpoint create \ --name nexus-messaging-nexus-endpoint \ --target-namespace nexus-messaging-handler-namespace \ --target-task-queue nexus-messaging-handler-task-queue @@ -68,10 +59,8 @@ In a third terminal, run the following command to start the example: Expected output: ``` -Attached approval context before the workflow existed: UserId One -Started remote greeting workflow: UserId One -Started remote greeting workflow: UserId Two -Attached approval context to the running workflow: UserId Two +started remote greeting workflow: UserId One +started remote greeting workflow: UserId Two Supported languages for UserId One: [CHINESE, ENGLISH] Supported languages for UserId Two: [CHINESE, ENGLISH] UserId One changed language: ENGLISH -> ARABIC diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/caller/CallerRemoteWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/caller/CallerRemoteWorkflowImpl.java index 9d7ae2725..d6767c04a 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/caller/CallerRemoteWorkflowImpl.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/caller/CallerRemoteWorkflowImpl.java @@ -54,28 +54,16 @@ public List run() { // There are examples for each of the three messaging types - // update, query, then signal. - // Attach information before the Workflow exists. Because attachApprovalContext is backed by - // Signal-with-Start on the handler, this call creates the Workflow and delivers the note to it. - greetingRemoteServiceOne.attachApprovalContext( - new NexusRemoteGreetingService.AttachApprovalContextInput( - "queued for localization review by the nightly batch", REMOTE_WORKFLOW_ONE)); - log.add("Attached approval context before the workflow existed: " + REMOTE_WORKFLOW_ONE); - logger.info("attached approval context for {}, creating the workflow", REMOTE_WORKFLOW_ONE); - - // This is an Async Nexus Operation — starts a Workflow on the handler and returns a handle. - // Unlike the sync Operations below (getLanguages, approve, etc.), this does not block until the - // Workflow completes. It is backed by TemporalNexusClient.startWorkflow on the handler side. - // - // The Workflow for this user is already running due to the call above. The handler sets the - // conflict policy to USE_EXISTING, so this call attaches the Operation's completion callback - // to the running execution. + // This is an Async Nexus operation — starts a workflow on the handler and returns a handle. + // Unlike the sync operations below (getLanguages, setLanguage, etc.), this does not block + // until the workflow completes. It is backed by WorkflowRunOperation on the handler side. NexusOperationHandle handleOne = Workflow.startNexusOperation( greetingRemoteServiceOne::runFromRemote, new NexusRemoteGreetingService.RunFromRemoteInput(REMOTE_WORKFLOW_ONE)); // Wait for the operation to be started (workflow is now running on the handler). handleOne.getExecution().get(); - log.add("Started remote greeting workflow: " + REMOTE_WORKFLOW_ONE); + log.add("started remote greeting workflow: " + REMOTE_WORKFLOW_ONE); logger.info("started remote greeting workflow {}", REMOTE_WORKFLOW_ONE); NexusOperationHandle handleTwo = @@ -84,18 +72,9 @@ public List run() { new NexusRemoteGreetingService.RunFromRemoteInput(REMOTE_WORKFLOW_TWO)); // Wait for the operation to be started (workflow is now running on the handler). handleTwo.getExecution().get(); - log.add("Started remote greeting workflow: " + REMOTE_WORKFLOW_TWO); + log.add("started remote greeting workflow: " + REMOTE_WORKFLOW_TWO); logger.info("started remote greeting workflow {}", REMOTE_WORKFLOW_TWO); - // This user's Workflow was created by runFromRemote just above, so here signalWithStart skips - // the start and only delivers the Signal. - greetingRemoteServiceTwo.attachApprovalContext( - new NexusRemoteGreetingService.AttachApprovalContextInput( - "translation approved by the localization team", REMOTE_WORKFLOW_TWO)); - log.add("Attached approval context to the running workflow: " + REMOTE_WORKFLOW_TWO); - logger.info( - "attached approval context for {}, messaging the existing workflow", REMOTE_WORKFLOW_TWO); - // Query the remote workflow for supported languages. NexusRemoteGreetingService.GetLanguagesOutput languagesOutput = greetingRemoteServiceOne.getLanguages( diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflow.java index ec0114b92..bd115c9b4 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflow.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflow.java @@ -19,10 +19,6 @@ @WorkflowInterface public interface GreetingWorkflow { - // The wire name of the setLanguageUsingActivity Update, needed by the Nexus handler when it - // starts the Update through TemporalNexusClient. - String SET_LANGUAGE_USING_ACTIVITY_UPDATE = "setLanguageUsingActivity"; - class ApproveInput { private final String name; @@ -37,20 +33,6 @@ public String getName() { } } - class AttachApprovalContextInput { - private final String note; - - @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) - public AttachApprovalContextInput(@JsonProperty("note") String note) { - this.note = note; - } - - @JsonProperty("note") - public String getNote() { - return note; - } - } - class GetLanguagesInput { private final boolean includeUnsupported; @@ -94,11 +76,6 @@ public Language getLanguage() { @SignalMethod void approve(ApproveInput input); - // Attaches supporting information for the eventual approval. Delivered with Signal-with-Start, - // so this may be the message that creates the Workflow. - @SignalMethod - void attachApprovalContext(AttachApprovalContextInput input); - // Changes the active language synchronously (only supports languages already in the greetings // map). @UpdateMethod diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflowImpl.java index d8a3c13c8..ca923da15 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflowImpl.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/GreetingWorkflowImpl.java @@ -20,7 +20,6 @@ public class GreetingWorkflowImpl implements GreetingWorkflow { private boolean approvedForRelease = false; private final Map greetings = new EnumMap<>(Language.class); private Language language = Language.ENGLISH; - private String approvalContext = null; private final GreetingActivity greetingActivity = Workflow.newActivityStub( @@ -59,16 +58,10 @@ public Language getLanguage() { @Override public void approve(ApproveInput input) { - logger.info("Approval signal received (context: {})", approvalContext); + logger.info("Approval signal received"); approvedForRelease = true; } - @Override - public void attachApprovalContext(GreetingWorkflow.AttachApprovalContextInput input) { - logger.info("attachApprovalContext signal received: {}", input.getNote()); - approvalContext = input.getNote(); - } - @Override public Language setLanguage(GreetingWorkflow.SetLanguageInput input) { logger.info("setLanguage update received"); diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java index 7c8a0c01c..b73979878 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java @@ -1,26 +1,20 @@ package io.temporal.samples.nexusmessaging.ondemandpattern.handler; -import io.nexusrpc.OperationException; +import io.nexusrpc.handler.OperationHandler; +import io.nexusrpc.handler.OperationImpl; import io.nexusrpc.handler.ServiceImpl; -import io.temporal.api.enums.v1.WorkflowIdConflictPolicy; -import io.temporal.client.BatchRequest; -import io.temporal.client.UpdateOptions; -import io.temporal.client.WorkflowClient; import io.temporal.client.WorkflowOptions; -import io.temporal.client.WorkflowUpdateStage; -import io.temporal.nexus.TemporalNexusClient; -import io.temporal.nexus.TemporalOperation; -import io.temporal.nexus.TemporalOperationResult; -import io.temporal.nexus.TemporalOperationStartContext; +import io.temporal.nexus.Nexus; +import io.temporal.nexus.WorkflowHandle; +import io.temporal.nexus.WorkflowRunOperation; import io.temporal.samples.nexusmessaging.ondemandpattern.service.Language; import io.temporal.samples.nexusmessaging.ondemandpattern.service.NexusRemoteGreetingService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** - * Nexus Operation handler for the on-demand pattern. Each Operation receives a userId, which {@link - * #getWorkflowId} maps to the target Workflow ID, and {@code runFromRemote} starts the - * GreetingWorkflow for that user. + * Nexus operation handler for the on-demand pattern. Each operation receives the target workflow ID + * in its input, and {@code runFromRemote} starts a brand-new GreetingWorkflow. */ @ServiceImpl(service = NexusRemoteGreetingService.class) public class NexusRemoteGreetingServiceImpl { @@ -38,125 +32,75 @@ public static String getWorkflowId(String userId) { return WORKFLOW_ID_PREFIX + userId; } - private GreetingWorkflow getWorkflowStub(TemporalNexusClient client, String userId) { - return client + private GreetingWorkflow getWorkflowStub(String userId) { + return Nexus.getOperationContext() .getWorkflowClient() .newWorkflowStub(GreetingWorkflow.class, getWorkflowId(userId)); } - // Starts the GreetingWorkflow for the given user, or attaches to one already running (see the - // conflict policy below). startWorkflow attaches a completion callback, so the Operation - // completes when the Workflow returns. - @TemporalOperation - public TemporalOperationResult runFromRemote( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.RunFromRemoteInput input) { - logger.info("RunFromRemote was received for userID {}", input.getUserId()); - return client.startWorkflow( - GreetingWorkflow.class, - GreetingWorkflow::run, - WorkflowOptions.newBuilder() - .setWorkflowId(getWorkflowId(input.getUserId())) - .setTaskQueue(HandlerWorker.TASK_QUEUE) - // By default, starting a Workflow whose ID is already running fails the - // Operation. Since attachApprovalContext below can create the GreetingWorkflow - // first, this Operation needs to attach to the running execution rather than - // fail. - .setWorkflowIdConflictPolicy( - WorkflowIdConflictPolicy.WORKFLOW_ID_CONFLICT_POLICY_USE_EXISTING) - .build()); + // Starts a new GreetingWorkflow with the caller-specified workflow ID. This is an async + // Nexus operation backed by WorkflowRunOperation. + @OperationImpl + public OperationHandler runFromRemote() { + return WorkflowRunOperation.fromWorkflowHandle( + (ctx, details, input) -> { + logger.info("RunFromRemote was received for userID {}", input.getUserId()); + return WorkflowHandle.fromWorkflowMethod( + Nexus.getOperationContext() + .getWorkflowClient() + .newWorkflowStub( + GreetingWorkflow.class, + WorkflowOptions.newBuilder() + .setWorkflowId(getWorkflowId(input.getUserId())) + .setTaskQueue(HandlerWorker.TASK_QUEUE) + .build()) + ::run); + }); } - @TemporalOperation - public TemporalOperationResult getLanguages( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.GetLanguagesInput input) { - logger.info("Query for GetLanguages was received for userId {}", input.getUserId()); - return TemporalOperationResult.sync( - getWorkflowStub(client, input.getUserId()) - .getLanguages(new GreetingWorkflow.GetLanguagesInput(input.isIncludeUnsupported()))); + @OperationImpl + public OperationHandler< + NexusRemoteGreetingService.GetLanguagesInput, + NexusRemoteGreetingService.GetLanguagesOutput> + getLanguages() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Query for GetLanguages was received for userId {}", input.getUserId()); + return getWorkflowStub(input.getUserId()) + .getLanguages(new GreetingWorkflow.GetLanguagesInput(input.isIncludeUnsupported())); + }); } - @TemporalOperation - public TemporalOperationResult getLanguage( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.GetLanguageInput input) { - logger.info("Query for GetLanguage was received for userId {}", input.getUserId()); - return TemporalOperationResult.sync(getWorkflowStub(client, input.getUserId()).getLanguage()); + @OperationImpl + public OperationHandler getLanguage() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Query for GetLanguage was received for userId {}", input.getUserId()); + return getWorkflowStub(input.getUserId()).getLanguage(); + }); } // Uses setLanguageUsingActivity so that new languages are fetched via an activity. - @TemporalOperation - public TemporalOperationResult setLanguage( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.SetLanguageInput input) - throws OperationException { - logger.info("Update for SetLanguage was received for userId {}", input.getUserId()); - return client.startWorkflowUpdate( - GreetingWorkflow.class, - getWorkflowId(input.getUserId()), - GreetingWorkflow::setLanguageUsingActivity, - new GreetingWorkflow.SetLanguageInput(input.getLanguage()), - UpdateOptions.newBuilder() - .setResultClass(Language.class) - // The Update to invoke has to be named explicitly; the method reference above - // supplies the argument and result types but not the wire name. - .setUpdateName(GreetingWorkflow.SET_LANGUAGE_USING_ACTIVITY_UPDATE) - // An Update-backed Operation must wait for the ACCEPTED stage. Any other stage - // is rejected with "nexus op workflow updates only support - // WorkflowUpdateStageAccepted for async updates". - .setWaitForStage(WorkflowUpdateStage.ACCEPTED) - .build()); + @OperationImpl + public OperationHandler setLanguage() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Update for SetLanguage was received for userId {}", input.getUserId()); + return getWorkflowStub(input.getUserId()) + .setLanguageUsingActivity(new GreetingWorkflow.SetLanguageInput(input.getLanguage())); + }); } - @TemporalOperation - public TemporalOperationResult approve( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.ApproveInput input) { - logger.info("Signal for Approve was received for userId {}", input.getUserId()); - getWorkflowStub(client, input.getUserId()) - .approve(new GreetingWorkflow.ApproveInput(input.getName())); - return TemporalOperationResult.sync(new NexusRemoteGreetingService.ApproveOutput()); - } - - // Signal-with-Start. Supporting information for an approval is often produced by a different - // system than the one requesting it, so the two messages can arrive in either order. This - // Operation is written so that either works, which means it may have to create the Workflow - // itself: signalWithStart delivers the Signal, starting the Workflow first if it is not already - // running. When the Workflow already exists, only the Signal is delivered. - // - // Both this and runFromRemote derive the same Workflow ID from the same userId, which is what - // lets them agree on which execution they mean regardless of which arrives first. - // - // Like approve, this is sync messaging: it completes during the handler call. - @TemporalOperation - public TemporalOperationResult attachApprovalContext( - TemporalOperationStartContext ctx, - TemporalNexusClient client, - NexusRemoteGreetingService.AttachApprovalContextInput input) { - logger.info( - "AttachApprovalContext was received for userId {}: {}", input.getUserId(), input.getNote()); - WorkflowClient workflowClient = client.getWorkflowClient(); - GreetingWorkflow stub = - workflowClient.newWorkflowStub( - GreetingWorkflow.class, - WorkflowOptions.newBuilder() - .setWorkflowId(getWorkflowId(input.getUserId())) - .setTaskQueue(HandlerWorker.TASK_QUEUE) - .build()); - - BatchRequest request = workflowClient.newSignalWithStartRequest(); - request.add( - stub::attachApprovalContext, - new GreetingWorkflow.AttachApprovalContextInput(input.getNote())); - request.add(stub::run); - workflowClient.signalWithStart(request); - - return TemporalOperationResult.sync(null); + @OperationImpl + public OperationHandler< + NexusRemoteGreetingService.ApproveInput, NexusRemoteGreetingService.ApproveOutput> + approve() { + return OperationHandler.sync( + (ctx, details, input) -> { + logger.info("Signal for Approve was received for userId {}", input.getUserId()); + getWorkflowStub(input.getUserId()) + .approve(new GreetingWorkflow.ApproveInput(input.getName())); + return new NexusRemoteGreetingService.ApproveOutput(); + }); } } diff --git a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/service/NexusRemoteGreetingService.java b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/service/NexusRemoteGreetingService.java index a17225ad8..6c6b35cb1 100644 --- a/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/service/NexusRemoteGreetingService.java +++ b/core/src/main/java/io/temporal/samples/nexusmessaging/ondemandpattern/service/NexusRemoteGreetingService.java @@ -111,28 +111,6 @@ public String getUserId() { } } - class AttachApprovalContextInput { - private final String note; - private final String userId; - - @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) - public AttachApprovalContextInput( - @JsonProperty("note") String note, @JsonProperty("userId") String userId) { - this.note = note; - this.userId = userId; - } - - @JsonProperty("note") - public String getNote() { - return note; - } - - @JsonProperty("userId") - public String getUserId() { - return userId; - } - } - class GetLanguagesOutput { private final List languages; @@ -173,11 +151,4 @@ public ApproveOutput() {} // Approves the specified workflow, allowing it to complete. @Operation ApproveOutput approve(ApproveInput input); - - // Attaches supporting information for the eventual approval. Unlike every Operation above, this - // one does not require the Workflow to already exist: the handler delivers it with - // Signal-with-Start, so the same call either messages a running GreetingWorkflow or creates one. - // That makes it safe to call before or after runFromRemote. - @Operation - void attachApprovalContext(AttachApprovalContextInput input); } diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/README.md b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/README.md new file mode 100644 index 000000000..f944dd4f6 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/README.md @@ -0,0 +1,19 @@ +This sample shows how to expose a long-running Workflow's queries, updates, and signals as Nexus operations. This +version uses the experimental `@TemporalOperation` annotation to declare Temporal-backed Nexus operations as ordinary +methods on `@ServiceImpl` classes instead of `@OperationImpl` factories that return `TemporalOperationHandler.create(...)`. + +There are two self-contained examples, each in its own directory: + +| | `callerpattern/` | `ondemandpattern/` | +|---|---|--------------------------------------------------------------| +| **Pattern** | Signal an existing Workflow | Create and run Workflows on demand, and send signals to them | +| **Who creates the Workflow?** | The handler worker starts it on boot | The caller starts it via a Nexus operation | +| **Who knows the Workflow ID?** | Only the handler | The caller chooses and passes it in every operation | +| **Nexus service** | `NexusGreetingService` | `NexusRemoteGreetingService` | + +Each directory is fully self-contained for clarity. The +`GreetingWorkflow`, `GreetingWorkflowImpl`, `GreetingActivity` and `GreetingActivityImpl` classes are pretty much the same between the two — only the +Nexus service interface and its implementation differ. This highlights that the same Workflow can be +exposed through Nexus in different ways depending on whether the caller needs lifecycle control. + +See each directory's README for running instructions. diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/README.md b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/README.md new file mode 100644 index 000000000..586c49142 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/README.md @@ -0,0 +1,69 @@ +## Caller pattern + +The handler worker starts a `GreetingWorkflow` for a User ID. +`NexusGreetingServiceImpl` derives the Workflow ID and routes every Nexus operation to it. +The caller's input does not have that Workflow ID as the caller doesn't know it - but the caller sends in the User ID, +and `NexusGreetingServiceImpl` knows how to get the desired Workflow ID from that User ID (see the `getWorkflowId` call). + +`HandlerWorker` is using the same `getWorkflowId` call to generate a Workflow ID from a User ID when it launches the Workflow. + +The caller Workflow: +1. Queries for supported languages (`getLanguages` — backed by a `@QueryMethod`) +2. Changes the language to Arabic (`setLanguage` — backed by an `@UpdateMethod` that calls an activity) +3. Confirms the change via a second query (`getLanguage`) +4. Approves the Workflow (`approve` — backed by a `@SignalMethod`) + +### Running + +This sample requires a Temporal dev server build that supports Workflow Update callbacks. Download the compatible +binary from the [Temporal CLI pre-release instructions](https://docs.temporal.io/standalone-nexus-operation#temporal-cli-support). + +Start the Temporal dev server with the required namespaces pre-created and Workflow Update callbacks enabled: + +```bash +./temporal server start-dev \ + --dynamic-config-value history.enableUpdateCallbacks=true \ + --dynamic-config-value history.enableCHASMSignalBacklinks=true \ + --namespace nexus-messaging-handler-namespace \ + --namespace nexus-messaging-caller-namespace +``` + +Create the Nexus endpoint: + +```bash +./temporal operator nexus endpoint create \ + --name nexus-messaging-nexus-endpoint \ + --target-namespace nexus-messaging-handler-namespace \ + --target-task-queue nexus-messaging-handler-task-queue +``` + +This sample loads connection settings from `ClientConfigProfile`. The +`nexus-messaging-handler` and `nexus-messaging-caller` profiles are defined in +`core/src/main/resources/config.toml`. You can override settings with environment +variables or by editing the TOML file (see the `envconfig` sample for details). + +In one terminal, start the handler worker: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler.HandlerWorker +``` + +In a second terminal, start the caller worker: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller.CallerWorker +``` + +In a third terminal, run the following command to start the example: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller.CallerStarter +``` + +Expected output: + +``` +Supported languages: [CHINESE, ENGLISH] +Language changed: ENGLISH -> ARABIC +Workflow approved +``` diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerStarter.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerStarter.java new file mode 100644 index 000000000..601408446 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerStarter.java @@ -0,0 +1,46 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller; + +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowOptions; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import java.nio.file.Paths; +import java.util.List; +import java.util.UUID; + +public class CallerStarter { + + public static void main(String[] args) { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(CallerStarter.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CallerWorker.CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + CallerWorkflow workflow = + client.newWorkflowStub( + CallerWorkflow.class, + WorkflowOptions.newBuilder() + .setWorkflowId("nexus-messaging-caller-" + UUID.randomUUID()) + .setTaskQueue(CallerWorker.TASK_QUEUE) + .build()); + + // Launch the worker, passing in an identifier which the Nexus service will use + // to find the matching workflow (See NexusGreetingServiceImpl::getWorkflowId) + List log = workflow.run("user-1"); + log.forEach(System.out::println); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorker.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorker.java new file mode 100644 index 000000000..0484eb198 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorker.java @@ -0,0 +1,58 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller; + +import io.temporal.client.WorkflowClient; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import io.temporal.worker.WorkflowImplementationOptions; +import io.temporal.workflow.NexusServiceOptions; +import java.nio.file.Paths; +import java.util.Collections; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class CallerWorker { + private static final Logger logger = LoggerFactory.getLogger(CallerWorker.class); + + static final String CONFIG_PROFILE = "nexus-messaging-caller"; + public static final String TASK_QUEUE = "nexus-messaging-caller-task-queue"; + static final String NEXUS_ENDPOINT = "nexus-messaging-nexus-endpoint"; + + public static void main(String[] args) throws InterruptedException { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(CallerWorker.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + Worker worker = factory.newWorker(TASK_QUEUE); + worker.registerWorkflowImplementationTypes( + WorkflowImplementationOptions.newBuilder() + .setNexusServiceOptions( + // The key must match the @Service-annotated interface name. + Collections.singletonMap( + "NexusGreetingService", + NexusServiceOptions.newBuilder().setEndpoint(NEXUS_ENDPOINT).build())) + .build(), + CallerWorkflowImpl.class); + + factory.start(); + logger.info("Caller worker started, ctrl+c to exit"); + Thread.currentThread().join(); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflow.java new file mode 100644 index 000000000..e9b60e2a3 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflow.java @@ -0,0 +1,11 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller; + +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; +import java.util.List; + +@WorkflowInterface +public interface CallerWorkflow { + @WorkflowMethod + List run(String userId); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflowImpl.java new file mode 100644 index 000000000..b9bb1112b --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/caller/CallerWorkflowImpl.java @@ -0,0 +1,73 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.caller; + +import io.temporal.failure.ApplicationFailure; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.NexusGreetingService; +import io.temporal.workflow.NexusOperationOptions; +import io.temporal.workflow.NexusServiceOptions; +import io.temporal.workflow.Workflow; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import org.slf4j.Logger; + +public class CallerWorkflowImpl implements CallerWorkflow { + + private static final Logger logger = Workflow.getLogger(CallerWorkflowImpl.class); + + // The endpoint is configured at the worker level in CallerWorker; only operation options are + // set here. + NexusGreetingService greetingService = + Workflow.newNexusServiceStub( + NexusGreetingService.class, + NexusServiceOptions.newBuilder() + .setOperationOptions( + NexusOperationOptions.newBuilder() + .setScheduleToCloseTimeout(Duration.ofSeconds(10)) + .build()) + .build()); + + @Override + public List run(String userId) { + + // Messages in the log array are passed back to the caller who will then log them to report what + // is happening. + // The same message is also logged for demo purposes, so that things are visible in the caller + // workflow output. + List log = new ArrayList<>(); + + // Call a Nexus operation backed by a query against the entity workflow. + // The workflow must already be running on the handler, otherwise you will + // get an error saying the workflow has already terminated. + NexusGreetingService.GetLanguagesOutput languagesOutput = + greetingService.getLanguages(new NexusGreetingService.GetLanguagesInput(false, userId)); + log.add("Supported languages: " + languagesOutput.getLanguages()); + logger.info("Supported languages: {}", languagesOutput.getLanguages()); + + // Following are examples for each of the three messaging types - + // update, query, then signal. + + // Call a Nexus operation backed by an update against the entity workflow. + Language previousLanguage = + greetingService.setLanguage( + new NexusGreetingService.SetLanguageInput(Language.ARABIC, userId)); + + // Call a Nexus operation backed by a query to confirm the language change. + Language currentLanguage = + greetingService.getLanguage(new NexusGreetingService.GetLanguageInput(userId)); + if (currentLanguage != Language.ARABIC) { + throw ApplicationFailure.newFailure( + "Expected language ARABIC, got " + currentLanguage, "AssertionError"); + } + + log.add("Language changed: " + previousLanguage.name() + " -> " + Language.ARABIC.name()); + logger.info("Language changed from {} to {}", previousLanguage, Language.ARABIC); + + // Call a Nexus operation backed by a signal against the entity workflow. + greetingService.approve(new NexusGreetingService.ApproveInput("caller", userId)); + log.add("Workflow approved"); + logger.info("Workflow approved"); + + return log; + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivity.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivity.java new file mode 100644 index 000000000..53e5e9aa2 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivity.java @@ -0,0 +1,12 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.temporal.activity.ActivityInterface; +import io.temporal.activity.ActivityMethod; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; + +@ActivityInterface +public interface GreetingActivity { + // Simulates a call to a remote greeting service. Returns null if the language is not supported. + @ActivityMethod + String callGreetingService(Language language); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivityImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivityImpl.java new file mode 100644 index 000000000..250997532 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingActivityImpl.java @@ -0,0 +1,25 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; +import java.util.EnumMap; +import java.util.Map; + +public class GreetingActivityImpl implements GreetingActivity { + + private static final Map GREETINGS = new EnumMap<>(Language.class); + + static { + GREETINGS.put(Language.ARABIC, "مرحبا بالعالم"); + GREETINGS.put(Language.CHINESE, "你好,世界"); + GREETINGS.put(Language.ENGLISH, "Hello, world"); + GREETINGS.put(Language.FRENCH, "Bonjour, monde"); + GREETINGS.put(Language.HINDI, "नमस्ते दुनिया"); + GREETINGS.put(Language.PORTUGUESE, "Olá mundo"); + GREETINGS.put(Language.SPANISH, "Hola mundo"); + } + + @Override + public String callGreetingService(Language language) { + return GREETINGS.get(language); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflow.java new file mode 100644 index 000000000..fd37328c3 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflow.java @@ -0,0 +1,51 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.NexusGreetingService; +import io.temporal.workflow.QueryMethod; +import io.temporal.workflow.SignalMethod; +import io.temporal.workflow.UpdateMethod; +import io.temporal.workflow.UpdateValidatorMethod; +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; + +/** + * A long-running "entity" workflow that backs the NexusGreetingService Nexus operations. The + * workflow exposes queries, an update, and a signal. These are private implementation details of + * the Nexus service: the caller only interacts via Nexus operations. + */ +@WorkflowInterface +public interface GreetingWorkflow { + + // The wire name of the setLanguageUsingActivity Update, needed by the Nexus handler when it + // starts the Update through TemporalNexusClient. + String SET_LANGUAGE_USING_ACTIVITY_UPDATE = "setLanguageUsingActivity"; + + @WorkflowMethod + String run(); + + // Returns the languages currently supported by the workflow. + @QueryMethod + NexusGreetingService.GetLanguagesOutput getLanguages( + NexusGreetingService.GetLanguagesInput input); + + // Returns the currently active language. + @QueryMethod + Language getLanguage(); + + // Approves the workflow, allowing it to complete. + @SignalMethod + void approve(NexusGreetingService.ApproveInput input); + + // Changes the active language synchronously (only supports languages already in the greetings + // map). + @UpdateMethod + Language setLanguage(NexusGreetingService.SetLanguageInput input); + + @UpdateValidatorMethod(updateName = "setLanguage") + void validateSetLanguage(NexusGreetingService.SetLanguageInput input); + + // Changes the active language, calling an activity to fetch a greeting for new languages. + @UpdateMethod + Language setLanguageUsingActivity(NexusGreetingService.SetLanguageInput input); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflowImpl.java new file mode 100644 index 000000000..cec73e442 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/GreetingWorkflowImpl.java @@ -0,0 +1,103 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.temporal.activity.ActivityOptions; +import io.temporal.failure.ApplicationFailure; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.NexusGreetingService; +import io.temporal.workflow.Workflow; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.EnumMap; +import java.util.List; +import java.util.Map; +import org.slf4j.Logger; + +public class GreetingWorkflowImpl implements GreetingWorkflow { + + private boolean approvedForRelease = false; + private final Map greetings = new EnumMap<>(Language.class); + private Language language = Language.ENGLISH; + + private static final Logger logger = Workflow.getLogger(GreetingWorkflowImpl.class); + + private final GreetingActivity greetingActivity = + Workflow.newActivityStub( + GreetingActivity.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(10)).build()); + + public GreetingWorkflowImpl() { + greetings.put(Language.CHINESE, "你好,世界"); + greetings.put(Language.ENGLISH, "Hello, world"); + } + + @Override + public String run() { + // Wait until approved and all in-flight update handlers have finished. + Workflow.await(() -> approvedForRelease && Workflow.isEveryHandlerFinished()); + return greetings.get(language); + } + + @Override + public NexusGreetingService.GetLanguagesOutput getLanguages( + NexusGreetingService.GetLanguagesInput input) { + List result; + if (input.isIncludeUnsupported()) { + result = new ArrayList<>(Arrays.asList(Language.values())); + } else { + result = new ArrayList<>(greetings.keySet()); + } + Collections.sort(result); + return new NexusGreetingService.GetLanguagesOutput(result); + } + + @Override + public Language getLanguage() { + return language; + } + + @Override + public void approve(NexusGreetingService.ApproveInput input) { + logger.info( + "Approval signal received for workflow {}", + NexusGreetingServiceImpl.getWorkflowId(input.getUserId())); + approvedForRelease = true; + } + + @Override + public Language setLanguage(NexusGreetingService.SetLanguageInput input) { + logger.info( + "setLanguage update received for workflow {}", + NexusGreetingServiceImpl.getWorkflowId(input.getUserId())); + Language previous = language; + language = input.getLanguage(); + return previous; + } + + @Override + public void validateSetLanguage(NexusGreetingService.SetLanguageInput input) { + logger.info( + "validateSetLanguage called for workflow {}", + NexusGreetingServiceImpl.getWorkflowId(input.getUserId())); + if (!greetings.containsKey(input.getLanguage())) { + throw new IllegalArgumentException(input.getLanguage().name() + " is not supported"); + } + } + + @Override + public Language setLanguageUsingActivity(NexusGreetingService.SetLanguageInput input) { + if (!greetings.containsKey(input.getLanguage())) { + String greeting = greetingActivity.callGreetingService(input.getLanguage()); + if (greeting == null) { + throw ApplicationFailure.newFailure( + "Greeting service does not support " + input.getLanguage().name(), + "UnsupportedLanguage"); + } + greetings.put(input.getLanguage(), greeting); + } + Language previous = language; + language = input.getLanguage(); + return previous; + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/HandlerWorker.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/HandlerWorker.java new file mode 100644 index 000000000..774118766 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/HandlerWorker.java @@ -0,0 +1,73 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowExecutionAlreadyStarted; +import io.temporal.client.WorkflowOptions; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import java.nio.file.Paths; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class HandlerWorker { + private static final Logger logger = LoggerFactory.getLogger(HandlerWorker.class); + + static final String CONFIG_PROFILE = "nexus-messaging-handler"; + public static final String TASK_QUEUE = "nexus-messaging-handler-task-queue"; + static final String USER_ID = "user-1"; + + public static void main(String[] args) throws InterruptedException { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(HandlerWorker.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + // Start the long-running entity workflow that backs the Nexus service, if not already running. + // Create a workflow ID derived from the given user ID. + // This would be for a process that would create a workflow for each UserID, + // if you had a single long running workflow for all users then you could + // remove all the USER_IDs from the inputs and just make everything refer + // to a single workflow ID. + String workflowId = NexusGreetingServiceImpl.getWorkflowId(USER_ID); + + GreetingWorkflow greetingWorkflow = + client.newWorkflowStub( + GreetingWorkflow.class, + WorkflowOptions.newBuilder() + .setWorkflowId(workflowId) + .setTaskQueue(TASK_QUEUE) + .build()); + try { + WorkflowClient.start(greetingWorkflow::run); + logger.info("Started greeting workflow: {}", workflowId); + } catch (WorkflowExecutionAlreadyStarted e) { + logger.info("Greeting workflow already running: {}", workflowId); + } + + WorkerFactory factory = WorkerFactory.newInstance(client); + Worker worker = factory.newWorker(TASK_QUEUE); + worker.registerWorkflowImplementationTypes(GreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivityImpl()); + worker.registerNexusServiceImplementation(new NexusGreetingServiceImpl()); + + factory.start(); + logger.info("Handler worker started, ctrl+c to exit"); + Thread.currentThread().join(); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/NexusGreetingServiceImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/NexusGreetingServiceImpl.java new file mode 100644 index 000000000..e16875a35 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/handler/NexusGreetingServiceImpl.java @@ -0,0 +1,96 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.handler; + +import io.nexusrpc.OperationException; +import io.nexusrpc.handler.ServiceImpl; +import io.temporal.client.UpdateOptions; +import io.temporal.client.WorkflowUpdateStage; +import io.temporal.nexus.TemporalNexusClient; +import io.temporal.nexus.TemporalOperation; +import io.temporal.nexus.TemporalOperationResult; +import io.temporal.nexus.TemporalOperationStartContext; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service.NexusGreetingService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Nexus operation handler implementation. Each operation receives a userId, which is mapped to a + * workflow ID using {@link #WORKFLOW_ID_PREFIX}. The operations are synchronous because queries and + * updates against a running workflow complete quickly. + */ +@ServiceImpl(service = NexusGreetingService.class) +public class NexusGreetingServiceImpl { + + private static final Logger logger = LoggerFactory.getLogger(NexusGreetingServiceImpl.class); + + static final String WORKFLOW_ID_PREFIX = "GreetingWorkflow_for_"; + + // This example assumes you might have multiple workflows, one for each user. + // If you had a single workflow for all users, then you could remove the + // getWorkflowId method, remove the user ID from each input, and just + // use the single workflow ID in the getWorkflowStub method below. + public static String getWorkflowId(String userId) { + return WORKFLOW_ID_PREFIX + userId; + } + + private GreetingWorkflow getWorkflowStub(TemporalNexusClient client, String userId) { + return client + .getWorkflowClient() + .newWorkflowStub(GreetingWorkflow.class, getWorkflowId(userId)); + } + + @TemporalOperation + public TemporalOperationResult getLanguages( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusGreetingService.GetLanguagesInput input) { + logger.info("Query for GetLanguages was received for user {}", input.getUserId()); + return TemporalOperationResult.sync( + getWorkflowStub(client, input.getUserId()).getLanguages(input)); + } + + @TemporalOperation + public TemporalOperationResult getLanguage( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusGreetingService.GetLanguageInput input) { + logger.info("Query for GetLanguage was received for user {}", input.getUserId()); + return TemporalOperationResult.sync(getWorkflowStub(client, input.getUserId()).getLanguage()); + } + + // Routes to setLanguageUsingActivity (not setLanguage) so that new languages not already in the + // greetings map can be fetched via an activity. + @TemporalOperation + public TemporalOperationResult setLanguage( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusGreetingService.SetLanguageInput input) + throws OperationException { + logger.info("Update for SetLanguage was received for user {}", input.getUserId()); + return client.startWorkflowUpdate( + GreetingWorkflow.class, + getWorkflowId(input.getUserId()), + GreetingWorkflow::setLanguageUsingActivity, + input, + UpdateOptions.newBuilder() + .setResultClass(Language.class) + // The Update to invoke has to be named explicitly; the method reference above + // supplies the argument and result types but not the wire name. + .setUpdateName(GreetingWorkflow.SET_LANGUAGE_USING_ACTIVITY_UPDATE) + // An Update-backed Operation must wait for the ACCEPTED stage. Any other stage + // is rejected with "nexus op workflow updates only support + // WorkflowUpdateStageAccepted for async updates". + .setWaitForStage(WorkflowUpdateStage.ACCEPTED) + .build()); + } + + @TemporalOperation + public TemporalOperationResult approve( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusGreetingService.ApproveInput input) { + logger.info("Signal for Approve was received for user {}", input.getUserId()); + getWorkflowStub(client, input.getUserId()).approve(input); + return TemporalOperationResult.sync(new NexusGreetingService.ApproveOutput()); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/Language.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/Language.java new file mode 100644 index 000000000..2818f2ece --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/Language.java @@ -0,0 +1,11 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service; + +public enum Language { + ARABIC, + CHINESE, + ENGLISH, + FRENCH, + HINDI, + PORTUGUESE, + SPANISH +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/NexusGreetingService.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/NexusGreetingService.java new file mode 100644 index 000000000..8fee77593 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/callerpattern/service/NexusGreetingService.java @@ -0,0 +1,132 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.callerpattern.service; + +import com.fasterxml.jackson.annotation.JsonAutoDetect; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.nexusrpc.Operation; +import io.nexusrpc.Service; +import java.util.List; + +/** + * Nexus service definition. Shared between the handler and caller. The caller uses this to create a + * type-safe Nexus client stub; the handler implements the operations. + */ +@Service +public interface NexusGreetingService { + + class GetLanguagesInput { + private final boolean includeUnsupported; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguagesInput( + @JsonProperty("includeUnsupported") boolean includeUnsupported, + @JsonProperty("userId") String userId) { + this.includeUnsupported = includeUnsupported; + this.userId = userId; + } + + @JsonProperty("includeUnsupported") + public boolean isIncludeUnsupported() { + return includeUnsupported; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class GetLanguagesOutput { + private final List languages; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguagesOutput(@JsonProperty("languages") List languages) { + this.languages = languages; + } + + @JsonProperty("languages") + public List getLanguages() { + return languages; + } + } + + class GetLanguageInput { + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguageInput(@JsonProperty("userId") String userId) { + this.userId = userId; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class ApproveInput { + private final String name; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public ApproveInput(@JsonProperty("name") String name, @JsonProperty("userId") String userId) { + this.name = name; + this.userId = userId; + } + + @JsonProperty("name") + public String getName() { + return name; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + @JsonAutoDetect(fieldVisibility = JsonAutoDetect.Visibility.ANY) + class ApproveOutput { + @JsonCreator + public ApproveOutput() {} + } + + class SetLanguageInput { + private final Language language; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public SetLanguageInput( + @JsonProperty("language") Language language, @JsonProperty("userId") String userId) { + this.language = language; + this.userId = userId; + } + + @JsonProperty("language") + public Language getLanguage() { + return language; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + // Returns the languages supported by the greeting workflow. + @Operation + GetLanguagesOutput getLanguages(GetLanguagesInput input); + + // Returns the currently active language. + @Operation + Language getLanguage(GetLanguageInput input); + + // Changes the active language, returning the previous one. + @Operation + Language setLanguage(SetLanguageInput input); + + // Approves the workflow, allowing it to complete. + @Operation + ApproveOutput approve(ApproveInput input); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/README.md b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/README.md new file mode 100644 index 000000000..5e16ca058 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/README.md @@ -0,0 +1,82 @@ +## On-demand pattern + +No Workflow is pre-started. The caller creates and controls Workflow instances through Nexus +operations. `NexusRemoteGreetingService` adds a `runFromRemote` operation that starts a new +`GreetingWorkflow`, and every other operation includes a `workflowId` so the handler knows which +instance to target. + +The caller Workflow: +1. Attaches approval context for the first user via `attachApprovalContext`, before anything has + started that user's Workflow +2. Starts or attaches to two remote `GreetingWorkflow` instances via `runFromRemote` (backed by a Workflow started + through `TemporalNexusClient.startWorkflow`) +3. Attaches approval context for the second user, whose Workflow now already exists +4. Queries each for supported languages +5. Changes the language on each (Arabic and Hindi) +6. Confirms the changes via queries +7. Approves both Workflows +8. Waits for each to complete and returns their results + +### Running + +This sample requires a Temporal dev server build that supports Workflow Update callbacks. Download the compatible +binary from the [Temporal CLI pre-release instructions](https://docs.temporal.io/standalone-nexus-operation#temporal-cli-support). + +Start the Temporal dev server with the required namespaces pre-created and Workflow Update callbacks enabled: + +```bash +./temporal server start-dev \ + --dynamic-config-value history.enableUpdateCallbacks=true \ + --dynamic-config-value history.enableCHASMSignalBacklinks=true \ + --dynamic-config-value history.enableSignalWithStartFromWorkflow=true \ + --namespace nexus-messaging-handler-namespace \ + --namespace nexus-messaging-caller-namespace +``` + +Create the Nexus endpoint: + +```bash +./temporal operator nexus endpoint create \ + --name nexus-messaging-nexus-endpoint \ + --target-namespace nexus-messaging-handler-namespace \ + --target-task-queue nexus-messaging-handler-task-queue +``` + +This sample loads connection settings from `ClientConfigProfile`. The +`nexus-messaging-handler` and `nexus-messaging-caller` profiles are defined in +`core/src/main/resources/config.toml`. You can override settings with environment +variables or by editing the TOML file (see the `envconfig` sample for details). + +In one terminal, start the handler worker: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler.HandlerWorker +``` + +In a second terminal, start the caller worker: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller.CallerRemoteWorker +``` + +In a third terminal, run the following command to start the example: + +```bash +./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller.CallerRemoteStarter +``` + +Expected output: + +``` +Attached approval context before the workflow existed: UserId One +Started remote greeting workflow: UserId One +Started remote greeting workflow: UserId Two +Attached approval context to the running workflow: UserId Two +Supported languages for UserId One: [CHINESE, ENGLISH] +Supported languages for UserId Two: [CHINESE, ENGLISH] +UserId One changed language: ENGLISH -> ARABIC +UserId Two changed language: ENGLISH -> HINDI +Workflows approved +Workflow one result: مرحبا بالعالم +Workflow two result: नमस्ते दुनिया +``` diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteStarter.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteStarter.java new file mode 100644 index 000000000..c43cbeeee --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteStarter.java @@ -0,0 +1,44 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller; + +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowOptions; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import java.nio.file.Paths; +import java.util.List; +import java.util.UUID; + +public class CallerRemoteStarter { + + public static void main(String[] args) { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(CallerRemoteStarter.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CallerRemoteWorker.CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + CallerRemoteWorkflow workflow = + client.newWorkflowStub( + CallerRemoteWorkflow.class, + WorkflowOptions.newBuilder() + .setWorkflowId("nexus-messaging-remote-caller-" + UUID.randomUUID()) + .setTaskQueue(CallerRemoteWorker.TASK_QUEUE) + .build()); + + List log = workflow.run(); + log.forEach(System.out::println); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorker.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorker.java new file mode 100644 index 000000000..63f29b241 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorker.java @@ -0,0 +1,58 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller; + +import io.temporal.client.WorkflowClient; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import io.temporal.worker.WorkflowImplementationOptions; +import io.temporal.workflow.NexusServiceOptions; +import java.nio.file.Paths; +import java.util.Collections; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class CallerRemoteWorker { + private static final Logger logger = LoggerFactory.getLogger(CallerRemoteWorker.class); + + static final String CONFIG_PROFILE = "nexus-messaging-caller"; + public static final String TASK_QUEUE = "nexus-messaging-caller-remote-task-queue"; + static final String NEXUS_ENDPOINT = "nexus-messaging-nexus-endpoint"; + + public static void main(String[] args) throws InterruptedException { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(CallerRemoteWorker.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + Worker worker = factory.newWorker(TASK_QUEUE); + worker.registerWorkflowImplementationTypes( + WorkflowImplementationOptions.newBuilder() + .setNexusServiceOptions( + // The key must match the @Service-annotated interface name. + Collections.singletonMap( + "NexusRemoteGreetingService", + NexusServiceOptions.newBuilder().setEndpoint(NEXUS_ENDPOINT).build())) + .build(), + CallerRemoteWorkflowImpl.class); + + factory.start(); + logger.info("Caller remote worker started, ctrl+c to exit"); + Thread.currentThread().join(); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflow.java new file mode 100644 index 000000000..3fcd5e214 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflow.java @@ -0,0 +1,11 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller; + +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; +import java.util.List; + +@WorkflowInterface +public interface CallerRemoteWorkflow { + @WorkflowMethod + List run(); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflowImpl.java new file mode 100644 index 000000000..0e65c62f8 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/caller/CallerRemoteWorkflowImpl.java @@ -0,0 +1,175 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.caller; + +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.NexusRemoteGreetingService; +import io.temporal.workflow.NexusOperationHandle; +import io.temporal.workflow.NexusOperationOptions; +import io.temporal.workflow.NexusServiceOptions; +import io.temporal.workflow.Workflow; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import org.slf4j.Logger; + +public class CallerRemoteWorkflowImpl implements CallerRemoteWorkflow { + + private static final Logger logger = Workflow.getLogger(CallerRemoteWorkflowImpl.class); + + // This is going to create two workflows and send messages to them. + // We need to have an ID to differentiate so that Nexus knows how to name + // a workflow and then how to know the correct destination workflow. + private static final String REMOTE_WORKFLOW_ONE = "UserId One"; + private static final String REMOTE_WORKFLOW_TWO = "UserId Two"; + + NexusRemoteGreetingService greetingRemoteServiceOne = + Workflow.newNexusServiceStub( + NexusRemoteGreetingService.class, + NexusServiceOptions.newBuilder() + .setOperationOptions( + NexusOperationOptions.newBuilder() + .setScheduleToCloseTimeout(Duration.ofSeconds(10)) + .build()) + .build()); + NexusRemoteGreetingService greetingRemoteServiceTwo = + Workflow.newNexusServiceStub( + NexusRemoteGreetingService.class, + NexusServiceOptions.newBuilder() + .setOperationOptions( + NexusOperationOptions.newBuilder() + .setScheduleToCloseTimeout(Duration.ofSeconds(10)) + .build()) + .build()); + + @Override + public List run() { + // Messages in the log array are passed back to the caller who will then log them to report what + // is happening. + // The same message is also logged for demo purposes, so that things are visible in the caller + // workflow output. + List log = new ArrayList<>(); + + // Each call is performed twice in this example. This assumes there are two users we want + // to process. The first call starts two workflows, one for each user. + // Subsequent calls perform different actions between the two users. + // There are examples for each of the three messaging types - + // update, query, then signal. + + // Attach information before the Workflow exists. Because attachApprovalContext is backed by + // Signal-with-Start on the handler, this call creates the Workflow and delivers the note to it. + greetingRemoteServiceOne.attachApprovalContext( + new NexusRemoteGreetingService.AttachApprovalContextInput( + "queued for localization review by the nightly batch", REMOTE_WORKFLOW_ONE)); + log.add("Attached approval context before the workflow existed: " + REMOTE_WORKFLOW_ONE); + logger.info("attached approval context for {}, creating the workflow", REMOTE_WORKFLOW_ONE); + + // This is an Async Nexus Operation — starts a Workflow on the handler and returns a handle. + // Unlike the sync Operations below (getLanguages, approve, etc.), this does not block until the + // Workflow completes. It is backed by TemporalNexusClient.startWorkflow on the handler side. + // + // The Workflow for this user is already running due to the call above. The handler sets the + // conflict policy to USE_EXISTING, so this call attaches the Operation's completion callback + // to the running execution. + NexusOperationHandle handleOne = + Workflow.startNexusOperation( + greetingRemoteServiceOne::runFromRemote, + new NexusRemoteGreetingService.RunFromRemoteInput(REMOTE_WORKFLOW_ONE)); + // Wait for the operation to be started (workflow is now running on the handler). + handleOne.getExecution().get(); + log.add("Started remote greeting workflow: " + REMOTE_WORKFLOW_ONE); + logger.info("started remote greeting workflow {}", REMOTE_WORKFLOW_ONE); + + NexusOperationHandle handleTwo = + Workflow.startNexusOperation( + greetingRemoteServiceTwo::runFromRemote, + new NexusRemoteGreetingService.RunFromRemoteInput(REMOTE_WORKFLOW_TWO)); + // Wait for the operation to be started (workflow is now running on the handler). + handleTwo.getExecution().get(); + log.add("Started remote greeting workflow: " + REMOTE_WORKFLOW_TWO); + logger.info("started remote greeting workflow {}", REMOTE_WORKFLOW_TWO); + + // This user's Workflow was created by runFromRemote just above, so here signalWithStart skips + // the start and only delivers the Signal. + greetingRemoteServiceTwo.attachApprovalContext( + new NexusRemoteGreetingService.AttachApprovalContextInput( + "translation approved by the localization team", REMOTE_WORKFLOW_TWO)); + log.add("Attached approval context to the running workflow: " + REMOTE_WORKFLOW_TWO); + logger.info( + "attached approval context for {}, messaging the existing workflow", REMOTE_WORKFLOW_TWO); + + // Query the remote workflow for supported languages. + NexusRemoteGreetingService.GetLanguagesOutput languagesOutput = + greetingRemoteServiceOne.getLanguages( + new NexusRemoteGreetingService.GetLanguagesInput(false, REMOTE_WORKFLOW_ONE)); + log.add( + "Supported languages for " + REMOTE_WORKFLOW_ONE + ": " + languagesOutput.getLanguages()); + logger.info( + "supported languages are {} for workflow {}", + languagesOutput.getLanguages(), + REMOTE_WORKFLOW_ONE); + + languagesOutput = + greetingRemoteServiceTwo.getLanguages( + new NexusRemoteGreetingService.GetLanguagesInput(false, REMOTE_WORKFLOW_TWO)); + log.add( + "Supported languages for " + REMOTE_WORKFLOW_TWO + ": " + languagesOutput.getLanguages()); + logger.info( + "supported languages are {} for workflow {}", + languagesOutput.getLanguages(), + REMOTE_WORKFLOW_TWO); + + // Update the language on the remote workflow. + Language previousLanguageOne = + greetingRemoteServiceOne.setLanguage( + new NexusRemoteGreetingService.SetLanguageInput(Language.ARABIC, REMOTE_WORKFLOW_ONE)); + + Language previousLanguageTwo = + greetingRemoteServiceTwo.setLanguage( + new NexusRemoteGreetingService.SetLanguageInput(Language.HINDI, REMOTE_WORKFLOW_TWO)); + + // Confirm the change by querying. + Language currentLanguage = + greetingRemoteServiceOne.getLanguage( + new NexusRemoteGreetingService.GetLanguageInput(REMOTE_WORKFLOW_ONE)); + log.add( + REMOTE_WORKFLOW_ONE + + " changed language: " + + previousLanguageOne.name() + + " -> " + + currentLanguage.name()); + logger.info( + "Language changed from {} to {} for workflow {}", + previousLanguageOne, + currentLanguage, + REMOTE_WORKFLOW_ONE); + + currentLanguage = + greetingRemoteServiceTwo.getLanguage( + new NexusRemoteGreetingService.GetLanguageInput(REMOTE_WORKFLOW_TWO)); + log.add( + REMOTE_WORKFLOW_TWO + + " changed language: " + + previousLanguageTwo.name() + + " -> " + + currentLanguage.name()); + logger.info( + "Language changed from {} to {} for workflow {}", + previousLanguageTwo, + currentLanguage, + REMOTE_WORKFLOW_TWO); + + // Approve the remote workflow so it can complete. + greetingRemoteServiceOne.approve( + new NexusRemoteGreetingService.ApproveInput("remote-caller", REMOTE_WORKFLOW_ONE)); + greetingRemoteServiceTwo.approve( + new NexusRemoteGreetingService.ApproveInput("remote-caller", REMOTE_WORKFLOW_TWO)); + log.add("Workflows approved"); + + // Wait for the remote workflow to finish and return its result. + String result = handleOne.getResult().get(); + log.add("Workflow one result: " + result); + + result = handleTwo.getResult().get(); + log.add("Workflow two result: " + result); + return log; + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivity.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivity.java new file mode 100644 index 000000000..6ab7b0023 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivity.java @@ -0,0 +1,12 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import io.temporal.activity.ActivityInterface; +import io.temporal.activity.ActivityMethod; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; + +@ActivityInterface +public interface GreetingActivity { + // Simulates a call to a remote greeting service. Returns null if the language is not supported. + @ActivityMethod + String callGreetingService(Language language); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivityImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivityImpl.java new file mode 100644 index 000000000..64b805a9b --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingActivityImpl.java @@ -0,0 +1,25 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; +import java.util.EnumMap; +import java.util.Map; + +public class GreetingActivityImpl implements GreetingActivity { + + private static final Map GREETINGS = new EnumMap<>(Language.class); + + static { + GREETINGS.put(Language.ARABIC, "مرحبا بالعالم"); + GREETINGS.put(Language.CHINESE, "你好,世界"); + GREETINGS.put(Language.ENGLISH, "Hello, world"); + GREETINGS.put(Language.FRENCH, "Bonjour, monde"); + GREETINGS.put(Language.HINDI, "नमस्ते दुनिया"); + GREETINGS.put(Language.PORTUGUESE, "Olá mundo"); + GREETINGS.put(Language.SPANISH, "Hola mundo"); + } + + @Override + public String callGreetingService(Language language) { + return GREETINGS.get(language); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflow.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflow.java new file mode 100644 index 000000000..990cb375d --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflow.java @@ -0,0 +1,113 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.NexusRemoteGreetingService; +import io.temporal.workflow.QueryMethod; +import io.temporal.workflow.SignalMethod; +import io.temporal.workflow.UpdateMethod; +import io.temporal.workflow.UpdateValidatorMethod; +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; + +/** + * A long-running "entity" workflow that backs the NexusRemoteGreetingService Nexus operations. The + * workflow exposes queries, an update, and a signal. These are private implementation details of + * the Nexus service: the caller only interacts via Nexus operations. + */ +@WorkflowInterface +public interface GreetingWorkflow { + + // The wire name of the setLanguageUsingActivity Update, needed by the Nexus handler when it + // starts the Update through TemporalNexusClient. + String SET_LANGUAGE_USING_ACTIVITY_UPDATE = "setLanguageUsingActivity"; + + class ApproveInput { + private final String name; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public ApproveInput(@JsonProperty("name") String name) { + this.name = name; + } + + @JsonProperty("name") + public String getName() { + return name; + } + } + + class AttachApprovalContextInput { + private final String note; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public AttachApprovalContextInput(@JsonProperty("note") String note) { + this.note = note; + } + + @JsonProperty("note") + public String getNote() { + return note; + } + } + + class GetLanguagesInput { + private final boolean includeUnsupported; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguagesInput(@JsonProperty("includeUnsupported") boolean includeUnsupported) { + this.includeUnsupported = includeUnsupported; + } + + @JsonProperty("includeUnsupported") + public boolean isIncludeUnsupported() { + return includeUnsupported; + } + } + + class SetLanguageInput { + private final Language language; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public SetLanguageInput(@JsonProperty("language") Language language) { + this.language = language; + } + + @JsonProperty("language") + public Language getLanguage() { + return language; + } + } + + @WorkflowMethod + String run(); + + // Returns the languages currently supported by the workflow. + @QueryMethod + NexusRemoteGreetingService.GetLanguagesOutput getLanguages(GetLanguagesInput input); + + // Returns the currently active language. + @QueryMethod + Language getLanguage(); + + // Approves the workflow, allowing it to complete. + @SignalMethod + void approve(ApproveInput input); + + // Attaches supporting information for the eventual approval. Delivered with Signal-with-Start, + // so this may be the message that creates the Workflow. + @SignalMethod + void attachApprovalContext(AttachApprovalContextInput input); + + // Changes the active language synchronously (only supports languages already in the greetings + // map). + @UpdateMethod + Language setLanguage(SetLanguageInput input); + + @UpdateValidatorMethod(updateName = "setLanguage") + void validateSetLanguage(SetLanguageInput input); + + // Changes the active language, calling an activity to fetch a greeting for new languages. + @UpdateMethod + Language setLanguageUsingActivity(SetLanguageInput input); +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflowImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflowImpl.java new file mode 100644 index 000000000..a29a86240 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/GreetingWorkflowImpl.java @@ -0,0 +1,103 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import io.temporal.activity.ActivityOptions; +import io.temporal.failure.ApplicationFailure; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.NexusRemoteGreetingService; +import io.temporal.workflow.Workflow; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.EnumMap; +import java.util.List; +import java.util.Map; +import org.slf4j.Logger; + +public class GreetingWorkflowImpl implements GreetingWorkflow { + + private static final Logger logger = Workflow.getLogger(GreetingWorkflowImpl.class); + private boolean approvedForRelease = false; + private final Map greetings = new EnumMap<>(Language.class); + private Language language = Language.ENGLISH; + private String approvalContext = null; + + private final GreetingActivity greetingActivity = + Workflow.newActivityStub( + GreetingActivity.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(10)).build()); + + public GreetingWorkflowImpl() { + greetings.put(Language.CHINESE, "你好,世界"); + greetings.put(Language.ENGLISH, "Hello, world"); + } + + @Override + public String run() { + // Wait until approved and all in-flight update handlers have finished. + Workflow.await(() -> approvedForRelease && Workflow.isEveryHandlerFinished()); + return greetings.get(language); + } + + @Override + public NexusRemoteGreetingService.GetLanguagesOutput getLanguages( + GreetingWorkflow.GetLanguagesInput input) { + List result; + if (input.isIncludeUnsupported()) { + result = new ArrayList<>(Arrays.asList(Language.values())); + } else { + result = new ArrayList<>(greetings.keySet()); + } + Collections.sort(result); + return new NexusRemoteGreetingService.GetLanguagesOutput(result); + } + + @Override + public Language getLanguage() { + return language; + } + + @Override + public void approve(ApproveInput input) { + logger.info("Approval signal received (context: {})", approvalContext); + approvedForRelease = true; + } + + @Override + public void attachApprovalContext(GreetingWorkflow.AttachApprovalContextInput input) { + logger.info("attachApprovalContext signal received: {}", input.getNote()); + approvalContext = input.getNote(); + } + + @Override + public Language setLanguage(GreetingWorkflow.SetLanguageInput input) { + logger.info("setLanguage update received"); + Language previous = language; + language = input.getLanguage(); + return previous; + } + + @Override + public void validateSetLanguage(GreetingWorkflow.SetLanguageInput input) { + logger.info("validateSetLanguage called"); + if (!greetings.containsKey(input.getLanguage())) { + throw new IllegalArgumentException(input.getLanguage().name() + " is not supported"); + } + } + + @Override + public Language setLanguageUsingActivity(GreetingWorkflow.SetLanguageInput input) { + if (!greetings.containsKey(input.getLanguage())) { + String greeting = greetingActivity.callGreetingService(input.getLanguage()); + if (greeting == null) { + throw ApplicationFailure.newFailure( + "Greeting service does not support " + input.getLanguage().name(), + "UnsupportedLanguage"); + } + greetings.put(input.getLanguage(), greeting); + } + Language previous = language; + language = input.getLanguage(); + return previous; + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/HandlerWorker.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/HandlerWorker.java new file mode 100644 index 000000000..dba1ef3c2 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/HandlerWorker.java @@ -0,0 +1,48 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import io.temporal.client.WorkflowClient; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.envconfig.LoadClientConfigProfileOptions; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import java.nio.file.Paths; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class HandlerWorker { + private static final Logger logger = LoggerFactory.getLogger(HandlerWorker.class); + + static final String CONFIG_PROFILE = "nexus-messaging-handler"; + public static final String TASK_QUEUE = "nexus-messaging-handler-task-queue"; + + public static void main(String[] args) throws InterruptedException { + ClientConfigProfile profile; + try { + String configFilePath = + Paths.get(HandlerWorker.class.getResource("/config.toml").toURI()).toString(); + profile = + ClientConfigProfile.load( + LoadClientConfigProfileOptions.newBuilder() + .setConfigFilePath(configFilePath) + .setConfigFileProfile(CONFIG_PROFILE) + .build()); + } catch (Exception e) { + throw new RuntimeException("Failed to load client configuration", e); + } + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs(profile.toWorkflowServiceStubsOptions()); + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + Worker worker = factory.newWorker(TASK_QUEUE); + worker.registerWorkflowImplementationTypes(GreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivityImpl()); + worker.registerNexusServiceImplementation(new NexusRemoteGreetingServiceImpl()); + + factory.start(); + logger.info("Handler worker started, ctrl+c to exit"); + Thread.currentThread().join(); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java new file mode 100644 index 000000000..fa9f5a65f --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/handler/NexusRemoteGreetingServiceImpl.java @@ -0,0 +1,162 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.handler; + +import io.nexusrpc.OperationException; +import io.nexusrpc.handler.ServiceImpl; +import io.temporal.api.enums.v1.WorkflowIdConflictPolicy; +import io.temporal.client.BatchRequest; +import io.temporal.client.UpdateOptions; +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowOptions; +import io.temporal.client.WorkflowUpdateStage; +import io.temporal.nexus.TemporalNexusClient; +import io.temporal.nexus.TemporalOperation; +import io.temporal.nexus.TemporalOperationResult; +import io.temporal.nexus.TemporalOperationStartContext; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.Language; +import io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service.NexusRemoteGreetingService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Nexus Operation handler for the on-demand pattern. Each Operation receives a userId, which {@link + * #getWorkflowId} maps to the target Workflow ID, and {@code runFromRemote} starts the + * GreetingWorkflow for that user. + */ +@ServiceImpl(service = NexusRemoteGreetingService.class) +public class NexusRemoteGreetingServiceImpl { + + private static final Logger logger = + LoggerFactory.getLogger(NexusRemoteGreetingServiceImpl.class); + + static final String WORKFLOW_ID_PREFIX = "GreetingWorkflow_for_"; + + // This example assumes you might have multiple workflows, one for each user. + // If you had a single workflow for all users, then you could remove the + // getWorkflowId method, remove the user ID from each input, and just + // use the single workflow ID in the getWorkflowStub method below. + public static String getWorkflowId(String userId) { + return WORKFLOW_ID_PREFIX + userId; + } + + private GreetingWorkflow getWorkflowStub(TemporalNexusClient client, String userId) { + return client + .getWorkflowClient() + .newWorkflowStub(GreetingWorkflow.class, getWorkflowId(userId)); + } + + // Starts the GreetingWorkflow for the given user, or attaches to one already running (see the + // conflict policy below). startWorkflow attaches a completion callback, so the Operation + // completes when the Workflow returns. + @TemporalOperation + public TemporalOperationResult runFromRemote( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.RunFromRemoteInput input) { + logger.info("RunFromRemote was received for userID {}", input.getUserId()); + return client.startWorkflow( + GreetingWorkflow.class, + GreetingWorkflow::run, + WorkflowOptions.newBuilder() + .setWorkflowId(getWorkflowId(input.getUserId())) + .setTaskQueue(HandlerWorker.TASK_QUEUE) + // By default, starting a Workflow whose ID is already running fails the + // Operation. Since attachApprovalContext below can create the GreetingWorkflow + // first, this Operation needs to attach to the running execution rather than + // fail. + .setWorkflowIdConflictPolicy( + WorkflowIdConflictPolicy.WORKFLOW_ID_CONFLICT_POLICY_USE_EXISTING) + .build()); + } + + @TemporalOperation + public TemporalOperationResult getLanguages( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.GetLanguagesInput input) { + logger.info("Query for GetLanguages was received for userId {}", input.getUserId()); + return TemporalOperationResult.sync( + getWorkflowStub(client, input.getUserId()) + .getLanguages(new GreetingWorkflow.GetLanguagesInput(input.isIncludeUnsupported()))); + } + + @TemporalOperation + public TemporalOperationResult getLanguage( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.GetLanguageInput input) { + logger.info("Query for GetLanguage was received for userId {}", input.getUserId()); + return TemporalOperationResult.sync(getWorkflowStub(client, input.getUserId()).getLanguage()); + } + + // Uses setLanguageUsingActivity so that new languages are fetched via an activity. + @TemporalOperation + public TemporalOperationResult setLanguage( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.SetLanguageInput input) + throws OperationException { + logger.info("Update for SetLanguage was received for userId {}", input.getUserId()); + return client.startWorkflowUpdate( + GreetingWorkflow.class, + getWorkflowId(input.getUserId()), + GreetingWorkflow::setLanguageUsingActivity, + new GreetingWorkflow.SetLanguageInput(input.getLanguage()), + UpdateOptions.newBuilder() + .setResultClass(Language.class) + // The Update to invoke has to be named explicitly; the method reference above + // supplies the argument and result types but not the wire name. + .setUpdateName(GreetingWorkflow.SET_LANGUAGE_USING_ACTIVITY_UPDATE) + // An Update-backed Operation must wait for the ACCEPTED stage. Any other stage + // is rejected with "nexus op workflow updates only support + // WorkflowUpdateStageAccepted for async updates". + .setWaitForStage(WorkflowUpdateStage.ACCEPTED) + .build()); + } + + @TemporalOperation + public TemporalOperationResult approve( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.ApproveInput input) { + logger.info("Signal for Approve was received for userId {}", input.getUserId()); + getWorkflowStub(client, input.getUserId()) + .approve(new GreetingWorkflow.ApproveInput(input.getName())); + return TemporalOperationResult.sync(new NexusRemoteGreetingService.ApproveOutput()); + } + + // Signal-with-Start. Supporting information for an approval is often produced by a different + // system than the one requesting it, so the two messages can arrive in either order. This + // Operation is written so that either works, which means it may have to create the Workflow + // itself: signalWithStart delivers the Signal, starting the Workflow first if it is not already + // running. When the Workflow already exists, only the Signal is delivered. + // + // Both this and runFromRemote derive the same Workflow ID from the same userId, which is what + // lets them agree on which execution they mean regardless of which arrives first. + // + // Like approve, this is sync messaging: it completes during the handler call. + @TemporalOperation + public TemporalOperationResult attachApprovalContext( + TemporalOperationStartContext ctx, + TemporalNexusClient client, + NexusRemoteGreetingService.AttachApprovalContextInput input) { + logger.info( + "AttachApprovalContext was received for userId {}: {}", input.getUserId(), input.getNote()); + WorkflowClient workflowClient = client.getWorkflowClient(); + GreetingWorkflow stub = + workflowClient.newWorkflowStub( + GreetingWorkflow.class, + WorkflowOptions.newBuilder() + .setWorkflowId(getWorkflowId(input.getUserId())) + .setTaskQueue(HandlerWorker.TASK_QUEUE) + .build()); + + BatchRequest request = workflowClient.newSignalWithStartRequest(); + request.add( + stub::attachApprovalContext, + new GreetingWorkflow.AttachApprovalContextInput(input.getNote())); + request.add(stub::run); + workflowClient.signalWithStart(request); + + return TemporalOperationResult.sync(null); + } +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/Language.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/Language.java new file mode 100644 index 000000000..36f013930 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/Language.java @@ -0,0 +1,11 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service; + +public enum Language { + ARABIC, + CHINESE, + ENGLISH, + FRENCH, + HINDI, + PORTUGUESE, + SPANISH +} diff --git a/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/NexusRemoteGreetingService.java b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/NexusRemoteGreetingService.java new file mode 100644 index 000000000..f3ce245c6 --- /dev/null +++ b/core/src/main/java/io/temporal/samples/nexusmessagingtemporaloperation/ondemandpattern/service/NexusRemoteGreetingService.java @@ -0,0 +1,183 @@ +package io.temporal.samples.nexusmessagingtemporaloperation.ondemandpattern.service; + +import com.fasterxml.jackson.annotation.JsonAutoDetect; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.nexusrpc.Operation; +import io.nexusrpc.Service; +import java.util.List; + +/** + * Nexus service definition for the on-demand pattern. Every operation includes a {@code userId} so + * the caller controls which workflow instance is targeted though that, while the Nexus service + * converts that UserId into a WorkflowId. This also exposes a {@code runFromRemote} operation that + * starts a new GreetingWorkflow. + */ +@Service +public interface NexusRemoteGreetingService { + + class RunFromRemoteInput { + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public RunFromRemoteInput(@JsonProperty("userId") String userId) { + this.userId = userId; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class GetLanguagesInput { + private final boolean includeUnsupported; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguagesInput( + @JsonProperty("includeUnsupported") boolean includeUnsupported, + @JsonProperty("userId") String userId) { + this.includeUnsupported = includeUnsupported; + this.userId = userId; + } + + @JsonProperty("includeUnsupported") + public boolean isIncludeUnsupported() { + return includeUnsupported; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + @JsonAutoDetect(fieldVisibility = JsonAutoDetect.Visibility.ANY) + class GetLanguageInput { + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguageInput(@JsonProperty("userId") String userId) { + this.userId = userId; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class SetLanguageInput { + private final Language language; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public SetLanguageInput( + @JsonProperty("language") Language language, @JsonProperty("userId") String userId) { + this.language = language; + this.userId = userId; + } + + @JsonProperty("language") + public Language getLanguage() { + return language; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class ApproveInput { + private final String name; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public ApproveInput(@JsonProperty("name") String name, @JsonProperty("userId") String userId) { + this.name = name; + this.userId = userId; + } + + @JsonProperty("name") + public String getName() { + return name; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class AttachApprovalContextInput { + private final String note; + private final String userId; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public AttachApprovalContextInput( + @JsonProperty("note") String note, @JsonProperty("userId") String userId) { + this.note = note; + this.userId = userId; + } + + @JsonProperty("note") + public String getNote() { + return note; + } + + @JsonProperty("userId") + public String getUserId() { + return userId; + } + } + + class GetLanguagesOutput { + private final List languages; + + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + public GetLanguagesOutput(@JsonProperty("languages") List languages) { + this.languages = languages; + } + + @JsonProperty("languages") + public List getLanguages() { + return languages; + } + } + + @JsonAutoDetect(fieldVisibility = JsonAutoDetect.Visibility.ANY) + class ApproveOutput { + @JsonCreator + public ApproveOutput() {} + } + + // Starts a new GreetingWorkflow for the given user ID. This is an asynchronous Nexus + // operation: the caller receives a handle and can wait for the workflow to complete. + @Operation + String runFromRemote(RunFromRemoteInput input); + + // Returns the languages supported by the specified workflow. + @Operation + GetLanguagesOutput getLanguages(GetLanguagesInput input); + + // Returns the currently active language of the specified workflow. + @Operation + Language getLanguage(GetLanguageInput input); + + // Changes the active language on the specified workflow, returning the previous one. + @Operation + Language setLanguage(SetLanguageInput input); + + // Approves the specified workflow, allowing it to complete. + @Operation + ApproveOutput approve(ApproveInput input); + + // Attaches supporting information for the eventual approval. Unlike every Operation above, this + // one does not require the Workflow to already exist: the handler delivers it with + // Signal-with-Start, so the same call either messages a running GreetingWorkflow or creates one. + // That makes it safe to call before or after runFromRemote. + @Operation + void attachApprovalContext(AttachApprovalContextInput input); +}