Repository navigation
Sample for NexusSerializationContext #802
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
b498f12
09927cd
c262de7
bc717ac
a30fe10
548fb13
d31656f
98cfa76
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,115 @@ | ||
| # Nexus serialization context | ||
|
|
||
| This sample calls a synchronous and an asynchronous Nexus operation through each of | ||
| two endpoints. Each endpoint routes to its own handler worker, which registers both | ||
| `SyncEchoService` and `AsyncEchoService`. The caller's `NexusCodec` uses | ||
| `NexusSerializationContext` to select the `PayloadCodec` registered for each endpoint. | ||
| The sample uses these keys: | ||
|
|
||
| - Key A: compresses with zlib, then encrypts synchronous and asynchronous payloads for `nexus-serialization-compressed-encrypted`. | ||
| - Key B: encrypts synchronous and asynchronous payloads for `nexus-serialization-encrypted`. | ||
| - Key C: encrypts the caller workflow's input and final result. | ||
|
|
||
| The caller configures a codec for each endpoint and a separate codec for its own | ||
| workflow input and result. Each handler worker uses the fixed key for its endpoint. | ||
|
|
||
| The caller schedules all four operations before waiting for their results, so each | ||
| result must be decoded using the context of its own endpoint. The caller's | ||
| `NexusCodec` uses Key C for its own workflow payloads, which have no Nexus endpoint. | ||
|
|
||
| For both Nexus endpoints, the outer encrypted payload stores the endpoint name in | ||
| `nexus-endpoint-name` metadata alongside `binary/nexus-aes-gcm` and a sample key ID | ||
| (`key-a` or `key-b`). A Codec Server can use the endpoint name to select the matching | ||
| key and decompression chain without SDK context. | ||
|
|
||
| For the compressed endpoint, the outer payload's decoded metadata looks like: | ||
|
|
||
| ```text | ||
| encoding: binary/nexus-aes-gcm | ||
| encryption-key-id: key-a | ||
| nexus-endpoint-name: nexus-serialization-compressed-encrypted | ||
| ``` | ||
|
|
||
| The asynchronous results have `encryption-key-id: key-a` or `key-b` and their respective | ||
| endpoint names in the outer payload metadata. | ||
| The caller workflow's input and final result use `encryption-key-id: key-c`; they do not | ||
| have endpoint metadata. The starter uses the same converter as the caller worker, so it | ||
| can decode the final result before printing it. | ||
|
|
||
| `NexusSerializationContext` works end to end for synchronous Nexus operations. | ||
| For an asynchronous operation, the handler's final result is serialized as a workflow | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I'd mention this is something we are going to improve
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Updated |
||
| result and does not receive `NexusSerializationContext`. We plan to add | ||
| `NexusSerializationContext` support for asynchronous operation results in the Java SDK. | ||
| Until then, the sample's `NexusEndpointInterceptor` captures the endpoint when the | ||
| handler starts the backing `EchoWorkflow`. `NexusEndpointContextPropagator` saves it in | ||
| the workflow headers and restores it on the workflow thread. Each handler's codec uses | ||
| its fixed key to encrypt the workflow result and includes the propagated endpoint name | ||
| in the outer payload metadata. This propagation is needed for the endpoint metadata; | ||
| the handler already knows which key to use. When the result reaches the caller, the SDK | ||
| supplies `NexusSerializationContext` to decode it. | ||
|
|
||
| This endpoint propagation covers asynchronous operations backed by workflows. If the | ||
| endpoint is not propagated, the result remains encrypted with the handler's fixed key | ||
| but lacks the endpoint name in its metadata. | ||
|
|
||
| The hard-coded keys are only for this local example. For production encryption, | ||
| use a secure key store, as in the | ||
| [AWS Encryption SDK sample](../keymanagementencryption/awsencryptionsdk/README.md). | ||
|
|
||
| Requires Java SDK 1.40.0 or later and Temporal Server 1.30.0 or later with Nexus enabled | ||
| so the handler can read the endpoint name. | ||
|
|
||
| ## Run locally | ||
|
|
||
| Start a Temporal dev server: | ||
|
|
||
| ```bash | ||
| temporal server start-dev | ||
| ``` | ||
|
|
||
| In another terminal, create the namespaces and endpoints: | ||
|
|
||
| ```bash | ||
| temporal operator namespace create --namespace nexus-serialization-caller | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I would lean towards separating this out more so it is a more realistic sample, we have gotten some feedback from product on our samples not being reflective of what the customer would actually do. So I would suggest having 3 namespaces here like 1 caller and two handler and each handler has its own endpoint and worker
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Created 3 in different namespaces |
||
| temporal operator namespace create --namespace nexus-serialization-key-a-handler | ||
| temporal operator namespace create --namespace nexus-serialization-key-b-handler | ||
| temporal operator nexus endpoint create \ | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Do we need 4 endpoints? I was thinking the async call could just be on the same nexus endpoint?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Made this to be 2 services for Async and Sync, and two endpoints for encrypt and compress+encrypt |
||
| --name nexus-serialization-compressed-encrypted \ | ||
| --target-namespace nexus-serialization-key-a-handler \ | ||
| --target-task-queue nexus-serialization-key-a-handler | ||
| temporal operator nexus endpoint create \ | ||
| --name nexus-serialization-encrypted \ | ||
| --target-namespace nexus-serialization-key-b-handler \ | ||
| --target-task-queue nexus-serialization-key-b-handler | ||
| ``` | ||
|
|
||
| Run each of the following in its own terminal from the repository root: | ||
|
|
||
| ```bash | ||
| ./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusserializationcontext.handler.CompressedEncryptedHandlerWorker \ | ||
| --args="-namespace nexus-serialization-key-a-handler" | ||
| ``` | ||
|
|
||
| ```bash | ||
| ./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusserializationcontext.handler.EncryptedHandlerWorker \ | ||
| --args="-namespace nexus-serialization-key-b-handler" | ||
| ``` | ||
|
|
||
| ```bash | ||
| ./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusserializationcontext.caller.CallerWorker \ | ||
| --args="-namespace nexus-serialization-caller" | ||
| ``` | ||
|
|
||
| ```bash | ||
| ./gradlew -q :core:execute -PmainClass=io.temporal.samples.nexusserializationcontext.caller.CallerStarter \ | ||
| --args="-namespace nexus-serialization-caller" | ||
| ``` | ||
|
|
||
| The starter will print: | ||
|
|
||
| ```text | ||
| Compressed and encrypted endpoint sync result: Hello from Nexus | ||
| Encrypted endpoint sync result: Hello from Nexus | ||
| Compressed and encrypted endpoint async result: Hello from Nexus | ||
| Encrypted endpoint async result: Hello from Nexus | ||
| ``` | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,18 @@ | ||
| package io.temporal.samples.nexusserializationcontext; | ||
|
|
||
| public final class SampleConfig { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I would remove most things here except for the endpoints since you wouldn't want to expose the handlers or callers task queue that breaks part of the abstraction of Nexus
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Updated |
||
| public static final String COMPRESSED_ENCRYPTED_ENDPOINT = | ||
| "nexus-serialization-compressed-encrypted"; | ||
| public static final String ENCRYPTED_ENDPOINT = "nexus-serialization-encrypted"; | ||
| public static final String ENDPOINT_METADATA_KEY = "nexus-endpoint-name"; | ||
| public static final String KEY_A_ID = "key-a"; | ||
| public static final String KEY_B_ID = "key-b"; | ||
| public static final String KEY_C_ID = "key-c"; | ||
|
|
||
| // Hard-coded keys are only for this local sample. | ||
| public static final String KEY_A_VALUE = "sample-key-A-123"; | ||
| public static final String KEY_B_VALUE = "sample-key-B-123"; | ||
| public static final String KEY_C_VALUE = "sample-key-C-123"; | ||
|
|
||
| private SampleConfig() {} | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,29 @@ | ||
| package io.temporal.samples.nexusserializationcontext.caller; | ||
|
|
||
| import io.temporal.client.WorkflowClient; | ||
| import io.temporal.client.WorkflowClientOptions; | ||
| import io.temporal.client.WorkflowOptions; | ||
| import io.temporal.samples.nexus.options.ClientOptions; | ||
|
|
||
| public class CallerStarter { | ||
| public static void main(String[] args) { | ||
| WorkflowClient client = | ||
| ClientOptions.getWorkflowClient( | ||
| args, | ||
| WorkflowClientOptions.newBuilder().setDataConverter(CallerWorker.dataConverter())); | ||
| CallerWorkflow workflow = | ||
| client.newWorkflowStub( | ||
| CallerWorkflow.class, | ||
| WorkflowOptions.newBuilder().setTaskQueue(CallerWorker.TASK_QUEUE).build()); | ||
|
|
||
| EndpointResults results = workflow.echoThroughEndpoints("Hello from Nexus"); | ||
| System.out.println( | ||
| "Compressed and encrypted endpoint sync result: " | ||
| + results.compressedEncryptedSyncResult()); | ||
| System.out.println("Encrypted endpoint sync result: " + results.encryptedSyncResult()); | ||
| System.out.println( | ||
| "Compressed and encrypted endpoint async result: " | ||
| + results.compressedEncryptedAsyncResult()); | ||
| System.out.println("Encrypted endpoint async result: " + results.encryptedAsyncResult()); | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,58 @@ | ||
| package io.temporal.samples.nexusserializationcontext.caller; | ||
|
|
||
| import io.temporal.client.WorkflowClient; | ||
| import io.temporal.client.WorkflowClientOptions; | ||
| import io.temporal.common.converter.CodecDataConverter; | ||
| import io.temporal.common.converter.DataConverter; | ||
| import io.temporal.common.converter.DefaultDataConverter; | ||
| import io.temporal.payload.codec.ChainCodec; | ||
| import io.temporal.payload.codec.PayloadCodec; | ||
| import io.temporal.samples.nexus.options.ClientOptions; | ||
| import io.temporal.samples.nexusserializationcontext.SampleConfig; | ||
| import io.temporal.samples.nexusserializationcontext.codec.AesGcmCodec; | ||
| import io.temporal.samples.nexusserializationcontext.codec.NexusCodec; | ||
| import io.temporal.samples.nexusserializationcontext.codec.ZlibCodec; | ||
| import io.temporal.worker.Worker; | ||
| import io.temporal.worker.WorkerFactory; | ||
| import java.nio.charset.StandardCharsets; | ||
| import java.util.List; | ||
| import java.util.Map; | ||
| import javax.crypto.SecretKey; | ||
| import javax.crypto.spec.SecretKeySpec; | ||
|
|
||
| public class CallerWorker { | ||
| static final String TASK_QUEUE = "nexus-serialization-caller"; | ||
|
|
||
| public static void main(String[] args) { | ||
| WorkflowClient client = | ||
| ClientOptions.getWorkflowClient( | ||
| args, WorkflowClientOptions.newBuilder().setDataConverter(dataConverter())); | ||
| WorkerFactory factory = WorkerFactory.newInstance(client); | ||
| Worker worker = factory.newWorker(TASK_QUEUE); | ||
| worker.registerWorkflowImplementationTypes(CallerWorkflowImpl.class); | ||
| factory.start(); | ||
| } | ||
|
|
||
| public static DataConverter dataConverter() { | ||
| SecretKey keyA = | ||
| new SecretKeySpec(SampleConfig.KEY_A_VALUE.getBytes(StandardCharsets.UTF_8), "AES"); | ||
| SecretKey keyB = | ||
| new SecretKeySpec(SampleConfig.KEY_B_VALUE.getBytes(StandardCharsets.UTF_8), "AES"); | ||
| SecretKey keyC = | ||
| new SecretKeySpec(SampleConfig.KEY_C_VALUE.getBytes(StandardCharsets.UTF_8), "AES"); | ||
| // ChainCodec encodes last to first: compress, then encrypt. | ||
| PayloadCodec compressedEncrypted = | ||
| new ChainCodec(List.of(new AesGcmCodec(SampleConfig.KEY_A_ID, keyA), new ZlibCodec())); | ||
| PayloadCodec encrypted = new AesGcmCodec(SampleConfig.KEY_B_ID, keyB); | ||
| return new CodecDataConverter( | ||
| DefaultDataConverter.newDefaultInstance(), | ||
| List.of( | ||
| new NexusCodec( | ||
| Map.of( | ||
| SampleConfig.COMPRESSED_ENCRYPTED_ENDPOINT, | ||
| compressedEncrypted, | ||
| SampleConfig.ENCRYPTED_ENDPOINT, | ||
| encrypted), | ||
| new AesGcmCodec(SampleConfig.KEY_C_ID, keyC)))); | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,10 @@ | ||
| package io.temporal.samples.nexusserializationcontext.caller; | ||
|
|
||
| import io.temporal.workflow.WorkflowInterface; | ||
| import io.temporal.workflow.WorkflowMethod; | ||
|
|
||
| @WorkflowInterface | ||
| public interface CallerWorkflow { | ||
| @WorkflowMethod | ||
| EndpointResults echoThroughEndpoints(String message); | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,51 @@ | ||
| package io.temporal.samples.nexusserializationcontext.caller; | ||
|
|
||
| import io.temporal.samples.nexusserializationcontext.SampleConfig; | ||
| import io.temporal.samples.nexusserializationcontext.service.AsyncEchoService; | ||
| import io.temporal.samples.nexusserializationcontext.service.SyncEchoService; | ||
| import io.temporal.workflow.NexusOperationHandle; | ||
| import io.temporal.workflow.NexusOperationOptions; | ||
| import io.temporal.workflow.NexusServiceOptions; | ||
| import io.temporal.workflow.Workflow; | ||
| import java.time.Duration; | ||
|
|
||
| public class CallerWorkflowImpl implements CallerWorkflow { | ||
| @Override | ||
| public EndpointResults echoThroughEndpoints(String message) { | ||
| SyncEchoService compressedEncryptedService = | ||
| serviceFor(SyncEchoService.class, SampleConfig.COMPRESSED_ENCRYPTED_ENDPOINT); | ||
| SyncEchoService encryptedService = | ||
| serviceFor(SyncEchoService.class, SampleConfig.ENCRYPTED_ENDPOINT); | ||
| AsyncEchoService compressedEncryptedAsyncService = | ||
| serviceFor(AsyncEchoService.class, SampleConfig.COMPRESSED_ENCRYPTED_ENDPOINT); | ||
| AsyncEchoService encryptedAsyncService = | ||
| serviceFor(AsyncEchoService.class, SampleConfig.ENCRYPTED_ENDPOINT); | ||
|
|
||
| // Start all operations before awaiting results. Each result keeps its endpoint context. | ||
| NexusOperationHandle<String> compressedEncryptedSync = | ||
| Workflow.startNexusOperation(compressedEncryptedService::echo, message); | ||
| NexusOperationHandle<String> encryptedSync = | ||
| Workflow.startNexusOperation(encryptedService::echo, message); | ||
| NexusOperationHandle<String> compressedEncryptedAsync = | ||
| Workflow.startNexusOperation(compressedEncryptedAsyncService::echoAsync, message); | ||
| NexusOperationHandle<String> encryptedAsync = | ||
| Workflow.startNexusOperation(encryptedAsyncService::echoAsync, message); | ||
| return new EndpointResults( | ||
| compressedEncryptedSync.getResult().get(), | ||
| encryptedSync.getResult().get(), | ||
| compressedEncryptedAsync.getResult().get(), | ||
| encryptedAsync.getResult().get()); | ||
| } | ||
|
|
||
| private static <T> T serviceFor(Class<T> serviceClass, String endpoint) { | ||
| return Workflow.newNexusServiceStub( | ||
| serviceClass, | ||
| NexusServiceOptions.newBuilder() | ||
| .setEndpoint(endpoint) | ||
| .setOperationOptions( | ||
| NexusOperationOptions.newBuilder() | ||
| .setScheduleToCloseTimeout(Duration.ofSeconds(30)) | ||
| .build()) | ||
| .build()); | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,7 @@ | ||
| package io.temporal.samples.nexusserializationcontext.caller; | ||
|
|
||
| public record EndpointResults( | ||
| String compressedEncryptedSyncResult, | ||
| String encryptedSyncResult, | ||
| String compressedEncryptedAsyncResult, | ||
| String encryptedAsyncResult) {} |
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This is a fix for unit tests, it was stalling after this SDK version change