Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This condition is going to check shortly after starting the workflow and possibly not give the task to the worker on the first few iterations. Should we wait for the workflow to close via timeout before checking there is a WFT started?

@maciejdudko maciejdudko Oct 2, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I thought streamHistory is going to keep waiting for more events. Apparently not. Going to replace it with direct GetWorkflowExecutionHistory RPC call streaming history after execution timeout.

The previous (flaky) version waited for the timeout, but sometimes the server reached the timeout before issuing the first task, and then queries didn't work. The idea here is to ensure the workflow task was issued before the timeout, and if not, try again.

.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<String> 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<String> getTrace() {
return trace;
public String query() {
return state;
}
}
}
Loading