From d1a5d3368b29a722956c2416af66a77142ada470 Mon Sep 17 00:00:00 2001 From: Shivam <57200924+Shivs11@users.noreply.github.com> Date: Thu, 18 Dec 2025 15:27:06 -0500 Subject: [PATCH] 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 --- > [!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. > > Written by [Cursor Bugbot](https://cursor.com/dashboard?tab=bugbot) for commit 0a32731d15ec65e943573762c68606b77d2a84f2. This will update automatically on new commits. Configure [here](https://cursor.com/dashboard?tab=bugbot). --- tests/versioning_3_test.go | 156 +++++++++++++++++++------------------ 1 file changed, 81 insertions(+), 75 deletions(-) diff --git a/tests/versioning_3_test.go b/tests/versioning_3_test.go index 7ceb3cbf5f..767430a6f7 100644 --- a/tests/versioning_3_test.go +++ b/tests/versioning_3_test.go @@ -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) }