Skip to content
Merged
Show file tree
Hide file tree
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
13 changes: 2 additions & 11 deletions tests/e2e/workflow_execution_suite_e2e_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/stretchr/testify/suite"
)

Expand Down Expand Up @@ -43,11 +42,7 @@ func (s *WorkflowExecutionSuite) triggerAndWaitFinished(schemaID string) {
t := s.T()
t.Helper()

wfID := TriggerExampleWorkflow(t, s.client, s.baseURL, schemaID)
require.NotEmpty(t, wfID, "trigger should return a workflow ID")

resp, err := WaitForWorkflowTerminal(s.client, s.baseURL, wfID, FastStatusTimeout)
require.NoError(t, err, "workflow %s should reach terminal state", schemaID)
wfID, resp := TriggerAndWaitTerminal(t, s.client, s.baseURL, schemaID, FastStatusTimeout)
assert.Equal(t, wfID, resp.WorkflowID, "response should echo back the workflow ID")
assert.Equal(t, "finished", resp.Status, "workflow %s should finish successfully", schemaID)
}
Expand Down Expand Up @@ -80,10 +75,6 @@ func (s *WorkflowExecutionSuite) TestDurableExecution_Finishes() {
// TestSumRandBranch_Finishes runs a workflow with a 3s timer, three parallel rands, sum, and conditional branching.
func (s *WorkflowExecutionSuite) TestSumRandBranch_Finishes() {
t := s.T()
wfID := TriggerExampleWorkflow(t, s.client, s.baseURL, "sum-rand-branch")
require.NotEmpty(t, wfID)

resp, err := WaitForWorkflowTerminal(s.client, s.baseURL, wfID, LongStatusTimeout)
require.NoError(t, err, "workflow should reach terminal state")
_, resp := TriggerAndWaitTerminal(t, s.client, s.baseURL, "sum-rand-branch", LongStatusTimeout)
assert.Equal(t, "finished", resp.Status, "sum-rand-branch should finish successfully")
}
36 changes: 36 additions & 0 deletions tests/e2e/workflow_fixture_e2e_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,46 @@ import (
"net/http"
"path/filepath"
"testing"
"time"

"github.com/stretchr/testify/require"
)

// untriggeredReTriggerAttempts bounds how many times TriggerAndWaitTerminal re-triggers a
// workflow that gets stuck in "untriggered".
const untriggeredReTriggerAttempts = 3

// TriggerAndWaitTerminal triggers schemaID and waits for a terminal state, returning the final
// status.
//
// It works around an intermittent e2e issue where a freshly-triggered workflow can remain in
// "untriggered" — the trigger is accepted (a workflowId is returned) but the workflow actor is
// not spawned/claimed in time on the multi-node HA stack. This is a pre-existing flake unrelated
// to workflow logic (observed on main across unrelated workflows, and on the same commit both
// passing and failing); tracked separately for a proper root-cause. Only the stuck-"untriggered"
// state is retried by re-triggering (each trigger yields a fresh workflowId); a workflow that
// reaches any terminal state — including "error" — is returned as-is so real failures still surface.
func TriggerAndWaitTerminal(t *testing.T, client *http.Client, baseURL, schemaID string, timeout time.Duration) (string, *WorkflowStatusResponse) {
t.Helper()
var wfID string
var resp *WorkflowStatusResponse
var err error
for attempt := 1; attempt <= untriggeredReTriggerAttempts; attempt++ {
wfID = TriggerExampleWorkflow(t, client, baseURL, schemaID)
resp, err = WaitForWorkflowTerminal(client, baseURL, wfID, timeout)
if err == nil {
return wfID, resp
}
if resp == nil || resp.Status != "untriggered" {
break
}
t.Logf("e2e: workflow %s for %q stuck in 'untriggered' (attempt %d/%d); re-triggering [tracked flake]",
wfID, schemaID, attempt, untriggeredReTriggerAttempts)
}
require.NoError(t, err, "workflow %s should reach terminal state", schemaID)
return wfID, resp
}

// UpsertSchema PUTs the JSON schema without triggering it.
// If an e2e overlay variant exists, it is used instead of the default schema.
func UpsertSchema(t *testing.T, client *http.Client, baseURL, workflowsDir, schemaID string) {
Expand Down
Loading