diff --git a/CHANGELOG.md b/CHANGELOG.md index 4ebe7bc903..fa8de3cab5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,11 @@ to docs, or any other relevant information. ## [Unreleased] +### :boom: Breaking Changes +- `WorkflowOutboundCallsInterceptor` has new `sleep(Duration, TimerOptions)` and + `await(Duration, TimerOptions, String, Supplier)` methods. Implementations that don't extend + `WorkflowOutboundCallsInterceptorBase` need to add them, and Base subclasses override them to see calls with options. + ### Added - Added experimental `ChildWorkflowOptions.Builder.setVersioningOverride` and `VersioningOverride.OneTimeVersioningOverride` for explicit pinned, auto-upgrade, and one-time @@ -27,6 +32,8 @@ to docs, or any other relevant information. Child workflow overrides and one-time routing require Temporal Server 1.32.0 or later. - `WorkerFactoryOptions.Builder.setLoggerTagPrefix` that can be used to customized structured logging tags (MDC keys) set by Temporal SDK in worker context. +- `Workflow.sleep(Duration, TimerOptions)` and `Workflow.await(Duration, TimerOptions, Supplier)` let workflows set a + summary on the timer behind a sleep or a timed await, like `Workflow.newTimer(Duration, TimerOptions)` already does. ### Changed - Release notes for all future releases are now in a single CHANGELOG.md file. `releases` directory with old release diff --git a/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptor.java b/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptor.java index 3b3d7f6f7d..a1b581a4e5 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptor.java +++ b/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptor.java @@ -784,8 +784,13 @@ public DynamicUpdateHandler getHandler() { void sleep(Duration duration); + void sleep(Duration duration, TimerOptions options); + boolean await(Duration timeout, String reason, Supplier unblockCondition); + boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition); + void await(String reason, Supplier unblockCondition); Promise newTimer(Duration duration); diff --git a/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptorBase.java b/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptorBase.java index 9d99d4c78b..6f55f2f4ae 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptorBase.java +++ b/temporal-sdk/src/main/java/io/temporal/common/interceptors/WorkflowOutboundCallsInterceptorBase.java @@ -65,11 +65,22 @@ public void sleep(Duration duration) { next.sleep(duration); } + @Override + public void sleep(Duration duration, TimerOptions options) { + next.sleep(duration, options); + } + @Override public boolean await(Duration timeout, String reason, Supplier unblockCondition) { return next.await(timeout, reason, unblockCondition); } + @Override + public boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition) { + return next.await(timeout, options, reason, unblockCondition); + } + @Override public void await(String reason, Supplier unblockCondition) { next.await(reason, unblockCondition); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflowContext.java b/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflowContext.java index e0b28a77e0..2d5cfcba33 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflowContext.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflowContext.java @@ -1369,11 +1369,22 @@ public SignalExternalOutput signalExternalWorkflow(SignalExternalInput input) { @Override public void sleep(Duration duration) { - newTimer(duration).get(); + sleep(duration, TimerOptions.newBuilder().build()); + } + + @Override + public void sleep(Duration duration, TimerOptions options) { + newTimer(duration, options).get(); } @Override public boolean await(Duration timeout, String reason, Supplier unblockCondition) { + return await(timeout, TimerOptions.newBuilder().build(), reason, unblockCondition); + } + + @Override + public boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition) { boolean cancelTimerOnCondition = replayContext.tryUseSdkFlag(SdkFlag.CANCEL_AWAIT_TIMER_ON_CONDITION); @@ -1385,7 +1396,7 @@ public boolean await(Duration timeout, String reason, Supplier unblockC // Create timer in a cancellation scope so we can cancel it when condition is satisfied CompletablePromise timer = Workflow.newPromise(); CancellationScope timerScope = - Workflow.newCancellationScope(() -> timer.completeFrom(newTimer(timeout))); + Workflow.newCancellationScope(() -> timer.completeFrom(newTimer(timeout, options))); timerScope.run(); WorkflowThread.await(reason, () -> (timer.isCompleted() || unblockCondition.get())); @@ -1397,7 +1408,7 @@ public boolean await(Duration timeout, String reason, Supplier unblockC return conditionSatisfied; } else { // Old behavior: timer is not cancelled when condition is satisfied - Promise timer = newTimer(timeout); + Promise timer = newTimer(timeout, options); WorkflowThread.await(reason, () -> (timer.isCompleted() || unblockCondition.get())); return !timer.isCompleted(); } 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 0a1c982f07..5a68f6e8e7 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 @@ -570,6 +570,25 @@ public static boolean await(Duration timeout, String reason, Supplier u }); } + public static boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition) + throws DestroyWorkflowThreadError { + assertNotReadOnly(reason); + return getWorkflowOutboundInterceptor() + .await( + timeout, + options, + reason, + () -> { + getRootWorkflowContext().setReadOnly(true); + try { + return unblockCondition.get(); + } finally { + getRootWorkflowContext().setReadOnly(false); + } + }); + } + public static R sideEffect(Class resultClass, Type resultType, Func func) { assertNotReadOnly("side effect"); return getWorkflowOutboundInterceptor().sideEffect(resultClass, resultType, func); @@ -710,6 +729,11 @@ public static void sleep(Duration duration) { getWorkflowOutboundInterceptor().sleep(duration); } + public static void sleep(Duration duration, TimerOptions options) { + assertNotReadOnly("sleep"); + getWorkflowOutboundInterceptor().sleep(duration, options); + } + public static boolean isWorkflowThread() { return WorkflowThreadMarker.isWorkflowThread(); } 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 40ecad495c..5f1c1f86f6 100644 --- a/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java +++ b/temporal-sdk/src/main/java/io/temporal/workflow/Workflow.java @@ -885,6 +885,17 @@ public static void sleep(long millis) { WorkflowInternal.sleep(Duration.ofMillis(millis)); } + /** + * Must be called instead of {@link Thread#sleep(long)} to guarantee determinism. + * + * @param duration time to sleep. + * @param options options for the underlying timer, such as its summary. + * @see #newTimer(Duration, TimerOptions) + */ + public static void sleep(Duration duration, TimerOptions options) { + WorkflowInternal.sleep(duration, options); + } + /** * Block current thread until unblockCondition is evaluated to true. * @@ -926,6 +937,31 @@ public static boolean await(Duration timeout, Supplier unblockCondition }); } + /** + * Block current workflow thread until unblockCondition is evaluated to true or timeout passes. + * + * @param timeout time to unblock even if unblockCondition is not satisfied. + * @param options options for the timer that implements the timeout, such as its summary. + * @param unblockCondition condition that should return true to indicate that thread should + * unblock. The condition is called on every state transition, so it should not contain any + * code that mutates any workflow state. It should also not contain any time based conditions. + * Use timeout parameter for those. + * @return false if timed out. + * @throws CanceledFailure if thread (or current {@link CancellationScope} was canceled). + * @see #newTimer(Duration, TimerOptions) + */ + public static boolean await( + Duration timeout, TimerOptions options, Supplier unblockCondition) { + return WorkflowInternal.await( + timeout, + options, + "await", + () -> { + CancellationScope.throwCanceled(); + return unblockCondition.get(); + }); + } + /** * Invokes function retrying in case of failures according to retry options. Synchronous variant. * Use {@link Async#retry(RetryOptions, Optional, Functions.Func)} for asynchronous functions. diff --git a/temporal-sdk/src/test/java/io/temporal/internal/sync/SyncWorkflowContextTest.java b/temporal-sdk/src/test/java/io/temporal/internal/sync/SyncWorkflowContextTest.java index acd1221819..10a7627c0e 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/sync/SyncWorkflowContextTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/sync/SyncWorkflowContextTest.java @@ -1,6 +1,8 @@ package io.temporal.internal.sync; import static org.junit.Assert.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -9,17 +11,22 @@ import io.temporal.api.command.v1.ContinueAsNewWorkflowExecutionCommandAttributes; import io.temporal.api.common.v1.SearchAttributes; import io.temporal.api.common.v1.WorkflowType; +import io.temporal.api.sdk.v1.UserMetadata; +import io.temporal.common.converter.DefaultDataConverter; import io.temporal.common.interceptors.Header; import io.temporal.common.interceptors.WorkflowOutboundCallsInterceptor.ContinueAsNewInput; +import io.temporal.internal.common.SdkFlag; import io.temporal.internal.common.SearchAttributesUtil; import io.temporal.internal.logging.PrefixedMdc; import io.temporal.internal.replay.ReplayWorkflowContext; import io.temporal.workflow.ContinueAsNewOptions; +import io.temporal.workflow.TimerOptions; import java.time.Duration; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.function.Supplier; import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentCaptor; @@ -51,6 +58,45 @@ public void testUpsertSearchAttributesException() { context.upsertSearchAttributes(attr); } + /** + * Live workflow tests take the CANCEL_AWAIT_TIMER_ON_CONDITION branch, so this forces the branch + * used without that flag through a mocked replay context. + */ + @Test + public void testAwaitPassesTimerSummaryWithoutCancelAwaitTimerFlag() { + when(mockReplayWorkflowContext.tryUseSdkFlag(SdkFlag.CANCEL_AWAIT_TIMER_ON_CONDITION)) + .thenReturn(false); + when(mockReplayWorkflowContext.getWorkflowType()) + .thenReturn(WorkflowType.newBuilder().setName("dummy-workflow").build()); + ExecutorService threadPool = Executors.newCachedThreadPool(); + Supplier neverTrue = () -> false; + DeterministicRunner runner = + DeterministicRunner.newRunner( + threadPool::submit, + context, + () -> + context.await( + Duration.ofHours(1), + TimerOptions.newBuilder().setSummary("await-summary").build(), + "await", + neverTrue)); + + try { + runner.runUntilAllBlocked(DeterministicRunner.DEFAULT_DEADLOCK_DETECTION_TIMEOUT_MS); + + ArgumentCaptor metadata = ArgumentCaptor.forClass(UserMetadata.class); + verify(mockReplayWorkflowContext) + .newTimer(eq(Duration.ofHours(1)), metadata.capture(), any()); + assertEquals( + "await-summary", + DefaultDataConverter.STANDARD_INSTANCE.fromPayload( + metadata.getValue().getSummary(), String.class, String.class)); + } finally { + runner.close(); + threadPool.shutdown(); + } + } + @Test public void testContinueAsNewBackoffStartInterval() { ExecutorService threadPool = Executors.newCachedThreadPool(); diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/TimerMetadataTest.java b/temporal-sdk/src/test/java/io/temporal/workflow/TimerMetadataTest.java new file mode 100644 index 0000000000..69725de2c3 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/workflow/TimerMetadataTest.java @@ -0,0 +1,136 @@ +package io.temporal.workflow; + +import static org.junit.Assert.assertEquals; + +import io.temporal.api.enums.v1.EventType; +import io.temporal.api.history.v1.HistoryEvent; +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowStub; +import io.temporal.testUtils.HistoryUtils; +import io.temporal.testing.WorkflowReplayer; +import io.temporal.testing.internal.SDKTestWorkflowRule; +import io.temporal.workflow.cancellationTests.WorkflowAwaitCancelTimerOnConditionTest.TestAwaitWorkflow; +import java.time.Duration; +import java.util.List; +import java.util.stream.Collectors; +import org.junit.Rule; +import org.junit.Test; + +public class TimerMetadataTest { + + private static final String SLEEP_SUMMARY = "sleep-summary"; + private static final String AWAIT_SUMMARY = "await-summary"; + + @Rule + public SDKTestWorkflowRule testWorkflowRule = + SDKTestWorkflowRule.newBuilder() + .setWorkflowTypes( + SleepWithSummaryWorkflowImpl.class, + AwaitTimeoutWithSummaryWorkflowImpl.class, + AwaitWithSummaryWorkflowImpl.class) + .build(); + + @Test + public void sleepSetsTimerSummary() { + SleepWorkflow workflow = testWorkflowRule.newWorkflowStub(SleepWorkflow.class); + workflow.execute(); + + List timers = + timerStartedEvents(WorkflowStub.fromTyped(workflow).getExecution().getWorkflowId()); + assertEquals(1, timers.size()); + HistoryUtils.assertEventMetadata(timers.get(0), SLEEP_SUMMARY, null); + } + + @Test + public void awaitSetsTimerSummaryWhenTimingOut() { + AwaitTimeoutWorkflow workflow = testWorkflowRule.newWorkflowStub(AwaitTimeoutWorkflow.class); + assertEquals("timed out", workflow.execute()); + + String workflowId = WorkflowStub.fromTyped(workflow).getExecution().getWorkflowId(); + List timers = timerStartedEvents(workflowId); + assertEquals(1, timers.size()); + HistoryUtils.assertEventMetadata(timers.get(0), AWAIT_SUMMARY, null); + testWorkflowRule.assertHistoryEvent(workflowId, EventType.EVENT_TYPE_TIMER_FIRED); + } + + @Test + public void awaitSetsTimerSummaryWhenConditionIsSatisfied() { + TestAwaitWorkflow workflow = testWorkflowRule.newWorkflowStub(TestAwaitWorkflow.class); + String workflowId = WorkflowClient.start(workflow::execute).getWorkflowId(); + testWorkflowRule.waitForTheEndOfWFT(workflowId); + workflow.unblock(); + assertEquals("condition satisfied", WorkflowStub.fromTyped(workflow).getResult(String.class)); + + List timers = timerStartedEvents(workflowId); + assertEquals(1, timers.size()); + HistoryUtils.assertEventMetadata(timers.get(0), AWAIT_SUMMARY, null); + testWorkflowRule.assertHistoryEvent(workflowId, EventType.EVENT_TYPE_TIMER_CANCELED); + } + + /** + * The history was recorded by an older SDK without CANCEL_AWAIT_TIMER_ON_CONDITION and without a + * timer summary. Adding a summary must not cause a nondeterminism error on replay. + */ + @Test + public void awaitWithSummaryReplaysHistoryRecordedWithoutSummary() throws Exception { + WorkflowReplayer.replayWorkflowExecutionFromResource( + "awaitTimerConditionOldBehavior.json", AwaitWithSummaryWorkflowImpl.class); + } + + private List timerStartedEvents(String workflowId) { + return testWorkflowRule.getWorkflowClient().fetchHistory(workflowId).getEvents().stream() + .filter(HistoryEvent::hasTimerStartedEventAttributes) + .collect(Collectors.toList()); + } + + @WorkflowInterface + public interface SleepWorkflow { + @WorkflowMethod + void execute(); + } + + @WorkflowInterface + public interface AwaitTimeoutWorkflow { + @WorkflowMethod + String execute(); + } + + public static class SleepWithSummaryWorkflowImpl implements SleepWorkflow { + @Override + public void execute() { + Workflow.sleep( + Duration.ofMillis(100), TimerOptions.newBuilder().setSummary(SLEEP_SUMMARY).build()); + } + } + + public static class AwaitTimeoutWithSummaryWorkflowImpl implements AwaitTimeoutWorkflow { + @Override + public String execute() { + boolean satisfied = + Workflow.await( + Duration.ofMillis(100), + TimerOptions.newBuilder().setSummary(AWAIT_SUMMARY).build(), + () -> false); + return satisfied ? "condition satisfied" : "timed out"; + } + } + + public static class AwaitWithSummaryWorkflowImpl implements TestAwaitWorkflow { + private boolean unblocked = false; + + @Override + public String execute() { + boolean result = + Workflow.await( + Duration.ofHours(1), + TimerOptions.newBuilder().setSummary(AWAIT_SUMMARY).build(), + () -> unblocked); + return result ? "condition satisfied" : "timed out"; + } + + @Override + public void unblock() { + unblocked = true; + } + } +} diff --git a/temporal-testing/src/main/java/io/temporal/testing/TestActivityEnvironmentInternal.java b/temporal-testing/src/main/java/io/temporal/testing/TestActivityEnvironmentInternal.java index dc2a0d3d3c..8217007744 100644 --- a/temporal-testing/src/main/java/io/temporal/testing/TestActivityEnvironmentInternal.java +++ b/temporal-testing/src/main/java/io/temporal/testing/TestActivityEnvironmentInternal.java @@ -405,11 +405,22 @@ public void sleep(Duration duration) { throw new UnsupportedOperationException("not implemented"); } + @Override + public void sleep(Duration duration, TimerOptions options) { + throw new UnsupportedOperationException("not implemented"); + } + @Override public boolean await(Duration timeout, String reason, Supplier unblockCondition) { throw new UnsupportedOperationException("not implemented"); } + @Override + public boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition) { + throw new UnsupportedOperationException("not implemented"); + } + @Override public void await(String reason, Supplier unblockCondition) { throw new UnsupportedOperationException("not implemented"); diff --git a/temporal-testing/src/main/java/io/temporal/testing/internal/TracingWorkerInterceptor.java b/temporal-testing/src/main/java/io/temporal/testing/internal/TracingWorkerInterceptor.java index 77d82bdaee..c807c4c2fe 100644 --- a/temporal-testing/src/main/java/io/temporal/testing/internal/TracingWorkerInterceptor.java +++ b/temporal-testing/src/main/java/io/temporal/testing/internal/TracingWorkerInterceptor.java @@ -217,6 +217,14 @@ public void sleep(Duration duration) { next.sleep(duration); } + @Override + public void sleep(Duration duration, TimerOptions options) { + if (!WorkflowUnsafe.isReplaying()) { + trace.add("sleep " + duration); + } + next.sleep(duration, options); + } + @Override public boolean await(Duration timeout, String reason, Supplier unblockCondition) { if (!WorkflowUnsafe.isReplaying()) { @@ -225,6 +233,15 @@ public boolean await(Duration timeout, String reason, Supplier unblockC return next.await(timeout, reason, unblockCondition); } + @Override + public boolean await( + Duration timeout, TimerOptions options, String reason, Supplier unblockCondition) { + if (!WorkflowUnsafe.isReplaying()) { + trace.add("await " + timeout + " " + reason); + } + return next.await(timeout, options, reason, unblockCondition); + } + @Override public void await(String reason, Supplier unblockCondition) { if (!WorkflowUnsafe.isReplaying()) {