mirror of
https://github.com/temporalio/temporal.git
synced 2026-08-30 18:41:49 -07:00
Flake fixer: use s.assertions and new assert struct within eventuallyWithT (#8867)
## What changed?
- WISOTT
- When the function verifyWorkflowVersioning is called from within an
eventuallyWithT, we have to be careful and pass in a new Assertions
object to it. This is because on encountering failures inside the
function, we want to collect those failures and let eventuallyWithT
retry again. This was not happening right now since we were reporting
failures back on the assertions object of the test suite itself.
- For all other tests that are not calling verifyWorkflowVersioning from
within the eventuallyWithT, just pass in the underlying test suite's
assertion object and fail fast.
## Why?
- To resolve a flake and also testing correctness.
## How did you test it?
- [x] built
- [ ] run locally and tested manually
- [ ] covered by existing tests
- [ ] added new unit test(s)
- [ ] added new functional test(s)
## Potential risks
- None
<!-- CURSOR_SUMMARY -->
---
> [!NOTE]
> Route assertions through a passed-in require.Assertions to allow
failure collection in EventuallyWithT; update all test call sites
accordingly.
>
> - **Tests (versioning_3_test.go)**:
> - Change `verifyWorkflowVersioning` signature to accept
`*require.Assertions` and replace internal `s.*` assertions with the
provided handle.
> - Update all callers to pass either `s.Assertions` or a new
`require.New(t)` (inside `EventuallyWithT`), ensuring failures are
collected/retried correctly.
> - Add localized creation of `require.Assertions` in several
`EventuallyWithT` blocks and forward to `verifyWorkflowVersioning`.
> - No functional logic changes to test flows; only assertion
plumbing/refactoring.
>
> <sup>Written by [Cursor
Bugbot](https://cursor.com/dashboard?tab=bugbot) for commit
0a32731d15. This will update automatically
on new commits. Configure
[here](https://cursor.com/dashboard?tab=bugbot).</sup>
<!-- /CURSOR_SUMMARY -->
This commit is contained in:
@@ -275,21 +275,21 @@ func (s *Versioning3Suite) testWorkflowWithPinnedOverride(sticky bool) {
|
||||
runID := s.startWorkflow(tv, tv.VersioningOverridePinned(s.useV32))
|
||||
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyVersioningSAs(tv, vbPinned, tv)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv.WithRunID(runID))
|
||||
}
|
||||
|
||||
s.WaitForChannel(ctx, actCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
|
||||
s.pollWftAndHandle(tv, sticky, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
return respondCompleteWorkflow(tv, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestQueryWithPinnedOverride_NoSticky() {
|
||||
@@ -362,7 +362,7 @@ func (s *Versioning3Suite) testPinnedQuery_DrainedVersion(pollersPresent bool, r
|
||||
|
||||
s.startWorkflow(tv, tv.VersioningOverridePinned(s.useV32))
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbPinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbPinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
|
||||
// create version v2 and make it current which shall make v1 go from current -> draining/drained
|
||||
idlePollerDone = make(chan struct{})
|
||||
@@ -445,7 +445,7 @@ func (s *Versioning3Suite) testQueryWithPinnedOverride(sticky bool) {
|
||||
runID := s.startWorkflow(tv, tv.VersioningOverridePinned(s.useV32))
|
||||
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), tv.VersioningOverridePinned(s.useV32), nil)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv.WithRunID(runID))
|
||||
}
|
||||
@@ -483,7 +483,7 @@ func (s *Versioning3Suite) testUnpinnedQuery(sticky bool) {
|
||||
s.pollWftAndHandle(tv, false, wftCompleted,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
return respondEmptyWft(tv, sticky, vbUnpinned), nil
|
||||
})
|
||||
|
||||
@@ -493,7 +493,7 @@ func (s *Versioning3Suite) testUnpinnedQuery(sticky bool) {
|
||||
runID := s.startWorkflow(tv, nil)
|
||||
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv.WithRunID(runID))
|
||||
}
|
||||
@@ -563,7 +563,7 @@ func (s *Versioning3Suite) testPinnedWorkflowWithLateActivityPoller() {
|
||||
s.startWorkflow(tv, override)
|
||||
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
// Wait long enough to make sure the activity is backlogged.
|
||||
s.validateBacklogCount(tv, tqTypeAct, 1)
|
||||
|
||||
@@ -574,7 +574,7 @@ func (s *Versioning3Suite) testPinnedWorkflowWithLateActivityPoller() {
|
||||
s.NotNil(task)
|
||||
return respondActivity(), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
s.validateBacklogCount(tv, tqTypeAct, 0)
|
||||
|
||||
s.pollWftAndHandle(tv, false, nil,
|
||||
@@ -582,7 +582,7 @@ func (s *Versioning3Suite) testPinnedWorkflowWithLateActivityPoller() {
|
||||
s.NotNil(task)
|
||||
return respondCompleteWorkflow(tv, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), override, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestUnpinnedWorkflow_Sticky() {
|
||||
@@ -615,7 +615,7 @@ func (s *Versioning3Suite) testUnpinnedWorkflow(sticky bool) {
|
||||
s.pollWftAndHandle(tv, false, wftCompleted,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
return respondWftWithActivities(tv, tv, sticky, vbUnpinned, "5"), nil
|
||||
})
|
||||
|
||||
@@ -631,21 +631,21 @@ func (s *Versioning3Suite) testUnpinnedWorkflow(sticky bool) {
|
||||
runID := s.startWorkflow(tv, nil)
|
||||
|
||||
s.WaitForChannel(ctx, wftCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyVersioningSAs(tv, vbUnpinned, tv)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv.WithRunID(runID))
|
||||
}
|
||||
|
||||
s.WaitForChannel(ctx, actCompleted)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
|
||||
s.pollWftAndHandle(tv, sticky, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
return respondCompleteWorkflow(tv, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
// drainWorkflowTaskAfterSetCurrent is a helper that sets the current deployment version,
|
||||
@@ -658,7 +658,7 @@ func (s *Versioning3Suite) drainWorkflowTaskAfterSetCurrent(
|
||||
s.pollWftAndHandle(tv, false, wftCompleted,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
return respondEmptyWft(tv, false, vbUnpinned), nil
|
||||
})
|
||||
s.waitForDeploymentDataPropagation(tv, versionStatusInactive, false, tqTypeWf)
|
||||
@@ -713,7 +713,7 @@ func (s *Versioning3Suite) TestUnpinnedWorkflow_SuccessfulUpdate_TransitionsToNe
|
||||
|
||||
// VersioningInfo should not have changed before the update has been processed by the poller.
|
||||
// Deployment version transition should also be nil since this is a speculative task.
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
|
||||
return &workflowservice.RespondWorkflowTaskCompletedRequest{
|
||||
Commands: s.UpdateAcceptCompleteCommands(tv2),
|
||||
@@ -755,7 +755,7 @@ func (s *Versioning3Suite) TestUnpinnedWorkflow_SuccessfulUpdate_TransitionsToNe
|
||||
// Since the poller accepted the update, the Worker Deployment Version that completed the last workflow task
|
||||
// of this workflow execution should have changed to the new version. However, the version transition should
|
||||
// still be nil.
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
|
||||
}
|
||||
|
||||
@@ -799,7 +799,7 @@ func (s *Versioning3Suite) TestUnpinnedWorkflow_FailedUpdate_DoesNotTransitionTo
|
||||
|
||||
// VersioningInfo should not have changed before the update has been processed by the poller.
|
||||
// Deployment version transition should also be nil since this is a speculative task.
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
|
||||
updRequestMsg := task.Messages[0]
|
||||
updRequest := protoutils.UnmarshalAny[*updatepb.Request](s.T(), updRequestMsg.GetBody())
|
||||
@@ -833,7 +833,7 @@ func (s *Versioning3Suite) TestUnpinnedWorkflow_FailedUpdate_DoesNotTransitionTo
|
||||
|
||||
// Since the poller rejected the update, the Worker Deployment Version that completed the last workflow task
|
||||
// of this workflow execution should not have changed.
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) sendUpdateNoError(tv *testvars.TestVars) <-chan *workflowservice.UpdateWorkflowExecutionResponse {
|
||||
@@ -1221,10 +1221,10 @@ func (s *Versioning3Suite) testTransitionFromWft(sticky bool, toUnversioned bool
|
||||
s.pollWftAndHandle(tv1, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
return respondWftWithActivities(tv1, tv1, sticky, vbUnpinned, "5"), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv1.WithRunID(runID))
|
||||
}
|
||||
@@ -1234,7 +1234,7 @@ func (s *Versioning3Suite) testTransitionFromWft(sticky bool, toUnversioned bool
|
||||
s.NotNil(task)
|
||||
return respondActivity(), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
|
||||
if toUnversioned {
|
||||
// unset A as current
|
||||
@@ -1251,10 +1251,10 @@ func (s *Versioning3Suite) testTransitionFromWft(sticky bool, toUnversioned bool
|
||||
s.unversionedPollWftAndHandle(tv1, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, &workflowpb.DeploymentVersionTransition{Version: "__unversioned__"})
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, &workflowpb.DeploymentVersionTransition{Version: "__unversioned__"})
|
||||
return respondCompleteWorkflowUnversioned(tv1), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, vbUnspecified, nil, nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnspecified, nil, nil, nil)
|
||||
} else {
|
||||
|
||||
// Set B as the current deployment
|
||||
@@ -1275,10 +1275,10 @@ func (s *Versioning3Suite) testTransitionFromWft(sticky bool, toUnversioned bool
|
||||
s.pollWftAndHandle(tv2, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
return respondCompleteWorkflow(tv2, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1346,7 +1346,7 @@ func (s *Versioning3Suite) testDoubleTransition(unversionedSrc bool, signal bool
|
||||
s.NotNil(task)
|
||||
return respondWftWithActivities(tv1, tv1, false, sourceVB, "5"), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, sourceVB, sourceV, nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, sourceVB, sourceV, nil, nil)
|
||||
|
||||
if signal {
|
||||
// Send a signal so a wf task is scheduled before we poll the activity
|
||||
@@ -1417,16 +1417,16 @@ func (s *Versioning3Suite) testDoubleTransition(unversionedSrc bool, signal bool
|
||||
s.doPollWftAndHandle(tv1, !unversionedSrc, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, sourceVB, sourceV, nil, sourceTransition)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, sourceVB, sourceV, nil, sourceTransition)
|
||||
return respondEmptyWft(tv1, false, sourceVB), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, sourceVB, sourceV, nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, sourceVB, sourceV, nil, nil)
|
||||
|
||||
// Activity should be unblocked now to sourceV poller
|
||||
s.doPollActivityAndHandle(tv1, !unversionedSrc, nil,
|
||||
func(task *workflowservice.PollActivityTaskQueueResponse) (*workflowservice.RespondActivityTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, sourceVB, sourceV, nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, sourceVB, sourceV, nil, nil)
|
||||
return respondActivity(), nil
|
||||
})
|
||||
|
||||
@@ -1448,10 +1448,10 @@ func (s *Versioning3Suite) testDoubleTransition(unversionedSrc bool, signal bool
|
||||
s.pollWftAndHandle(tv2, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv2, sourceVB, sourceV, nil, tv2.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, sourceVB, sourceV, nil, tv2.DeploymentVersionTransition())
|
||||
return respondCompleteWorkflow(tv2, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestNexusTask_StaysOnCurrentDeployment() {
|
||||
@@ -1556,12 +1556,12 @@ func (s *Versioning3Suite) TestEagerActivity() {
|
||||
poller, resp := s.pollWftAndHandle(tv, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnspecified, nil, nil, tv.DeploymentVersionTransition())
|
||||
resp := respondWftWithActivities(tv, tv, true, vbUnpinned, "5")
|
||||
resp.Commands[0].GetScheduleActivityTaskCommandAttributes().RequestEagerExecution = true
|
||||
return resp, nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
|
||||
s.NotEmpty(resp.GetActivityTasks())
|
||||
|
||||
@@ -1571,14 +1571,14 @@ func (s *Versioning3Suite) TestEagerActivity() {
|
||||
return respondActivity(), nil
|
||||
})
|
||||
s.NoError(err)
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
|
||||
s.pollWftAndHandle(tv, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
return respondCompleteWorkflow(tv, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv, vbUnpinned, tv.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestTransitionFromActivity_Sticky() {
|
||||
@@ -1628,10 +1628,10 @@ func (s *Versioning3Suite) testTransitionFromActivity(sticky bool) {
|
||||
s.pollWftAndHandle(tv1, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
return respondWftWithActivities(tv1, tv1, sticky, vbUnpinned, "5", "6", "7", "8"), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv1.WithRunID(runID))
|
||||
}
|
||||
@@ -1671,7 +1671,7 @@ func (s *Versioning3Suite) testTransitionFromActivity(sticky bool) {
|
||||
})
|
||||
|
||||
s.WaitForChannel(ctx, act2Started)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
|
||||
// 2. Set d2 as the current deployment
|
||||
if s.useNewDeploymentData {
|
||||
@@ -1717,7 +1717,7 @@ func (s *Versioning3Suite) testTransitionFromActivity(sticky bool) {
|
||||
s.pollWftAndHandle(tv2, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
close(transitionStarted)
|
||||
s.Logger.Info("Transition wft started")
|
||||
// 8. Complete the transition after act1 completes and act2's first attempt fails.
|
||||
@@ -1727,7 +1727,7 @@ func (s *Versioning3Suite) testTransitionFromActivity(sticky bool) {
|
||||
s.Logger.Info("Transition wft completed")
|
||||
return respondEmptyWft(tv2, sticky, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
if sticky {
|
||||
s.verifyWorkflowStickyQueue(tv2)
|
||||
}
|
||||
@@ -1740,7 +1740,7 @@ func (s *Versioning3Suite) testTransitionFromActivity(sticky bool) {
|
||||
s.Logger.Info("Final wft completed")
|
||||
return respondCompleteWorkflow(tv2, vbUnpinned), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv2.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestIndependentVersionedActivity_Pinned() {
|
||||
@@ -1801,11 +1801,11 @@ func (s *Versioning3Suite) testIndependentActivity(behavior enumspb.VersioningBe
|
||||
s.pollWftAndHandle(tvWf, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
s.verifyWorkflowVersioning(tvWf, vbUnspecified, nil, nil, tvWf.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tvWf, vbUnspecified, nil, nil, tvWf.DeploymentVersionTransition())
|
||||
s.Logger.Info("First wf task completed")
|
||||
return respondWftWithActivities(tvWf, tvAct, false, behavior, "5"), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
|
||||
if unversionedActivity {
|
||||
s.unversionedPollActivityAndHandle(tvAct, nil,
|
||||
@@ -1822,14 +1822,14 @@ func (s *Versioning3Suite) testIndependentActivity(behavior enumspb.VersioningBe
|
||||
return respondActivity(), nil
|
||||
})
|
||||
}
|
||||
s.verifyWorkflowVersioning(tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
|
||||
s.pollWftAndHandle(tvWf, false, nil,
|
||||
func(task *workflowservice.PollWorkflowTaskQueueResponse) (*workflowservice.RespondWorkflowTaskCompletedRequest, error) {
|
||||
s.NotNil(task)
|
||||
return respondCompleteWorkflow(tvWf, behavior), nil
|
||||
})
|
||||
s.verifyWorkflowVersioning(tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tvWf, behavior, tvWf.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestChildWorkflowInheritance_PinnedParent() {
|
||||
@@ -1876,7 +1876,7 @@ func (s *Versioning3Suite) testChildWorkflowInheritance_ExpectInherit(crossTq bo
|
||||
currentChanged := make(chan struct{}, 1)
|
||||
|
||||
childv1 := func(ctx workflow.Context) (string, error) {
|
||||
s.verifyWorkflowVersioning(tv1Child, vbPinned, tv1Child.Deployment(), override, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1Child, vbPinned, tv1Child.Deployment(), override, nil)
|
||||
return "v1", nil
|
||||
}
|
||||
wf1 := func(ctx workflow.Context) (string, error) {
|
||||
@@ -1892,7 +1892,7 @@ func (s *Versioning3Suite) testChildWorkflowInheritance_ExpectInherit(crossTq bo
|
||||
var val1 string
|
||||
s.NoError(fut1.Get(ctx, &val1))
|
||||
|
||||
s.verifyWorkflowVersioning(tv1, parentRegistrationBehavior, tv1.Deployment(), override, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, parentRegistrationBehavior, tv1.Deployment(), override, nil)
|
||||
return val1, nil
|
||||
}
|
||||
|
||||
@@ -2030,7 +2030,7 @@ func (s *Versioning3Suite) testChildWorkflowInheritance_ExpectNoInherit(crossTq
|
||||
var val1 string
|
||||
s.NoError(fut1.Get(ctx, &val1))
|
||||
|
||||
s.verifyWorkflowVersioning(tv1, parentBehavior, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, parentBehavior, tv1.Deployment(), nil, nil)
|
||||
return val1, nil
|
||||
}
|
||||
|
||||
@@ -2121,11 +2121,11 @@ func (s *Versioning3Suite) testChildWorkflowInheritance_ExpectNoInherit(crossTq
|
||||
s.Equal("v2", out)
|
||||
|
||||
if parentBehavior == vbPinned {
|
||||
s.verifyWorkflowVersioning(tv1, parentBehavior, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, parentBehavior, tv1.Deployment(), nil, nil)
|
||||
} else {
|
||||
s.verifyWorkflowVersioning(tv1, parentBehavior, tv2.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, parentBehavior, tv2.Deployment(), nil, nil)
|
||||
}
|
||||
s.verifyWorkflowVersioning(tv2Child, vbPinned, tv2Child.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2Child, vbPinned, tv2Child.Deployment(), nil, nil)
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) TestPinnedCaN_SameTQ() {
|
||||
@@ -2167,13 +2167,13 @@ func (s *Versioning3Suite) testCan(crossTq bool, behavior enumspb.VersioningBeha
|
||||
if crossTq {
|
||||
newCtx = workflow.WithWorkflowTaskQueue(newCtx, canxTq)
|
||||
}
|
||||
s.verifyWorkflowVersioning(tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbUnspecified, nil, nil, tv1.DeploymentVersionTransition())
|
||||
wfStarted <- struct{}{}
|
||||
// wait for current version to change.
|
||||
<-currentChanged
|
||||
return "", workflow.NewContinueAsNewError(newCtx, "wf", attempt+1)
|
||||
case 1:
|
||||
s.verifyWorkflowVersioning(tv1, vbPinned, tv1.Deployment(), nil, nil)
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv1, vbPinned, tv1.Deployment(), nil, nil)
|
||||
return "v1", nil
|
||||
}
|
||||
s.FailNow("workflow should not get to this point")
|
||||
@@ -2183,9 +2183,9 @@ func (s *Versioning3Suite) testCan(crossTq bool, behavior enumspb.VersioningBeha
|
||||
wf2 := func(ctx workflow.Context, attempt int) (string, error) {
|
||||
if behavior == vbUnpinned && s.deploymentWorkflowVersion >= workerdeployment.AsyncSetCurrentAndRamping {
|
||||
// Unpinned CaN should inherit parent deployment version and behaviour
|
||||
s.verifyWorkflowVersioning(tv2, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnpinned, tv1.Deployment(), nil, tv2.DeploymentVersionTransition())
|
||||
} else {
|
||||
s.verifyWorkflowVersioning(tv2, vbUnspecified, nil, nil, tv2.DeploymentVersionTransition())
|
||||
s.verifyWorkflowVersioning(s.Assertions, tv2, vbUnspecified, nil, nil, tv2.DeploymentVersionTransition())
|
||||
}
|
||||
return "v2", nil
|
||||
}
|
||||
@@ -2968,6 +2968,7 @@ func (s *Versioning3Suite) forgetTaskQueueDeploymentVersion(
|
||||
}
|
||||
|
||||
func (s *Versioning3Suite) verifyWorkflowVersioning(
|
||||
a *require.Assertions,
|
||||
tv *testvars.TestVars,
|
||||
behavior enumspb.VersioningBehavior,
|
||||
deployment *deploymentpb.Deployment,
|
||||
@@ -2982,23 +2983,23 @@ func (s *Versioning3Suite) verifyWorkflowVersioning(
|
||||
},
|
||||
},
|
||||
)
|
||||
s.NoError(err)
|
||||
a.NoError(err)
|
||||
|
||||
versioningInfo := dwf.WorkflowExecutionInfo.GetVersioningInfo()
|
||||
s.Equal(behavior.String(), versioningInfo.GetBehavior().String())
|
||||
a.Equal(behavior.String(), versioningInfo.GetBehavior().String())
|
||||
var v *deploymentspb.WorkerDeploymentVersion
|
||||
if versioningInfo.GetVersion() != "" { //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
//nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
v, err = worker_versioning.WorkerDeploymentVersionFromStringV31(versioningInfo.GetVersion())
|
||||
s.NoError(err)
|
||||
s.NotNil(versioningInfo.GetDeploymentVersion()) // make sure we are always populating this whenever Version string is populated
|
||||
a.NoError(err)
|
||||
a.NotNil(versioningInfo.GetDeploymentVersion()) // make sure we are always populating this whenever Version string is populated
|
||||
}
|
||||
if dv := versioningInfo.GetDeploymentVersion(); dv != nil {
|
||||
v = worker_versioning.DeploymentVersionFromDeployment(worker_versioning.DeploymentFromExternalDeploymentVersion(dv))
|
||||
}
|
||||
actualDeployment := worker_versioning.DeploymentFromDeploymentVersion(v)
|
||||
if !deployment.Equal(actualDeployment) {
|
||||
s.Fail(fmt.Sprintf("deployment version mismatch. expected: {%s}, actual: {%s}",
|
||||
a.Fail(fmt.Sprintf("deployment version mismatch. expected: {%s}, actual: {%s}",
|
||||
deployment,
|
||||
actualDeployment,
|
||||
))
|
||||
@@ -3006,30 +3007,30 @@ func (s *Versioning3Suite) verifyWorkflowVersioning(
|
||||
|
||||
if s.useV32 {
|
||||
// v0.32 override
|
||||
s.Equal(override.GetAutoUpgrade(), versioningInfo.GetVersioningOverride().GetAutoUpgrade())
|
||||
s.Equal(override.GetPinned().GetVersion().GetBuildId(), versioningInfo.GetVersioningOverride().GetPinned().GetVersion().GetBuildId())
|
||||
s.Equal(override.GetPinned().GetVersion().GetDeploymentName(), versioningInfo.GetVersioningOverride().GetPinned().GetVersion().GetDeploymentName())
|
||||
s.Equal(override.GetPinned().GetBehavior(), versioningInfo.GetVersioningOverride().GetPinned().GetBehavior())
|
||||
a.Equal(override.GetAutoUpgrade(), versioningInfo.GetVersioningOverride().GetAutoUpgrade())
|
||||
a.Equal(override.GetPinned().GetVersion().GetBuildId(), versioningInfo.GetVersioningOverride().GetPinned().GetVersion().GetBuildId())
|
||||
a.Equal(override.GetPinned().GetVersion().GetDeploymentName(), versioningInfo.GetVersioningOverride().GetPinned().GetVersion().GetDeploymentName())
|
||||
a.Equal(override.GetPinned().GetBehavior(), versioningInfo.GetVersioningOverride().GetPinned().GetBehavior())
|
||||
if worker_versioning.OverrideIsPinned(override) {
|
||||
s.Equal(override.GetPinned().GetVersion().GetDeploymentName(), dwf.WorkflowExecutionInfo.GetWorkerDeploymentName())
|
||||
a.Equal(override.GetPinned().GetVersion().GetDeploymentName(), dwf.WorkflowExecutionInfo.GetWorkerDeploymentName())
|
||||
}
|
||||
} else {
|
||||
// v0.31 override
|
||||
s.Equal(override.GetBehavior().String(), versioningInfo.GetVersioningOverride().GetBehavior().String()) //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
a.Equal(override.GetBehavior().String(), versioningInfo.GetVersioningOverride().GetBehavior().String()) //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
if actualOverrideDeployment := versioningInfo.GetVersioningOverride().GetPinnedVersion(); override.GetPinnedVersion() != actualOverrideDeployment { //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
s.Fail(fmt.Sprintf("pinned override mismatch. expected: {%s}, actual: {%s}",
|
||||
a.Fail(fmt.Sprintf("pinned override mismatch. expected: {%s}, actual: {%s}",
|
||||
override.GetPinnedVersion(), //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
actualOverrideDeployment,
|
||||
))
|
||||
}
|
||||
if worker_versioning.OverrideIsPinned(override) {
|
||||
d, _ := worker_versioning.WorkerDeploymentVersionFromStringV31(override.GetPinnedVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
|
||||
s.Equal(d.GetDeploymentName(), dwf.WorkflowExecutionInfo.GetWorkerDeploymentName())
|
||||
a.Equal(d.GetDeploymentName(), dwf.WorkflowExecutionInfo.GetWorkerDeploymentName())
|
||||
}
|
||||
}
|
||||
|
||||
if !versioningInfo.GetVersionTransition().Equal(transition) {
|
||||
s.Fail(fmt.Sprintf("version transition mismatch. expected: {%s}, actual: {%s}",
|
||||
a.Fail(fmt.Sprintf("version transition mismatch. expected: {%s}, actual: {%s}",
|
||||
transition,
|
||||
versioningInfo.GetVersionTransition(),
|
||||
))
|
||||
@@ -3684,7 +3685,8 @@ func (s *Versioning3Suite) TestAutoUpgradeWorkflows_NoBouncingBetweenVersions()
|
||||
|
||||
// Verify that the workflow is running on v1
|
||||
s.EventuallyWithT(func(t *assert.CollectT) {
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
a := require.New(t)
|
||||
s.verifyWorkflowVersioning(a, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
}, 10*time.Second, 100*time.Millisecond)
|
||||
|
||||
// Start v0 workers to ensure they never receive a task
|
||||
@@ -3754,7 +3756,8 @@ func (s *Versioning3Suite) TestWorkflowTQLags_DependentActivityStartsTransition(
|
||||
|
||||
// Verify that the workflow is running on v1.
|
||||
s.EventuallyWithT(func(t *assert.CollectT) {
|
||||
s.verifyWorkflowVersioning(tv0, vbUnpinned, tv0.Deployment(), nil, nil)
|
||||
a := require.New(t)
|
||||
s.verifyWorkflowVersioning(a, tv0, vbUnpinned, tv0.Deployment(), nil, nil)
|
||||
}, 10*time.Second, 100*time.Millisecond)
|
||||
|
||||
// Update the userData for the activity TQ by setting the current version to v1.
|
||||
@@ -3792,7 +3795,8 @@ func (s *Versioning3Suite) TestWorkflowTQLags_DependentActivityStartsTransition(
|
||||
|
||||
// Verify that the workflow is running on v1.
|
||||
s.EventuallyWithT(func(t *assert.CollectT) {
|
||||
s.verifyWorkflowVersioning(tv0, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
a := require.New(t)
|
||||
s.verifyWorkflowVersioning(a, tv0, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
}, 10*time.Second, 100*time.Millisecond)
|
||||
}
|
||||
|
||||
@@ -3850,7 +3854,8 @@ func (s *Versioning3Suite) TestActivityTQLags_DependentActivityCompletesOnTheNew
|
||||
|
||||
// Verify that the workflow is running on v0.
|
||||
s.EventuallyWithT(func(t *assert.CollectT) {
|
||||
s.verifyWorkflowVersioning(tv0, vbUnpinned, tv0.Deployment(), nil, nil)
|
||||
a := require.New(t)
|
||||
s.verifyWorkflowVersioning(a, tv0, vbUnpinned, tv0.Deployment(), nil, nil)
|
||||
}, 10*time.Second, 100*time.Millisecond)
|
||||
|
||||
// Update the userData for the workflow TQ *only* by setting the current version to v1
|
||||
@@ -3892,7 +3897,8 @@ func (s *Versioning3Suite) TestActivityTQLags_DependentActivityCompletesOnTheNew
|
||||
|
||||
// Verify that the workflow is still running on v1.
|
||||
s.EventuallyWithT(func(t *assert.CollectT) {
|
||||
s.verifyWorkflowVersioning(tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
a := require.New(t)
|
||||
s.verifyWorkflowVersioning(a, tv1, vbUnpinned, tv1.Deployment(), nil, nil)
|
||||
}, 10*time.Second, 100*time.Millisecond)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user