Skip reactivation signals for current/ramping/draining versions (#9778)

## Summary

- Extends `CheckTaskQueueVersionMembership` response with two new
fields: `is_version_active_or_draining` (bool) and `revision_number`
(int64). Matching populates both from its deployment data.
- Reactivation signals are **skipped** when matching reports the target
version as CURRENT/RAMPING/DRAINING; otherwise they are sent with a
**deterministic UUID v5 RequestId** derived from `revision_number`.
- Replaces the old TTL-based `ReactivationSignalCache` with a per-pod
**revision-based dedup LRU** on the worker-deployment client: each entry
records the highest revision this pod has successfully signaled for a
given version, so older or equal signals are skipped.

## What changed on the wire (matching → history)

`CheckTaskQueueVersionMembershipResponse` now has two new flat fields
(no wrapper message):

```proto
bool is_version_active_or_draining = 2;  // true when status is CURRENT/RAMPING/DRAINING
int64 revision_number = 3;               // from WorkerDeploymentVersionData.revision_number; 0 if unknown / legacy
```

Matching's `CheckTaskQueueVersionMembership` fills both via the helper
`worker_versioning.IsVersionActiveOrDraining(deploymentData, dep, build)
(bool, int64)`.

Naming choice — we picked `is_version_active_or_draining` (negative
polarity) rather than something like `supports_reactivation` so the
proto zero value (`false`) maps to the safe default ("send the signal").
Old matching binaries and runtime "version not found" both produce the
zero value, and history correctly falls through.

## Where `revision_number` flows

- **Matching**: populates the response field from the version's tracked
revision.
- **History-side helper/caches**:
`ValidateVersioningOverrideAndGetReactivationEligibility` returns
`(isVersionActiveOrDraining bool, revisionNumber int64, err)`.
`VersionMembershipAndReactivationStatusCache` stores both.
- **History signaler plumbing**: `VersionReactivationSignalerFn`,
`ReactivateVersionWorkflowIfPinned`, and all five call sites
(`startworkflow`, `signalwithstartworkflow`, `updateworkflowoptions`,
`resetworkflow`, `multioperation`) carry `revisionNumber int64`.
`resetworkflow.validatePostResetOperationInputs` returns parallel slices
`([]bool, []int64, error)` for per-operation inputs.
- **Signal RequestId**: `ClientImpl.SignalVersionReactivation` composes
`requestID = uuid.NewSHA1(uuid.NameSpaceOID,
[]byte("reactivation-signal:" + revisionNumber)).String()` — a
deterministic UUID v5 derived from the revision alone. Cassandra's
`signal_requested set<uuid>` column requires UUID-formatted RequestIds.

## Why revision-based dedup

History is sharded on `(namespaceID, workflowID)`. N concurrent
`StartWorkflow` calls pinned to the same drained version fan out across
potentially every history pod in the fleet. Before this PR each pod
independently fired a reactivation signal at the version workflow,
producing up to N `WorkflowExecutionSignaled` events — directly at odds
with the version workflow's design (it intentionally keeps history
minimal and CaNs aggressively, see `version_workflow.go:68-74`).

Per-pod caches alone can't fix this because they don't coordinate. What
we need is a **cluster-wide-deterministic dedup key** so all pods
converge on the same value for the same reactivation cycle. The
version's `revision_number` — incremented in `syncTaskQueuesAsync` on
every status change — is exactly that signal. Every pod reads the same
revision from matching, every pod composes the same UUID RequestId, and
Temporal's built-in `mutableState.pendingSignalRequestedIDs` dedup (see
`service/history/api/signalworkflow/api.go:40`) collapses concurrent
signals into exactly one event on the version workflow.

The per-pod map is a local optimization on top of that: it prevents a
single pod from re-sending the same-or-older-revision signal once it has
successfully sent one, cutting RPC volume.

## How the new caches look

### 1. `VersionMembershipAndReactivationStatusCache` (read-side,
per-pod)
Caches matching's `CheckTaskQueueVersionMembership` response so repeated
pinned-override validations on the same task queue don't re-hit
matching.

- **Key**: `(namespaceID, taskQueue, taskQueueType, deploymentName,
buildID)`
- **Value**: `(isMember bool, isVersionActiveOrDraining bool,
revisionNumber int64)`
- **Eviction**: `VersionMembershipCacheTTL` (1s default; 5s in
functional tests).

### 2. `highestRevSignaledToVersionWf` (write-side dedup, per-pod)
A field on `ClientImpl` in `service/worker/workerdeployment/client.go`.
For each target version workflow, stores the highest revision this pod
has successfully signaled. Subsequent calls at the same-or-lower
revision skip the RPC.

- **Key**: `reactivationVersionKey{namespaceID, deploymentName,
buildID}`
- **Value**: `int64` (highest revision successfully signaled)
- **Eviction**: LRU, bounded by `VersionReactivationSignalCacheMaxSize`.

The previous TTL-based `ReactivationSignalCache` module (in
`common/worker_versioning/`) has been deleted along with its provider
and `VersionReactivationSignalCacheTTL` config.

## Backwards/forwards compatibility

- **Old matching → new history**: old binaries don't set
`is_version_active_or_draining` or `revision_number`; both default to
proto zero values. `false` on the active bool → history falls through →
signal fires (safe default). `revisionNumber = 0` flows through as-is.
- **New matching → old history**: new fields on the response are ignored
by old history → identical to pre-PR behavior.
- **New matching → new history**: signal fires only when the version is
not active/draining; cross-pod fires converge on one UUID RequestId and
fold into one `WorkflowExecutionSignaled` event.

## Test plan

- [x] Unit tests for `IsVersionActiveOrDraining` covering all status
cases (CURRENT, RAMPING, DRAINING, DRAINED, INACTIVE, UNSPECIFIED), new
vs. old format, deleted and not-found versions.
- [x] Unit tests for
`ValidateVersioningOverrideAndGetReactivationEligibility` (cache
hit/miss, RPC with/without eligibility, Unimplemented fallback).
- [x] Unit tests for the per-pod dedup on
`ClientImpl.SignalVersionReactivation`: same-rev dedups, newer-rev
fires, older-rev skipped, different version isolated, signal-failure
allows retry.
- [x] Unit test for RequestId format (UUID v5, deterministic across
calls with the same revision).
- [x] Functional tests (all pass on SQLite and cass-es):
  - `TestStartWorkflowExecution_ReactivateVersionOnPinned`
-
`TestStartWorkflowExecution_ReactivateVersionOnPinned_WithConflictPolicy`
  - `TestSignalWithStartWorkflowExecution_ReactivateVersionOnPinned`
  - `TestUpdateWorkflowExecutionOptions_ReactivateVersionOnPinned`
  - `TestResetWorkflowExecution_ReactivateVersionOnPinned`

(The four `TestReactivationSignalCache_Deduplication_*` functional tests
from an earlier iteration were deleted — their coverage moved to unit
tests.)

<!-- CURSOR_SUMMARY -->
---

> [!NOTE]
> **Medium Risk**
> Changes matching↔history API and reactivation signaling semantics by
skipping signals for active/draining versions and introducing
revision-based dedup via deterministic RequestIds; issues could affect
version workflow state transitions or signal fan-out during upgrades.
> 
> **Overview**
> Matching’s `CheckTaskQueueVersionMembershipResponse` is extended with
`should_skip_reactivation` and `revision_number`, and matching now
populates both from per-task-queue deployment data.
> 
> History-side versioning validation is refactored to return and cache
reactivation eligibility + revision, and reactivation signaling paths
(`StartWorkflow`, `SignalWithStart`, `UpdateWorkflowExecutionOptions`,
`ResetWorkflow`, multi-op) now **skip signals** when matching reports
the version as *CURRENT/RAMPING/DRAINING*.
> 
> The old TTL-based `ReactivationSignalCache` is removed
(configs/metrics/providers updated), and the worker-deployment client
now performs **revision-based per-pod dedup** plus receiver-side dedup
by sending signals with a deterministic UUIDv5-like `RequestId` derived
from `revision_number`. Tests are updated/added to cover status
evaluation, new plumbing, and dedup behavior.
> 
> <sup>Reviewed by [Cursor Bugbot](https://cursor.com/bugbot) for commit
a1ec5e93db. Bugbot is set up for automated
code reviews on this repo. Configure
[here](https://www.cursor.com/dashboard/bugbot).</sup>
<!-- /CURSOR_SUMMARY -->

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Shivam
2026-04-23 10:39:36 -04:00
committed by GitHub
parent 28fd8c47fd
commit 562d26fc83
30 changed files with 850 additions and 1014 deletions

View File

@@ -5374,10 +5374,25 @@ func (x *CheckTaskQueueVersionMembershipRequest) GetVersion() *v110.WorkerDeploy
}
type CheckTaskQueueVersionMembershipResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
IsMember bool `protobuf:"varint,1,opt,name=is_member,json=isMember,proto3" json:"is_member,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
state protoimpl.MessageState `protogen:"open.v1"`
IsMember bool `protobuf:"varint,1,opt,name=is_member,json=isMember,proto3" json:"is_member,omitempty"`
// True when a reactivation signal to this version would be redundant — i.e., matching
// determined the version is already in a state where it does not need to be reactivated
// (today: CURRENT, RAMPING, or DRAINING). History uses this to suppress such signals.
// The zero value (false) is the safe default; it applies when matching has no definitive
// answer (version not present in matching's deployment data, or old matching servers
// that do not set this field) and tells history to send the signal.
ShouldSkipReactivation bool `protobuf:"varint,2,opt,name=should_skip_reactivation,json=shouldSkipReactivation,proto3" json:"should_skip_reactivation,omitempty"`
// revision_number is the version's current revision as tracked in matching's per-TQ
// deployment data. It is returned so history can compose a stable, cluster-wide-deterministic
// RequestId on the reactivation signal. All history pods querying the same version at the
// same point in time converge on the same revision_number, so Temporal's built-in
// SignalRequestedIds dedup (see signalworkflow/api.go) collapses the N-pod signal fan-out
// into exactly one event on the version workflow. Zero when unknown (old matching server or
// legacy DeploymentVersionData format that does not carry revision_number).
RevisionNumber int64 `protobuf:"varint,3,opt,name=revision_number,json=revisionNumber,proto3" json:"revision_number,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *CheckTaskQueueVersionMembershipResponse) Reset() {
@@ -5417,6 +5432,20 @@ func (x *CheckTaskQueueVersionMembershipResponse) GetIsMember() bool {
return false
}
func (x *CheckTaskQueueVersionMembershipResponse) GetShouldSkipReactivation() bool {
if x != nil {
return x.ShouldSkipReactivation
}
return false
}
func (x *CheckTaskQueueVersionMembershipResponse) GetRevisionNumber() int64 {
if x != nil {
return x.RevisionNumber
}
return 0
}
// PollConditions are extra conditions to set on the poll. Only supported with new matcher.
type PollConditions struct {
state protoimpl.MessageState `protogen:"open.v1"`
@@ -6160,9 +6189,11 @@ const file_temporal_server_api_matchingservice_v1_request_response_proto_rawDesc
"\n" +
"task_queue\x18\x02 \x01(\tR\ttaskQueue\x12L\n" +
"\x0ftask_queue_type\x18\x03 \x01(\x0e2$.temporal.api.enums.v1.TaskQueueTypeR\rtaskQueueType\x12T\n" +
"\aversion\x18\x04 \x01(\v2:.temporal.server.api.deployment.v1.WorkerDeploymentVersionR\aversion\"F\n" +
"\aversion\x18\x04 \x01(\v2:.temporal.server.api.deployment.v1.WorkerDeploymentVersionR\aversion\"\xa9\x01\n" +
"'CheckTaskQueueVersionMembershipResponse\x12\x1b\n" +
"\tis_member\x18\x01 \x01(\bR\bisMember\"L\n" +
"\tis_member\x18\x01 \x01(\bR\bisMember\x128\n" +
"\x18should_skip_reactivation\x18\x02 \x01(\bR\x16shouldSkipReactivation\x12'\n" +
"\x0frevision_number\x18\x03 \x01(\x03R\x0erevisionNumber\"L\n" +
"\x0ePollConditions\x12!\n" +
"\fmin_priority\x18\x01 \x01(\x05R\vminPriority\x12\x17\n" +
"\ano_wait\x18\x02 \x01(\bR\x06noWaitB>Z<go.temporal.io/server/api/matchingservice/v1;matchingserviceb\x06proto3"

View File

@@ -2921,17 +2921,12 @@ instead of the previous HSM backed implementation.`,
`Maximum number of entries in the version membership cache.`,
)
VersionReactivationSignalCacheTTL = NewGlobalDurationSetting(
"history.versionReactivationSignalCacheTTL",
10*time.Second,
`TTL for caching drainage reactivation signals to version workflows. These signals are sent from the history service to update the version workflow's
draining status to DRAINING from DRAINED/INACTIVE states.`,
)
VersionReactivationSignalCacheMaxSize = NewGlobalIntSetting(
"history.versionReactivationSignalCacheMaxSize",
ReactivationSignalDedupCacheMaxSize = NewGlobalIntSetting(
"worker.reactivationSignalDedupCacheMaxSize",
10000,
`Maximum number of entries in the version reactivation signal cache.`,
`Maximum number of entries in the per-pod reactivation-signal dedup cache on the
worker deployment client. Each entry tracks the highest revision signaled for one
target version workflow.`,
)
EnableVersionReactivationSignals = NewGlobalBoolSetting(

View File

@@ -47,7 +47,7 @@ const (
MutableStateCacheTypeTagValue = "mutablestate"
EventsCacheTypeTagValue = "events"
VersionMembershipCacheTypeTagValue = "version_membership"
VersionReactivationSignalCacheTypeTagValue = "version_reactivation_signal"
ReactivationSignalDedupCacheTypeTagValue = "reactivation_signal_dedup"
RoutingInfoCacheTypeTagValue = "routing_info"
NexusEndpointRegistryReadThroughCacheTypeTagValue = "nexus_endpoint_registry_readthrough"
@@ -459,8 +459,9 @@ const (
VersionMembershipCacheGetScope = "VersionMembershipCacheGet"
// VersionMembershipCachePutScope is the scope used by version membership cache
VersionMembershipCachePutScope = "VersionMembershipCachePut"
// VersionReactivationSignalCacheShouldSendScope is the scope used by version reactivation signal cache
VersionReactivationSignalCacheShouldSendScope = "VersionReactivationSignalCacheShouldSend"
// ReactivationSignalDedupScope is the scope used by the per-pod reactivation-signal
// dedup cache on the worker-deployment client.
ReactivationSignalDedupScope = "ReactivationSignalDedup"
// RoutingInfoCacheGetScope is the scope used by routing info cache
RoutingInfoCacheGetScope = "RoutingInfoCacheGet"
// RoutingInfoCachePutScope is the scope used by routing info cache

View File

@@ -7,20 +7,27 @@ import (
"go.temporal.io/server/common/metrics"
)
// VersionMembershipCache is used to cache results of Matching's CheckTaskQueueVersionMembership
// calls (used internally by the worker versioning pinned override validation).
// VersionMembershipAndReactivationStatusCache caches results of Matching's
// CheckTaskQueueVersionMembership calls. It stores three pieces of information per version:
// - isMember: whether the task queue exists in the version (used for pinned override validation).
// - shouldSkipReactivation: true when the version's current status is CURRENT, RAMPING,
// or DRAINING, in which case the caller skips the reactivation signal. The zero value
// (false) is the safe default and covers unknown / not-found / old-matching cases.
// - revisionNumber: the version's current revision per matching's view. Used as part of the
// reactivation signal's RequestId so that all history pods targeting the same version at
// the same revision compose the same dedup key. Zero means unknown (old matching server or
// legacy DeploymentVersionData format with no revision_number field).
//
// Implementations are expected to be safe for concurrent use.
type (
VersionMembershipCache interface {
// Get returns (isMember, ok). ok=false means there was no cached value.
VersionMembershipAndReactivationStatusCache interface {
Get(
namespaceID string,
taskQueue string,
taskQueueType enumspb.TaskQueueType,
deploymentName string,
buildID string,
) (isMember bool, ok bool)
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, ok bool)
Put(
namespaceID string,
@@ -29,6 +36,8 @@ type (
deploymentName string,
buildID string,
isMember bool,
shouldSkipReactivation bool,
revisionNumber int64,
)
}
@@ -40,28 +49,34 @@ type (
buildID string
}
VersionMembershipCacheImpl struct {
versionTaskQueueInfoCacheValue struct {
isMember bool
shouldSkipReactivation bool // false = unknown / not-found / eligible-for-reactivation
revisionNumber int64 // 0 = unknown (old matching server or legacy format)
}
VersionMembershipAndReactivationStatusCacheImpl struct {
cache.Cache
metricsHandler metrics.Handler
}
)
// NewVersionMembershipCache wraps the provided cache with a typed API and metrics.
func NewVersionMembershipCache(c cache.Cache, metricsHandler metrics.Handler) VersionMembershipCache {
// NewVersionMembershipAndReactivationStatusCache wraps the provided cache with a typed API and metrics.
func NewVersionMembershipAndReactivationStatusCache(c cache.Cache, metricsHandler metrics.Handler) VersionMembershipAndReactivationStatusCache {
h := metricsHandler.WithTags(metrics.CacheTypeTag(metrics.VersionMembershipCacheTypeTagValue))
return &VersionMembershipCacheImpl{
return &VersionMembershipAndReactivationStatusCacheImpl{
Cache: c,
metricsHandler: h,
}
}
func (c *VersionMembershipCacheImpl) Get(
func (c *VersionMembershipAndReactivationStatusCacheImpl) Get(
namespaceID string,
taskQueue string,
taskQueueType enumspb.TaskQueueType,
deploymentName string,
buildID string,
) (isMember bool, ok bool) {
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, ok bool) {
handler := c.metricsHandler.WithTags(metrics.OperationTag(metrics.VersionMembershipCacheGetScope), metrics.NamespaceIDTag(namespaceID))
metrics.CacheRequests.With(handler).Record(1)
@@ -75,24 +90,26 @@ func (c *VersionMembershipCacheImpl) Get(
v := c.Cache.Get(key)
if v == nil {
metrics.CacheMissCounter.With(handler).Record(1)
return false, false
return false, false, 0, false
}
isMember, ok = v.(bool)
value, ok := v.(versionTaskQueueInfoCacheValue)
if !ok {
// Unexpected type: treat as miss to avoid false positives.
metrics.CacheMissCounter.With(handler).Record(1)
return false, false
return false, false, 0, false
}
return isMember, true
return value.isMember, value.shouldSkipReactivation, value.revisionNumber, true
}
func (c *VersionMembershipCacheImpl) Put(
func (c *VersionMembershipAndReactivationStatusCacheImpl) Put(
namespaceID string,
taskQueue string,
taskQueueType enumspb.TaskQueueType,
deploymentName string,
buildID string,
isMember bool,
shouldSkipReactivation bool,
revisionNumber int64,
) {
handler := c.metricsHandler.WithTags(metrics.OperationTag(metrics.VersionMembershipCachePutScope), metrics.NamespaceIDTag(namespaceID))
metrics.CacheRequests.With(handler).Record(1)
@@ -104,5 +121,9 @@ func (c *VersionMembershipCacheImpl) Put(
deploymentName: deploymentName,
buildID: buildID,
}
c.Cache.Put(key, isMember)
c.Cache.Put(key, versionTaskQueueInfoCacheValue{
isMember: isMember,
shouldSkipReactivation: shouldSkipReactivation,
revisionNumber: revisionNumber,
})
}

View File

@@ -1,65 +0,0 @@
//nolint:staticcheck
package worker_versioning
import (
"go.temporal.io/server/common/cache"
"go.temporal.io/server/common/metrics"
)
// ReactivationSignalCache deduplicates reactivation signals to version workflows.
//
// Implementations are expected to be safe for concurrent use.
type (
ReactivationSignalCache interface {
// ShouldSendSignal returns true if signal should be sent (not recently sent)
// and atomically marks it as sent. Returns false if recently sent.
ShouldSendSignal(namespaceID, deploymentName, buildID string) bool
}
reactivationSignalCacheKey struct {
namespaceID string
deploymentName string
buildID string
}
ReactivationSignalCacheImpl struct {
cache.Cache
metricsHandler metrics.Handler
}
)
// NewReactivationSignalCache wraps the provided cache with a typed API and metrics.
func NewReactivationSignalCache(c cache.Cache, metricsHandler metrics.Handler) ReactivationSignalCache {
h := metricsHandler.WithTags(metrics.CacheTypeTag(metrics.VersionReactivationSignalCacheTypeTagValue))
return &ReactivationSignalCacheImpl{
Cache: c,
metricsHandler: h,
}
}
func (c *ReactivationSignalCacheImpl) ShouldSendSignal(
namespaceID, deploymentName, buildID string,
) bool {
handler := c.metricsHandler.WithTags(
metrics.OperationTag(metrics.VersionReactivationSignalCacheShouldSendScope),
metrics.NamespaceIDTag(namespaceID),
)
metrics.CacheRequests.With(handler).Record(1)
key := reactivationSignalCacheKey{
namespaceID: namespaceID,
deploymentName: deploymentName,
buildID: buildID,
}
// Check if we recently sent a signal for this version
if c.Cache.Get(key) != nil {
// Entry exists, signal was recently sent - deduplicate
return false
}
// No recent signal, mark as sent and return true
metrics.CacheMissCounter.With(handler).Record(1)
c.Cache.Put(key, true)
return true
}

View File

@@ -288,13 +288,13 @@ func MakeDirectiveForWorkflowTask(
type IsWFTaskQueueInVersionDetector = func(ctx context.Context, namespaceID, tq string, version *deploymentpb.WorkerDeploymentVersion) (bool, error)
func GetIsWFTaskQueueInVersionDetector(matchingClient resource.MatchingClient, versionMembershipCache VersionMembershipCache) IsWFTaskQueueInVersionDetector {
func GetIsWFTaskQueueInVersionDetector(matchingClient resource.MatchingClient, versionCache VersionMembershipAndReactivationStatusCache) IsWFTaskQueueInVersionDetector {
return func(ctx context.Context,
namespaceID, tq string,
version *deploymentpb.WorkerDeploymentVersion) (bool, error) {
// Check cache first.
if isMember, ok := versionMembershipCache.Get(
if isMember, _, _, ok := versionCache.Get(
namespaceID, tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW,
version.GetDeploymentName(), version.GetBuildId(),
); ok {
@@ -302,31 +302,33 @@ func GetIsWFTaskQueueInVersionDetector(matchingClient resource.MatchingClient, v
}
// Cache miss — resolve via matching RPC.
isMember, err := checkTaskQueueVersionMembership(ctx, matchingClient, namespaceID, tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW, version)
isMember, shouldSkipReactivation, revisionNumber, err := checkVersionMembershipAndReactivationEligibility(ctx, matchingClient, namespaceID, tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW, version)
if err != nil {
return false, err
}
// Add result to cache
versionMembershipCache.Put(
versionCache.Put(
namespaceID, tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW,
version.GetDeploymentName(), version.GetBuildId(),
isMember,
isMember, shouldSkipReactivation, revisionNumber,
)
return isMember, nil
}
}
// checkTaskQueueVersionMembership calls matching to check if a task queue belongs to a version,
// falling back to fetching the full user data if the CheckTaskQueueVersionMembership RPC is not implemented. (this can happen
// during rolling upgrades of matching and history services where in history would be on a higher version than matching)
func checkTaskQueueVersionMembership(
// checkVersionMembershipAndReactivationEligibility calls matching to check if a task queue belongs to a version
// and whether the version is currently active-or-draining (see ShouldSkipReactivation for the
// exact status set). Falls back to fetching the full user data if the CheckTaskQueueVersionMembership
// RPC is not implemented (this can happen during rolling upgrades where history is on a higher
// version than matching).
func checkVersionMembershipAndReactivationEligibility(
ctx context.Context,
matchingClient resource.MatchingClient,
namespaceID, tq string,
tqType enumspb.TaskQueueType,
version *deploymentpb.WorkerDeploymentVersion,
) (bool, error) {
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, err error) {
resp, err := matchingClient.CheckTaskQueueVersionMembership(ctx, &matchingservice.CheckTaskQueueVersionMembershipRequest{
NamespaceId: namespaceID,
TaskQueue: tq,
@@ -340,9 +342,9 @@ func checkTaskQueueVersionMembership(
if errors.As(err, &unimplErr) {
return checkVersionMembershipViaUserData(ctx, matchingClient, namespaceID, tq, tqType, version)
}
return false, err
return false, false, 0, err
}
return resp.GetIsMember(), nil
return resp.GetIsMember(), resp.GetShouldSkipReactivation(), resp.GetRevisionNumber(), nil
}
// checkVersionMembershipViaUserData is the fallback for when matching doesn't support
@@ -355,7 +357,7 @@ func checkVersionMembershipViaUserData(
tq string,
tqType enumspb.TaskQueueType,
version *deploymentpb.WorkerDeploymentVersion,
) (bool, error) {
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, err error) {
resp, err := matchingClient.GetTaskQueueUserData(ctx,
&matchingservice.GetTaskQueueUserDataRequest{
NamespaceId: namespaceID,
@@ -363,13 +365,16 @@ func checkVersionMembershipViaUserData(
TaskQueueType: tqType,
})
if err != nil {
return false, err
return false, false, 0, err
}
tqData, ok := resp.GetUserData().GetData().GetPerType()[int32(tqType)]
if !ok {
return false, nil
return false, false, 0, nil
}
return HasDeploymentVersion(tqData.GetDeploymentData(), DeploymentVersionFromDeployment(DeploymentFromExternalDeploymentVersion(version))), nil
deploymentData := tqData.GetDeploymentData()
isMember = HasDeploymentVersion(deploymentData, DeploymentVersionFromDeployment(DeploymentFromExternalDeploymentVersion(version)))
shouldSkipReactivation, revisionNumber = ShouldSkipReactivation(deploymentData, version.GetDeploymentName(), version.GetBuildId())
return isMember, shouldSkipReactivation, revisionNumber, nil
}
func FindOldDeploymentVersion(deployments *persistencespb.DeploymentData, v *deploymentspb.WorkerDeploymentVersion) int {
@@ -403,6 +408,46 @@ func HasDeploymentVersion(deployments *persistencespb.DeploymentData, v *deploym
return false
}
// ShouldSkipReactivation reports whether a reactivation signal to the given version would
// be redundant. Returns true when the version's status is CURRENT, RAMPING, or DRAINING.
// Returns false for DRAINED and INACTIVE (the two statuses the reactivation handler in
// version_workflow.go acts on) and when the version is not present in the deployment data.
// (UNSPECIFIED also yields false; in practice the deployment workflow sets a status at
// construction so this branch should not trigger.)
//
// The returned revisionNumber is the version's revision as tracked in the new deployment
// data format (WorkerDeploymentVersionData.revision_number). It is 0 for the legacy
// DeploymentVersionData format (which does not carry a revision number) and when the
// version is not found at all.
//
//nolint:staticcheck
func ShouldSkipReactivation(
deployments *persistencespb.DeploymentData,
deploymentName string,
buildID string,
) (bool, int64) {
// Check old format first (deprecated versions list).
for _, vd := range deployments.GetVersions() {
if vd.GetVersion().GetDeploymentName() == deploymentName && vd.GetVersion().GetBuildId() == buildID {
return isStatusSkippableFromReactivation(vd.GetStatus()), 0
}
}
// Check new format (deployments_data map).
deploymentData := deployments.GetDeploymentsData()[deploymentName]
versionData := deploymentData.GetVersions()[buildID]
if versionData == nil || versionData.GetDeleted() {
return false, 0
}
return isStatusSkippableFromReactivation(versionData.GetStatus()), versionData.GetRevisionNumber()
}
func isStatusSkippableFromReactivation(s enumspb.WorkerDeploymentVersionStatus) bool {
return s == enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT ||
s == enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_RAMPING ||
s == enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING
}
func CountDeploymentVersions(deployments *persistencespb.DeploymentData) int {
//nolint:staticcheck // SA1019
res := len(deployments.GetVersions())
@@ -623,16 +668,16 @@ func ExtractVersioningBehaviorFromOverride(override *workflowpb.VersioningOverri
return override.GetBehavior()
}
func validatePinnedVersionInTaskQueue(ctx context.Context,
func validateVersionAndGetReactivationEligibility(ctx context.Context,
pinnedVersion *deploymentpb.WorkerDeploymentVersion,
matchingClient resource.MatchingClient,
versionMembershipCache VersionMembershipCache,
versionCache VersionMembershipAndReactivationStatusCache,
tq string,
tqType enumspb.TaskQueueType,
namespaceID string) error {
namespaceID string) (shouldSkipReactivation bool, revisionNumber int64, err error) {
// Check if we have recently queried matching to validate if this version exists in the task queue.
if isMember, ok := versionMembershipCache.Get(
if isMember, cachedActiveOrDraining, cachedRevision, ok := versionCache.Get(
namespaceID,
tq,
tqType,
@@ -640,88 +685,90 @@ func validatePinnedVersionInTaskQueue(ctx context.Context,
pinnedVersion.BuildId,
); ok {
if isMember {
return nil
return cachedActiveOrDraining, cachedRevision, nil
}
return serviceerror.NewFailedPrecondition(
return false, 0, serviceerror.NewFailedPrecondition(
FormatPinnedVersionNotInTaskQueueError(pinnedVersion.GetDeploymentName(), pinnedVersion.GetBuildId(), tq, tqType),
)
}
isMember, err := checkTaskQueueVersionMembership(ctx, matchingClient, namespaceID, tq, tqType, pinnedVersion)
isMember, shouldSkipReactivation, revisionNumber, err := checkVersionMembershipAndReactivationEligibility(ctx, matchingClient, namespaceID, tq, tqType, pinnedVersion)
if err != nil {
return err
return false, 0, err
}
// Add result to cache
versionMembershipCache.Put(
versionCache.Put(
namespaceID,
tq,
tqType,
pinnedVersion.DeploymentName,
pinnedVersion.BuildId,
isMember,
shouldSkipReactivation,
revisionNumber,
)
if !isMember {
return serviceerror.NewFailedPrecondition(
return false, 0, serviceerror.NewFailedPrecondition(
FormatPinnedVersionNotInTaskQueueError(pinnedVersion.GetDeploymentName(), pinnedVersion.GetBuildId(), tq, tqType),
)
}
return nil
return shouldSkipReactivation, revisionNumber, nil
}
func ValidateVersioningOverride(ctx context.Context,
func ValidateVersioningOverrideAndGetReactivationEligibility(ctx context.Context,
override *workflowpb.VersioningOverride,
matchingClient resource.MatchingClient,
versionMembershipCache VersionMembershipCache,
versionCache VersionMembershipAndReactivationStatusCache,
tq string,
tqType enumspb.TaskQueueType,
namespaceID string) error {
namespaceID string) (shouldSkipReactivation bool, revisionNumber int64, err error) {
if override == nil {
return nil
return false, 0, nil
}
if override.GetAutoUpgrade() { // v0.32
return nil
return false, 0, nil
} else if p := override.GetPinned(); p != nil {
if p.GetVersion() == nil {
return serviceerror.NewInvalidArgument("must provide version if override is pinned.")
return false, 0, serviceerror.NewInvalidArgument("must provide version if override is pinned.")
}
if p.GetBehavior() == workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_UNSPECIFIED {
return serviceerror.NewInvalidArgument("must specify pinned override behavior if override is pinned.")
return false, 0, serviceerror.NewInvalidArgument("must specify pinned override behavior if override is pinned.")
}
return validatePinnedVersionInTaskQueue(ctx, p.GetVersion(), matchingClient, versionMembershipCache, tq, tqType, namespaceID)
return validateVersionAndGetReactivationEligibility(ctx, p.GetVersion(), matchingClient, versionCache, tq, tqType, namespaceID)
}
//nolint:staticcheck // SA1019: worker versioning v0.31
switch override.GetBehavior() {
case enumspb.VERSIONING_BEHAVIOR_PINNED:
if override.GetDeployment() != nil {
return ValidateDeployment(override.GetDeployment())
return false, 0, ValidateDeployment(override.GetDeployment())
} else if override.GetPinnedVersion() != "" {
_, err := ValidateDeploymentVersionStringV31(override.GetPinnedVersion())
if err != nil {
return err
return false, 0, err
}
return validatePinnedVersionInTaskQueue(ctx, ExternalWorkerDeploymentVersionFromStringV31(override.GetPinnedVersion()), matchingClient, versionMembershipCache, tq, tqType, namespaceID)
return validateVersionAndGetReactivationEligibility(ctx, ExternalWorkerDeploymentVersionFromStringV31(override.GetPinnedVersion()), matchingClient, versionCache, tq, tqType, namespaceID)
} else {
return serviceerror.NewInvalidArgument("must provide deployment (deprecated) or pinned version if behavior is 'PINNED'")
return false, 0, serviceerror.NewInvalidArgument("must provide deployment (deprecated) or pinned version if behavior is 'PINNED'")
}
case enumspb.VERSIONING_BEHAVIOR_AUTO_UPGRADE:
if override.GetDeployment() != nil {
return serviceerror.NewInvalidArgument("only provide deployment if behavior is 'PINNED'")
return false, 0, serviceerror.NewInvalidArgument("only provide deployment if behavior is 'PINNED'")
}
if override.GetPinnedVersion() != "" {
return serviceerror.NewInvalidArgument("only provide pinned version if behavior is 'PINNED'")
return false, 0, serviceerror.NewInvalidArgument("only provide pinned version if behavior is 'PINNED'")
}
case enumspb.VERSIONING_BEHAVIOR_UNSPECIFIED:
return serviceerror.NewInvalidArgument("override behavior is required")
return false, 0, serviceerror.NewInvalidArgument("override behavior is required")
default:
//nolint:staticcheck // SA1019 deprecated stamp will clean up later
return serviceerror.NewInvalidArgumentf("override behavior %s not recognized", override.GetBehavior())
return false, 0, serviceerror.NewInvalidArgumentf("override behavior %s not recognized", override.GetBehavior())
}
return nil
return false, 0, nil
}
// FindTargetDeploymentVersionAndRevisionNumberForWorkflowID returns the deployment version and revision number (if applicable) for

View File

@@ -42,16 +42,16 @@ var (
type testVersionMembershipCache struct {
mu sync.Mutex
m map[versionMembershipCacheKey]bool
m map[versionMembershipCacheKey]versionTaskQueueInfoCacheValue
getCalls int
putCalls int
}
func newTestVersionMembershipCache() *testVersionMembershipCache {
return &testVersionMembershipCache{m: make(map[versionMembershipCacheKey]bool)}
return &testVersionMembershipCache{m: make(map[versionMembershipCacheKey]versionTaskQueueInfoCacheValue)}
}
var _ VersionMembershipCache = (*testVersionMembershipCache)(nil)
var _ VersionMembershipAndReactivationStatusCache = (*testVersionMembershipCache)(nil)
func (c *testVersionMembershipCache) Get(
namespaceID string,
@@ -59,7 +59,7 @@ func (c *testVersionMembershipCache) Get(
taskQueueType enumspb.TaskQueueType,
deploymentName string,
buildID string,
) (isMember bool, ok bool) {
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, ok bool) {
c.mu.Lock()
defer c.mu.Unlock()
c.getCalls++
@@ -70,7 +70,7 @@ func (c *testVersionMembershipCache) Get(
deploymentName: deploymentName,
buildID: buildID,
}]
return v, ok
return v.isMember, v.shouldSkipReactivation, v.revisionNumber, ok
}
func (c *testVersionMembershipCache) Put(
@@ -80,6 +80,8 @@ func (c *testVersionMembershipCache) Put(
deploymentName string,
buildID string,
isMember bool,
shouldSkipReactivation bool,
revisionNumber int64,
) {
c.mu.Lock()
defer c.mu.Unlock()
@@ -90,7 +92,159 @@ func (c *testVersionMembershipCache) Put(
taskQueueType: taskQueueType,
deploymentName: deploymentName,
buildID: buildID,
}] = isMember
}] = versionTaskQueueInfoCacheValue{isMember: isMember, shouldSkipReactivation: shouldSkipReactivation, revisionNumber: revisionNumber}
}
func TestShouldSkipReactivation(t *testing.T) {
tests := []struct {
name string
deployments *persistencespb.DeploymentData
deploymentName string
buildID string
expected bool
expectedRevisionNumber int64
}{
{
name: "nil deployments returns false (not found → caller should send signal)",
deployments: nil,
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 0,
},
{
name: "version not found returns false",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{},
},
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 0,
},
{
name: "new format: DRAINED returns false (caller should send signal) with revision_number",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, RevisionNumber: 42},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 42,
},
{
name: "new format: INACTIVE returns false with revision_number",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, RevisionNumber: 7},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 7,
},
{
name: "new format: CURRENT returns true with revision_number",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT, RevisionNumber: 3},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: true,
expectedRevisionNumber: 3,
},
{
name: "new format: RAMPING returns true",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_RAMPING},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: true,
expectedRevisionNumber: 0,
},
{
name: "new format: DRAINING returns true",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, RevisionNumber: 11},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: true,
expectedRevisionNumber: 11,
},
{
name: "new format: deleted version returns false (treated as not-found)",
deployments: &persistencespb.DeploymentData{
DeploymentsData: map[string]*persistencespb.WorkerDeploymentData{
"dep": {Versions: map[string]*deploymentspb.WorkerDeploymentVersionData{
"build": {Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, Deleted: true, RevisionNumber: 99},
}},
},
},
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 0,
},
{
name: "old format: DRAINED returns false with revision_number 0 (no field in old format)",
deployments: &persistencespb.DeploymentData{
Versions: []*deploymentspb.DeploymentVersionData{
{
Version: &deploymentspb.WorkerDeploymentVersion{DeploymentName: "dep", BuildId: "build"},
Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
},
},
},
deploymentName: "dep",
buildID: "build",
expected: false,
expectedRevisionNumber: 0,
},
{
name: "old format: CURRENT returns true with revision_number 0",
deployments: &persistencespb.DeploymentData{
Versions: []*deploymentspb.DeploymentVersionData{
{
Version: &deploymentspb.WorkerDeploymentVersion{DeploymentName: "dep", BuildId: "build"},
Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT,
},
},
},
deploymentName: "dep",
buildID: "build",
expected: true,
expectedRevisionNumber: 0,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
result, revisionNumber := ShouldSkipReactivation(tt.deployments, tt.deploymentName, tt.buildID)
require.Equal(t, tt.expected, result)
require.Equal(t, tt.expectedRevisionNumber, revisionNumber)
})
}
}
func TestCalculateTaskQueueVersioningInfo(t *testing.T) {
@@ -815,7 +969,7 @@ func TestWorkerDeploymentVersionFromStringV32(t *testing.T) {
}
}
func TestValidateVersioningOverride(t *testing.T) {
func TestValidateVersioningOverrideAndGetReactivationEligibility(t *testing.T) {
testNamespaceID := "test-namespace-id"
testTaskQueue := "test-task-queue"
testVersion := &deploymentpb.WorkerDeploymentVersion{
@@ -829,13 +983,15 @@ func TestValidateVersioningOverride(t *testing.T) {
}
tests := []struct {
name string
override *workflowpb.VersioningOverride
taskQueueType enumspb.TaskQueueType
setupCache func(c *testVersionMembershipCache)
setupMock func(m *matchingservicemock.MockMatchingServiceClient)
expectError bool
errorContains string
name string
override *workflowpb.VersioningOverride
taskQueueType enumspb.TaskQueueType
setupCache func(c *testVersionMembershipCache)
setupMock func(m *matchingservicemock.MockMatchingServiceClient)
expectError bool
errorContains string
expectedShouldSkipReactivation bool
expectedRevisionNumber int64
}{
{
name: "nil override returns nil",
@@ -856,7 +1012,7 @@ func TestValidateVersioningOverride(t *testing.T) {
expectError: false,
},
{
name: "v0.32: Pinned override, with cache hit, returns nil",
name: "v0.32: Pinned override, with cache hit (drained), returns isDrainedOrInactive=true and cached revision",
override: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
@@ -866,12 +1022,36 @@ func TestValidateVersioningOverride(t *testing.T) {
},
},
setupCache: func(c *testVersionMembershipCache) {
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, true)
// Cache says the version is drained/inactive (shouldSkipReactivation=false)
// with revision 42.
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, true, false, 42)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(gomock.Any(), gomock.Any()).Times(0) // No RPC call expected!
},
expectError: false,
expectError: false,
expectedShouldSkipReactivation: false,
expectedRevisionNumber: 42,
},
{
name: "v0.32: Pinned override, with cache hit (active), returns isDrainedOrInactive=false",
override: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: testVersion,
},
},
},
setupCache: func(c *testVersionMembershipCache) {
// Cache says the version is active (shouldSkipReactivation=true).
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, true, true, 0)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(gomock.Any(), gomock.Any()).Times(0) // No RPC call expected!
},
expectError: false,
expectedShouldSkipReactivation: true,
},
{
name: "v0.32: Pinned override, with cache hit, returns error (since version is not present in the task queue)",
@@ -884,7 +1064,7 @@ func TestValidateVersioningOverride(t *testing.T) {
},
},
setupCache: func(c *testVersionMembershipCache) {
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, false)
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, false, false, 0)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(gomock.Any(), gomock.Any()).Times(0) // No RPC call expected!
@@ -904,7 +1084,7 @@ func TestValidateVersioningOverride(t *testing.T) {
},
},
setupCache: func(c *testVersionMembershipCache) {
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, true)
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId, true, false, 0)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(
@@ -939,7 +1119,55 @@ func TestValidateVersioningOverride(t *testing.T) {
errorContains: getPinnedVersionErrorMsg(testVersion, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW),
},
{
name: "v0.32: Pinned override, with cache miss, calls RPC and caches true",
name: "v0.32: Pinned override, with cache miss, RPC returns member and drained with revision",
override: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: testVersion,
},
},
},
setupCache: func(c *testVersionMembershipCache) {},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(
gomock.Any(),
gomock.Any(),
).Return(&matchingservice.CheckTaskQueueVersionMembershipResponse{
IsMember: true,
ShouldSkipReactivation: false, // drained/inactive on matching's side
RevisionNumber: 7,
}, nil)
},
expectError: false,
expectedShouldSkipReactivation: false,
expectedRevisionNumber: 7,
},
{
name: "v0.32: Pinned override, with cache miss, RPC returns member and active",
override: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: testVersion,
},
},
},
setupCache: func(c *testVersionMembershipCache) {},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(
gomock.Any(),
gomock.Any(),
).Return(&matchingservice.CheckTaskQueueVersionMembershipResponse{
IsMember: true,
ShouldSkipReactivation: true, // active on matching's side
}, nil)
},
expectError: false,
expectedShouldSkipReactivation: true,
},
{
name: "v0.32: Pinned override, with cache miss, RPC returns member without eligibility (old matching)",
override: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
@@ -957,7 +1185,8 @@ func TestValidateVersioningOverride(t *testing.T) {
IsMember: true,
}, nil)
},
expectError: false,
expectError: false,
expectedShouldSkipReactivation: false, // old matching server — zero-value treated as "don't know, send signal"
},
{
name: "v0.32: Pinned override, without version, returns error",
@@ -1030,7 +1259,7 @@ func TestValidateVersioningOverride(t *testing.T) {
PinnedVersion: "test-deployment.test-build-id",
},
setupCache: func(c *testVersionMembershipCache) {
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, "test-deployment", "test-build-id", true)
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, "test-deployment", "test-build-id", true, false, 0)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(gomock.Any(), gomock.Any()).Times(0)
@@ -1044,7 +1273,7 @@ func TestValidateVersioningOverride(t *testing.T) {
PinnedVersion: "test-deployment.test-build-id",
},
setupCache: func(c *testVersionMembershipCache) {
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, "test-deployment", "test-build-id", false)
c.Put(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, "test-deployment", "test-build-id", false, false, 0)
},
setupMock: func(m *matchingservicemock.MockMatchingServiceClient) {
m.EXPECT().CheckTaskQueueVersionMembership(gomock.Any(), gomock.Any()).Times(0)
@@ -1154,7 +1383,8 @@ func TestValidateVersioningOverride(t *testing.T) {
},
}, nil)
},
expectError: false,
expectError: false,
expectedShouldSkipReactivation: false, // version data exists with UNSPECIFIED status — not in {CURRENT, RAMPING, DRAINING}
},
{
name: "v0.32: Pinned override, Unimplemented fallback, version is not member",
@@ -1203,7 +1433,7 @@ func TestValidateVersioningOverride(t *testing.T) {
if tqType == enumspb.TASK_QUEUE_TYPE_UNSPECIFIED {
tqType = enumspb.TASK_QUEUE_TYPE_WORKFLOW
}
err := ValidateVersioningOverride(
shouldSkipReactivation, revisionNumber, err := ValidateVersioningOverrideAndGetReactivationEligibility(
context.Background(),
tt.override,
mockMatchingClient,
@@ -1220,6 +1450,8 @@ func TestValidateVersioningOverride(t *testing.T) {
}
} else {
require.NoError(t, err)
require.Equal(t, tt.expectedShouldSkipReactivation, shouldSkipReactivation)
require.Equal(t, tt.expectedRevisionNumber, revisionNumber)
}
})
}
@@ -1369,7 +1601,7 @@ func TestGetIsWFTaskQueueInVersionDetector(t *testing.T) {
taskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
deploymentName: testVersion.DeploymentName,
buildID: testVersion.BuildId,
}] = true
}] = versionTaskQueueInfoCacheValue{isMember: true}
function := GetIsWFTaskQueueInVersionDetector(mockClient, cache)
isMember, err := function(context.Background(), testNamespaceID, testTaskQueue, testVersion)
@@ -1391,7 +1623,7 @@ func TestGetIsWFTaskQueueInVersionDetector(t *testing.T) {
taskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
deploymentName: testVersion.DeploymentName,
buildID: testVersion.BuildId,
}] = false
}] = versionTaskQueueInfoCacheValue{isMember: false}
function := GetIsWFTaskQueueInVersionDetector(mockClient, cache)
isMember, err := function(context.Background(), testNamespaceID, testTaskQueue, testVersion)
@@ -1415,7 +1647,7 @@ func TestGetIsWFTaskQueueInVersionDetector(t *testing.T) {
assert.Equal(t, 1, cache.putCalls)
// Verify the value was actually stored.
cached, ok := cache.Get(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId)
cached, _, _, ok := cache.Get(testNamespaceID, testTaskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, testVersion.DeploymentName, testVersion.BuildId)
assert.True(t, ok)
assert.True(t, cached)
})

View File

@@ -738,6 +738,21 @@ message CheckTaskQueueVersionMembershipRequest {
message CheckTaskQueueVersionMembershipResponse {
bool is_member = 1;
// True when a reactivation signal to this version would be redundant — i.e., matching
// determined the version is already in a state where it does not need to be reactivated
// (today: CURRENT, RAMPING, or DRAINING). History uses this to suppress such signals.
// The zero value (false) is the safe default; it applies when matching has no definitive
// answer (version not present in matching's deployment data, or old matching servers
// that do not set this field) and tells history to send the signal.
bool should_skip_reactivation = 2;
// revision_number is the version's current revision as tracked in matching's per-TQ
// deployment data. It is returned so history can compose a stable, cluster-wide-deterministic
// RequestId on the reactivation signal. All history pods querying the same version at the
// same point in time converge on the same revision_number, so Temporal's built-in
// SignalRequestedIds dedup (see signalworkflow/api.go) collapses the N-pod signal fan-out
// into exactly one event on the version workflow. Zero when unknown (old matching server or
// legacy DeploymentVersionData format that does not carry revision_number).
int64 revision_number = 3;
}
// PollConditions are extra conditions to set on the poll. Only supported with new matcher.

View File

@@ -60,8 +60,7 @@ func Invoke(
workflowConsistencyChecker api.WorkflowConsistencyChecker,
tokenSerializer *tasktoken.Serializer,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
reactivationSignaler api.VersionReactivationSignalerFn,
testHooks testhooks.TestHooks,
) (*historyservice.ExecuteMultiOperationResponse, error) {
@@ -102,8 +101,7 @@ func Invoke(
tokenSerializer,
startReq,
matchingClient,
versionMembershipCache,
reactivationSignalCache,
versionCache,
reactivationSignaler,
uws.workflowLeaseCallback(ctx),
)

View File

@@ -30,8 +30,7 @@ func Invoke(
shardContext historyi.ShardContext,
workflowConsistencyChecker api.WorkflowConsistencyChecker,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
reactivationSignaler api.VersionReactivationSignalerFn,
) (_ *historyservice.ResetWorkflowExecutionResponse, retError error) {
namespaceID := namespace.ID(resetRequest.GetNamespaceId())
@@ -77,7 +76,7 @@ func Invoke(
}
// Validate versioning override, if any.
err = validatePostResetOperationInputs(ctx, request.GetPostResetOperations(), matchingClient, versionMembershipCache,
shouldSkipReactivationPerOp, revisionNumberPerOp, err := validatePostResetOperationInputs(ctx, request.GetPostResetOperations(), matchingClient, versionCache,
baseMutableState.GetExecutionInfo().GetTaskQueue(), namespaceID.String())
if err != nil {
return nil, err
@@ -193,10 +192,10 @@ func Invoke(
}
// Notify version workflow if we're pinning to a potentially drained version via post-reset operations
for _, operation := range request.GetPostResetOperations() {
for i, operation := range request.GetPostResetOperations() {
if updateOpts, ok := operation.GetVariant().(*workflowpb.PostResetOperation_UpdateWorkflowOptions_); ok {
api.ReactivateVersionWorkflowIfPinned(ctx, namespaceEntry,
updateOpts.UpdateWorkflowOptions.GetWorkflowExecutionOptions().GetVersioningOverride(), reactivationSignalCache, reactivationSignaler, shardContext.GetConfig().EnableVersionReactivationSignals())
updateOpts.UpdateWorkflowOptions.GetWorkflowExecutionOptions().GetVersioningOverride(), reactivationSignaler, shardContext.GetConfig().EnableVersionReactivationSignals(), shouldSkipReactivationPerOp[i], revisionNumberPerOp[i])
}
}
@@ -234,22 +233,37 @@ func GetResetReapplyExcludeTypes(
}
// validatePostResetOperationInputs validates the optional post reset operation inputs.
// Returns parallel slices (one entry per operation) carrying the reactivation-signal inputs
// derived from the operation's versioning override:
// - shouldSkipReactivationPerOp: whether each operation's pinned version is active or
// still draining per matching (true → no need to send a reactivation signal; false
// covers drained/inactive, unknown, not-found, and old-matching cases).
// - revisionNumberPerOp: the pinned version's revision number per matching's view, used to
// compose a stable RequestId on the reactivation signal for receiver-side dedup.
//
// Both are only populated for UpdateWorkflowOptions operations; other operation types default
// to zero values.
func validatePostResetOperationInputs(ctx context.Context,
postResetOperations []*workflowpb.PostResetOperation,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
taskQueue string,
namespaceID string) error {
for _, operation := range postResetOperations {
namespaceID string) ([]bool, []int64, error) {
shouldSkipReactivationPerOp := make([]bool, len(postResetOperations))
revisionNumberPerOp := make([]int64, len(postResetOperations))
for i, operation := range postResetOperations {
switch op := operation.GetVariant().(type) {
case *workflowpb.PostResetOperation_UpdateWorkflowOptions_:
opts := op.UpdateWorkflowOptions.GetWorkflowExecutionOptions()
if err := worker_versioning.ValidateVersioningOverride(ctx, opts.GetVersioningOverride(), matchingClient, versionMembershipCache, taskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, namespaceID); err != nil {
return err
shouldSkipReactivation, revisionNumber, err := worker_versioning.ValidateVersioningOverrideAndGetReactivationEligibility(ctx, opts.GetVersioningOverride(), matchingClient, versionCache, taskQueue, enumspb.TASK_QUEUE_TYPE_WORKFLOW, namespaceID)
if err != nil {
return nil, nil, err
}
shouldSkipReactivationPerOp[i] = shouldSkipReactivation
revisionNumberPerOp[i] = revisionNumber
default:
return serviceerror.NewInvalidArgumentf("unsupported post reset operation: %T", op)
return nil, nil, serviceerror.NewInvalidArgumentf("unsupported post reset operation: %T", op)
}
}
return nil
return shouldSkipReactivationPerOp, revisionNumberPerOp, nil
}

View File

@@ -64,7 +64,7 @@ type (
commandHandlerRegistry *workflow.CommandHandlerRegistry
chasmWorkflowRegistry *chasmworkflow.Registry
matchingClient matchingservice.MatchingServiceClient
versionMembershipCache worker_versioning.VersionMembershipCache
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache
}
)
@@ -78,7 +78,7 @@ func NewWorkflowTaskCompletedHandler(
visibilityManager manager.VisibilityManager,
workflowConsistencyChecker api.WorkflowConsistencyChecker,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
) *WorkflowTaskCompletedHandler {
return &WorkflowTaskCompletedHandler{
config: shardContext.GetConfig(),
@@ -102,7 +102,7 @@ func NewWorkflowTaskCompletedHandler(
commandHandlerRegistry: commandHandlerRegistry,
chasmWorkflowRegistry: chasmWorkflowRegistry,
matchingClient: matchingClient,
versionMembershipCache: versionMembershipCache,
versionCache: versionCache,
}
}
@@ -410,7 +410,7 @@ func (handler *WorkflowTaskCompletedHandler) Invoke(
handler.commandHandlerRegistry,
handler.chasmWorkflowRegistry,
handler.matchingClient,
handler.versionMembershipCache,
handler.versionCache,
)
if responseMutations, err = workflowTaskHandler.handleCommands(

View File

@@ -84,7 +84,7 @@ type (
commandHandlerRegistry *workflow.CommandHandlerRegistry
chasmWorkflowRegistry *chasmworkflow.Registry
matchingClient matchingservice.MatchingServiceClient
versionMembershipCache worker_versioning.VersionMembershipCache
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache
}
workflowTaskFailedCause struct {
@@ -126,7 +126,7 @@ func newWorkflowTaskCompletedHandler(
commandHandlerRegistry *workflow.CommandHandlerRegistry,
chasmWorkflowRegistry *chasmworkflow.Registry,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
) *workflowTaskCompletedHandler {
return &workflowTaskCompletedHandler{
identity: identity,
@@ -161,7 +161,7 @@ func newWorkflowTaskCompletedHandler(
commandHandlerRegistry: commandHandlerRegistry,
chasmWorkflowRegistry: chasmWorkflowRegistry,
matchingClient: matchingClient,
versionMembershipCache: versionMembershipCache,
versionCache: versionCache,
}
}
@@ -1122,7 +1122,7 @@ func (handler *workflowTaskCompletedHandler) handleCommandContinueAsNewWorkflow(
handler.workflowTaskCompletedID,
parentNamespace,
attr,
worker_versioning.GetIsWFTaskQueueInVersionDetector(handler.matchingClient, handler.versionMembershipCache),
worker_versioning.GetIsWFTaskQueueInVersionDetector(handler.matchingClient, handler.versionCache),
)
if err != nil {
return nil, err

View File

@@ -22,8 +22,7 @@ func Invoke(
shard historyi.ShardContext,
workflowConsistencyChecker api.WorkflowConsistencyChecker,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
reactivationSignaler api.VersionReactivationSignalerFn,
) (_ *historyservice.SignalWithStartWorkflowExecutionResponse, retError error) {
namespaceEntry, err := api.GetActiveNamespace(shard, namespace.ID(signalWithStartRequest.GetNamespaceId()), signalWithStartRequest.SignalWithStartRequest.WorkflowId)
@@ -71,7 +70,7 @@ func Invoke(
}
// Validation for versioning override, if any.
err = worker_versioning.ValidateVersioningOverride(ctx, request.GetVersioningOverride(), matchingClient, versionMembershipCache, request.GetTaskQueue().GetName(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, namespaceID.String())
shouldSkipReactivation, revisionNumber, err := worker_versioning.ValidateVersioningOverrideAndGetReactivationEligibility(ctx, request.GetVersioningOverride(), matchingClient, versionCache, request.GetTaskQueue().GetName(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, namespaceID.String())
if err != nil {
return nil, err
}
@@ -90,7 +89,7 @@ func Invoke(
// Notify version workflow if we're starting a new workflow pinned to a potentially drained version
if started {
api.ReactivateVersionWorkflowIfPinned(ctx, namespaceEntry, request.GetVersioningOverride(), reactivationSignalCache, reactivationSignaler, shard.GetConfig().EnableVersionReactivationSignals())
api.ReactivateVersionWorkflowIfPinned(ctx, namespaceEntry, request.GetVersioningOverride(), reactivationSignaler, shard.GetConfig().EnableVersionReactivationSignals(), shouldSkipReactivation, revisionNumber)
}
return &historyservice.SignalWithStartWorkflowExecutionResponse{

View File

@@ -58,9 +58,10 @@ type Starter struct {
request *historyservice.StartWorkflowExecutionRequest
namespace *namespace.Namespace
createOrUpdateLeaseFn api.CreateOrUpdateLeaseFunc
versionMembershipCache worker_versioning.VersionMembershipCache
reactivationSignalCache worker_versioning.ReactivationSignalCache
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache
reactivationSignaler api.VersionReactivationSignalerFn
shouldSkipReactivation bool
revisionNumber int64
}
// creationParams is a container for all information obtained from creating the uncommitted execution.
@@ -89,8 +90,7 @@ func NewStarter(
tokenSerializer *tasktoken.Serializer,
request *historyservice.StartWorkflowExecutionRequest,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
reactivationSignaler api.VersionReactivationSignalerFn,
createLeaseFn api.CreateOrUpdateLeaseFunc,
) (*Starter, error) {
@@ -109,8 +109,7 @@ func NewStarter(
request: request,
namespace: namespaceEntry,
createOrUpdateLeaseFn: createLeaseFn,
versionMembershipCache: versionMembershipCache,
reactivationSignalCache: reactivationSignalCache,
versionCache: versionCache,
reactivationSignaler: reactivationSignaler,
}, nil
}
@@ -136,7 +135,7 @@ func (s *Starter) prepare(ctx context.Context) error {
}
// Validation for versioning override, if any.
err = worker_versioning.ValidateVersioningOverride(ctx, request.GetVersioningOverride(), s.matchingClient, s.versionMembershipCache, request.GetTaskQueue().GetName(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, s.namespace.ID().String())
s.shouldSkipReactivation, s.revisionNumber, err = worker_versioning.ValidateVersioningOverrideAndGetReactivationEligibility(ctx, request.GetVersioningOverride(), s.matchingClient, s.versionCache, request.GetTaskQueue().GetName(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, s.namespace.ID().String())
if err != nil {
return err
}
@@ -213,7 +212,7 @@ func (s *Starter) Invoke(
// Notify version workflow if we are starting a workflow execution on a potentially drained version.
// Only signal when a new workflow was actually created (StartNew), not for deduped retries
// (StartDeduped) or reused existing workflows (StartReused) where the pinned override is not applied.
api.ReactivateVersionWorkflowIfPinned(ctx, s.namespace, s.request.StartRequest.GetVersioningOverride(), s.reactivationSignalCache, s.reactivationSignaler, s.shardContext.GetConfig().EnableVersionReactivationSignals())
api.ReactivateVersionWorkflowIfPinned(ctx, s.namespace, s.request.StartRequest.GetVersioningOverride(), s.reactivationSignaler, s.shardContext.GetConfig().EnableVersionReactivationSignals(), s.shouldSkipReactivation, s.revisionNumber)
}
return resp, outcome, conflictErr
}
@@ -221,7 +220,7 @@ func (s *Starter) Invoke(
}
// Notify version workflow if we're pinning to a potentially drained version
api.ReactivateVersionWorkflowIfPinned(ctx, s.namespace, s.request.StartRequest.GetVersioningOverride(), s.reactivationSignalCache, s.reactivationSignaler, s.shardContext.GetConfig().EnableVersionReactivationSignals())
api.ReactivateVersionWorkflowIfPinned(ctx, s.namespace, s.request.StartRequest.GetVersioningOverride(), s.reactivationSignaler, s.shardContext.GetConfig().EnableVersionReactivationSignals(), s.shouldSkipReactivation, s.revisionNumber)
resp, err = s.generateResponse(
creationParams.runID,

View File

@@ -27,8 +27,7 @@ func Invoke(
shardCtx historyi.ShardContext,
workflowConsistencyChecker api.WorkflowConsistencyChecker,
matchingClient matchingservice.MatchingServiceClient,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
reactivationSignaler api.VersionReactivationSignalerFn,
) (*historyservice.UpdateWorkflowExecutionOptionsResponse, error) {
req := request.GetUpdateRequest()
@@ -40,6 +39,8 @@ func Invoke(
// Store versioning override to send reactivation signal after successful persistence
var versioningOverrideForReactivation *workflowpb.VersioningOverride
var shouldSkipReactivation bool
var revisionNumber int64
err = api.GetAndUpdateWorkflowWithNew(
ctx,
@@ -79,7 +80,7 @@ func Invoke(
}
// Validate versioning override, if any.
err = worker_versioning.ValidateVersioningOverride(ctx, requestedOptions.GetVersioningOverride(), matchingClient, versionMembershipCache, mutableState.GetExecutionInfo().GetTaskQueue(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, ns.ID().String())
shouldSkipReactivation, revisionNumber, err = worker_versioning.ValidateVersioningOverrideAndGetReactivationEligibility(ctx, requestedOptions.GetVersioningOverride(), matchingClient, versionCache, mutableState.GetExecutionInfo().GetTaskQueue(), enumspb.TASK_QUEUE_TYPE_WORKFLOW, ns.ID().String())
if err != nil {
return nil, err
}
@@ -118,7 +119,7 @@ func Invoke(
// Notify version workflow if we're pinning to a potentially drained version.
// This is done after successful persistence to avoid signaling if the update fails.
api.ReactivateVersionWorkflowIfPinned(ctx, ns, versioningOverrideForReactivation, reactivationSignalCache, reactivationSignaler, shardCtx.GetConfig().EnableVersionReactivationSignals())
api.ReactivateVersionWorkflowIfPinned(ctx, ns, versioningOverrideForReactivation, reactivationSignaler, shardCtx.GetConfig().EnableVersionReactivationSignals(), shouldSkipReactivation, revisionNumber)
return ret, nil
}

View File

@@ -33,36 +33,32 @@ import (
"google.golang.org/protobuf/types/known/timestamppb"
)
type noopVersionMembershipCache struct{}
type noopVersionCache struct{}
func (noopVersionMembershipCache) Get(
func (noopVersionCache) Get(
_ string,
_ string,
_ enumspb.TaskQueueType,
_ string,
_ string,
) (isMember bool, ok bool) {
return false, false
) (isMember bool, shouldSkipReactivation bool, revisionNumber int64, ok bool) {
return false, false, 0, false
}
func (noopVersionMembershipCache) Put(
func (noopVersionCache) Put(
_ string,
_ string,
_ enumspb.TaskQueueType,
_ string,
_ string,
_ bool,
_ bool,
_ int64,
) {
}
type noopReactivationSignalCache struct{}
func (noopReactivationSignalCache) ShouldSendSignal(_, _, _ string) bool {
return false // Always return false to skip sending signals in tests
}
// noopReactivationSignaler is a no-op signaler function for tests
func noopReactivationSignaler(_ context.Context, _ *namespace.Namespace, _, _ string) error {
func noopReactivationSignaler(_ context.Context, _ *namespace.Namespace, _, _ string, _ int64) error {
return nil
}
@@ -330,9 +326,8 @@ func (s *updateWorkflowOptionsSuite) TestInvoke_Success() {
s.shardContext,
s.workflowConsistencyChecker,
s.mockMatchingClient,
noopVersionMembershipCache{}, // cache not meant to be used in this test
noopReactivationSignalCache{}, // cache not meant to be used in this test
noopReactivationSignaler, // signaler not meant to be used in this test
noopVersionCache{}, // cache not meant to be used in this test
noopReactivationSignaler, // signaler not meant to be used in this test
)
s.NoError(err)
s.NotNil(resp)

View File

@@ -11,32 +11,42 @@ import (
// VersionReactivationSignalerFn is a function type for sending reactivation signals to version workflows.
// This abstraction allows the history API layer to use the deployment client without importing it directly,
// avoiding import cycles between history/api and worker/workerdeployment packages.
// revisionNumber is the version's current revision per matching's view and is used by the signaler
// to compose a cluster-wide-deterministic RequestId on the signal for receiver-side dedup.
type VersionReactivationSignalerFn func(
ctx context.Context,
namespaceEntry *namespace.Namespace,
deploymentName, buildID string,
revisionNumber int64,
) error
// ReactivateVersionWorkflowIfPinned sends a reactivation signal to the version workflow
// when workflows are pinned to a potentially DRAINED/INACTIVE version. It also deduplicates
// signals within the cache TTL window.
// when workflows are pinned to a potentially DRAINED/INACTIVE version.
// This is a fire-and-forget operation - the signal is sent asynchronously and errors are
// logged by the signaler implementation.
// logged by the signaler implementation. The signaler itself is responsible for per-pod
// dedup by revision number; cross-pod duplicates fold at the receiver via a deterministic
// UUID RequestId.
//
//nolint:revive,errcheck
func ReactivateVersionWorkflowIfPinned(
ctx context.Context,
namespaceEntry *namespace.Namespace,
override *workflowpb.VersioningOverride,
signalCache worker_versioning.ReactivationSignalCache,
signaler VersionReactivationSignalerFn,
enabled bool,
shouldSkipReactivation bool,
revisionNumber int64,
) {
// Check if signals are enabled globally
if !enabled {
return
}
// Skip signal if matching confirmed the version is active or still draining.
if shouldSkipReactivation {
return
}
// Only process if we're pinning to a specific version
if !worker_versioning.OverrideIsPinned(override) {
return
@@ -47,19 +57,10 @@ func ReactivateVersionWorkflowIfPinned(
return
}
// Check cache - skip if signal was recently sent
if signalCache != nil && !signalCache.ShouldSendSignal(
namespaceEntry.ID().String(),
pinnedVersion.GetDeploymentName(),
pinnedVersion.GetBuildId(),
) {
return
}
// Send the signal asynchronously to avoid adding latency to the caller's request.
// Errors are logged by the signaler implementation (e.g. via convertAndRecordError). However,
// errors are not propagated to the caller as this is a fire-and-forget operation.
go func() {
signaler(context.Background(), namespaceEntry, pinnedVersion.GetDeploymentName(), pinnedVersion.GetBuildId()) //nolint:errcheck
signaler(context.Background(), namespaceEntry, pinnedVersion.GetDeploymentName(), pinnedVersion.GetBuildId(), revisionNumber) //nolint:errcheck
}()
}

View File

@@ -416,16 +416,14 @@ type Config struct {
NumConsecutiveWorkflowTaskProblemsToTriggerSearchAttribute dynamicconfig.IntPropertyFnWithNamespaceFilter
// Worker-Versioning related settings
EnableSuggestCaNOnNewTargetVersion dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnableSendTargetVersionChanged dynamicconfig.BoolPropertyFnWithNamespaceFilter
UseRevisionNumberForWorkerVersioning dynamicconfig.BoolPropertyFnWithNamespaceFilter
VersionMembershipCacheTTL dynamicconfig.DurationPropertyFn
VersionMembershipCacheMaxSize dynamicconfig.IntPropertyFn
VersionReactivationSignalCacheTTL dynamicconfig.DurationPropertyFn
VersionReactivationSignalCacheMaxSize dynamicconfig.IntPropertyFn
EnableVersionReactivationSignals dynamicconfig.BoolPropertyFn
RoutingInfoCacheTTL dynamicconfig.DurationPropertyFn
RoutingInfoCacheMaxSize dynamicconfig.IntPropertyFn
EnableSuggestCaNOnNewTargetVersion dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnableSendTargetVersionChanged dynamicconfig.BoolPropertyFnWithNamespaceFilter
UseRevisionNumberForWorkerVersioning dynamicconfig.BoolPropertyFnWithNamespaceFilter
VersionMembershipCacheTTL dynamicconfig.DurationPropertyFn
VersionMembershipCacheMaxSize dynamicconfig.IntPropertyFn
EnableVersionReactivationSignals dynamicconfig.BoolPropertyFn
RoutingInfoCacheTTL dynamicconfig.DurationPropertyFn
RoutingInfoCacheMaxSize dynamicconfig.IntPropertyFn
}
// NewConfig returns new service config with default values
@@ -797,16 +795,14 @@ func NewConfig(
NumConsecutiveWorkflowTaskProblemsToTriggerSearchAttribute: dynamicconfig.NumConsecutiveWorkflowTaskProblemsToTriggerSearchAttribute.Get(dc),
// Worker-Versioning related
UseRevisionNumberForWorkerVersioning: dynamicconfig.UseRevisionNumberForWorkerVersioning.Get(dc),
EnableSuggestCaNOnNewTargetVersion: dynamicconfig.EnableSuggestCaNOnNewTargetVersion.Get(dc),
EnableSendTargetVersionChanged: dynamicconfig.EnableSendTargetVersionChanged.Get(dc),
VersionMembershipCacheTTL: dynamicconfig.VersionMembershipCacheTTL.Get(dc),
VersionMembershipCacheMaxSize: dynamicconfig.VersionMembershipCacheMaxSize.Get(dc),
VersionReactivationSignalCacheTTL: dynamicconfig.VersionReactivationSignalCacheTTL.Get(dc),
VersionReactivationSignalCacheMaxSize: dynamicconfig.VersionReactivationSignalCacheMaxSize.Get(dc),
EnableVersionReactivationSignals: dynamicconfig.EnableVersionReactivationSignals.Get(dc),
RoutingInfoCacheTTL: dynamicconfig.RoutingInfoCacheTTL.Get(dc),
RoutingInfoCacheMaxSize: dynamicconfig.RoutingInfoCacheMaxSize.Get(dc),
UseRevisionNumberForWorkerVersioning: dynamicconfig.UseRevisionNumberForWorkerVersioning.Get(dc),
EnableSuggestCaNOnNewTargetVersion: dynamicconfig.EnableSuggestCaNOnNewTargetVersion.Get(dc),
EnableSendTargetVersionChanged: dynamicconfig.EnableSendTargetVersionChanged.Get(dc),
VersionMembershipCacheTTL: dynamicconfig.VersionMembershipCacheTTL.Get(dc),
VersionMembershipCacheMaxSize: dynamicconfig.VersionMembershipCacheMaxSize.Get(dc),
EnableVersionReactivationSignals: dynamicconfig.EnableVersionReactivationSignals.Get(dc),
RoutingInfoCacheTTL: dynamicconfig.RoutingInfoCacheTTL.Get(dc),
RoutingInfoCacheMaxSize: dynamicconfig.RoutingInfoCacheMaxSize.Get(dc),
}
return cfg

View File

@@ -89,7 +89,6 @@ var Module = fx.Options(
fx.Provide(NewService),
fx.Provide(ReplicationProgressCacheProvider),
fx.Provide(VersionMembershipCacheProvider),
fx.Provide(ReactivationSignalCacheProvider),
workerdeployment.ClientModule,
fx.Provide(RoutingInfoCacheProvider),
fx.Invoke(ServiceLifetimeHooks),
@@ -406,7 +405,7 @@ func VersionMembershipCacheProvider(
lc fx.Lifecycle,
serviceConfig *configs.Config,
metricsHandler metrics.Handler,
) worker_versioning.VersionMembershipCache {
) worker_versioning.VersionMembershipAndReactivationStatusCache {
c := commoncache.New(serviceConfig.VersionMembershipCacheMaxSize(), &commoncache.Options{
TTL: max(1*time.Second, serviceConfig.VersionMembershipCacheTTL()),
})
@@ -416,24 +415,7 @@ func VersionMembershipCacheProvider(
return nil
},
})
return worker_versioning.NewVersionMembershipCache(c, metricsHandler)
}
func ReactivationSignalCacheProvider(
lc fx.Lifecycle,
serviceConfig *configs.Config,
metricsHandler metrics.Handler,
) worker_versioning.ReactivationSignalCache {
c := commoncache.New(serviceConfig.VersionReactivationSignalCacheMaxSize(), &commoncache.Options{
TTL: max(1*time.Second, serviceConfig.VersionReactivationSignalCacheTTL()),
})
lc.Append(fx.Hook{
OnStop: func(context.Context) error {
c.Stop()
return nil
},
})
return worker_versioning.NewReactivationSignalCache(c, metricsHandler)
return worker_versioning.NewVersionMembershipAndReactivationStatusCache(c, metricsHandler)
}
func RoutingInfoCacheProvider(

View File

@@ -136,8 +136,7 @@ type (
workflowConsistencyChecker api.WorkflowConsistencyChecker
chasmEngine chasm.Engine
versionChecker headers.VersionChecker
versionMembershipCache worker_versioning.VersionMembershipCache
reactivationSignalCache worker_versioning.ReactivationSignalCache
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache
workerDeploymentClient workerdeployment.Client
routingInfoCache worker_versioning.RoutingInfoCache
tracer trace.Tracer
@@ -160,8 +159,7 @@ func NewEngineWithShardContext(
sdkClientFactory sdk.ClientFactory,
eventNotifier events.Notifier,
config *configs.Config,
versionMembershipCache worker_versioning.VersionMembershipCache,
reactivationSignalCache worker_versioning.ReactivationSignalCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
workerDeploymentClient workerdeployment.Client,
routingInfoCache worker_versioning.RoutingInfoCache,
rawMatchingClient matchingservice.MatchingServiceClient,
@@ -235,8 +233,7 @@ func NewEngineWithShardContext(
outboundQueueCBPool: outboundQueueCBPool,
testHooks: testHooks,
chasmEngine: chasmEngine,
versionMembershipCache: versionMembershipCache,
reactivationSignalCache: reactivationSignalCache,
versionCache: versionCache,
workerDeploymentClient: workerDeploymentClient,
routingInfoCache: routingInfoCache,
}
@@ -416,8 +413,7 @@ func (e *historyEngineImpl) StartWorkflowExecution(
e.tokenSerializer,
startRequest,
e.matchingClient,
e.versionMembershipCache,
e.reactivationSignalCache,
e.versionCache,
e.workerDeploymentClient.SignalVersionReactivation,
api.NewWorkflowLeaseAndContext,
)
@@ -440,8 +436,7 @@ func (e *historyEngineImpl) ExecuteMultiOperation(
e.workflowConsistencyChecker,
e.tokenSerializer,
e.matchingClient,
e.versionMembershipCache,
e.reactivationSignalCache,
e.versionCache,
e.workerDeploymentClient.SignalVersionReactivation,
e.testHooks,
)
@@ -589,7 +584,7 @@ func (e *historyEngineImpl) RespondWorkflowTaskCompleted(
e.persistenceVisibilityMgr,
e.workflowConsistencyChecker,
e.matchingClient,
e.versionMembershipCache,
e.versionCache,
)
return h.Invoke(ctx, req)
}
@@ -658,7 +653,7 @@ func (e *historyEngineImpl) SignalWithStartWorkflowExecution(
ctx context.Context,
req *historyservice.SignalWithStartWorkflowExecutionRequest,
) (_ *historyservice.SignalWithStartWorkflowExecutionResponse, retError error) {
return signalwithstartworkflow.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionMembershipCache, e.reactivationSignalCache, e.workerDeploymentClient.SignalVersionReactivation)
return signalwithstartworkflow.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionCache, e.workerDeploymentClient.SignalVersionReactivation)
}
func (e *historyEngineImpl) UpdateWorkflowExecution(
@@ -859,7 +854,7 @@ func (e *historyEngineImpl) ResetWorkflowExecution(
ctx context.Context,
req *historyservice.ResetWorkflowExecutionRequest,
) (*historyservice.ResetWorkflowExecutionResponse, error) {
return resetworkflow.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionMembershipCache, e.reactivationSignalCache, e.workerDeploymentClient.SignalVersionReactivation)
return resetworkflow.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionCache, e.workerDeploymentClient.SignalVersionReactivation)
}
// UpdateWorkflowExecutionOptions updates the options of a specific workflow execution.
@@ -868,7 +863,7 @@ func (e *historyEngineImpl) UpdateWorkflowExecutionOptions(
ctx context.Context,
req *historyservice.UpdateWorkflowExecutionOptionsRequest,
) (*historyservice.UpdateWorkflowExecutionOptionsResponse, error) {
return updateworkflowoptions.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionMembershipCache, e.reactivationSignalCache, e.workerDeploymentClient.SignalVersionReactivation)
return updateworkflowoptions.Invoke(ctx, req, e.shardContext, e.workflowConsistencyChecker, e.matchingClient, e.versionCache, e.workerDeploymentClient.SignalVersionReactivation)
}
func (e *historyEngineImpl) NotifyNewHistoryEvent(

View File

@@ -104,7 +104,7 @@ type (
// by the history engine as a function value.
type noopWorkerDeploymentClient struct{ workerdeployment.Client }
func (noopWorkerDeploymentClient) SignalVersionReactivation(context.Context, *namespace.Namespace, string, string) error {
func (noopWorkerDeploymentClient) SignalVersionReactivation(context.Context, *namespace.Namespace, string, string, int64) error {
return nil
}

View File

@@ -52,8 +52,7 @@ type (
PersistenceRateLimiter replication.PersistenceRateLimiter
TestHooks testhooks.TestHooks
ChasmEngine chasm.Engine
VersionMembershipCache worker_versioning.VersionMembershipCache
ReactivationSignalCache worker_versioning.ReactivationSignalCache
VersionMembershipCache worker_versioning.VersionMembershipAndReactivationStatusCache
WorkerDeploymentClient workerdeployment.Client
RoutingInfoCache worker_versioning.RoutingInfoCache
}
@@ -74,7 +73,6 @@ func (f *historyEngineFactory) CreateEngine(
f.EventNotifier,
f.Config,
f.VersionMembershipCache,
f.ReactivationSignalCache,
f.WorkerDeploymentClient,
f.RoutingInfoCache,
f.RawMatchingClient,

View File

@@ -159,5 +159,5 @@ type unusedDependencies struct {
cache.Cache
chasm.Engine
ChasmRegistry *chasm.Registry
worker_versioning.VersionMembershipCache
worker_versioning.VersionMembershipAndReactivationStatusCache
}

View File

@@ -54,7 +54,7 @@ type (
workflowResetter ndc.WorkflowResetter
parentClosePolicyClient parentclosepolicy.Client
versionMembershipCache worker_versioning.VersionMembershipCache
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache
}
)
@@ -69,7 +69,7 @@ func newTransferQueueActiveTaskExecutor(
matchingRawClient resource.MatchingRawClient,
visibilityManager manager.VisibilityManager,
chasmEngine chasm.Engine,
versionMembershipCache worker_versioning.VersionMembershipCache,
versionCache worker_versioning.VersionMembershipAndReactivationStatusCache,
) queues.Executor {
return &transferQueueActiveTaskExecutor{
transferQueueTaskExecutorBase: newTransferQueueTaskExecutorBase(
@@ -93,7 +93,7 @@ func newTransferQueueActiveTaskExecutor(
sdkClientFactory,
config.NumParentClosePolicySystemWorkflows(),
),
versionMembershipCache: versionMembershipCache,
versionCache: versionCache,
}
}
@@ -925,7 +925,7 @@ func (t *transferQueueActiveTaskExecutor) processStartChildExecution(
if attributes.GetNamespaceId() != mutableState.GetExecutionInfo().GetNamespaceId() { // don't inherit pinned version if child is in a different namespace
inheritedPinnedVersion = nil
} else if newTQ != mutableState.GetExecutionInfo().GetTaskQueue() {
newTQInPinnedVersion, err = worker_versioning.GetIsWFTaskQueueInVersionDetector(t.matchingRawClient, t.versionMembershipCache)(ctx, attributes.GetNamespaceId(), newTQ, inheritedPinnedVersion)
newTQInPinnedVersion, err = worker_versioning.GetIsWFTaskQueueInVersionDetector(t.matchingRawClient, t.versionCache)(ctx, attributes.GetNamespaceId(), newTQ, inheritedPinnedVersion)
if err != nil {
return fmt.Errorf("error determining child task queue presence in inherited version: %w", err)
}
@@ -964,7 +964,7 @@ func (t *transferQueueActiveTaskExecutor) processStartChildExecution(
if attributes.GetNamespaceId() != mutableState.GetExecutionInfo().GetNamespaceId() { // don't inherit auto upgrade info if child is in a different namespace
inheritedAutoUpgradeInfo = nil
} else if newTQ != mutableState.GetExecutionInfo().GetTaskQueue() {
TQInSourceDeploymentVersion, err := worker_versioning.GetIsWFTaskQueueInVersionDetector(t.matchingRawClient, t.versionMembershipCache)(ctx, attributes.GetNamespaceId(), newTQ, inheritedAutoUpgradeInfo.GetSourceDeploymentVersion())
TQInSourceDeploymentVersion, err := worker_versioning.GetIsWFTaskQueueInVersionDetector(t.matchingRawClient, t.versionCache)(ctx, attributes.GetNamespaceId(), newTQ, inheritedAutoUpgradeInfo.GetSourceDeploymentVersion())
if err != nil {
return fmt.Errorf("error determining child task queue presence in inherited version: %w", err)
}

View File

@@ -32,7 +32,7 @@ type (
HistoryRawClient resource.HistoryRawClient
MatchingRawClient resource.MatchingRawClient
VisibilityManager manager.VisibilityManager
VersionMembershipCache worker_versioning.VersionMembershipCache
VersionMembershipCache worker_versioning.VersionMembershipAndReactivationStatusCache
}
transferQueueFactory struct {

View File

@@ -3462,8 +3462,25 @@ func (e *matchingEngineImpl) CheckTaskQueueVersionMembership(
}
typedUserData := userData.GetData().GetPerType()[int32(request.GetTaskQueueType())]
present := worker_versioning.HasDeploymentVersion(typedUserData.GetDeploymentData(), request.GetVersion())
return &matchingservice.CheckTaskQueueVersionMembershipResponse{IsMember: present}, nil
deploymentData := typedUserData.GetDeploymentData()
present := worker_versioning.HasDeploymentVersion(deploymentData, request.GetVersion())
// Report whether the version is active-or-draining so callers can skip sending
// reactivation signals to versions that don't need one (CURRENT/RAMPING/DRAINING —
// see worker_versioning.ShouldSkipReactivation). The revision number flows back
// so history can compose a cluster-wide-deterministic RequestId on the reactivation
// signal for receiver-side dedup.
shouldSkipReactivation, revisionNumber := worker_versioning.ShouldSkipReactivation(
deploymentData,
request.GetVersion().GetDeploymentName(),
request.GetVersion().GetBuildId(),
)
return &matchingservice.CheckTaskQueueVersionMembershipResponse{
IsMember: present,
ShouldSkipReactivation: shouldSkipReactivation,
RevisionNumber: revisionNumber,
}, nil
}
func (e *matchingEngineImpl) UpdateTaskQueueConfig(

View File

@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"sort"
"strconv"
"strings"
"time"
@@ -26,6 +27,7 @@ import (
"go.temporal.io/server/api/matchingservice/v1"
"go.temporal.io/server/common"
"go.temporal.io/server/common/backoff"
"go.temporal.io/server/common/cache"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/log/tag"
@@ -205,10 +207,14 @@ type Client interface {
// SignalVersionReactivation sends a reactivation signal to a version workflow.
// Used when workflows are pinned to a potentially DRAINED/INACTIVE version.
// This is a fire-and-forget operation - errors are logged but returned for caller handling.
// revisionNumber is used to compose a stable RequestId on the signal so concurrent signals
// for the same (deployment, buildID, revisionNumber) tuple are deduplicated by the receiver
// via Temporal's built-in pendingSignalRequestedIDs mechanism.
SignalVersionReactivation(
ctx context.Context,
namespaceEntry *namespace.Namespace,
deploymentName, buildID string,
revisionNumber int64,
) error
ValidateComputeConfig(
@@ -226,6 +232,13 @@ type ErrRegister struct{ error }
var retryPolicy = backoff.NewExponentialRetryPolicy(100 * time.Millisecond).WithExpirationInterval(1 * time.Minute)
// ClientImpl implements Client
// reactivationVersionKey identifies one target version workflow for reactivation-signal dedup.
type reactivationVersionKey struct {
namespaceID string
deploymentName string
buildID string
}
type ClientImpl struct {
logger log.Logger
historyClient historyservice.HistoryServiceClient
@@ -238,6 +251,14 @@ type ClientImpl struct {
maxDeployments dynamicconfig.IntPropertyFnWithNamespaceFilter
testHooks testhooks.TestHooks
metricsHandler metrics.Handler
// highestRevSignaledToVersionWf is a per-pod LRU that holds, for each target version
// workflow, the highest revision number this pod has successfully signaled for
// reactivation. A reactivation signal for the same or older revision is skipped.
// Size-bounded by ReactivationSignalDedupCacheMaxSize; LRU eviction is pure memory
// hygiene, not part of the dedup contract. Cross-pod deduplication is handled
// separately by the deterministic UUID RequestId on each signal.
highestRevSignaledToVersionWf cache.Cache
}
func (d *ClientImpl) SetManager(
@@ -1988,12 +2009,48 @@ func (d *ClientImpl) SignalVersionReactivation(
ctx context.Context,
namespaceEntry *namespace.Namespace,
deploymentName, buildID string,
revisionNumber int64,
) (retErr error) {
//revive:disable-next-line:defer
defer d.convertAndRecordError("SignalVersionReactivation", deploymentName, &retErr, buildID)()
key := reactivationVersionKey{
namespaceID: namespaceEntry.ID().String(),
deploymentName: deploymentName,
buildID: buildID,
}
metricsHandler := d.metricsHandler.WithTags(
metrics.CacheTypeTag(metrics.ReactivationSignalDedupCacheTypeTagValue),
metrics.OperationTag(metrics.ReactivationSignalDedupScope),
metrics.NamespaceIDTag(namespaceEntry.ID().String()),
)
metrics.CacheRequests.With(metricsHandler).Record(1)
if stored := d.highestRevSignaledToVersionWf.Get(key); stored != nil {
if storedRev, ok := stored.(int64); ok && revisionNumber <= storedRev {
// This pod has already signaled the target version workflow at this revision
// or a newer one; skip to avoid a redundant RPC. Another pod may still send;
// the receiver dedups via the deterministic UUID RequestId.
return nil
}
}
metrics.CacheMissCounter.With(metricsHandler).Record(1)
workflowID := GenerateVersionWorkflowID(deploymentName, buildID)
// Deterministic UUID v5 RequestId derived from the revision number. Multiple history
// pods that independently decide to reactivate the same version at the same revision
// compute the same RequestId and fold into a single signal delivery, since the receiver
// dedups on RequestId.
//
// Revision alone is enough input because the dedup is scoped to this one version
// workflow (addressed by workflowID); different (deployment, build) pairs at the same
// revision don't collide because they're different workflows.
requestID := uuid.NewSHA1(
uuid.Nil,
[]byte("reactivation-signal:"+strconv.FormatInt(revisionNumber, 10)),
).String()
signalRequest := &historyservice.SignalWorkflowExecutionRequest{
NamespaceId: namespaceEntry.ID().String(),
SignalRequest: &workflowservice.SignalWorkflowExecutionRequest{
@@ -2004,11 +2061,19 @@ func (d *ClientImpl) SignalVersionReactivation(
SignalName: ReactivateVersionSignalName,
Input: nil,
Identity: "history-service",
RequestId: requestID,
},
}
_, err := d.historyClient.SignalWorkflowExecution(ctx, signalRequest)
return err
if err != nil {
return err
}
// Record success so subsequent calls to the same version workflow for this or older
// revisions skip the RPC.
d.highestRevSignaledToVersionWf.Put(key, revisionNumber)
return nil
}
func (d *ClientImpl) getSyncBatchSize() int32 {

View File

@@ -1,12 +1,14 @@
package workerdeployment
import (
"context"
"time"
wciclient "go.temporal.io/auto-scaled-workers/wci/client"
sdkworker "go.temporal.io/sdk/worker"
"go.temporal.io/sdk/workflow"
deploymentspb "go.temporal.io/server/api/deployment/v1"
"go.temporal.io/server/common/cache"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/metrics"
@@ -67,6 +69,7 @@ var Module = fx.Options(
)
func ClientProvider(
lc fx.Lifecycle,
logger log.Logger,
historyClient resource.HistoryClient,
matchingClient resource.MatchingClient,
@@ -76,6 +79,13 @@ func ClientProvider(
testHooks testhooks.TestHooks,
metricsHandler metrics.Handler,
) Client {
highestRevSignaledToVersionWf := cache.New(dynamicconfig.ReactivationSignalDedupCacheMaxSize.Get(dc)(), nil)
lc.Append(fx.Hook{
OnStop: func(context.Context) error {
highestRevSignaledToVersionWf.Stop()
return nil
},
})
return &ClientImpl{
logger: logger,
historyClient: historyClient,
@@ -88,6 +98,7 @@ func ClientProvider(
maxDeployments: dynamicconfig.MatchingMaxDeployments.Get(dc),
testHooks: testHooks,
metricsHandler: metricsHandler,
highestRevSignaledToVersionWf: highestRevSignaledToVersionWf,
}
}

View File

@@ -1,8 +1,10 @@
package workerdeployment
import (
"context"
"fmt"
"os"
"regexp"
"strings"
"sync"
"testing"
@@ -13,14 +15,19 @@ import (
deploymentpb "go.temporal.io/api/deployment/v1"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/api/serviceerror"
"go.temporal.io/server/api/historyservice/v1"
"go.temporal.io/server/api/historyservicemock/v1"
persistencespb "go.temporal.io/server/api/persistence/v1"
"go.temporal.io/server/common/cache"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/namespace"
"go.temporal.io/server/common/persistence/visibility/manager"
"go.temporal.io/server/common/worker_versioning"
"go.temporal.io/server/service/history/consts"
update2 "go.temporal.io/server/service/history/workflow/update"
"go.uber.org/mock/gomock"
"google.golang.org/grpc"
)
// testMaxIDLengthLimit is the current default value used by dynamic config for
@@ -68,8 +75,11 @@ func (d *deploymentWorkflowClientSuite) SetupTest() {
BuildId: testBuildID,
}
d.deploymentClient = &ClientImpl{
historyClient: d.mockHistoryClient,
visibilityManager: d.VisibilityManager,
logger: log.NewNoopLogger(),
historyClient: d.mockHistoryClient,
visibilityManager: d.VisibilityManager,
metricsHandler: metrics.NoopMetricsHandler,
highestRevSignaledToVersionWf: cache.New(128, nil),
}
}
@@ -382,3 +392,146 @@ func TestIsRetryableQueryError(t *testing.T) {
require.False(t, isRetryableQueryError(err))
})
}
// TestSignalVersionReactivation_RequestIdFormat verifies that SignalVersionReactivation
// produces a deterministic UUID v5 RequestId derived from the revision number. Every
// history pod that observes the same revisionNumber must produce the same RequestId so
// the receiver collapses concurrent signals on the same dedup key.
func (d *deploymentWorkflowClientSuite) TestSignalVersionReactivation_RequestIdFormat() {
testCases := []struct {
name string
deploymentName string
buildID string
revisionNumber int64
}{
{
name: "non-zero revision",
deploymentName: "my-deployment",
buildID: "build-42",
revisionNumber: 7,
},
{
// Distinct (deployment, build) from other cases so per-pod dedup in the client
// doesn't skip this call. Only the RequestId format is under test here.
name: "zero revision (legacy format / never synced)",
deploymentName: "legacy-deployment",
buildID: "legacy-build",
revisionNumber: 0,
},
{
name: "same revision, different deployment/build — same RequestId (workflow scoping comes from WorkflowID)",
deploymentName: "other-deployment",
buildID: "build-1",
revisionNumber: 7,
},
}
uuidRe := regexp.MustCompile(`^[0-9a-f]{8}-[0-9a-f]{4}-5[0-9a-f]{3}-[0-9a-f]{4}-[0-9a-f]{12}$`)
reqIDsByRevision := map[int64]string{}
for _, tc := range testCases {
d.Run(tc.name, func() {
var capturedReqID string
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
DoAndReturn(func(_ context.Context, req *historyservice.SignalWorkflowExecutionRequest, _ ...grpc.CallOption) (*historyservice.SignalWorkflowExecutionResponse, error) {
capturedReqID = req.GetSignalRequest().GetRequestId()
d.Equal(ReactivateVersionSignalName, req.GetSignalRequest().GetSignalName())
d.Equal(GenerateVersionWorkflowID(tc.deploymentName, tc.buildID), req.GetSignalRequest().GetWorkflowExecution().GetWorkflowId())
return &historyservice.SignalWorkflowExecutionResponse{}, nil
})
err := d.deploymentClient.SignalVersionReactivation(
context.Background(),
d.ns,
tc.deploymentName,
tc.buildID,
tc.revisionNumber,
)
d.NoError(err)
// RequestId must be a UUID v5 (the "5" in the third group marks name-based SHA-1).
d.Regexp(uuidRe, capturedReqID)
// Determinism: identical revisionNumber across calls must produce identical RequestIds,
// regardless of deploymentName/buildID.
if prev, ok := reqIDsByRevision[tc.revisionNumber]; ok {
d.Equal(prev, capturedReqID)
} else {
reqIDsByRevision[tc.revisionNumber] = capturedReqID
}
})
}
}
// TestSignalVersionReactivation_DedupsByRevision verifies the per-pod dedup in
// SignalVersionReactivation: once a revision (or higher) has been signaled for a given
// version workflow, subsequent calls at the same-or-lower revision skip the RPC.
func (d *deploymentWorkflowClientSuite) TestSignalVersionReactivation_DedupsByRevision() {
const (
dep = "my-deployment"
build = "build-1"
)
// First call fires.
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
Return(&historyservice.SignalWorkflowExecutionResponse{}, nil).
Times(1)
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 5))
// Same revision → skip (no mock call expected; gomock fails if called).
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 5))
// Lower revision → skip.
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 3))
// Higher revision → fires again.
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
Return(&historyservice.SignalWorkflowExecutionResponse{}, nil).
Times(1)
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 7))
// Same-or-lower after the higher rev → still skipped.
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 7))
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 5))
}
// TestSignalVersionReactivation_DedupIsolatedPerVersion verifies the dedup key includes
// (namespaceID, deploymentName, buildID) — different version workflows do not share state.
func (d *deploymentWorkflowClientSuite) TestSignalVersionReactivation_DedupIsolatedPerVersion() {
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
Return(&historyservice.SignalWorkflowExecutionResponse{}, nil).
Times(3)
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, "dep-a", "build-1", 5))
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, "dep-b", "build-1", 5))
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, "dep-a", "build-2", 5))
}
// TestSignalVersionReactivation_FailureAllowsRetry verifies that a failed signal does NOT
// record the revision in the cache — a subsequent call for the same revision retries.
func (d *deploymentWorkflowClientSuite) TestSignalVersionReactivation_FailureAllowsRetry() {
const (
dep = "my-deployment"
build = "build-1"
)
// First attempt fails.
failErr := serviceerror.NewUnavailable("downstream unavailable")
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
Return(nil, failErr).
Times(1)
err := d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 5)
d.Error(err)
// Retry at the same revision must hit the wire again (cache not recorded on failure).
d.mockHistoryClient.EXPECT().
SignalWorkflowExecution(gomock.Any(), gomock.Any()).
Return(&historyservice.SignalWorkflowExecutionResponse{}, nil).
Times(1)
d.NoError(d.deploymentClient.SignalVersionReactivation(context.Background(), d.ns, dep, build, 5))
}

View File

@@ -49,7 +49,6 @@ const (
testLongVersionDrainageRefreshInterval = 10 * time.Second
testLongVersionDrainageVisibilityGracePeriod = 10 * time.Second
testVersionMembershipCacheTTL = 5 * time.Second
testLongVersionReactivationCacheTTL = 5 * time.Minute
testMaxVersionsInDeployment = 4
)
@@ -88,9 +87,6 @@ func (s *DeploymentVersionSuite) SetupSuite() {
dynamicconfig.VersionDrainageStatusRefreshInterval.Key(): testVersionDrainageRefreshInterval,
dynamicconfig.VersionDrainageStatusVisibilityGracePeriod.Key(): testVersionDrainageVisibilityGracePeriod,
dynamicconfig.VersionMembershipCacheTTL.Key(): testVersionMembershipCacheTTL,
// Large TTL for deduplication test. Must be set at suite level for cache initialization to work.
dynamicconfig.VersionReactivationSignalCacheTTL.Key(): testLongVersionReactivationCacheTTL,
}))
}
@@ -2213,36 +2209,24 @@ func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetPinnedSet
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_ReactivateVersionOnPinned() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testLongVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
// Include workflow version in deployment name to avoid conflicts in parallel tests
deploymentName := fmt.Sprintf("test-reactivate-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-task-queue") // First version (will become DRAINED)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-task-queue") // Second version (will become CURRENT)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-task-queue") // Pinned target (INACTIVE)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-task-queue") // Current version
// set version 1 as current
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(ctx, tv1)
err := s.setCurrent(tv1, true)
s.NoError(err)
// set version 2 as current which shall cause version 1 to be in 'Draining' state
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(ctx, tv2)
err = s.setCurrent(tv2, true)
err := s.setCurrent(tv2, true)
s.NoError(err)
// Wait for version 1 to become DRAINED
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testLongVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval) // wait for grace period and refresh interval
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
@@ -2250,7 +2234,7 @@ func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_ReactivateVe
}
}
// Register a worker for version 1 (DRAINED) so it can accept workflows when
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// UpdateWorkflowExecutionOptions is called to pin the workflow to version 1.
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
@@ -2289,7 +2273,7 @@ func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_ReactivateVe
}, "waitingWorkflow")
s.NoError(err)
// Pin the running workflow to version 1 (DRAINED) using UpdateWorkflowExecutionOptions.
// Pin the running workflow to version 1 (INACTIVE) using UpdateWorkflowExecutionOptions.
pinnedOverride := &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
@@ -2347,35 +2331,15 @@ func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_ReactivateVe
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnPinned() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testLongVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
// Include workflow version in deployment name to avoid conflicts in parallel tests
deploymentName := fmt.Sprintf("test-start-reactivate-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-start-task-queue") // First version (will become DRAINED)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-start-task-queue") // Second version (will become CURRENT)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-start-task-queue")
// set version 1 as current
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(ctx, tv1)
err := s.setCurrent(tv1, true)
s.NoError(err)
// set version 2 as current which shall cause version 1 to be in 'Draining' state
s.startVersionWorkflow(ctx, tv2)
err = s.setCurrent(tv2, true)
s.NoError(err)
// Wait for version 1 to become DRAINED
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testLongVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval) // wait for grace period and refresh interval
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
@@ -2384,7 +2348,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
}
}
// Register a worker for version 1 (DRAINED) so it can accept workflows when
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// StartWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
@@ -2399,7 +2363,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
s.NoError(w1.Start())
defer w1.Stop()
// Start a new workflow with the pinned override pointing to version 1 (DRAINED).
// Start a new workflow with the pinned override pointing to version 1 (INACTIVE).
wfTV := testvars.New(s)
var run sdkclient.WorkflowRun
s.Eventually(func() bool {
@@ -2429,7 +2393,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
},
})
// Wait for version 1 to show up as DRAINING (reactivated from DRAINED)
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
@@ -2451,9 +2415,6 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnPinned_WithConflictPolicy() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testLongVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
@@ -2461,24 +2422,15 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-conflict-task-queue")
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-conflict-task-queue")
// set version 1 as current
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(ctx, tv1)
err := s.setCurrent(tv1, true)
s.NoError(err)
// set version 2 as current which shall cause version 1 to be in 'Draining' state
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(ctx, tv2)
err = s.setCurrent(tv2, true)
err := s.setCurrent(tv2, true)
s.NoError(err)
// Wait for version 1 to become DRAINED
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testLongVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval)
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
@@ -2486,7 +2438,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
}
}
// Register a worker for version 1 (DRAINED) so it can accept workflows
// Register a worker for version 1 (INACTIVE) so it can accept workflows
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
@@ -2524,7 +2476,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Now start a second workflow with the SAME workflow ID, pinned to v1 (DRAINED),
// Now start a second workflow with the SAME workflow ID, pinned to v1 (INACTIVE),
// using TERMINATE_EXISTING conflict policy. This goes through the handleConflict method in api.go.
var run sdkclient.WorkflowRun
s.Eventually(func() bool {
@@ -2555,7 +2507,7 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
},
})
// Wait for version 1 to show up as DRAINING (reactivated from DRAINED)
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
@@ -2577,35 +2529,15 @@ func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnP
}
func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_ReactivateVersionOnPinned() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testLongVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
// Include workflow version in deployment name to avoid conflicts in parallel tests
deploymentName := fmt.Sprintf("test-sws-reactivate-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-sws-task-queue") // First version (will become DRAINED)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-sws-task-queue") // Second version (will become CURRENT)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-sws-task-queue")
// set version 1 as current
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(ctx, tv1)
err := s.setCurrent(tv1, true)
s.NoError(err)
// set version 2 as current which shall cause version 1 to be in 'Draining' state
s.startVersionWorkflow(ctx, tv2)
err = s.setCurrent(tv2, true)
s.NoError(err)
// Wait for version 1 to become DRAINED
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testLongVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval) // wait for grace period and refresh interval
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
@@ -2614,7 +2546,7 @@ func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_Reactivate
}
}
// Register a worker for version 1 (DRAINED) so it can accept workflows when
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// SignalWithStartWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
@@ -2629,7 +2561,7 @@ func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_Reactivate
s.NoError(w1.Start())
defer w1.Stop()
// Use SignalWithStart with the pinned override pointing to version 1 (DRAINED).
// Use SignalWithStart with the pinned override pointing to version 1 (INACTIVE).
// This should START a new workflow (not signal an existing one) since no workflow exists yet.
wfTV := testvars.New(s)
var run sdkclient.WorkflowRun
@@ -2665,7 +2597,7 @@ func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_Reactivate
},
})
// Wait for version 1 to show up as DRAINING (reactivated from DRAINED)
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
@@ -2687,36 +2619,24 @@ func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_Reactivate
}
func (s *DeploymentVersionSuite) TestResetWorkflowExecution_ReactivateVersionOnPinned() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testLongVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
// Include workflow version in deployment name to avoid conflicts in parallel tests
deploymentName := fmt.Sprintf("test-reset-reactivate-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-reset-task-queue") // First version (will become DRAINED)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-reset-task-queue") // Second version (will become CURRENT)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-reset-task-queue") // Pinned target (INACTIVE)
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-reset-task-queue") // Current version
// set version 1 as current
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(ctx, tv1)
err := s.setCurrent(tv1, true)
s.NoError(err)
// set version 2 as current which shall cause version 1 to be in 'Draining' state
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(ctx, tv2)
err = s.setCurrent(tv2, true)
err := s.setCurrent(tv2, true)
s.NoError(err)
// Wait for version 1 to become DRAINED
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testLongVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval) // wait for grace period and refresh interval
// Workflow that waits for a signal, used for both versions.
// Returns a string indicating which version completed it.
wf := func(version string) func(ctx workflow.Context) (string, error) {
@@ -2726,7 +2646,7 @@ func (s *DeploymentVersionSuite) TestResetWorkflowExecution_ReactivateVersionOnP
}
}
// Register a worker for version 1 (DRAINED) so it can accept workflows when
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// ResetWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
@@ -2793,7 +2713,7 @@ func (s *DeploymentVersionSuite) TestResetWorkflowExecution_ReactivateVersionOnP
}
s.Positive(resetEventID, "Should have found a workflow task complete event")
// Reset the workflow with PostResetOperations containing a versioning override pinned to v1 (which is currently DRAINED)
// Reset the workflow with PostResetOperations containing a versioning override pinned to v1 (which is currently INACTIVE)
var resetResp *workflowservice.ResetWorkflowExecutionResponse
s.Eventually(func() bool {
var resetErr error
@@ -2848,7 +2768,7 @@ func (s *DeploymentVersionSuite) TestResetWorkflowExecution_ReactivateVersionOnP
},
})
// Wait for version 1 to show up as DRAINING (reactivated from DRAINED)
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
@@ -3158,591 +3078,6 @@ func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_WithUnpinn
s.checkDescribeWorkflowAfterOverride(ctx, wf, override)
}
// The following tests verify that the version_reactivation_signal_cache works as intended.
// Setup note: v1 is left INACTIVE (never made current) rather than forced through a full
// CURRENT → DRAINING → DRAINED transition. The cache and the reactivation-signal handler
// both treat INACTIVE and DRAINED identically, and this shaves the drainage wait off every
// test. DRAINED-side coverage lives in the unit test Test_ReactivateVersion_FromDrained.
func (s *DeploymentVersionSuite) TestReactivationSignalCache_Deduplication_StartWorkflow() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
deploymentName := fmt.Sprintf("test-cache-dedup-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-cache-dedup-tq")
s.startVersionWorkflow(ctx, tv1)
// Workflow that waits for a signal before completing
wf := func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done", nil
}
// Register worker for version 1
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// === First workflow run: Should trigger reactivation signal to be sent (cache miss) ===
wfTV1 := testvars.New(s)
var run1 sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run1, startErr = s.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV1.WorkflowID(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
}, "waitingWorkflow")
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Version 1 should transition to DRAINING
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Signal workflow to complete.
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV1.WorkflowID(), run1.GetRunID(), "complete", nil))
// Wait for workflow to complete.
var result string
s.NoError(run1.Get(ctx, &result))
// Wait for version 1 to become Drained again
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval)
// === Second workflow run: Should NOT trigger reactivation signal to be sent (cache hit) ===
wfTV2 := testvars.New(s)
var run2 sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run2, startErr = s.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV2.WorkflowID(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
}, "waitingWorkflow")
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify version stays Drained for several checks (even though there is a workflow running)
// Use Eventually with a counter to check multiple times that the reactivation signal was cached
drainedCheckCount := 0
s.Eventually(func() bool {
resp, err := s.describeVersion(tv1)
s.NoError(err)
s.Equalf(enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus(),
"Version should remain DRAINED because reactivation signal was cached (check %d)", drainedCheckCount)
s.Equalf(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetStatus(),
"Version status should remain DRAINED because reactivation signal was cached (check %d)", drainedCheckCount)
drainedCheckCount++
return drainedCheckCount >= 5
}, 10*time.Second, 1*time.Second)
// Signal the workflow to complete.
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV2.WorkflowID(), run2.GetRunID(), "complete", nil))
s.NoError(run2.Get(ctx, &result))
}
func (s *DeploymentVersionSuite) TestReactivationSignalCache_Deduplication_SignalWithStart() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
deploymentName := fmt.Sprintf("test-sws-cache-dedup-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-sws-cache-dedup-tq")
s.startVersionWorkflow(ctx, tv1)
// Workflow that waits for a signal before completing
wf := func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done", nil
}
// Register worker for version 1
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// === FIRST SIGNAL WITH START: Should trigger reactivation (cache miss) ===
wfTV1 := testvars.New(s)
var run1 sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run1, startErr = s.SdkClient().SignalWithStartWorkflow(ctx,
wfTV1.WorkflowID(),
"start-signal",
nil,
sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
},
"waitingWorkflow",
)
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Version 1 should transition to DRAINING (reactivated)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Signal workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV1.WorkflowID(), run1.GetRunID(), "complete", nil))
// Wait for workflow to complete
var result string
s.NoError(run1.Get(ctx, &result))
// Wait for version 1 to become DRAINED again (workflow completed)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval)
// === SECOND SIGNAL WITH START: Should NOT trigger reactivation (cache hit) ===
wfTV2 := testvars.New(s)
var run2 sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run2, startErr = s.SdkClient().SignalWithStartWorkflow(ctx,
wfTV2.WorkflowID(),
"start-signal",
nil,
sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
},
"waitingWorkflow",
)
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify version stays DRAINED for several checks (workflow is still running)
// Use Eventually with a counter to check multiple times that the signal was cached
drainedCheckCount := 0
s.Eventually(func() bool {
resp, err := s.describeVersion(tv1)
s.NoError(err)
s.Equal(enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus(),
"Version should remain DRAINED because reactivation signal was cached")
s.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetStatus(),
"Version status should remain DRAINED because reactivation signal was cached")
drainedCheckCount++
return drainedCheckCount >= 5
}, 10*time.Second, 1*time.Second)
// Signal the workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV2.WorkflowID(), run2.GetRunID(), "complete", nil))
s.NoError(run2.Get(ctx, &result))
}
func (s *DeploymentVersionSuite) TestReactivationSignalCache_Deduplication_UpdateOptions() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
deploymentName := fmt.Sprintf("test-opts-cache-dedup-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-opts-cache-dedup-tq")
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-opts-cache-dedup-tq")
s.startVersionWorkflow(ctx, tv1)
// Set version 2 as current so workflows can start on it before being pinned/updated to v1
s.startVersionWorkflow(ctx, tv2)
err := s.setCurrent(tv2, true)
s.NoError(err)
// Workflow that waits for a signal before completing
wf := func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done", nil
}
// Register worker for version 1 (INACTIVE)
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Register worker for version 2 (CURRENT) on the same task queue
w2 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv2.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w2.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w2.Start())
defer w2.Stop()
s.waitForPollers(ctx, tv1, tv2)
pinnedOverride := &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
}
// === FIRST UPDATE OPTIONS: Should trigger reactivation (cache miss) ===
// Start workflow on v2 (current version) - waits for signal to complete
wfTV1 := testvars.New(s)
run1, err := s.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV1.WorkflowID(),
}, "waitingWorkflow")
s.NoError(err)
// Pin the workflow to v1 using UpdateWorkflowExecutionOptions
s.Eventually(func() bool {
_, err = s.FrontendClient().UpdateWorkflowExecutionOptions(ctx,
&workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: s.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV1.WorkflowID(),
RunId: run1.GetRunID(),
},
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: pinnedOverride,
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
})
return err == nil
}, 10*time.Second, 500*time.Millisecond)
// Signal workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV1.WorkflowID(), run1.GetRunID(), "complete", nil))
// Wait for workflow to complete
var result string
s.NoError(run1.Get(ctx, &result))
// Version 1 should transition to DRAINING (reactivated)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Wait for version 1 to become DRAINED again (workflow completed)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval)
// === SECOND UPDATE OPTIONS: Should NOT trigger reactivation (cache hit) ===
// Start another workflow on v2 - waits for signal to complete
wfTV2 := testvars.New(s)
run2, err := s.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV2.WorkflowID(),
}, "waitingWorkflow")
s.NoError(err)
// Pin this workflow to v1 using UpdateWorkflowExecutionOptions (should be cached)
s.Eventually(func() bool {
_, err = s.FrontendClient().UpdateWorkflowExecutionOptions(ctx,
&workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: s.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV2.WorkflowID(),
RunId: run2.GetRunID(),
},
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: pinnedOverride,
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
})
return err == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify version stays DRAINED for several checks (workflow is still running)
// Use Eventually with a counter to check multiple times that the signal was cached
drainedCheckCount := 0
s.Eventually(func() bool {
resp, err := s.describeVersion(tv1)
s.NoError(err)
s.Equal(enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus(),
"Version should remain DRAINED because reactivation signal was cached")
s.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetStatus(),
"Version status should remain DRAINED because reactivation signal was cached")
drainedCheckCount++
return drainedCheckCount >= 5
}, 10*time.Second, 1*time.Second)
// Signal the workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV2.WorkflowID(), run2.GetRunID(), "complete", nil))
s.NoError(run2.Get(ctx, &result))
}
// TestReactivationSignalCache_Deduplication_Reset verifies that the version reactivation signal cache
// deduplicates signals when ResetWorkflowExecution is called multiple times with a pinned override
// to a non-current (INACTIVE or DRAINED) version.
func (s *DeploymentVersionSuite) TestReactivationSignalCache_Deduplication_Reset() {
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testVersionDrainageVisibilityGracePeriod)
s.OverrideDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testLongVersionDrainageRefreshInterval)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
// Use shorter, explicit deployment series names to avoid truncation issues
deploymentName := fmt.Sprintf("test-reset-cache-dedup-wfv%d", s.workflowVersion)
tv1 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v1").WithTaskQueue("test-reset-cache-dedup-tq")
tv2 := testvars.New(s).WithDeploymentSeries(deploymentName).WithBuildID(deploymentName + "-v2").WithTaskQueue("test-reset-cache-dedup-tq")
s.startVersionWorkflow(ctx, tv1)
// Set version 2 as current so workflows can start on it before being reset with a pinned override to v1
s.startVersionWorkflow(ctx, tv2)
err := s.setCurrent(tv2, true)
s.NoError(err)
// Workflow that waits for a signal before completing
wf := func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done", nil
}
// Register worker for version 1 (INACTIVE)
w1 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Register worker for version 2 (CURRENT) on the same task queue
w2 := worker.New(s.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv2.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w2.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w2.Start())
defer w2.Stop()
s.waitForPollers(ctx, tv1, tv2)
pinnedOverride := &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
}
// Helper function to start a workflow, wait for task completion, and get reset event ID
startAndGetResetEventID := func(wfID string) (sdkclient.WorkflowRun, int64) {
run, err := s.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfID,
}, "waitingWorkflow")
s.NoError(err)
// Wait for workflow task completion (creates a reset point)
s.Eventually(func() bool {
hist := s.SdkClient().GetWorkflowHistory(ctx, wfID, run.GetRunID(), false, enumspb.HISTORY_EVENT_FILTER_TYPE_ALL_EVENT)
for hist.HasNext() {
event, err := hist.Next()
if err != nil {
return false
}
if event.EventType == enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED {
return true
}
}
return false
}, 10*time.Second, 200*time.Millisecond)
// Find the reset event ID
var resetEventID int64
hist := s.SdkClient().GetWorkflowHistory(ctx, wfID, run.GetRunID(), false, enumspb.HISTORY_EVENT_FILTER_TYPE_ALL_EVENT)
for hist.HasNext() {
event, err := hist.Next()
s.NoError(err)
if event.EventType == enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED {
resetEventID = event.EventId
break
}
}
s.Positive(resetEventID)
return run, resetEventID
}
// === FIRST RESET: Should trigger reactivation (cache miss) ===
wfTV1 := testvars.New(s)
run1, resetEventID1 := startAndGetResetEventID(wfTV1.WorkflowID())
// Reset with pinned override to v1 (DRAINED)
var resetResp1 *workflowservice.ResetWorkflowExecutionResponse
s.Eventually(func() bool {
var resetErr error
resetResp1, resetErr = s.FrontendClient().ResetWorkflowExecution(ctx, &workflowservice.ResetWorkflowExecutionRequest{
Namespace: s.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV1.WorkflowID(),
RunId: run1.GetRunID(),
},
Reason: "testing-reset-cache-dedup-1",
RequestId: uuid.NewString(),
WorkflowTaskFinishEventId: resetEventID1,
PostResetOperations: []*workflowpb.PostResetOperation{
{
Variant: &workflowpb.PostResetOperation_UpdateWorkflowOptions_{
UpdateWorkflowOptions: &workflowpb.PostResetOperation_UpdateWorkflowOptions{
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: pinnedOverride,
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
},
},
},
},
})
return resetErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Signal the reset workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV1.WorkflowID(), resetResp1.RunId, "complete", nil))
// Wait for workflow to complete
resetRun1 := s.SdkClient().GetWorkflow(ctx, wfTV1.WorkflowID(), resetResp1.RunId)
var result string
s.NoError(resetRun1.Get(ctx, &result))
// Version 1 should transition to DRAINING (reactivated)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Wait for version 1 to become DRAINED again (workflow completed)
s.checkVersionDrainageAndVersionStatus(ctx, tv1,
&deploymentpb.VersionDrainageInfo{Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED,
testVersionDrainageVisibilityGracePeriod+testLongVersionDrainageRefreshInterval)
// === SECOND RESET: Should NOT trigger reactivation (cache hit) ===
wfTV2 := testvars.New(s)
run2, resetEventID2 := startAndGetResetEventID(wfTV2.WorkflowID())
// Reset with pinned override to v1 (should be cached)
var resetResp2 *workflowservice.ResetWorkflowExecutionResponse
s.Eventually(func() bool {
var resetErr error
resetResp2, resetErr = s.FrontendClient().ResetWorkflowExecution(ctx, &workflowservice.ResetWorkflowExecutionRequest{
Namespace: s.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV2.WorkflowID(),
RunId: run2.GetRunID(),
},
Reason: "testing-reset-cache-dedup-2",
RequestId: uuid.NewString(),
WorkflowTaskFinishEventId: resetEventID2,
PostResetOperations: []*workflowpb.PostResetOperation{
{
Variant: &workflowpb.PostResetOperation_UpdateWorkflowOptions_{
UpdateWorkflowOptions: &workflowpb.PostResetOperation_UpdateWorkflowOptions{
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: pinnedOverride,
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
},
},
},
},
})
return resetErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify version stays DRAINED for several checks (reset workflow is still running)
// Use Eventually with a counter to check multiple times that the signal was cached
drainedCheckCount := 0
s.Eventually(func() bool {
resp, err := s.describeVersion(tv1)
s.NoError(err)
s.Equal(enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus(),
"Version should remain DRAINED because reactivation signal was cached")
s.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetStatus(),
"Version status should remain DRAINED because reactivation signal was cached")
drainedCheckCount++
return drainedCheckCount >= 5
}, 10*time.Second, 1*time.Second)
// Signal the reset workflow to complete
s.NoError(s.SdkClient().SignalWorkflow(ctx, wfTV2.WorkflowID(), resetResp2.RunId, "complete", nil))
// Wait for workflow to complete
resetRun2 := s.SdkClient().GetWorkflow(ctx, wfTV2.WorkflowID(), resetResp2.RunId)
var result2 string
s.NoError(resetRun2.Get(ctx, &result2))
}
func (s *DeploymentVersionSuite) TestDeleteVersion_ThenRecreateByPolling() {
s.skipBeforeVersion(workerdeployment.VersionDataRevisionNumber)
s.OverrideDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond)