temporalio / temporalio/sdk-java

Workflow execution with Workflow.await(condition) times out in unit tests with enabled time skipping

Open
#1,291 3 comments 1 reaction 0 assignees View on GitHub

Nobody has claimed this yet.

bug test server
Dominant language
Java
Stars
433
Forks
249
Avg merge
5d 6h
Merged PRs (30d)
26

Description

Expected Behavior

The unit test below should always pass

Actual Behavior

Sometimes the test fails with io.temporal.client.WorkflowNotFoundException. Changing Workflow.await(condition) to Workflow.await(Duration.ofSeconds(100), condition) in TestWorkflowImpl seems to fix the problem, but not sure why.
Attached are the TRACE logs for io.temporal for when the issue reproduces: bug.log

Steps to Reproduce the Problem

Run the following test:

import io.temporal.activity.ActivityOptions;
import io.temporal.client.WorkflowClient;
import io.temporal.client.WorkflowOptions;
import io.temporal.client.WorkflowStub;
import io.temporal.testing.TestWorkflowEnvironment;
import io.temporal.testing.TestWorkflowExtension;
import io.temporal.worker.Worker;
import io.temporal.worker.WorkflowImplementationOptions;
import io.temporal.workflow.*;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;

import java.time.Duration;
import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.stream.IntStream;

import static org.junit.jupiter.api.Assertions.assertEquals;

public class WorkflowExecutionTimeoutTest {
    @RegisterExtension
    public static final TestWorkflowExtension TEST_WORKFLOW_EXTENSION =
            TestWorkflowExtension.newBuilder()
                    .setDoNotStart(true)
                    .build();

    private ProcessEventsWorkflow processWorkflowStub;
    private TestWorkflow testWorkflowStub;

    @BeforeEach
    public void setUpTemporal(TestWorkflowEnvironment testEnv,
                              Worker worker,
                              WorkflowClient workflowClient,
                              WorkflowOptions workflowOptions) {
        worker.registerWorkflowImplementationTypes(
                WorkflowImplementationOptions.newBuilder()
                        .setDefaultActivityOptions(ActivityOptions.newBuilder()
                                .setStartToCloseTimeout(Duration.ofSeconds(10))
                                .build())
                        .build(),
                ProcessEventsWorkflowImpl.class,
                TestWorkflowImpl.class);

        testEnv.start();

        processWorkflowStub = workflowClient.newWorkflowStub(ProcessEventsWorkflow.class,
                WorkflowOptions.newBuilder(workflowOptions)
                        .setWorkflowId("ProcessEventsWorkflow")
                        .build());
        testWorkflowStub = workflowClient.newWorkflowStub(TestWorkflow.class,
                WorkflowOptions.newBuilder(workflowOptions)
                        .setWorkflowId("TestWorkflow")
                        .build());
    }

    @Test
    public void testBug() throws TimeoutException {
        // create artificial load to reproduce the bug - seems to help, but still the bug does not always reproduce
        IntStream.range(0, 20).forEach(index -> new Thread(this::busyWork).start());

        WorkflowClient.start(testWorkflowStub::execute);

        WorkflowStub.fromTyped(processWorkflowStub).signalWithStart("addEvent",
                new Object[]{"testEvent"},
                new Object[]{"TestWorkflow", Duration.ofSeconds(1)});
        WorkflowStub.fromTyped(processWorkflowStub).getResult(10, TimeUnit.SECONDS, Object.class);

        testWorkflowStub.stop(); // <----- fails here
        WorkflowStub.fromTyped(testWorkflowStub).getResult(10, TimeUnit.SECONDS, Object.class);

        assertEquals(Arrays.asList("testEvent"), testWorkflowStub.getEvents());
    }

    private void busyWork() {
        int count = 100000000;
        int sleepIndex = (int) (Math.random() * count);
        while(count-- > 0) {
            if (sleepIndex == count) { // yield at random intervals
                try {
                    Thread.sleep(100);
                } catch (InterruptedException e) {
                    throw new RuntimeException(e);
                }
            }

            Math.sqrt(count);
        }
    }


    @WorkflowInterface
    public interface ProcessEventsWorkflow {
        @WorkflowMethod
        void execute(String targetWorkflowId, Duration keepAliveTimeout);
        @SignalMethod
        void addEvent(String event);
    }

    public static class ProcessEventsWorkflowImpl implements ProcessEventsWorkflow {
        private final Queue<String> events = new LinkedList<>();

        @Override
        public void execute(String targetWorkflowId, Duration keepAliveTimeout) {
            while (true) {
                while (!events.isEmpty()) {
                    String event = events.poll();
                    Workflow.newExternalWorkflowStub(TestWorkflow.class, targetWorkflowId).onEvent(event);
                }

                Workflow.await(keepAliveTimeout, () -> !events.isEmpty());
                if (events.isEmpty()) {
                    return;
                }
            }
        }

        @Override
        public void addEvent(String event) {
            events.add(event);
        }
    }

    @WorkflowInterface
    public interface TestWorkflow {
        @WorkflowMethod
        void execute();

        @SignalMethod
        void stop();

        @SignalMethod
        void onEvent(String event);

        @QueryMethod
        List<String> getEvents();
    }

    public static class TestWorkflowImpl implements TestWorkflow {
        private boolean stop = false;
        private final List<String> events = new ArrayList<>();

        @Override
        public void execute() {
            Workflow.await(() -> stop);
        }

        @Override
        public void stop() {
            stop = true;
        }

        @Override
        public void onEvent(String event) {
            events.add(event);
        }

        @Override
        public List<String> getEvents() {
            return events;
        }
    }
}

Specifications

  • Version: Temporal Java SDK 1.11.0, 1.12.0, 1.13.0 (reproduces on all of these), Temporal Server 1.16.2
  • Platform: Java

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start by running the provided WorkflowExecutionTimeoutTest with time skipping enabled and review the attached bug.log TRACE output. Focus on the interaction between Workflow.await(condition), test-environment time skipping, and the stop signal; done means the reproduction no longer intermittently raises WorkflowNotFoundException across the listed SDK versions.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
distributed-systems, testing-qa
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.