diff --git a/temporal-sdk/src/test/java/io/temporal/workflow/TerminatedWorkflowTest.java b/temporal-sdk/src/test/java/io/temporal/workflow/TerminatedWorkflowTest.java index 6e8b7676b0..914d314adc 100644 --- a/temporal-sdk/src/test/java/io/temporal/workflow/TerminatedWorkflowTest.java +++ b/temporal-sdk/src/test/java/io/temporal/workflow/TerminatedWorkflowTest.java @@ -2,18 +2,15 @@ import static org.junit.Assert.*; +import io.temporal.api.common.v1.WorkflowExecution; +import io.temporal.api.enums.v1.EventType; import io.temporal.client.*; import io.temporal.failure.TerminatedFailure; +import io.temporal.failure.TimeoutFailure; import io.temporal.testing.internal.SDKTestOptions; import io.temporal.testing.internal.SDKTestWorkflowRule; -import io.temporal.workflow.shared.TestActivities.TestActivitiesImpl; -import io.temporal.workflow.shared.TestActivities.VariousTestActivities; -import io.temporal.workflow.shared.TestWorkflows.TestTraceWorkflow; +import io.temporal.workflow.shared.TestWorkflows.TestWorkflowWithQuery; import java.time.Duration; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; import org.junit.Assert; import org.junit.Rule; import org.junit.Test; @@ -22,79 +19,80 @@ * Tests verifying the correct behavior of the SDK if the workflow is in unsuccessful final states */ public class TerminatedWorkflowTest { - private final TestActivitiesImpl activitiesImpl = new TestActivitiesImpl(); - @Rule public SDKTestWorkflowRule testWorkflowRule = - SDKTestWorkflowRule.newBuilder() - .setWorkflowTypes(TraceTimingOutWorkflowImpl.class) - .setActivityImplementations(activitiesImpl) - .setWorkflowClientOptions(WorkflowClientOptions.newBuilder().build()) - .build(); + SDKTestWorkflowRule.newBuilder().setWorkflowTypes(NeverEndingWorkflowImpl.class).build(); @Test public void testShouldReturnQueryResultAfterWorkflowTimeout() { - WorkflowOptions options = - SDKTestOptions.newWorkflowOptionsWithTimeouts(testWorkflowRule.getTaskQueue()).toBuilder() - .setWorkflowRunTimeout(Duration.ofSeconds(1)) - .build(); - TestTraceWorkflow workflow = - testWorkflowRule.getWorkflowClient().newWorkflowStub(TestTraceWorkflow.class, options); - - Assert.assertThrows( - "Workflow should throw because of timeout", - WorkflowFailedException.class, - workflow::execute); + // Queries only work if workflow task was started at least once before closing. + // Without time skipping, we can't ensure this happens before the timeout. + // To keep the test execution time low but to avoid making it time-sensitive, we are retrying + // the workflow with exponentially growing timeouts until we get an execution that we can query. + TestWorkflowWithQuery workflow; + for (int i = 0; ; i++) { + WorkflowOptions options = + SDKTestOptions.newWorkflowOptionsWithTimeouts(testWorkflowRule.getTaskQueue()).toBuilder() + .setWorkflowRunTimeout(Duration.ofSeconds((int) Math.pow(2, i))) + .build(); + + workflow = + testWorkflowRule + .getWorkflowClient() + .newWorkflowStub(TestWorkflowWithQuery.class, options); + + WorkflowFailedException e = + Assert.assertThrows( + "Workflow should throw because of timeout", + WorkflowFailedException.class, + workflow::execute); + Assert.assertTrue(e.getCause() instanceof TimeoutFailure); + + WorkflowExecution execution = WorkflowStub.fromTyped(workflow).getExecution(); + if (testWorkflowRule + .getWorkflowClient() + .streamHistory(execution.getWorkflowId(), execution.getRunId()) + .anyMatch(evt -> evt.getEventType() == EventType.EVENT_TYPE_WORKFLOW_TASK_STARTED)) { + break; + } + } - Assert.assertEquals(1, workflow.getTrace().size()); - Assert.assertEquals("started", workflow.getTrace().get(0)); + Assert.assertEquals("started", workflow.query()); } @Test public void getResultShouldThrowAfterTerminationOfWorkflow() { WorkflowOptions options = WorkflowOptions.newBuilder().setTaskQueue(testWorkflowRule.getTaskQueue()).build(); - WorkflowStub workflow = - testWorkflowRule.getWorkflowClient().newUntypedWorkflowStub("execute", options); + testWorkflowRule + .getWorkflowClient() + .newUntypedWorkflowStub("TestWorkflowWithQuery", options); workflow.start(); - workflow.terminate("testing"); - WorkflowFailedException exception = null; - try { - workflow.getResult(1000, TimeUnit.MILLISECONDS, String.class); - fail("getResult should throw WorkflowFailedException because the workflow was terminated"); - } catch (WorkflowFailedException e) { - // This is expected - exception = e; - } catch (TimeoutException e) { - fail( - "getResult shouldn't wait all 5 seconds till the end of the workflow because it was already terminated"); - } - assertNotNull(exception); - assertTrue(exception.getCause() instanceof TerminatedFailure); + WorkflowFailedException e = + Assert.assertThrows( + "Workflow should throw WorkflowFailedException because the workflow was terminated", + WorkflowFailedException.class, + () -> workflow.getResult(String.class)); + assertTrue(e.getCause() instanceof TerminatedFailure); } - public static class TraceTimingOutWorkflowImpl implements TestTraceWorkflow { - private final List trace = new ArrayList<>(); + public static class NeverEndingWorkflowImpl implements TestWorkflowWithQuery { + private String state = "not started"; @Override public String execute() { - VariousTestActivities localActivities = - Workflow.newLocalActivityStub( - VariousTestActivities.class, SDKTestOptions.newLocalActivityOptions()); - - trace.add("started"); - localActivities.sleepActivity(5000, 123); - trace.add("finished"); - return ""; + state = "started"; + Workflow.await(() -> false); + return "should never happen"; } @Override - public List getTrace() { - return trace; + public String query() { + return state; } } }