From cd21e5b599ae8084dc606392e8ec907ff36161fa Mon Sep 17 00:00:00 2001 From: Patrik Beqo Date: Tue, 15 Sep 2026 19:00:22 -0400 Subject: [PATCH] Add named Workflow random streams and read-only detection --- .../replay/ReplayWorkflowContext.java | 3 + .../replay/ReplayWorkflowContextImpl.java | 5 + .../statemachines/WorkflowRandomStreams.java | 31 +++ .../statemachines/WorkflowStateMachines.java | 8 + .../internal/sync/WorkflowInternal.java | 8 +- .../java/io/temporal/workflow/Workflow.java | 18 ++ .../workflow/unsafe/WorkflowUnsafe.java | 16 ++ .../WorkflowRandomStreamsTest.java | 102 +++++++ .../workflow/WorkflowRandomStreamTest.java | 258 ++++++++++++++++++ .../workflow/WorkflowUnsafeReadOnlyTest.java | 162 +++++++++++ .../sync/DummySyncWorkflowContext.java | 5 + 11 files changed, 614 insertions(+), 2 deletions(-) create mode 100644 temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowRandomStreams.java create mode 100644 temporal-sdk/src/test/java/io/temporal/internal/statemachines/WorkflowRandomStreamsTest.java create mode 100644 temporal-sdk/src/test/java/io/temporal/workflow/WorkflowRandomStreamTest.java create mode 100644 temporal-sdk/src/test/java/io/temporal/workflow/WorkflowUnsafeReadOnlyTest.java diff --git a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContext.java b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContext.java index 982737dbee..41854b9229 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContext.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContext.java @@ -292,6 +292,9 @@ Integer getVersion( /** Replay safe random. */ Random newRandom(); + /** Replay safe named random stream. */ + Random getRandomStream(String name); + /** * @return scope to be used for metrics reporting. */ diff --git a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContextImpl.java b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContextImpl.java index 2f600b20aa..97c7e20ff0 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContextImpl.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/replay/ReplayWorkflowContextImpl.java @@ -81,6 +81,11 @@ public Random newRandom() { return workflowStateMachines.newRandom(); } + @Override + public Random getRandomStream(String name) { + return workflowStateMachines.getRandomStream(name); + } + @Override public Scope getMetricsScope() { return replayAwareWorkflowMetricsScope; diff --git a/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowRandomStreams.java b/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowRandomStreams.java new file mode 100644 index 0000000000..dd2815a8aa --- /dev/null +++ b/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowRandomStreams.java @@ -0,0 +1,31 @@ +package io.temporal.internal.statemachines; + +import com.google.common.hash.HashCode; +import com.google.common.hash.Hashing; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; +import java.util.Random; +import javax.annotation.Nonnull; + +final class WorkflowRandomStreams { + private static final String SEED_VERSION = "temporal.sdk.random.v1"; + + private final Map streams = new HashMap<>(); + + long deriveSeed(@Nonnull String runId, @Nonnull String name) { + // The separators keep ("ab", "c") from colliding with ("a", "bc") + String seed = String.join("\0", SEED_VERSION, runId, name); + HashCode hash = Hashing.sha256().hashString(seed, StandardCharsets.UTF_8); + return ByteBuffer.wrap(hash.asBytes()).getLong(); + } + + Random get(@Nonnull String runId, @Nonnull String name) { + return streams.computeIfAbsent(name, key -> new Random(deriveSeed(runId, key))); + } + + void reseed(@Nonnull String runId) { + streams.forEach((name, stream) -> stream.setSeed(deriveSeed(runId, name))); + } +} diff --git a/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowStateMachines.java b/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowStateMachines.java index 2de6b6ea15..8f79cbbc10 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowStateMachines.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/statemachines/WorkflowStateMachines.java @@ -115,6 +115,8 @@ enum HandleEventStatus { /** Used Workflow.newRandom and randomUUID together with currentRunId. */ private long idCounter; + private final WorkflowRandomStreams randomStreams = new WorkflowRandomStreams(); + /** Current workflow time. */ private long currentTimeMillis = -1; @@ -1195,6 +1197,11 @@ public Random newRandom() { return new Random(randomUUID().getLeastSignificantBits()); } + public Random getRandomStream(String name) { + checkEventLoopExecuting(); + return randomStreams.get(currentRunId, name); + } + public void sideEffect( Functions.Func> func, UserMetadata userMetadata, @@ -1548,6 +1555,7 @@ public void workflowTaskStarted( @Override public void updateRunId(String currentRunId) { WorkflowStateMachines.this.currentRunId = currentRunId; + randomStreams.reseed(currentRunId); } } diff --git a/temporal-sdk/src/main/java/io/temporal/internal/sync/WorkflowInternal.java b/temporal-sdk/src/main/java/io/temporal/internal/sync/WorkflowInternal.java index 84b1e91fd3..07b71769f2 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/sync/WorkflowInternal.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/sync/WorkflowInternal.java @@ -744,6 +744,10 @@ public static Random newRandom() { return getRootWorkflowContext().newRandom(); } + public static Random getRandomStream(String name) { + return getRootWorkflowContext().getReplayContext().getRandomStream(name); + } + public static Logger getLogger(Class clazz) { Logger logger = LoggerFactory.getLogger(clazz); return new ReplayAwareLogger( @@ -919,11 +923,11 @@ static SyncWorkflowContext getRootWorkflowContext() { return DeterministicRunnerImpl.currentThreadInternal().getWorkflowContext(); } - static boolean isReadOnly() { + public static boolean isReadOnly() { return getRootWorkflowContext().isReadOnly(); } - static void assertNotReadOnly(String action) { + public static void assertNotReadOnly(String action) { if (isReadOnly()) { throw new ReadOnlyException(action); } diff --git a/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java b/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java index d04a617d5d..561c2221e2 100644 --- a/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java +++ b/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java @@ -711,6 +711,24 @@ public static Random newRandom() { return WorkflowInternal.newRandom(); } + /** + * Returns a deterministic pseudorandom stream private to {@code name}. + * + *

Calling this method again with the same name returns the same logical stream where earlier + * draws left it. A Workflow Reset replays the same values up to the reset point, then reseeds the + * stream for the new Run. Each Continue-As-New Run gets a new sequence. + * + *

Draws are not recorded in Workflow History, so do not draw in read-only code. Use {@link + * WorkflowUnsafe#isReadOnly()} to gate draws. + * + *

Use a stable package-style name. Stream names are retained for the life of the Workflow Run. + * The stream is deterministic pseudorandomness and is not cryptographically secure. + */ + @Experimental + public static Random getRandomStream(String name) { + return WorkflowInternal.getRandomStream(name); + } + /** * True if workflow code is being replayed. * diff --git a/temporal-sdk/src/main/java/io/temporal/workflow/unsafe/WorkflowUnsafe.java b/temporal-sdk/src/main/java/io/temporal/workflow/unsafe/WorkflowUnsafe.java index 1a67b4e8a6..da635d9727 100644 --- a/temporal-sdk/src/main/java/io/temporal/workflow/unsafe/WorkflowUnsafe.java +++ b/temporal-sdk/src/main/java/io/temporal/workflow/unsafe/WorkflowUnsafe.java @@ -1,5 +1,6 @@ package io.temporal.workflow.unsafe; +import io.temporal.common.Experimental; import io.temporal.internal.sync.WorkflowInternal; import io.temporal.workflow.Functions; @@ -46,6 +47,21 @@ public static boolean isReplaying() { return WorkflowInternal.isReplaying(); } + /** + * Reports whether the current code is running where Workflow state cannot be mutated. + * + *

Read-only code includes Query handlers, Update validators, Side Effect functions, Await + * conditions, and other SDK callbacks that must not mutate Workflow state. + * + *

Must be called from Workflow code. + * + * @return true in a read-only Workflow context + */ + @Experimental + public static boolean isReadOnly() { + return WorkflowInternal.isReadOnly(); + } + /** * Runs the supplied procedure in the calling thread with disabled deadlock detection if called * from the workflow thread. Does nothing except the procedure execution if called from a diff --git a/temporal-sdk/src/test/java/io/temporal/internal/statemachines/WorkflowRandomStreamsTest.java b/temporal-sdk/src/test/java/io/temporal/internal/statemachines/WorkflowRandomStreamsTest.java new file mode 100644 index 0000000000..1752d17218 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/internal/statemachines/WorkflowRandomStreamsTest.java @@ -0,0 +1,102 @@ +package io.temporal.internal.statemachines; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertSame; + +import com.google.common.io.BaseEncoding; +import java.util.ArrayList; +import java.util.List; +import java.util.Random; +import org.junit.Test; + +public class WorkflowRandomStreamsTest { + private static final String RUN_ID = "runID"; + private static final String NAME = "io.temporal.test"; + + @Test + public void deriveSeed() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + long seed = randoms.deriveSeed(RUN_ID, NAME); + assertNotEquals(seed, randoms.deriveSeed("other", NAME)); + assertNotEquals(seed, randoms.deriveSeed(RUN_ID, "other")); + assertNotEquals(seed, randoms.deriveSeed("other", "other")); + } + + /** Pins the seed derivation and resulting byte stream. Changing either breaks replay. */ + @Test + public void getRandomStreamGolden() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + assertEquals(8181915698088084985L, randoms.deriveSeed(RUN_ID, NAME)); + + byte[] bytes = new byte[32]; + randoms.get(RUN_ID, NAME).nextBytes(bytes); + assertEquals( + "1c1d4dd36999ff851d72aa41f660ecd3220d83499109ba3ba24e455a1b776c21", + BaseEncoding.base16().lowerCase().encode(bytes)); + } + + @Test + public void deriveSeedSeparators() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + assertNotEquals(randoms.deriveSeed("ab", "c"), randoms.deriveSeed("a", "bc")); + } + + /** A second lookup under the same name continues the sequence rather than restarting it. */ + @Test + public void getRandomStreamMemoizes() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + + Random first = randoms.get(RUN_ID, NAME); + long firstDraw = first.nextLong(); + + Random second = randoms.get(RUN_ID, NAME); + long secondDraw = second.nextLong(); + + assertSame(first, second); + assertNotEquals(firstDraw, secondDraw); + } + + /** + * Interleaving draws across two names yields the same sequence per name as drawing from each on + * its own, so how often a workflow draws from one name cannot shift another. + */ + @Test + public void getRandomStreamNamesAreIndependent() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + Random first = randoms.get(RUN_ID, NAME); + Random second = randoms.get(RUN_ID, "other"); + + List interleavedA = new ArrayList<>(); + List interleavedB = new ArrayList<>(); + for (int i = 0; i < 3; i++) { + interleavedA.add(first.nextLong()); + interleavedB.add(second.nextLong()); + } + + assertEquals(solo(NAME, 3), interleavedA); + assertEquals(solo("other", 3), interleavedB); + assertNotEquals(interleavedA, interleavedB); + } + + @Test + public void reseedRandomsInPlace() { + WorkflowRandomStreams randoms = new WorkflowRandomStreams(); + + Random first = randoms.get(RUN_ID, NAME); + randoms.reseed("other"); + Random second = randoms.get("other", NAME); + + assertSame(first, second); + assertEquals(new WorkflowRandomStreams().get("other", NAME).nextLong(), second.nextLong()); + } + + private static List solo(String name, int draws) { + Random random = new WorkflowRandomStreams().get(RUN_ID, name); + List result = new ArrayList<>(); + for (int i = 0; i < draws; i++) { + result.add(random.nextLong()); + } + return result; + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowRandomStreamTest.java b/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowRandomStreamTest.java new file mode 100644 index 0000000000..b9463a4cd7 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowRandomStreamTest.java @@ -0,0 +1,258 @@ +package io.temporal.workflow; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assume.assumeTrue; + +import io.temporal.activity.ActivityInterface; +import io.temporal.activity.ActivityMethod; +import io.temporal.activity.ActivityOptions; +import io.temporal.api.common.v1.WorkflowExecution; +import io.temporal.api.workflowservice.v1.ResetWorkflowExecutionRequest; +import io.temporal.api.workflowservice.v1.ResetWorkflowExecutionResponse; +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowStub; +import io.temporal.client.WorkflowTargetOptions; +import io.temporal.testing.internal.SDKTestWorkflowRule; +import io.temporal.worker.WorkerOptions; +import java.time.Duration; +import java.util.Random; +import java.util.UUID; +import org.junit.Rule; +import org.junit.Test; + +public class WorkflowRandomStreamTest { + private static final String STREAM_NAME = "io.temporal.test"; + + @Rule + public SDKTestWorkflowRule testWorkflowRule = + SDKTestWorkflowRule.newBuilder() + .setWorkflowTypes( + SimpleWorkflowImpl.class, + ReplayWorkflowImpl.class, + ResetWorkflowImpl.class, + ResetLateSourceWorkflowImpl.class, + ContinueAsNewWorkflowImpl.class, + ParentWorkflowImpl.class) + .setActivityImplementations(new SimpleActivityImpl()) + .setWorkerOptions( + WorkerOptions.newBuilder() + .setStickyQueueScheduleToStartTimeout(Duration.ZERO) + .build()) + .build(); + + @Test + public void noCollisionAcrossRuns() { + SimpleWorkflow first = testWorkflowRule.newWorkflowStubTimeoutOptions(SimpleWorkflow.class); + SimpleWorkflow second = testWorkflowRule.newWorkflowStubTimeoutOptions(SimpleWorkflow.class); + + assertNotEquals(first.run(), second.run()); + } + + @Test + public void deterministicReplay() { + ReplayWorkflow workflow = testWorkflowRule.newWorkflowStubTimeoutOptions(ReplayWorkflow.class); + + long result = workflow.run(); + + assertEquals(result, workflow.currentState()); + } + + @Test + public void resetReproducesValues() { + assumeTrue( + "Test Server doesn't support reset workflow", SDKTestWorkflowRule.useExternalService); + assertResetValues(ResetWorkflow.class); + } + + @Test + public void resetReseedsSourceCreatedAfterResetPoint() { + assumeTrue( + "Test Server doesn't support reset workflow", SDKTestWorkflowRule.useExternalService); + assertResetValues(ResetLateSourceWorkflow.class); + } + + @Test + public void continueAsNewDrawsNewValues() { + ContinueAsNewWorkflow workflow = + testWorkflowRule.newWorkflowStubTimeoutOptions(ContinueAsNewWorkflow.class); + + long[] values = workflow.run(null); + + assertEquals(2, values.length); + assertNotEquals(values[0], values[1]); + } + + @Test + public void childContinueAsNewDrawsNewValues() { + ParentWorkflow workflow = testWorkflowRule.newWorkflowStubTimeoutOptions(ParentWorkflow.class); + + long[] values = workflow.run(); + + assertEquals(3, values.length); + assertNotEquals(values[0], values[1]); + assertNotEquals(values[0], values[2]); + assertNotEquals(values[1], values[2]); + } + + private void assertResetValues(Class workflowType) { + WorkflowClient client = testWorkflowRule.getWorkflowClient(); + WorkflowStub stub = + WorkflowStub.fromTyped(testWorkflowRule.newWorkflowStubTimeoutOptions(workflowType)); + WorkflowExecution execution = stub.start(); + long[] original = stub.getResult(long[].class); + assertEquals(2, original.length); + assertNotEquals(original[0], original[1]); + + // The reset targets the second Workflow Task (id=10), so the first draw is replayed and the + // second draw is redrawn + ResetWorkflowExecutionResponse response = + client + .getWorkflowServiceStubs() + .blockingStub() + .resetWorkflowExecution( + ResetWorkflowExecutionRequest.newBuilder() + .setNamespace(SDKTestWorkflowRule.NAMESPACE) + .setWorkflowExecution(execution) + .setWorkflowTaskFinishEventId(10) + .setReason("Integration test") + .setRequestId(UUID.randomUUID().toString()) + .build()); + + long[] afterReset = + client + .newUntypedWorkflowStub( + WorkflowTargetOptions.newBuilder() + .setWorkflowId(execution.getWorkflowId()) + .setRunId(response.getRunId()) + .build()) + .getResult(long[].class); + assertEquals(2, afterReset.length); + assertNotEquals(afterReset[0], afterReset[1]); + + assertEquals(original[0], afterReset[0]); + assertNotEquals(original[1], afterReset[1]); + } + + @WorkflowInterface + public interface SimpleWorkflow { + @WorkflowMethod + long run(); + } + + @WorkflowInterface + public interface ReplayWorkflow { + @WorkflowMethod + long run(); + + @QueryMethod + long currentState(); + } + + @WorkflowInterface + public interface ResetWorkflow { + @WorkflowMethod + long[] run(); + } + + @WorkflowInterface + public interface ResetLateSourceWorkflow { + @WorkflowMethod + long[] run(); + } + + @WorkflowInterface + public interface ContinueAsNewWorkflow { + @WorkflowMethod + long[] run(Long previous); + } + + @WorkflowInterface + public interface ParentWorkflow { + @WorkflowMethod + long[] run(); + } + + @ActivityInterface + public interface SimpleActivity { + @ActivityMethod + void run(); + } + + public static class SimpleWorkflowImpl implements SimpleWorkflow { + @Override + public long run() { + return Workflow.getRandomStream(STREAM_NAME).nextLong(); + } + } + + public static class ReplayWorkflowImpl implements ReplayWorkflow { + private long state; + + @Override + public long run() { + Random random = Workflow.getRandomStream(STREAM_NAME); + state = random.nextLong(); + newSimpleActivity().run(); + state = random.nextLong(); + return state; + } + + @Override + public long currentState() { + return state; + } + } + + public static class ResetWorkflowImpl implements ResetWorkflow { + @Override + public long[] run() { + Random random = Workflow.getRandomStream(STREAM_NAME); + long first = random.nextLong(); + newSimpleActivity().run(); + long second = random.nextLong(); + return new long[] {first, second}; + } + } + + public static class ResetLateSourceWorkflowImpl implements ResetLateSourceWorkflow { + @Override + public long[] run() { + long first = Workflow.getRandomStream("other").nextLong(); + newSimpleActivity().run(); + long second = Workflow.getRandomStream(STREAM_NAME).nextLong(); + return new long[] {first, second}; + } + } + + public static class ContinueAsNewWorkflowImpl implements ContinueAsNewWorkflow { + @Override + public long[] run(Long previous) { + long current = Workflow.getRandomStream(STREAM_NAME).nextLong(); + if (previous == null) { + Workflow.continueAsNew(current); + } + return new long[] {previous, current}; + } + } + + public static class ParentWorkflowImpl implements ParentWorkflow { + @Override + public long[] run() { + long parent = Workflow.getRandomStream(STREAM_NAME).nextLong(); + long[] child = Workflow.newChildWorkflowStub(ContinueAsNewWorkflow.class).run(null); + return new long[] {parent, child[0], child[1]}; + } + } + + public static class SimpleActivityImpl implements SimpleActivity { + @Override + public void run() {} + } + + private static SimpleActivity newSimpleActivity() { + return Workflow.newActivityStub( + SimpleActivity.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofMinutes(1)).build()); + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowUnsafeReadOnlyTest.java b/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowUnsafeReadOnlyTest.java new file mode 100644 index 0000000000..5c12165d37 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/workflow/WorkflowUnsafeReadOnlyTest.java @@ -0,0 +1,162 @@ +package io.temporal.workflow; + +import static org.junit.Assert.assertEquals; + +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowStub; +import io.temporal.common.interceptors.WorkerInterceptorBase; +import io.temporal.common.interceptors.WorkflowInboundCallsInterceptor; +import io.temporal.common.interceptors.WorkflowInboundCallsInterceptorBase; +import io.temporal.testing.internal.SDKTestWorkflowRule; +import io.temporal.worker.WorkerFactoryOptions; +import io.temporal.workflow.unsafe.WorkflowUnsafe; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; + +public class WorkflowUnsafeReadOnlyTest { + private static final Map calls = new ConcurrentHashMap<>(); + + @Rule + public SDKTestWorkflowRule testWorkflowRule = + SDKTestWorkflowRule.newBuilder() + .setWorkflowTypes(ReadOnlyWorkflowImpl.class) + .setWorkerFactoryOptions( + WorkerFactoryOptions.newBuilder() + .setWorkerInterceptors(new ReadOnlyRecordingInterceptor()) + .build()) + .build(); + + @Before + public void setUp() { + calls.clear(); + } + + @Test + public void isReadOnly() { + ReadOnlyWorkflow workflow = + testWorkflowRule.newWorkflowStubTimeoutOptions(ReadOnlyWorkflow.class); + WorkflowClient.start(workflow::run); + + workflow.query(); + workflow.update(); + workflow.finish(); + WorkflowStub.fromTyped(workflow).getResult(Void.class); + + Map expected = new ConcurrentHashMap<>(); + expected.put("ExecuteWorkflow", false); + expected.put("workflowTask", false); + expected.put("ExecuteUpdate", false); + expected.put("updateHandler", false); + expected.put("HandleSignal", false); + expected.put("sideEffect", true); + expected.put("await", true); + expected.put("HandleQuery", true); + expected.put("query", true); + expected.put("ValidateUpdate", true); + expected.put("validator", true); + assertEquals(expected, calls); + } + + private static void record(String name) { + calls.put(name, WorkflowUnsafe.isReadOnly()); + } + + @WorkflowInterface + public interface ReadOnlyWorkflow { + @WorkflowMethod + void run(); + + @QueryMethod + boolean query(); + + @UpdateMethod + void update(); + + @UpdateValidatorMethod(updateName = "update") + void validateUpdate(); + + @SignalMethod + void finish(); + } + + public static class ReadOnlyWorkflowImpl implements ReadOnlyWorkflow { + private boolean finished; + + @Override + public void run() { + record("workflowTask"); + Workflow.sideEffect( + Void.class, + () -> { + record("sideEffect"); + return null; + }); + Workflow.await( + () -> { + record("await"); + return finished; + }); + } + + @Override + public boolean query() { + record("query"); + return true; + } + + @Override + public void update() { + record("updateHandler"); + } + + @Override + public void validateUpdate() { + record("validator"); + } + + @Override + public void finish() { + finished = true; + } + } + + private static class ReadOnlyRecordingInterceptor extends WorkerInterceptorBase { + @Override + public WorkflowInboundCallsInterceptor interceptWorkflow(WorkflowInboundCallsInterceptor next) { + return new WorkflowInboundCallsInterceptorBase(next) { + @Override + public WorkflowOutput execute(WorkflowInput input) { + record("ExecuteWorkflow"); + return super.execute(input); + } + + @Override + public void handleSignal(SignalInput input) { + record("HandleSignal"); + super.handleSignal(input); + } + + @Override + public QueryOutput handleQuery(QueryInput input) { + record("HandleQuery"); + return super.handleQuery(input); + } + + @Override + public void validateUpdate(UpdateInput input) { + record("ValidateUpdate"); + super.validateUpdate(input); + } + + @Override + public UpdateOutput executeUpdate(UpdateInput input) { + record("ExecuteUpdate"); + return super.executeUpdate(input); + } + }; + } + } +} diff --git a/temporal-testing/src/main/java/io/temporal/internal/sync/DummySyncWorkflowContext.java b/temporal-testing/src/main/java/io/temporal/internal/sync/DummySyncWorkflowContext.java index f89e61c64b..c62081a4a1 100644 --- a/temporal-testing/src/main/java/io/temporal/internal/sync/DummySyncWorkflowContext.java +++ b/temporal-testing/src/main/java/io/temporal/internal/sync/DummySyncWorkflowContext.java @@ -288,6 +288,11 @@ public Random newRandom() { throw new UnsupportedOperationException("not implemented"); } + @Override + public Random getRandomStream(String name) { + throw new UnsupportedOperationException("not implemented"); + } + @Override public Scope getMetricsScope() { return new NoopScope();