diff --git a/api/deployment/v1/message.go-helpers.pb.go b/api/deployment/v1/message.go-helpers.pb.go index dc4fbf5f10f..47b291a42ae 100644 --- a/api/deployment/v1/message.go-helpers.pb.go +++ b/api/deployment/v1/message.go-helpers.pb.go @@ -2224,3 +2224,40 @@ func (this *DemoteVersionSignalArgs) Equal(that interface{}) bool { return proto.Equal(this, that1) } + +// Marshal an object of type TaskQueueFamilySummary to the protobuf v3 wire format +func (val *TaskQueueFamilySummary) Marshal() ([]byte, error) { + return proto.Marshal(val) +} + +// Unmarshal an object of type TaskQueueFamilySummary from the protobuf v3 wire format +func (val *TaskQueueFamilySummary) Unmarshal(buf []byte) error { + return proto.Unmarshal(buf, val) +} + +// Size returns the size of the object, in bytes, once serialized +func (val *TaskQueueFamilySummary) Size() int { + return proto.Size(val) +} + +// Equal returns whether two TaskQueueFamilySummary values are equivalent by recursively +// comparing the message's fields. +// For more information see the documentation for +// https://pkg.go.dev/google.golang.org/protobuf/proto#Equal +func (this *TaskQueueFamilySummary) Equal(that interface{}) bool { + if that == nil { + return this == nil + } + + var that1 *TaskQueueFamilySummary + switch t := that.(type) { + case *TaskQueueFamilySummary: + that1 = t + case TaskQueueFamilySummary: + that1 = &t + default: + return false + } + + return proto.Equal(this, that1) +} diff --git a/api/deployment/v1/message.pb.go b/api/deployment/v1/message.pb.go index 34c10c21998..2ab90a631c6 100644 --- a/api/deployment/v1/message.pb.go +++ b/api/deployment/v1/message.pb.go @@ -335,8 +335,14 @@ type VersionLocalState struct { ComputeConfig *v12.ComputeConfigSummary `protobuf:"bytes,18,opt,name=compute_config,json=computeConfig,proto3" json:"compute_config,omitempty"` // Cached compute status, updated when WCI signals this version workflow. ComputeStatus *v11.ComputeStatus `protobuf:"bytes,19,opt,name=compute_status,json=computeStatus,proto3" json:"compute_status,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // Active Current or Ramping v2 Version workflows do not publish task queue family Bloom filter snapshots. + // After a successful task queue registration changes their state, they continue as new into v3. The new run + // must publish a snapshot of this Bloom filter so that the Deployment workflow can learn about them and use it + // for validation. Note: this field, once set to true, is carried across every continues-as-new run of the + // version workflow. It only becomes false when the version workflow is deleted and is then later recreated. + TaskQueueFamilySummarySignalSent bool `protobuf:"varint,20,opt,name=task_queue_family_summary_signal_sent,json=taskQueueFamilySummarySignalSent,proto3" json:"task_queue_family_summary_signal_sent,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *VersionLocalState) Reset() { @@ -503,6 +509,13 @@ func (x *VersionLocalState) GetComputeStatus() *v11.ComputeStatus { return nil } +func (x *VersionLocalState) GetTaskQueueFamilySummarySignalSent() bool { + if x != nil { + return x.TaskQueueFamilySummarySignalSent + } + return false +} + // Data specific to a task queue, from the perspective of a worker deployment version. type TaskQueueVersionData struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -869,8 +882,11 @@ type WorkerDeploymentVersionSummary struct { ComputeConfig *v12.ComputeConfigSummary `protobuf:"bytes,13,opt,name=compute_config,json=computeConfig,proto3" json:"compute_config,omitempty"` // Compute status for this version. Synced from the version workflow when WCI signals a status change. ComputeStatus *v11.ComputeStatus `protobuf:"bytes,14,opt,name=compute_status,json=computeStatus,proto3" json:"compute_status,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // Snapshot of registered task queue families published by the Version workflow. + // It is refreshed periodically or after relevant state changes and may lag the Version workflow's exact state. + TaskQueueFamilySummary *TaskQueueFamilySummary `protobuf:"bytes,15,opt,name=task_queue_family_summary,json=taskQueueFamilySummary,proto3" json:"task_queue_family_summary,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *WorkerDeploymentVersionSummary) Reset() { @@ -1002,6 +1018,13 @@ func (x *WorkerDeploymentVersionSummary) GetComputeStatus() *v11.ComputeStatus { return nil } +func (x *WorkerDeploymentVersionSummary) GetTaskQueueFamilySummary() *TaskQueueFamilySummary { + if x != nil { + return x.TaskQueueFamilySummary + } + return nil +} + // used as Worker Deployment Version workflow update input: type RegisterWorkerInVersionArgs struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -3937,6 +3960,78 @@ func (x *DemoteVersionSignalArgs) GetRoutingConfig() *v11.RoutingConfig { return nil } +type TaskQueueFamilySummary struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Exact number of task queue families registered in the Version. This is n when sizing the Bloom filter. + Count int32 `protobuf:"varint,1,opt,name=count,proto3" json:"count,omitempty"` + // Number of bits in the Bloom filter. + BloomFilterSize int64 `protobuf:"varint,2,opt,name=bloom_filter_size,json=bloomFilterSize,proto3" json:"bloom_filter_size,omitempty"` + // Number of hash functions used by the Bloom filter. + BloomFilterHashCount int32 `protobuf:"varint,3,opt,name=bloom_filter_hash_count,json=bloomFilterHashCount,proto3" json:"bloom_filter_hash_count,omitempty"` + // Bloom filter words containing the registered task queue family names. + BloomFilterWords []int64 `protobuf:"varint,4,rep,packed,name=bloom_filter_words,json=bloomFilterWords,proto3" json:"bloom_filter_words,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *TaskQueueFamilySummary) Reset() { + *x = TaskQueueFamilySummary{} + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[60] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *TaskQueueFamilySummary) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*TaskQueueFamilySummary) ProtoMessage() {} + +func (x *TaskQueueFamilySummary) ProtoReflect() protoreflect.Message { + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[60] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use TaskQueueFamilySummary.ProtoReflect.Descriptor instead. +func (*TaskQueueFamilySummary) Descriptor() ([]byte, []int) { + return file_temporal_server_api_deployment_v1_message_proto_rawDescGZIP(), []int{60} +} + +func (x *TaskQueueFamilySummary) GetCount() int32 { + if x != nil { + return x.Count + } + return 0 +} + +func (x *TaskQueueFamilySummary) GetBloomFilterSize() int64 { + if x != nil { + return x.BloomFilterSize + } + return 0 +} + +func (x *TaskQueueFamilySummary) GetBloomFilterHashCount() int32 { + if x != nil { + return x.BloomFilterHashCount + } + return 0 +} + +func (x *TaskQueueFamilySummary) GetBloomFilterWords() []int64 { + if x != nil { + return x.BloomFilterWords + } + return nil +} + type VersionLocalState_TaskQueueFamilyData struct { state protoimpl.MessageState `protogen:"open.v1"` // Key: Task Queue Type @@ -3947,7 +4042,7 @@ type VersionLocalState_TaskQueueFamilyData struct { func (x *VersionLocalState_TaskQueueFamilyData) Reset() { *x = VersionLocalState_TaskQueueFamilyData{} - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[61] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[62] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3959,7 +4054,7 @@ func (x *VersionLocalState_TaskQueueFamilyData) String() string { func (*VersionLocalState_TaskQueueFamilyData) ProtoMessage() {} func (x *VersionLocalState_TaskQueueFamilyData) ProtoReflect() protoreflect.Message { - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[61] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[62] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3993,7 +4088,7 @@ type SyncDeploymentVersionUserDataRequest_SyncUserData struct { func (x *SyncDeploymentVersionUserDataRequest_SyncUserData) Reset() { *x = SyncDeploymentVersionUserDataRequest_SyncUserData{} - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[65] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[66] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4005,7 +4100,7 @@ func (x *SyncDeploymentVersionUserDataRequest_SyncUserData) String() string { func (*SyncDeploymentVersionUserDataRequest_SyncUserData) ProtoMessage() {} func (x *SyncDeploymentVersionUserDataRequest_SyncUserData) ProtoReflect() protoreflect.Message { - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[65] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[66] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4051,7 +4146,7 @@ type CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes struct { func (x *CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes) Reset() { *x = CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes{} - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[71] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[72] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4063,7 +4158,7 @@ func (x *CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes) String() string func (*CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes) ProtoMessage() {} func (x *CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes) ProtoReflect() protoreflect.Message { - mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[71] + mi := &file_temporal_server_api_deployment_v1_message_proto_msgTypes[72] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4106,7 +4201,7 @@ const file_temporal_server_api_deployment_v1_message_proto_rawDesc = "" + "\vupdate_time\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\n" + "updateTime\x12\x18\n" + "\adeleted\x18\x03 \x01(\bR\adeleted\x12L\n" + - "\x06status\x18\x06 \x01(\x0e24.temporal.api.enums.v1.WorkerDeploymentVersionStatusR\x06status\"\x92\x0e\n" + + "\x06status\x18\x06 \x01(\x0e24.temporal.api.enums.v1.WorkerDeploymentVersionStatusR\x06status\"\xe3\x0e\n" + "\x11VersionLocalState\x12T\n" + "\aversion\x18\x01 \x01(\v2:.temporal.server.api.deployment.v1.WorkerDeploymentVersionR\aversion\x12;\n" + "\vcreate_time\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\n" + @@ -4128,7 +4223,8 @@ const file_temporal_server_api_deployment_v1_message_proto_rawDesc = "" + "\x0frevision_number\x18\x0f \x01(\x03R\x0erevisionNumber\x124\n" + "\x16last_modifier_identity\x18\x11 \x01(\tR\x14lastModifierIdentity\x12T\n" + "\x0ecompute_config\x18\x12 \x01(\v2-.temporal.api.compute.v1.ComputeConfigSummaryR\rcomputeConfig\x12P\n" + - "\x0ecompute_status\x18\x13 \x01(\v2).temporal.api.deployment.v1.ComputeStatusR\rcomputeStatus\x1a\x8e\x01\n" + + "\x0ecompute_status\x18\x13 \x01(\v2).temporal.api.deployment.v1.ComputeStatusR\rcomputeStatus\x12O\n" + + "%task_queue_family_summary_signal_sent\x18\x14 \x01(\bR taskQueueFamilySummarySignalSent\x1a\x8e\x01\n" + "\x16TaskQueueFamiliesEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12^\n" + "\x05value\x18\x02 \x01(\v2H.temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyDataR\x05value:\x028\x01\x1a\x88\x02\n" + @@ -4166,7 +4262,7 @@ const file_temporal_server_api_deployment_v1_message_proto_rawDesc = "" + "\x03key\x18\x01 \x01(\tR\x03key\x12M\n" + "\x05value\x18\x02 \x01(\v27.temporal.server.api.deployment.v1.PropagatingRevisionsR\x05value:\x028\x01\"A\n" + "\x14PropagatingRevisions\x12)\n" + - "\x10revision_numbers\x18\x01 \x03(\x03R\x0frevisionNumbers\"\x94\b\n" + + "\x10revision_numbers\x18\x01 \x03(\x03R\x0frevisionNumbers\"\x8a\t\n" + "\x1eWorkerDeploymentVersionSummary\x12\x18\n" + "\aversion\x18\x01 \x01(\tR\aversion\x12;\n" + "\vcreate_time\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\n" + @@ -4183,7 +4279,8 @@ const file_temporal_server_api_deployment_v1_message_proto_rawDesc = "" + " \x01(\x0e24.temporal.api.enums.v1.WorkerDeploymentVersionStatusR\x06status\x12*\n" + "\x11create_request_id\x18\f \x01(\tR\x0fcreateRequestId\x12T\n" + "\x0ecompute_config\x18\r \x01(\v2-.temporal.api.compute.v1.ComputeConfigSummaryR\rcomputeConfig\x12P\n" + - "\x0ecompute_status\x18\x0e \x01(\v2).temporal.api.deployment.v1.ComputeStatusR\rcomputeStatus\"\xa7\x02\n" + + "\x0ecompute_status\x18\x0e \x01(\v2).temporal.api.deployment.v1.ComputeStatusR\rcomputeStatus\x12t\n" + + "\x19task_queue_family_summary\x18\x0f \x01(\v29.temporal.server.api.deployment.v1.TaskQueueFamilySummaryR\x16taskQueueFamilySummary\"\xa7\x02\n" + "\x1bRegisterWorkerInVersionArgs\x12&\n" + "\x0ftask_queue_name\x18\x01 \x01(\tR\rtaskQueueName\x12L\n" + "\x0ftask_queue_type\x18\x02 \x01(\x0e2$.temporal.api.enums.v1.TaskQueueTypeR\rtaskQueueType\x12&\n" + @@ -4406,7 +4503,12 @@ const file_temporal_server_api_deployment_v1_message_proto_rawDesc = "" + "\x19ForceCANVersionSignalArgs\x12[\n" + "\x0eoverride_state\x18\x01 \x01(\v24.temporal.server.api.deployment.v1.VersionLocalStateR\roverrideState\"k\n" + "\x17DemoteVersionSignalArgs\x12P\n" + - "\x0erouting_config\x18\x01 \x01(\v2).temporal.api.deployment.v1.RoutingConfigR\rroutingConfigB4Z2go.temporal.io/server/api/deployment/v1;deploymentb\x06proto3" + "\x0erouting_config\x18\x01 \x01(\v2).temporal.api.deployment.v1.RoutingConfigR\rroutingConfig\"\xbf\x01\n" + + "\x16TaskQueueFamilySummary\x12\x14\n" + + "\x05count\x18\x01 \x01(\x05R\x05count\x12*\n" + + "\x11bloom_filter_size\x18\x02 \x01(\x03R\x0fbloomFilterSize\x125\n" + + "\x17bloom_filter_hash_count\x18\x03 \x01(\x05R\x14bloomFilterHashCount\x12,\n" + + "\x12bloom_filter_words\x18\x04 \x03(\x03R\x10bloomFilterWordsB4Z2go.temporal.io/server/api/deployment/v1;deploymentb\x06proto3" var ( file_temporal_server_api_deployment_v1_message_proto_rawDescOnce sync.Once @@ -4420,7 +4522,7 @@ func file_temporal_server_api_deployment_v1_message_proto_rawDescGZIP() []byte { return file_temporal_server_api_deployment_v1_message_proto_rawDescData } -var file_temporal_server_api_deployment_v1_message_proto_msgTypes = make([]protoimpl.MessageInfo, 75) +var file_temporal_server_api_deployment_v1_message_proto_msgTypes = make([]protoimpl.MessageInfo, 76) var file_temporal_server_api_deployment_v1_message_proto_goTypes = []any{ (*WorkerDeploymentVersion)(nil), // 0: temporal.server.api.deployment.v1.WorkerDeploymentVersion (*DeploymentVersionData)(nil), // 1: temporal.server.api.deployment.v1.DeploymentVersionData @@ -4482,147 +4584,149 @@ var file_temporal_server_api_deployment_v1_message_proto_goTypes = []any{ (*ForceCANDeploymentSignalArgs)(nil), // 57: temporal.server.api.deployment.v1.ForceCANDeploymentSignalArgs (*ForceCANVersionSignalArgs)(nil), // 58: temporal.server.api.deployment.v1.ForceCANVersionSignalArgs (*DemoteVersionSignalArgs)(nil), // 59: temporal.server.api.deployment.v1.DemoteVersionSignalArgs - nil, // 60: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry - (*VersionLocalState_TaskQueueFamilyData)(nil), // 61: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData - nil, // 62: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry - nil, // 63: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry - nil, // 64: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry - (*SyncDeploymentVersionUserDataRequest_SyncUserData)(nil), // 65: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData - nil, // 66: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.TaskQueueMaxVersionsEntry - nil, // 67: temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.TaskQueueMaxVersionsEntry - nil, // 68: temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.TaskQueueMaxVersionsEntry - nil, // 69: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry - nil, // 70: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry - (*CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes)(nil), // 71: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes - nil, // 72: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry - nil, // 73: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry - nil, // 74: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry - (*timestamppb.Timestamp)(nil), // 75: google.protobuf.Timestamp - (v1.WorkerDeploymentVersionStatus)(0), // 76: temporal.api.enums.v1.WorkerDeploymentVersionStatus - (*v11.VersionDrainageInfo)(nil), // 77: temporal.api.deployment.v1.VersionDrainageInfo - (*v11.VersionMetadata)(nil), // 78: temporal.api.deployment.v1.VersionMetadata - (*v12.ComputeConfigSummary)(nil), // 79: temporal.api.compute.v1.ComputeConfigSummary - (*v11.ComputeStatus)(nil), // 80: temporal.api.deployment.v1.ComputeStatus - (*v11.RoutingConfig)(nil), // 81: temporal.api.deployment.v1.RoutingConfig - (v1.VersionDrainageStatus)(0), // 82: temporal.api.enums.v1.VersionDrainageStatus - (v1.TaskQueueType)(0), // 83: temporal.api.enums.v1.TaskQueueType - (*v11.WorkerDeploymentVersionInfo_VersionTaskQueueInfo)(nil), // 84: temporal.api.deployment.v1.WorkerDeploymentVersionInfo.VersionTaskQueueInfo - (*v12.ComputeConfig)(nil), // 85: temporal.api.compute.v1.ComputeConfig - (*v11.WorkerDeploymentInfo_WorkerDeploymentVersionSummary)(nil), // 86: temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - (*v11.WorkerDeploymentVersion)(nil), // 87: temporal.api.deployment.v1.WorkerDeploymentVersion - (*v13.Payload)(nil), // 88: temporal.api.common.v1.Payload - (*v12.ComputeConfigScalingGroup)(nil), // 89: temporal.api.compute.v1.ComputeConfigScalingGroup - (*v12.ComputeConfigScalingGroupUpdate)(nil), // 90: temporal.api.compute.v1.ComputeConfigScalingGroupUpdate + (*TaskQueueFamilySummary)(nil), // 60: temporal.server.api.deployment.v1.TaskQueueFamilySummary + nil, // 61: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry + (*VersionLocalState_TaskQueueFamilyData)(nil), // 62: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData + nil, // 63: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry + nil, // 64: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry + nil, // 65: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry + (*SyncDeploymentVersionUserDataRequest_SyncUserData)(nil), // 66: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData + nil, // 67: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.TaskQueueMaxVersionsEntry + nil, // 68: temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.TaskQueueMaxVersionsEntry + nil, // 69: temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.TaskQueueMaxVersionsEntry + nil, // 70: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry + nil, // 71: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry + (*CheckTaskQueuesHavePollersActivityArgs_TaskQueueTypes)(nil), // 72: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes + nil, // 73: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry + nil, // 74: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry + nil, // 75: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry + (*timestamppb.Timestamp)(nil), // 76: google.protobuf.Timestamp + (v1.WorkerDeploymentVersionStatus)(0), // 77: temporal.api.enums.v1.WorkerDeploymentVersionStatus + (*v11.VersionDrainageInfo)(nil), // 78: temporal.api.deployment.v1.VersionDrainageInfo + (*v11.VersionMetadata)(nil), // 79: temporal.api.deployment.v1.VersionMetadata + (*v12.ComputeConfigSummary)(nil), // 80: temporal.api.compute.v1.ComputeConfigSummary + (*v11.ComputeStatus)(nil), // 81: temporal.api.deployment.v1.ComputeStatus + (*v11.RoutingConfig)(nil), // 82: temporal.api.deployment.v1.RoutingConfig + (v1.VersionDrainageStatus)(0), // 83: temporal.api.enums.v1.VersionDrainageStatus + (v1.TaskQueueType)(0), // 84: temporal.api.enums.v1.TaskQueueType + (*v11.WorkerDeploymentVersionInfo_VersionTaskQueueInfo)(nil), // 85: temporal.api.deployment.v1.WorkerDeploymentVersionInfo.VersionTaskQueueInfo + (*v12.ComputeConfig)(nil), // 86: temporal.api.compute.v1.ComputeConfig + (*v11.WorkerDeploymentInfo_WorkerDeploymentVersionSummary)(nil), // 87: temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + (*v11.WorkerDeploymentVersion)(nil), // 88: temporal.api.deployment.v1.WorkerDeploymentVersion + (*v13.Payload)(nil), // 89: temporal.api.common.v1.Payload + (*v12.ComputeConfigScalingGroup)(nil), // 90: temporal.api.compute.v1.ComputeConfigScalingGroup + (*v12.ComputeConfigScalingGroupUpdate)(nil), // 91: temporal.api.compute.v1.ComputeConfigScalingGroupUpdate } var file_temporal_server_api_deployment_v1_message_proto_depIdxs = []int32{ 0, // 0: temporal.server.api.deployment.v1.DeploymentVersionData.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion - 75, // 1: temporal.server.api.deployment.v1.DeploymentVersionData.routing_update_time:type_name -> google.protobuf.Timestamp - 75, // 2: temporal.server.api.deployment.v1.DeploymentVersionData.current_since_time:type_name -> google.protobuf.Timestamp - 75, // 3: temporal.server.api.deployment.v1.DeploymentVersionData.ramping_since_time:type_name -> google.protobuf.Timestamp - 76, // 4: temporal.server.api.deployment.v1.DeploymentVersionData.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus - 75, // 5: temporal.server.api.deployment.v1.WorkerDeploymentVersionData.update_time:type_name -> google.protobuf.Timestamp - 76, // 6: temporal.server.api.deployment.v1.WorkerDeploymentVersionData.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus + 76, // 1: temporal.server.api.deployment.v1.DeploymentVersionData.routing_update_time:type_name -> google.protobuf.Timestamp + 76, // 2: temporal.server.api.deployment.v1.DeploymentVersionData.current_since_time:type_name -> google.protobuf.Timestamp + 76, // 3: temporal.server.api.deployment.v1.DeploymentVersionData.ramping_since_time:type_name -> google.protobuf.Timestamp + 77, // 4: temporal.server.api.deployment.v1.DeploymentVersionData.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus + 76, // 5: temporal.server.api.deployment.v1.WorkerDeploymentVersionData.update_time:type_name -> google.protobuf.Timestamp + 77, // 6: temporal.server.api.deployment.v1.WorkerDeploymentVersionData.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus 0, // 7: temporal.server.api.deployment.v1.VersionLocalState.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion - 75, // 8: temporal.server.api.deployment.v1.VersionLocalState.create_time:type_name -> google.protobuf.Timestamp - 75, // 9: temporal.server.api.deployment.v1.VersionLocalState.routing_update_time:type_name -> google.protobuf.Timestamp - 75, // 10: temporal.server.api.deployment.v1.VersionLocalState.current_since_time:type_name -> google.protobuf.Timestamp - 75, // 11: temporal.server.api.deployment.v1.VersionLocalState.ramping_since_time:type_name -> google.protobuf.Timestamp - 75, // 12: temporal.server.api.deployment.v1.VersionLocalState.first_activation_time:type_name -> google.protobuf.Timestamp - 75, // 13: temporal.server.api.deployment.v1.VersionLocalState.last_current_time:type_name -> google.protobuf.Timestamp - 75, // 14: temporal.server.api.deployment.v1.VersionLocalState.last_deactivation_time:type_name -> google.protobuf.Timestamp - 77, // 15: temporal.server.api.deployment.v1.VersionLocalState.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo - 78, // 16: temporal.server.api.deployment.v1.VersionLocalState.metadata:type_name -> temporal.api.deployment.v1.VersionMetadata - 60, // 17: temporal.server.api.deployment.v1.VersionLocalState.task_queue_families:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry - 76, // 18: temporal.server.api.deployment.v1.VersionLocalState.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus - 79, // 19: temporal.server.api.deployment.v1.VersionLocalState.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary - 80, // 20: temporal.server.api.deployment.v1.VersionLocalState.compute_status:type_name -> temporal.api.deployment.v1.ComputeStatus + 76, // 8: temporal.server.api.deployment.v1.VersionLocalState.create_time:type_name -> google.protobuf.Timestamp + 76, // 9: temporal.server.api.deployment.v1.VersionLocalState.routing_update_time:type_name -> google.protobuf.Timestamp + 76, // 10: temporal.server.api.deployment.v1.VersionLocalState.current_since_time:type_name -> google.protobuf.Timestamp + 76, // 11: temporal.server.api.deployment.v1.VersionLocalState.ramping_since_time:type_name -> google.protobuf.Timestamp + 76, // 12: temporal.server.api.deployment.v1.VersionLocalState.first_activation_time:type_name -> google.protobuf.Timestamp + 76, // 13: temporal.server.api.deployment.v1.VersionLocalState.last_current_time:type_name -> google.protobuf.Timestamp + 76, // 14: temporal.server.api.deployment.v1.VersionLocalState.last_deactivation_time:type_name -> google.protobuf.Timestamp + 78, // 15: temporal.server.api.deployment.v1.VersionLocalState.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo + 79, // 16: temporal.server.api.deployment.v1.VersionLocalState.metadata:type_name -> temporal.api.deployment.v1.VersionMetadata + 61, // 17: temporal.server.api.deployment.v1.VersionLocalState.task_queue_families:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry + 77, // 18: temporal.server.api.deployment.v1.VersionLocalState.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus + 80, // 19: temporal.server.api.deployment.v1.VersionLocalState.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary + 81, // 20: temporal.server.api.deployment.v1.VersionLocalState.compute_status:type_name -> temporal.api.deployment.v1.ComputeStatus 3, // 21: temporal.server.api.deployment.v1.WorkerDeploymentVersionWorkflowArgs.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState 7, // 22: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowArgs.state:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState - 75, // 23: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.create_time:type_name -> google.protobuf.Timestamp - 81, // 24: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 63, // 25: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.versions:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry - 64, // 26: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.propagating_revisions:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry - 75, // 27: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.create_time:type_name -> google.protobuf.Timestamp - 82, // 28: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.drainage_status:type_name -> temporal.api.enums.v1.VersionDrainageStatus - 77, // 29: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo - 75, // 30: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.routing_update_time:type_name -> google.protobuf.Timestamp - 75, // 31: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.current_since_time:type_name -> google.protobuf.Timestamp - 75, // 32: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.ramping_since_time:type_name -> google.protobuf.Timestamp - 75, // 33: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.first_activation_time:type_name -> google.protobuf.Timestamp - 75, // 34: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.last_current_time:type_name -> google.protobuf.Timestamp - 75, // 35: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.last_deactivation_time:type_name -> google.protobuf.Timestamp - 76, // 36: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus - 79, // 37: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary - 80, // 38: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.compute_status:type_name -> temporal.api.deployment.v1.ComputeStatus - 83, // 39: temporal.server.api.deployment.v1.RegisterWorkerInVersionArgs.task_queue_type:type_name -> temporal.api.enums.v1.TaskQueueType - 81, // 40: temporal.server.api.deployment.v1.RegisterWorkerInVersionArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 83, // 41: temporal.server.api.deployment.v1.RegisterWorkerInWorkerDeploymentArgs.task_queue_type:type_name -> temporal.api.enums.v1.TaskQueueType - 0, // 42: temporal.server.api.deployment.v1.RegisterWorkerInWorkerDeploymentArgs.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion - 84, // 43: temporal.server.api.deployment.v1.DescribeVersionFromWorkerDeploymentActivityResult.task_queue_infos:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersionInfo.VersionTaskQueueInfo - 75, // 44: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.routing_update_time:type_name -> google.protobuf.Timestamp - 75, // 45: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.current_since_time:type_name -> google.protobuf.Timestamp - 75, // 46: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.ramping_since_time:type_name -> google.protobuf.Timestamp - 81, // 47: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 3, // 48: temporal.server.api.deployment.v1.SyncVersionStateResponse.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState - 9, // 49: temporal.server.api.deployment.v1.SyncVersionStateResponse.summary:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary - 75, // 50: temporal.server.api.deployment.v1.AddVersionUpdateArgs.create_time:type_name -> google.protobuf.Timestamp - 77, // 51: temporal.server.api.deployment.v1.SyncDrainageInfoSignalArgs.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo - 82, // 52: temporal.server.api.deployment.v1.SyncDrainageStatusSignalArgs.drainage_status:type_name -> temporal.api.enums.v1.VersionDrainageStatus - 3, // 53: temporal.server.api.deployment.v1.QueryDescribeVersionResponse.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState - 7, // 54: temporal.server.api.deployment.v1.QueryDescribeWorkerDeploymentResponse.state:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState - 79, // 55: temporal.server.api.deployment.v1.StartWorkerDeploymentVersionRequest.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary - 0, // 56: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion - 65, // 57: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.sync:type_name -> temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData - 81, // 58: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.update_routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 2, // 59: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.upsert_version_data:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionData - 66, // 60: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.TaskQueueMaxVersionsEntry - 67, // 61: temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.TaskQueueMaxVersionsEntry - 14, // 62: temporal.server.api.deployment.v1.SyncUnversionedRampActivityArgs.update_args:type_name -> temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs - 68, // 63: temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.TaskQueueMaxVersionsEntry - 69, // 64: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.upsert_entries:type_name -> temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry - 78, // 65: temporal.server.api.deployment.v1.UpdateVersionMetadataResponse.metadata:type_name -> temporal.api.deployment.v1.VersionMetadata - 85, // 66: temporal.server.api.deployment.v1.CreateWorkerDeploymentVersionArgs.compute_config:type_name -> temporal.api.compute.v1.ComputeConfig - 70, // 67: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.task_queues_and_types:type_name -> temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry - 0, // 68: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.worker_deployment_version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion - 14, // 69: temporal.server.api.deployment.v1.SyncVersionStateActivityArgs.update_args:type_name -> temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs - 3, // 70: temporal.server.api.deployment.v1.SyncVersionStateActivityResult.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState - 9, // 71: temporal.server.api.deployment.v1.SyncVersionStateActivityResult.summary:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary - 75, // 72: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.create_time:type_name -> google.protobuf.Timestamp - 81, // 73: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 86, // 74: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.latest_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 86, // 75: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.current_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 86, // 76: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.ramping_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 75, // 77: temporal.server.api.deployment.v1.WorkerDeploymentSummary.create_time:type_name -> google.protobuf.Timestamp - 81, // 78: temporal.server.api.deployment.v1.WorkerDeploymentSummary.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 86, // 79: temporal.server.api.deployment.v1.WorkerDeploymentSummary.latest_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 86, // 80: temporal.server.api.deployment.v1.WorkerDeploymentSummary.current_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 86, // 81: temporal.server.api.deployment.v1.WorkerDeploymentSummary.ramping_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary - 72, // 82: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.scaling_groups:type_name -> temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry - 87, // 83: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion - 73, // 84: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.upsert_scaling_groups:type_name -> temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry - 87, // 85: temporal.server.api.deployment.v1.DeleteWorkerControllerInstanceInput.version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion - 74, // 86: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.upsert_scaling_groups:type_name -> temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry - 7, // 87: temporal.server.api.deployment.v1.ForceCANDeploymentSignalArgs.override_state:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState - 3, // 88: temporal.server.api.deployment.v1.ForceCANVersionSignalArgs.override_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState - 81, // 89: temporal.server.api.deployment.v1.DemoteVersionSignalArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig - 61, // 90: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry.value:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData - 62, // 91: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.task_queues:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry - 4, // 92: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry.value:type_name -> temporal.server.api.deployment.v1.TaskQueueVersionData - 9, // 93: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry.value:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary - 8, // 94: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry.value:type_name -> temporal.server.api.deployment.v1.PropagatingRevisions - 83, // 95: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData.types:type_name -> temporal.api.enums.v1.TaskQueueType - 1, // 96: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData.data:type_name -> temporal.server.api.deployment.v1.DeploymentVersionData - 88, // 97: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry.value:type_name -> temporal.api.common.v1.Payload - 71, // 98: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry.value:type_name -> temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes - 83, // 99: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes.types:type_name -> temporal.api.enums.v1.TaskQueueType - 89, // 100: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroup - 90, // 101: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroupUpdate - 90, // 102: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroupUpdate - 103, // [103:103] is the sub-list for method output_type - 103, // [103:103] is the sub-list for method input_type - 103, // [103:103] is the sub-list for extension type_name - 103, // [103:103] is the sub-list for extension extendee - 0, // [0:103] is the sub-list for field type_name + 76, // 23: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.create_time:type_name -> google.protobuf.Timestamp + 82, // 24: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 64, // 25: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.versions:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry + 65, // 26: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.propagating_revisions:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry + 76, // 27: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.create_time:type_name -> google.protobuf.Timestamp + 83, // 28: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.drainage_status:type_name -> temporal.api.enums.v1.VersionDrainageStatus + 78, // 29: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo + 76, // 30: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.routing_update_time:type_name -> google.protobuf.Timestamp + 76, // 31: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.current_since_time:type_name -> google.protobuf.Timestamp + 76, // 32: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.ramping_since_time:type_name -> google.protobuf.Timestamp + 76, // 33: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.first_activation_time:type_name -> google.protobuf.Timestamp + 76, // 34: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.last_current_time:type_name -> google.protobuf.Timestamp + 76, // 35: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.last_deactivation_time:type_name -> google.protobuf.Timestamp + 77, // 36: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.status:type_name -> temporal.api.enums.v1.WorkerDeploymentVersionStatus + 80, // 37: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary + 81, // 38: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.compute_status:type_name -> temporal.api.deployment.v1.ComputeStatus + 60, // 39: temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary.task_queue_family_summary:type_name -> temporal.server.api.deployment.v1.TaskQueueFamilySummary + 84, // 40: temporal.server.api.deployment.v1.RegisterWorkerInVersionArgs.task_queue_type:type_name -> temporal.api.enums.v1.TaskQueueType + 82, // 41: temporal.server.api.deployment.v1.RegisterWorkerInVersionArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 84, // 42: temporal.server.api.deployment.v1.RegisterWorkerInWorkerDeploymentArgs.task_queue_type:type_name -> temporal.api.enums.v1.TaskQueueType + 0, // 43: temporal.server.api.deployment.v1.RegisterWorkerInWorkerDeploymentArgs.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion + 85, // 44: temporal.server.api.deployment.v1.DescribeVersionFromWorkerDeploymentActivityResult.task_queue_infos:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersionInfo.VersionTaskQueueInfo + 76, // 45: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.routing_update_time:type_name -> google.protobuf.Timestamp + 76, // 46: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.current_since_time:type_name -> google.protobuf.Timestamp + 76, // 47: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.ramping_since_time:type_name -> google.protobuf.Timestamp + 82, // 48: temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 3, // 49: temporal.server.api.deployment.v1.SyncVersionStateResponse.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState + 9, // 50: temporal.server.api.deployment.v1.SyncVersionStateResponse.summary:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary + 76, // 51: temporal.server.api.deployment.v1.AddVersionUpdateArgs.create_time:type_name -> google.protobuf.Timestamp + 78, // 52: temporal.server.api.deployment.v1.SyncDrainageInfoSignalArgs.drainage_info:type_name -> temporal.api.deployment.v1.VersionDrainageInfo + 83, // 53: temporal.server.api.deployment.v1.SyncDrainageStatusSignalArgs.drainage_status:type_name -> temporal.api.enums.v1.VersionDrainageStatus + 3, // 54: temporal.server.api.deployment.v1.QueryDescribeVersionResponse.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState + 7, // 55: temporal.server.api.deployment.v1.QueryDescribeWorkerDeploymentResponse.state:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState + 80, // 56: temporal.server.api.deployment.v1.StartWorkerDeploymentVersionRequest.compute_config:type_name -> temporal.api.compute.v1.ComputeConfigSummary + 0, // 57: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion + 66, // 58: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.sync:type_name -> temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData + 82, // 59: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.update_routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 2, // 60: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.upsert_version_data:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionData + 67, // 61: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataResponse.TaskQueueMaxVersionsEntry + 68, // 62: temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.CheckWorkerDeploymentUserDataPropagationRequest.TaskQueueMaxVersionsEntry + 14, // 63: temporal.server.api.deployment.v1.SyncUnversionedRampActivityArgs.update_args:type_name -> temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs + 69, // 64: temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.task_queue_max_versions:type_name -> temporal.server.api.deployment.v1.SyncUnversionedRampActivityResponse.TaskQueueMaxVersionsEntry + 70, // 65: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.upsert_entries:type_name -> temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry + 79, // 66: temporal.server.api.deployment.v1.UpdateVersionMetadataResponse.metadata:type_name -> temporal.api.deployment.v1.VersionMetadata + 86, // 67: temporal.server.api.deployment.v1.CreateWorkerDeploymentVersionArgs.compute_config:type_name -> temporal.api.compute.v1.ComputeConfig + 71, // 68: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.task_queues_and_types:type_name -> temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry + 0, // 69: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.worker_deployment_version:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersion + 14, // 70: temporal.server.api.deployment.v1.SyncVersionStateActivityArgs.update_args:type_name -> temporal.server.api.deployment.v1.SyncVersionStateUpdateArgs + 3, // 71: temporal.server.api.deployment.v1.SyncVersionStateActivityResult.version_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState + 9, // 72: temporal.server.api.deployment.v1.SyncVersionStateActivityResult.summary:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary + 76, // 73: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.create_time:type_name -> google.protobuf.Timestamp + 82, // 74: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 87, // 75: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.latest_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 87, // 76: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.current_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 87, // 77: temporal.server.api.deployment.v1.WorkerDeploymentWorkflowMemo.ramping_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 76, // 78: temporal.server.api.deployment.v1.WorkerDeploymentSummary.create_time:type_name -> google.protobuf.Timestamp + 82, // 79: temporal.server.api.deployment.v1.WorkerDeploymentSummary.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 87, // 80: temporal.server.api.deployment.v1.WorkerDeploymentSummary.latest_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 87, // 81: temporal.server.api.deployment.v1.WorkerDeploymentSummary.current_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 87, // 82: temporal.server.api.deployment.v1.WorkerDeploymentSummary.ramping_version_summary:type_name -> temporal.api.deployment.v1.WorkerDeploymentInfo.WorkerDeploymentVersionSummary + 73, // 83: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.scaling_groups:type_name -> temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry + 88, // 84: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion + 74, // 85: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.upsert_scaling_groups:type_name -> temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry + 88, // 86: temporal.server.api.deployment.v1.DeleteWorkerControllerInstanceInput.version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion + 75, // 87: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.upsert_scaling_groups:type_name -> temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry + 7, // 88: temporal.server.api.deployment.v1.ForceCANDeploymentSignalArgs.override_state:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentLocalState + 3, // 89: temporal.server.api.deployment.v1.ForceCANVersionSignalArgs.override_state:type_name -> temporal.server.api.deployment.v1.VersionLocalState + 82, // 90: temporal.server.api.deployment.v1.DemoteVersionSignalArgs.routing_config:type_name -> temporal.api.deployment.v1.RoutingConfig + 62, // 91: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamiliesEntry.value:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData + 63, // 92: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.task_queues:type_name -> temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry + 4, // 93: temporal.server.api.deployment.v1.VersionLocalState.TaskQueueFamilyData.TaskQueuesEntry.value:type_name -> temporal.server.api.deployment.v1.TaskQueueVersionData + 9, // 94: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.VersionsEntry.value:type_name -> temporal.server.api.deployment.v1.WorkerDeploymentVersionSummary + 8, // 95: temporal.server.api.deployment.v1.WorkerDeploymentLocalState.PropagatingRevisionsEntry.value:type_name -> temporal.server.api.deployment.v1.PropagatingRevisions + 84, // 96: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData.types:type_name -> temporal.api.enums.v1.TaskQueueType + 1, // 97: temporal.server.api.deployment.v1.SyncDeploymentVersionUserDataRequest.SyncUserData.data:type_name -> temporal.server.api.deployment.v1.DeploymentVersionData + 89, // 98: temporal.server.api.deployment.v1.UpdateVersionMetadataArgs.UpsertEntriesEntry.value:type_name -> temporal.api.common.v1.Payload + 72, // 99: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueuesAndTypesEntry.value:type_name -> temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes + 84, // 100: temporal.server.api.deployment.v1.CheckTaskQueuesHavePollersActivityArgs.TaskQueueTypes.types:type_name -> temporal.api.enums.v1.TaskQueueType + 90, // 101: temporal.server.api.deployment.v1.ValidateWorkerControllerInstanceSpecInput.ScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroup + 91, // 102: temporal.server.api.deployment.v1.UpdateWorkerControllerInstanceInput.UpsertScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroupUpdate + 91, // 103: temporal.server.api.deployment.v1.UpdateComputeConfigArgs.UpsertScalingGroupsEntry.value:type_name -> temporal.api.compute.v1.ComputeConfigScalingGroupUpdate + 104, // [104:104] is the sub-list for method output_type + 104, // [104:104] is the sub-list for method input_type + 104, // [104:104] is the sub-list for extension type_name + 104, // [104:104] is the sub-list for extension extendee + 0, // [0:104] is the sub-list for field type_name } func init() { file_temporal_server_api_deployment_v1_message_proto_init() } @@ -4636,7 +4740,7 @@ func file_temporal_server_api_deployment_v1_message_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_temporal_server_api_deployment_v1_message_proto_rawDesc), len(file_temporal_server_api_deployment_v1_message_proto_rawDesc)), NumEnums: 0, - NumMessages: 75, + NumMessages: 76, NumExtensions: 0, NumServices: 0, }, diff --git a/common/metrics/metric_defs.go b/common/metrics/metric_defs.go index b2a1298bbdb..c5ebbbfb103 100644 --- a/common/metrics/metric_defs.go +++ b/common/metrics/metric_defs.go @@ -1676,6 +1676,7 @@ var ( WorkerDeploymentVersioningOverrideCounter = NewCounterDef("worker_deployment_versioning_override_count") WorkerDeploymentVersioningOneTimeOverrideCounter = NewCounterDef("worker_deployment_versioning_one_time_override_count") WorkerDeploymentVersionDeletePropagationFailure = NewCounterDef("worker_deployment_version_delete_propagation_failure") + WorkerDeploymentTaskQueueFamilyBloomFilterOutcome = NewCounterDef("worker_deployment_task_queue_family_bloom_filter_outcome") StartDeploymentTransitionCounter = NewCounterDef("start_deployment_transition_count") VersioningDataPropagationLatency = NewTimerDef("versioning_data_propagation_latency") SlowVersioningDataPropagationCounter = NewCounterDef("slow_versioning_data_propagation") diff --git a/go.mod b/go.mod index 55125f0e90f..70fab6516e5 100644 --- a/go.mod +++ b/go.mod @@ -16,6 +16,7 @@ require ( github.com/aws/aws-sdk-go-v2/credentials v1.19.15 github.com/aws/aws-sdk-go-v2/service/s3 v1.99.1 github.com/aws/smithy-go v1.25.0 + github.com/bits-and-blooms/bloom/v3 v3.7.1 github.com/blang/semver/v4 v4.0.0 github.com/cactus/go-statsd-client/v5 v5.1.0 github.com/caio/go-tdigest/v5 v5.0.0 @@ -90,6 +91,7 @@ require ( require ( github.com/aclements/go-moremath v0.0.0-20210112150236-f10218a38794 // indirect + github.com/bits-and-blooms/bitset v1.24.2 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/clipperhouse/uax29/v2 v2.7.0 // indirect github.com/go-openapi/swag/cmdutils v0.26.0 // indirect diff --git a/go.sum b/go.sum index 1d241d3c204..cc8c11fee5c 100644 --- a/go.sum +++ b/go.sum @@ -99,6 +99,10 @@ github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/bitly/go-hostpool v0.0.0-20171023180738-a3a6125de932 h1:mXoPYz/Ul5HYEDvkta6I8/rnYM5gSdSV2tJ6XbZuEtY= github.com/bitly/go-hostpool v0.0.0-20171023180738-a3a6125de932/go.mod h1:NOuUCSz6Q9T7+igc/hlvDOUdtWKryOrtFyIVABv/p7k= +github.com/bits-and-blooms/bitset v1.24.2 h1:M7/NzVbsytmtfHbumG+K2bremQPMJuqv1JD3vOaFxp0= +github.com/bits-and-blooms/bitset v1.24.2/go.mod h1:7hO7Gc7Pp1vODcmWvKMRA9BNmbv6a/7QIWpPxHddWR8= +github.com/bits-and-blooms/bloom/v3 v3.7.1 h1:WXovk4TRKZttAMJfoQx6K2DM0zNIt8w+c67UqO+etV0= +github.com/bits-and-blooms/bloom/v3 v3.7.1/go.mod h1:rZzYLLje2dfzXfAkJNxQQHsKurAyK55KUnL43Euk0hU= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= github.com/blang/semver/v4 v4.0.0/go.mod h1:IbckMUScFkM3pff0VJDNKRiT6TG/YpiHIM2yvyW5YoQ= github.com/bmizerany/assert v0.0.0-20160611221934-b7ed37b82869 h1:DDGfHa7BWjL4YnC6+E63dPcxHo2sUxDIu8g3QgEJdRY= diff --git a/proto/internal/temporal/server/api/deployment/v1/message.proto b/proto/internal/temporal/server/api/deployment/v1/message.proto index bfe3be0429c..44a64b37d73 100644 --- a/proto/internal/temporal/server/api/deployment/v1/message.proto +++ b/proto/internal/temporal/server/api/deployment/v1/message.proto @@ -153,6 +153,13 @@ message VersionLocalState { // Cached compute status, updated when WCI signals this version workflow. temporal.api.deployment.v1.ComputeStatus compute_status = 19; + + // Active Current or Ramping v2 Version workflows do not publish task queue family Bloom filter snapshots. + // After a successful task queue registration changes their state, they continue as new into v3. The new run + // must publish a snapshot of this Bloom filter so that the Deployment workflow can learn about them and use it + // for validation. Note: this field, once set to true, is carried across every continues-as-new run of the + // version workflow. It only becomes false when the version workflow is deleted and is then later recreated. + bool task_queue_family_summary_signal_sent = 20; } // Data specific to a task queue, from the perspective of a worker deployment version. @@ -239,6 +246,10 @@ message WorkerDeploymentVersionSummary { // Compute status for this version. Synced from the version workflow when WCI signals a status change. temporal.api.deployment.v1.ComputeStatus compute_status = 14; + + // Snapshot of registered task queue families published by the Version workflow. + // It is refreshed periodically or after relevant state changes and may lag the Version workflow's exact state. + TaskQueueFamilySummary task_queue_family_summary = 15; } // used as Worker Deployment Version workflow update input: @@ -625,3 +636,14 @@ message ForceCANVersionSignalArgs { message DemoteVersionSignalArgs { temporal.api.deployment.v1.RoutingConfig routing_config = 1; } + +message TaskQueueFamilySummary { + // Exact number of task queue families registered in the Version. This is n when sizing the Bloom filter. + int32 count = 1; + // Number of bits in the Bloom filter. + int64 bloom_filter_size = 2; + // Number of hash functions used by the Bloom filter. + int32 bloom_filter_hash_count = 3; + // Bloom filter words containing the registered task queue family names. + repeated int64 bloom_filter_words = 4; +} diff --git a/service/worker/workerdeployment/fx.go b/service/worker/workerdeployment/fx.go index e1040d567ae..4e61c4d37ca 100644 --- a/service/worker/workerdeployment/fx.go +++ b/service/worker/workerdeployment/fx.go @@ -33,6 +33,8 @@ const ( AsyncSetCurrentAndRamping // Version Data has its own revision number with TaskQueue registration being async as well VersionDataRevisionNumber + // Version summaries include task queue family membership information. + TaskQueueFamilySummary ) type ( diff --git a/service/worker/workerdeployment/replaytester/replay_test.go b/service/worker/workerdeployment/replaytester/replay_test.go index 8cb66270922..6a9955f3e82 100644 --- a/service/worker/workerdeployment/replaytester/replay_test.go +++ b/service/worker/workerdeployment/replaytester/replay_test.go @@ -53,7 +53,7 @@ func TestReplays(t *testing.T) { func testReplays(t *testing.T, versionDemotionSignalEnabled bool) { // For each workflow implementation version we run all the replay tests for snapshots created by that version or older versions - for wv := workerdeployment.InitialVersion; wv <= workerdeployment.VersionDataRevisionNumber; wv++ { + for wv := workerdeployment.InitialVersion; wv <= workerdeployment.TaskQueueFamilySummary; wv++ { replayer := worker.NewWorkflowReplayer() // Create version workflow wrapper to match production registration diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/expected_counts.txt b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/expected_counts.txt new file mode 100644 index 00000000000..84fe614ee17 --- /dev/null +++ b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/expected_counts.txt @@ -0,0 +1,6 @@ +# Expected workflow counts for replay testing +# Generated by generate_history.sh on Mon Sep 7 20:13:49 EDT 2026 +EXPECTED_DEPLOYMENT_WORKFLOWS=25 +EXPECTED_VERSION_WORKFLOWS=14 +ACTUAL_DEPLOYMENT_WORKFLOWS=25 +ACTUAL_VERSION_WORKFLOWS=14 diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-1956-7a9a-8c0d-c83073f4ba70.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-1956-7a9a-8c0d-c83073f4ba70.json.gz new file mode 100644 index 00000000000..a6ccd4cbad2 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-1956-7a9a-8c0d-c83073f4ba70.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-7f7a-73d8-8c3d-4a3eec8995eb.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-7f7a-73d8-8c3d-4a3eec8995eb.json.gz new file mode 100644 index 00000000000..6eaf98c8a68 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_01a07e5c-7f7a-73d8-8c3d-4a3eec8995eb.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_0c9740ea-747f-422a-a799-7775d9a80a1a.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_0c9740ea-747f-422a-a799-7775d9a80a1a.json.gz new file mode 100644 index 00000000000..4a35b7e9387 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_0c9740ea-747f-422a-a799-7775d9a80a1a.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_1b67e3dd-4868-41c5-ab1f-b81d63166b5a.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_1b67e3dd-4868-41c5-ab1f-b81d63166b5a.json.gz new file mode 100644 index 00000000000..30be5a15dc4 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_1b67e3dd-4868-41c5-ab1f-b81d63166b5a.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_7c164999-cf9d-44e3-9801-d5fb39916aa8.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_7c164999-cf9d-44e3-9801-d5fb39916aa8.json.gz new file mode 100644 index 00000000000..d6ae790e58e Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_7c164999-cf9d-44e3-9801-d5fb39916aa8.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_939beb07-7568-411f-b109-b06af4345447.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_939beb07-7568-411f-b109-b06af4345447.json.gz new file mode 100644 index 00000000000..d5cf4e5d8a1 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_939beb07-7568-411f-b109-b06af4345447.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_a8839567-2c49-48a7-beee-0c7cb9dd2305.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_a8839567-2c49-48a7-beee-0c7cb9dd2305.json.gz new file mode 100644 index 00000000000..df60b631d53 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_a8839567-2c49-48a7-beee-0c7cb9dd2305.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_abfbef38-048b-4147-a2e0-3ea821476c6b.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_abfbef38-048b-4147-a2e0-3ea821476c6b.json.gz new file mode 100644 index 00000000000..c0d1189ba19 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_abfbef38-048b-4147-a2e0-3ea821476c6b.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b0ec7662-5281-4349-81dc-3a27aff3f2d1.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b0ec7662-5281-4349-81dc-3a27aff3f2d1.json.gz new file mode 100644 index 00000000000..36124a8ea77 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b0ec7662-5281-4349-81dc-3a27aff3f2d1.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b1db4eea-3907-4729-99cf-c26a417f1a36.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b1db4eea-3907-4729-99cf-c26a417f1a36.json.gz new file mode 100644 index 00000000000..0601bb41dbf Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_b1db4eea-3907-4729-99cf-c26a417f1a36.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_bf7dab99-6968-4f71-a574-c85953d0fd0b.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_bf7dab99-6968-4f71-a574-c85953d0fd0b.json.gz new file mode 100644 index 00000000000..bce64fa41d2 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_bf7dab99-6968-4f71-a574-c85953d0fd0b.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_e827c05d-e583-4c76-8b08-703ceccc09e0.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_e827c05d-e583-4c76-8b08-703ceccc09e0.json.gz new file mode 100644 index 00000000000..950471cf035 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_e827c05d-e583-4c76-8b08-703ceccc09e0.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f23ebfaa-51df-424d-a96d-3951799d8b94.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f23ebfaa-51df-424d-a96d-3951799d8b94.json.gz new file mode 100644 index 00000000000..3ae48511cdb Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f23ebfaa-51df-424d-a96d-3951799d8b94.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f640badf-a2f2-40ec-a3a6-245bf570b154.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f640badf-a2f2-40ec-a3a6-245bf570b154.json.gz new file mode 100644 index 00000000000..356054cc212 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_version_wf_run_f640badf-a2f2-40ec-a3a6-245bf570b154.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01a07e5c-194f-740f-a8eb-8cca71c1309b.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01a07e5c-194f-740f-a8eb-8cca71c1309b.json.gz new file mode 100644 index 00000000000..897b0e9170d Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01a07e5c-194f-740f-a8eb-8cca71c1309b.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01dc1d46-f82e-4d2d-9b8e-32127563d042.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01dc1d46-f82e-4d2d-9b8e-32127563d042.json.gz new file mode 100644 index 00000000000..a9fcefdadc4 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_01dc1d46-f82e-4d2d-9b8e-32127563d042.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0ae7a908-9947-44c6-98e1-06c9baf2c5c2.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0ae7a908-9947-44c6-98e1-06c9baf2c5c2.json.gz new file mode 100644 index 00000000000..d1b13aa82d3 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0ae7a908-9947-44c6-98e1-06c9baf2c5c2.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0e61ba91-58e4-4bdc-855a-035359f88f34.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0e61ba91-58e4-4bdc-855a-035359f88f34.json.gz new file mode 100644 index 00000000000..81071da8652 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_0e61ba91-58e4-4bdc-855a-035359f88f34.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_105f20ca-ab55-4e36-b930-ff9e8b25d5a2.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_105f20ca-ab55-4e36-b930-ff9e8b25d5a2.json.gz new file mode 100644 index 00000000000..27700d1ee47 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_105f20ca-ab55-4e36-b930-ff9e8b25d5a2.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_1714d86b-6741-4c04-8978-c47a617ab8e5.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_1714d86b-6741-4c04-8978-c47a617ab8e5.json.gz new file mode 100644 index 00000000000..fcf2d23fe32 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_1714d86b-6741-4c04-8978-c47a617ab8e5.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_24a003fe-052e-464c-84f9-969348d638ee.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_24a003fe-052e-464c-84f9-969348d638ee.json.gz new file mode 100644 index 00000000000..b01b1006ffe Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_24a003fe-052e-464c-84f9-969348d638ee.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_3fa5c2da-f555-475b-8a2e-dd9dff6f9801.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_3fa5c2da-f555-475b-8a2e-dd9dff6f9801.json.gz new file mode 100644 index 00000000000..d01cf32bbfe Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_3fa5c2da-f555-475b-8a2e-dd9dff6f9801.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_5de0c438-cb25-474d-bfd0-96e753b23e00.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_5de0c438-cb25-474d-bfd0-96e753b23e00.json.gz new file mode 100644 index 00000000000..41b4271cd1d Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_5de0c438-cb25-474d-bfd0-96e753b23e00.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_60d32118-bb14-4bb7-87e9-c00e5d497c28.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_60d32118-bb14-4bb7-87e9-c00e5d497c28.json.gz new file mode 100644 index 00000000000..7bb177f8122 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_60d32118-bb14-4bb7-87e9-c00e5d497c28.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_6549f32d-1b2f-4c34-b65e-b3af905de594.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_6549f32d-1b2f-4c34-b65e-b3af905de594.json.gz new file mode 100644 index 00000000000..24c39386553 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_6549f32d-1b2f-4c34-b65e-b3af905de594.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_7b4f797a-451a-4c04-aa1e-5284f6461539.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_7b4f797a-451a-4c04-aa1e-5284f6461539.json.gz new file mode 100644 index 00000000000..73a13f9007d Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_7b4f797a-451a-4c04-aa1e-5284f6461539.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_98f77bcd-5d0a-4fef-837d-1d8aea483a53.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_98f77bcd-5d0a-4fef-837d-1d8aea483a53.json.gz new file mode 100644 index 00000000000..26b2af29808 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_98f77bcd-5d0a-4fef-837d-1d8aea483a53.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_9a1fda0f-8d27-42c4-b0b5-4ed169623a74.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_9a1fda0f-8d27-42c4-b0b5-4ed169623a74.json.gz new file mode 100644 index 00000000000..05d841c9b1d Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_9a1fda0f-8d27-42c4-b0b5-4ed169623a74.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_b1a626bc-7d38-499f-bfe5-b5ce6786bb67.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_b1a626bc-7d38-499f-bfe5-b5ce6786bb67.json.gz new file mode 100644 index 00000000000..47a0afe91c2 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_b1a626bc-7d38-499f-bfe5-b5ce6786bb67.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_bf3c3466-4c69-48ab-aba8-7f1404c81b74.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_bf3c3466-4c69-48ab-aba8-7f1404c81b74.json.gz new file mode 100644 index 00000000000..3e6b6c30883 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_bf3c3466-4c69-48ab-aba8-7f1404c81b74.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c6f0fdc2-dfec-4142-8e6c-2f66698cd776.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c6f0fdc2-dfec-4142-8e6c-2f66698cd776.json.gz new file mode 100644 index 00000000000..7550d822758 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c6f0fdc2-dfec-4142-8e6c-2f66698cd776.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c928d0cd-f57b-4aaf-b212-6b7bbebff4c9.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c928d0cd-f57b-4aaf-b212-6b7bbebff4c9.json.gz new file mode 100644 index 00000000000..c9aaff57c87 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_c928d0cd-f57b-4aaf-b212-6b7bbebff4c9.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_d0b9a06c-86bd-4640-9bcd-56accef8e882.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_d0b9a06c-86bd-4640-9bcd-56accef8e882.json.gz new file mode 100644 index 00000000000..0803fa333ec Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_d0b9a06c-86bd-4640-9bcd-56accef8e882.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ddc2caa0-6da1-49ee-a6d6-18911fa23fec.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ddc2caa0-6da1-49ee-a6d6-18911fa23fec.json.gz new file mode 100644 index 00000000000..8c370a8dac3 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ddc2caa0-6da1-49ee-a6d6-18911fa23fec.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e17c2e23-c43d-4520-b37f-e6e47119f0c9.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e17c2e23-c43d-4520-b37f-e6e47119f0c9.json.gz new file mode 100644 index 00000000000..126350c963e Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e17c2e23-c43d-4520-b37f-e6e47119f0c9.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e8045edd-283f-49a0-bb43-e35468f0b8e5.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e8045edd-283f-49a0-bb43-e35468f0b8e5.json.gz new file mode 100644 index 00000000000..fb3d8781bd6 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_e8045edd-283f-49a0-bb43-e35468f0b8e5.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ebd1ae48-ac64-4525-85aa-c608b039f6f2.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ebd1ae48-ac64-4525-85aa-c608b039f6f2.json.gz new file mode 100644 index 00000000000..25a01eb2d12 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_ebd1ae48-ac64-4525-85aa-c608b039f6f2.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_edaee500-e84a-4cd1-acec-4d63b6c0f0d0.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_edaee500-e84a-4cd1-acec-4d63b6c0f0d0.json.gz new file mode 100644 index 00000000000..0ab471c07cf Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_edaee500-e84a-4cd1-acec-4d63b6c0f0d0.json.gz differ diff --git a/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_fcbf9813-945a-457d-94b9-bb69b886885e.json.gz b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_fcbf9813-945a-457d-94b9-bb69b886885e.json.gz new file mode 100644 index 00000000000..9505d620215 Binary files /dev/null and b/service/worker/workerdeployment/replaytester/testdata/v3/run_1788826427/replay_worker_deployment_wf_run_fcbf9813-945a-457d-94b9-bb69b886885e.json.gz differ diff --git a/service/worker/workerdeployment/task_queue_family_summary.go b/service/worker/workerdeployment/task_queue_family_summary.go new file mode 100644 index 00000000000..e5afb60ea8b --- /dev/null +++ b/service/worker/workerdeployment/task_queue_family_summary.go @@ -0,0 +1,52 @@ +package workerdeployment + +import ( + "github.com/bits-and-blooms/bloom/v3" + "go.temporal.io/sdk/workflow" + deploymentspb "go.temporal.io/server/api/deployment/v1" +) + +const taskQueueFamilyBloomFalsePositiveRate = 0.01 + +func buildTaskQueueFamilySummary( + taskQueueFamilies map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData, +) *deploymentspb.TaskQueueFamilySummary { + summary := &deploymentspb.TaskQueueFamilySummary{Count: int32(len(taskQueueFamilies))} + if len(taskQueueFamilies) == 0 { + return summary + } + + filter := bloom.NewWithEstimates(uint(len(taskQueueFamilies)), taskQueueFamilyBloomFalsePositiveRate) + for _, taskQueueName := range workflow.DeterministicKeys(taskQueueFamilies) { + filter.AddString(taskQueueName) + } + + summary.BloomFilterSize = int64(filter.Cap()) + summary.BloomFilterHashCount = int32(filter.K()) + summary.BloomFilterWords = make([]int64, len(filter.BitSet().Words())) + for index, word := range filter.BitSet().Words() { + summary.BloomFilterWords[index] = int64(word) + } + return summary +} + +func taskQueueFamilyMayExist(summary *deploymentspb.TaskQueueFamilySummary, taskQueueName string) bool { + count := summary.GetCount() + if count <= 0 { + return true + } + + filterWords := summary.GetBloomFilterWords() + // TODO: Validate the Bloom filter metadata before reconstructing it. + // The proto uses int64 words to satisfy proto lint, while the Bloom library requires uint64. + // This conversion preserves the existing bitset; it does not rebuild the filter from task queue names. + words := make([]uint64, len(filterWords)) + for index, word := range filterWords { + words[index] = uint64(word) + } + return bloom.FromWithM( + words, + uint(summary.GetBloomFilterSize()), + uint(summary.GetBloomFilterHashCount()), + ).TestString(taskQueueName) +} diff --git a/service/worker/workerdeployment/task_queue_family_summary_test.go b/service/worker/workerdeployment/task_queue_family_summary_test.go new file mode 100644 index 00000000000..35c6ea02e0e --- /dev/null +++ b/service/worker/workerdeployment/task_queue_family_summary_test.go @@ -0,0 +1,210 @@ +package workerdeployment + +import ( + "testing" + + "github.com/bits-and-blooms/bloom/v3" + "github.com/stretchr/testify/require" + enumspb "go.temporal.io/api/enums/v1" + sdkclient "go.temporal.io/sdk/client" + "go.temporal.io/sdk/temporal" + deploymentspb "go.temporal.io/server/api/deployment/v1" + "go.temporal.io/server/common/worker_versioning" +) + +func TestBuildTaskQueueFamilySummary(t *testing.T) { + t.Parallel() + + testCases := []struct { + name string + taskQueueFamilies map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData + wantCount int32 + }{ + { + name: "families", + taskQueueFamilies: map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData{ + "queue-a": { + TaskQueues: map[int32]*deploymentspb.TaskQueueVersionData{ + int32(enumspb.TASK_QUEUE_TYPE_WORKFLOW): {}, + int32(enumspb.TASK_QUEUE_TYPE_NEXUS): {}, + }, + }, + "queue-b": { + TaskQueues: map[int32]*deploymentspb.TaskQueueVersionData{ + int32(enumspb.TASK_QUEUE_TYPE_ACTIVITY): {}, + }, + }, + }, + wantCount: 2, + }, + { + name: "no families", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + summary := buildTaskQueueFamilySummary(tc.taskQueueFamilies) + require.Equal(t, tc.wantCount, summary.GetCount()) + if tc.wantCount == 0 { + require.Zero(t, summary.GetBloomFilterSize()) + require.Zero(t, summary.GetBloomFilterHashCount()) + require.Empty(t, summary.GetBloomFilterWords()) + return + } + + require.NotZero(t, summary.GetBloomFilterSize()) + require.NotZero(t, summary.GetBloomFilterHashCount()) + require.NotEmpty(t, summary.GetBloomFilterWords()) + words := make([]uint64, len(summary.GetBloomFilterWords())) + for index, word := range summary.GetBloomFilterWords() { + words[index] = uint64(word) + } + filter := bloom.FromWithM(words, uint(summary.GetBloomFilterSize()), uint(summary.GetBloomFilterHashCount())) + require.True(t, filter.TestString("queue-a")) + require.True(t, filter.TestString("queue-b")) + require.False(t, filter.TestString("queue-c")) + }) + } +} + +func TestVersionStateToSummaryTaskQueueFamilySummary(t *testing.T) { + t.Parallel() + + state := &deploymentspb.VersionLocalState{ + Version: &deploymentspb.WorkerDeploymentVersion{ + DeploymentName: "deployment", + BuildId: "build-id", + }, + TaskQueueFamilies: map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData{ + "queue": {}, + }, + } + + testCases := []struct { + name string + workflowVersion DeploymentWorkflowVersion + wantSummary bool + }{ + { + name: "v2 omits summary", + workflowVersion: VersionDataRevisionNumber, + }, + { + name: "v3 includes summary", + workflowVersion: TaskQueueFamilySummary, + wantSummary: true, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + summary := versionStateToSummary( + state, + tc.workflowVersion >= TaskQueueFamilySummary, + ).GetTaskQueueFamilySummary() + if !tc.wantSummary { + require.Nil(t, summary) + return + } + require.Equal(t, int32(1), summary.GetCount()) + }) + } +} + +func TestValidateRegisterWorkerTaskQueueFamilySummary(t *testing.T) { + t.Parallel() + + version := &deploymentspb.WorkerDeploymentVersion{ + DeploymentName: "deployment", + BuildId: "build-id", + } + versionString := worker_versioning.WorkerDeploymentVersionToStringV31(version) + completeSummary := buildTaskQueueFamilySummary(map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData{ + "existing-queue": {}, + }) + + testCases := []struct { + name string + workflowVersion DeploymentWorkflowVersion + summary *deploymentspb.TaskQueueFamilySummary + taskQueueName string + maxTaskQueues int32 + wantLimitError bool + wantBloomPass bool + }{ + { + name: "definite miss at limit", + workflowVersion: TaskQueueFamilySummary, + summary: completeSummary, + taskQueueName: "new-queue", + maxTaskQueues: 1, + wantLimitError: true, + }, + { + name: "existing family at limit", + workflowVersion: TaskQueueFamilySummary, + summary: completeSummary, + taskQueueName: "existing-queue", + maxTaskQueues: 1, + wantBloomPass: true, + }, + { + name: "below limit", + workflowVersion: TaskQueueFamilySummary, + summary: completeSummary, + taskQueueName: "new-queue", + maxTaskQueues: 2, + }, + { + name: "missing summary fails open", + workflowVersion: TaskQueueFamilySummary, + taskQueueName: "new-queue", + maxTaskQueues: 1, + }, + { + name: "old workflow version fails open", + workflowVersion: VersionDataRevisionNumber, + summary: completeSummary, + taskQueueName: "new-queue", + maxTaskQueues: 1, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + runner := &WorkflowRunner{ + WorkerDeploymentWorkflowArgs: &deploymentspb.WorkerDeploymentWorkflowArgs{ + State: &deploymentspb.WorkerDeploymentLocalState{ + Versions: map[string]*deploymentspb.WorkerDeploymentVersionSummary{ + versionString: {TaskQueueFamilySummary: tc.summary}, + }, + }, + }, + metrics: sdkclient.MetricsNopHandler, + workflowVersion: tc.workflowVersion, + } + bloomFilterPassed, err := runner.validateRegisterWorkerWithBloomFilterResult(&deploymentspb.RegisterWorkerInWorkerDeploymentArgs{ + TaskQueueName: tc.taskQueueName, + TaskQueueType: enumspb.TASK_QUEUE_TYPE_NEXUS, + MaxTaskQueues: tc.maxTaskQueues, + Version: version, + }) + require.Equal(t, tc.wantBloomPass, bloomFilterPassed) + + if !tc.wantLimitError { + require.NoError(t, err) + return + } + var applicationError *temporal.ApplicationError + require.ErrorAs(t, err, &applicationError) + require.Equal(t, errMaxTaskQueuesInVersionType, applicationError.Type()) + }) + } +} diff --git a/service/worker/workerdeployment/util.go b/service/worker/workerdeployment/util.go index b87691fc1b6..9b55aa5939e 100644 --- a/service/worker/workerdeployment/util.go +++ b/service/worker/workerdeployment/util.go @@ -91,6 +91,10 @@ const ( errVersionIsDraining = "errVersionIsDraining" errVersionHasPollers = "errVersionHasPollersSuffix" + taskQueueFamilyBloomFilterOutcomeAccepted = "accepted" + taskQueueFamilyBloomFilterOutcomeFalsePositive = "false_positive" + taskQueueFamilyBloomFilterOutcomeRejected = "rejected" + errFailedPrecondition = "FailedPrecondition" errInvalidComputeConfig = "errInvalidComputeConfig" @@ -122,6 +126,11 @@ var ( ) ) +func isMaxTaskQueuesInVersionError(err error) bool { + applicationError, ok := errors.AsType[*temporal.ApplicationError](err) + return ok && applicationError.Type() == errMaxTaskQueuesInVersionType +} + var ( defaultActivityOptions = workflow.ActivityOptions{ StartToCloseTimeout: 1 * time.Minute, diff --git a/service/worker/workerdeployment/version_workflow.go b/service/worker/workerdeployment/version_workflow.go index c949d427ee6..edb8d4db5e2 100644 --- a/service/worker/workerdeployment/version_workflow.go +++ b/service/worker/workerdeployment/version_workflow.go @@ -350,6 +350,12 @@ func (d *VersionWorkflowRunner) run(ctx workflow.Context) error { if err := d.syncVersionDataToComputeStatus(ctx); err != nil { return err } + if workflow.GetInfo(ctx).ContinuedExecutionRunID != "" && + d.hasMinVersion(TaskQueueFamilySummary) && + len(d.VersionState.GetTaskQueueFamilies()) > 0 && + !d.VersionState.GetTaskQueueFamilySummarySignalSent() { + d.syncSummary(ctx) + } // Listen to signals in a different goroutine to make business logic clearer workflow.Go(ctx, d.listenToSignals) @@ -777,6 +783,7 @@ func (d *VersionWorkflowRunner) handleRegisterWorker(ctx workflow.Context, args if d.VersionState.TaskQueueFamilies == nil { d.VersionState.TaskQueueFamilies = make(map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData) } + _, taskQueueFamilyExists := d.VersionState.TaskQueueFamilies[args.TaskQueueName] if d.VersionState.TaskQueueFamilies[args.TaskQueueName] == nil { d.VersionState.TaskQueueFamilies[args.TaskQueueName] = &deploymentspb.VersionLocalState_TaskQueueFamilyData{} } @@ -796,6 +803,9 @@ func (d *VersionWorkflowRunner) handleRegisterWorker(ctx workflow.Context, args d.VersionState.Status = enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE // deployment workflow updates the status in version summary to INACTIVE } + if !taskQueueFamilyExists && d.hasMinVersion(TaskQueueFamilySummary) { + d.syncSummary(ctx) + } if withRevisionNumbers && args.GetRoutingConfig() != nil { // Still need to check RoutingConfig not being nil because of edge cases during enabling dynamic config. @@ -984,7 +994,7 @@ func (d *VersionWorkflowRunner) handleSyncState(ctx workflow.Context, args *depl } return &deploymentspb.SyncVersionStateResponse{ - Summary: versionStateToSummary(state), + Summary: versionStateToSummary(state, d.hasMinVersion(TaskQueueFamilySummary)), }, nil } @@ -1065,19 +1075,30 @@ func (d *VersionWorkflowRunner) newUUID(ctx workflow.Context) string { // Sync version summary with the WorkerDeployment workflow. func (d *VersionWorkflowRunner) syncSummary(ctx workflow.Context) { + summary := versionStateToSummary( + d.GetVersionState(), + d.hasMinVersion(TaskQueueFamilySummary), + ) err := workflow.SignalExternalWorkflow(ctx, GenerateDeploymentWorkflowID(d.VersionState.Version.DeploymentName), "", SyncVersionSummarySignal, - versionStateToSummary(d.GetVersionState()), + summary, ).Get(ctx, nil) if err != nil { d.logger.Error("could not sync version summary to deployment workflow", "error", err) + return + } + if summary.GetTaskQueueFamilySummary() != nil { + d.VersionState.TaskQueueFamilySummarySignalSent = true } } -func versionStateToSummary(s *deploymentspb.VersionLocalState) *deploymentspb.WorkerDeploymentVersionSummary { - return &deploymentspb.WorkerDeploymentVersionSummary{ +func versionStateToSummary( + s *deploymentspb.VersionLocalState, + includeTaskQueueFamilySummary bool, +) *deploymentspb.WorkerDeploymentVersionSummary { + summary := &deploymentspb.WorkerDeploymentVersionSummary{ Version: worker_versioning.WorkerDeploymentVersionToStringV31(s.Version), CreateTime: s.CreateTime, DrainageStatus: s.DrainageInfo.GetStatus(), // deprecated. @@ -1092,6 +1113,10 @@ func versionStateToSummary(s *deploymentspb.VersionLocalState) *deploymentspb.Wo ComputeConfig: s.ComputeConfig, ComputeStatus: s.ComputeStatus, } + if includeTaskQueueFamilySummary { + summary.TaskQueueFamilySummary = buildTaskQueueFamilySummary(s.GetTaskQueueFamilies()) + } + return summary } func (d *VersionWorkflowRunner) refreshDrainageInfo(ctx workflow.Context) { diff --git a/service/worker/workerdeployment/version_workflow_test.go b/service/worker/workerdeployment/version_workflow_test.go index efea6f2791f..8eb2385746c 100644 --- a/service/worker/workerdeployment/version_workflow_test.go +++ b/service/worker/workerdeployment/version_workflow_test.go @@ -8,6 +8,7 @@ import ( "time" "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" computepb "go.temporal.io/api/compute/v1" deploymentpb "go.temporal.io/api/deployment/v1" @@ -17,6 +18,7 @@ import ( "go.temporal.io/sdk/testsuite" "go.temporal.io/sdk/workflow" deploymentspb "go.temporal.io/server/api/deployment/v1" + "go.temporal.io/server/common/payloads" "go.temporal.io/server/common/testing/protorequire" "go.temporal.io/server/common/testing/testvars" "go.temporal.io/server/common/worker_versioning" @@ -36,7 +38,7 @@ type VersionWorkflowSuite struct { func TestVersionWorkflowSuite(t *testing.T) { t.Parallel() - suite.Run(t, &VersionWorkflowSuite{workflowVersion: VersionDataRevisionNumber}) + suite.Run(t, &VersionWorkflowSuite{workflowVersion: TaskQueueFamilySummary}) } func (s *VersionWorkflowSuite) SetupTest() { @@ -68,6 +70,108 @@ func (s *VersionWorkflowSuite) TearDownTest() { s.env.AssertExpectations(s.T()) } +func TestVersionWorkflowTaskQueueFamilySummaryBootstrap(t *testing.T) { + t.Parallel() + + testCases := []struct { + name string + summarySent bool + wantSignal bool + }{ + { + name: "first v3 run after continue as new", + wantSignal: true, + }, + { + name: "later v3 run", + summarySent: true, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + tv := testvars.New(t) + var workflowSuite testsuite.WorkflowTestSuite + env := workflowSuite.NewTestWorkflowEnvironment() + versionWorkflow := func(ctx workflow.Context, args *deploymentspb.WorkerDeploymentVersionWorkflowArgs) error { + return VersionWorkflow( + ctx, + func() DeploymentWorkflowVersion { return TaskQueueFamilySummary }, + func() time.Duration { return 5 * time.Minute }, + func() time.Duration { return 3 * time.Minute }, + args, + ) + } + env.RegisterWorkflowWithOptions(versionWorkflow, workflow.RegisterOptions{Name: WorkerDeploymentVersionWorkflowType}) + env.SetContinuedExecutionRunID("previous-run-id") + + var capturedSummary *deploymentspb.WorkerDeploymentVersionSummary + signalCall := env.OnSignalExternalWorkflow( + mock.Anything, + GenerateDeploymentWorkflowID(tv.DeploymentSeries()), + "", + SyncVersionSummarySignal, + mock.Anything, + ).Return(func(namespace string, workflowID string, runID string, signalName string, arg any) error { + capturedSummary = arg.(*deploymentspb.WorkerDeploymentVersionSummary) + return nil + }) + if tc.wantSignal { + signalCall.Once() + } else { + signalCall.Maybe() + } + + env.RegisterDelayedCallback(func() { + env.SignalWorkflow(ForceCANSignalName, &deploymentspb.ForceCANVersionSignalArgs{}) + }, time.Millisecond) + + env.ExecuteWorkflow(WorkerDeploymentVersionWorkflowType, &deploymentspb.WorkerDeploymentVersionWorkflowArgs{ + NamespaceName: tv.NamespaceName().String(), + NamespaceId: tv.NamespaceID().String(), + VersionState: &deploymentspb.VersionLocalState{ + Version: &deploymentspb.WorkerDeploymentVersion{ + DeploymentName: tv.DeploymentSeries(), + BuildId: tv.BuildID(), + }, + TaskQueueFamilies: map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData{ + tv.TaskQueue().GetName(): { + TaskQueues: map[int32]*deploymentspb.TaskQueueVersionData{ + int32(enumspb.TASK_QUEUE_TYPE_WORKFLOW): {}, + }, + }, + }, + TaskQueueFamilySummarySignalSent: tc.summarySent, + }, + }) + + require.True(t, env.IsWorkflowCompleted()) + workflowErr := env.GetWorkflowError() + require.Error(t, workflowErr) + var executionErr *temporal.WorkflowExecutionError + require.ErrorAs(t, workflowErr, &executionErr) + var continueAsNewErr *workflow.ContinueAsNewError + require.ErrorAs(t, executionErr.Unwrap(), &continueAsNewErr) + var nextArgs deploymentspb.WorkerDeploymentVersionWorkflowArgs + require.NoError(t, payloads.Decode(continueAsNewErr.Input, &nextArgs)) + require.True(t, nextArgs.GetVersionState().GetTaskQueueFamilySummarySignalSent()) + + require.Equal(t, tc.wantSignal, capturedSummary != nil) + if !tc.wantSignal { + require.Nil(t, capturedSummary) + } else { + require.Equal(t, int32(1), capturedSummary.GetTaskQueueFamilySummary().GetCount()) + require.NotZero(t, capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterSize()) + require.NotZero(t, capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterHashCount()) + require.NotEmpty(t, capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterWords()) + } + env.AssertExpectations(t) + }) + } +} + // Test_SyncState_BatchSize verifies if the right number of batches are created during the SyncDeploymentVersionUserData activity func (s *VersionWorkflowSuite) Test_SyncState_Batch_SingleTaskQueue() { // TODO: refactor this test so it creates a version with the TQ already added to it and then @@ -1095,6 +1199,18 @@ func (s *VersionWorkflowSuite) Test_RegisterWorker_ResetRevisionNumber_WhenReviv // Make propagation check take long enough so register worker happens before workflow exits s.env.OnActivity(a.CheckWorkerDeploymentUserDataPropagation, mock.Anything, mock.Anything).After(100 * time.Millisecond).Return(nil).Maybe() + var capturedSummary *deploymentspb.WorkerDeploymentVersionSummary + s.env.OnSignalExternalWorkflow( + mock.Anything, + GenerateDeploymentWorkflowID(tv.DeploymentSeries()), + "", + SyncVersionSummarySignal, + mock.Anything, + ).Return(func(namespace string, workflowID string, runID string, signalName string, arg any) error { + capturedSummary = arg.(*deploymentspb.WorkerDeploymentVersionSummary) + return nil + }).Maybe() + // Delete the version s.env.RegisterDelayedCallback(func() { deleteArgs := &deploymentspb.DeleteVersionArgs{ @@ -1193,6 +1309,8 @@ func (s *VersionWorkflowSuite) Test_RegisterWorker_ResetRevisionNumber_WhenReviv }) s.True(s.env.IsWorkflowCompleted()) + s.Require().NotNil(capturedSummary) + s.Equal(int32(1), capturedSummary.GetTaskQueueFamilySummary().GetCount()) } // Test_SyncState_IncrementsRevisionNumber_InAsyncMode tests that revision numbers are tracked @@ -1617,7 +1735,20 @@ func (s *VersionWorkflowSuite) Test_RegisterWorker_DoesNotSignalPropagationCompl s.env.OnActivity(a.CheckWorkerDeploymentUserDataPropagation, mock.Anything, mock.Anything).Return(nil).Maybe() - s.env.OnSignalExternalWorkflow(mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Run(func(args mock.Arguments) { + var capturedSummary *deploymentspb.WorkerDeploymentVersionSummary + expectedWorkflowID := GenerateDeploymentWorkflowID(tv.DeploymentSeries()) + s.env.OnSignalExternalWorkflow( + mock.Anything, + expectedWorkflowID, + "", + SyncVersionSummarySignal, + mock.Anything, + ).Return(func(namespace string, workflowID string, runID string, signalName string, arg any) error { + capturedSummary = arg.(*deploymentspb.WorkerDeploymentVersionSummary) + return nil + }).Maybe() + + s.env.OnSignalExternalWorkflow(mock.Anything, mock.Anything, mock.Anything, PropagationCompleteSignal, mock.Anything).Run(func(args mock.Arguments) { s.Fail("Should not signal propagation complete for worker registration") }).Maybe() @@ -1663,6 +1794,11 @@ func (s *VersionWorkflowSuite) Test_RegisterWorker_DoesNotSignalPropagationCompl }) s.True(s.env.IsWorkflowCompleted()) + s.Require().NotNil(capturedSummary) + s.Equal(int32(1), capturedSummary.GetTaskQueueFamilySummary().GetCount()) + s.NotZero(capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterSize()) + s.NotZero(capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterHashCount()) + s.NotEmpty(capturedSummary.GetTaskQueueFamilySummary().GetBloomFilterWords()) } // Test_BatchTaskQueuesForSync_SingleBatch tests batching when task queues fit in one batch diff --git a/service/worker/workerdeployment/workflow.go b/service/worker/workerdeployment/workflow.go index 135914e85e1..1d27f677efb 100644 --- a/service/worker/workerdeployment/workflow.go +++ b/service/worker/workerdeployment/workflow.go @@ -363,10 +363,13 @@ func (d *WorkflowRunner) run(ctx workflow.Context) error { return err } - if err := workflow.SetUpdateHandler( + if err := workflow.SetUpdateHandlerWithOptions( ctx, RegisterWorkerInWorkerDeployment, d.handleRegisterWorker, + workflow.UpdateHandlerOptions{ + Validator: d.validateRegisterWorker, + }, ); err != nil { return err } @@ -714,6 +717,11 @@ func (d *WorkflowRunner) handleRegisterWorker(ctx workflow.Context, args *deploy d.setStateChanged() d.lock.Unlock() }() + // Revalidate after acquiring the lock because the summary may have changed after update validation. + bloomFilterPassed, err := d.validateRegisterWorkerWithBloomFilterResult(args) + if err != nil { + return err + } version := worker_versioning.WorkerDeploymentVersionToStringV31(args.Version) @@ -740,16 +748,20 @@ func (d *WorkflowRunner) handleRegisterWorker(ctx workflow.Context, args *deploy RoutingConfig: routingConfigToSync, }).Get(ctx, nil) if err != nil { - if appError, ok := errors.AsType[*temporal.ApplicationError](err); ok { - if appError.Type() == errMaxTaskQueuesInVersionType { - return temporal.NewApplicationError( - fmt.Sprintf("cannot add task queue %v since maximum number of task queues (%d) have been registered in deployment", args.TaskQueueName, args.MaxTaskQueues), - errMaxTaskQueuesInVersionType, - ) + if isMaxTaskQueuesInVersionError(err) { + if bloomFilterPassed { + d.recordTaskQueueFamilyBloomFilterOutcome(taskQueueFamilyBloomFilterOutcomeFalsePositive) } + return temporal.NewApplicationError( + fmt.Sprintf("cannot add task queue %v since maximum number of task queues (%d) have been registered in deployment", args.TaskQueueName, args.MaxTaskQueues), + errMaxTaskQueuesInVersionType, + ) } return err } + if bloomFilterPassed { + d.recordTaskQueueFamilyBloomFilterOutcome(taskQueueFamilyBloomFilterOutcomeAccepted) + } if d.State.Versions[version].Status == enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CREATED { // now that a poller is seen, we should update the status to INACTIVE @@ -760,6 +772,40 @@ func (d *WorkflowRunner) handleRegisterWorker(ctx workflow.Context, args *deploy return d.updateMemo(ctx) } +func (d *WorkflowRunner) validateRegisterWorker(args *deploymentspb.RegisterWorkerInWorkerDeploymentArgs) error { + _, err := d.validateRegisterWorkerWithBloomFilterResult(args) + return err +} + +func (d *WorkflowRunner) validateRegisterWorkerWithBloomFilterResult(args *deploymentspb.RegisterWorkerInWorkerDeploymentArgs) (bool, error) { + if !d.hasMinVersion(TaskQueueFamilySummary) { + return false, nil + } + + version := worker_versioning.WorkerDeploymentVersionToStringV31(args.GetVersion()) + versionSummary := d.GetState().GetVersions()[version] + taskQueueFamilySummary := versionSummary.GetTaskQueueFamilySummary() + if taskQueueFamilySummary == nil || + taskQueueFamilySummary.GetCount() < args.GetMaxTaskQueues() { + return false, nil + } + if taskQueueFamilyMayExist(taskQueueFamilySummary, args.GetTaskQueueName()) { + return true, nil + } + + // The bloom filter thinks that adding this task queue would exceed the currently set limit + d.recordTaskQueueFamilyBloomFilterOutcome(taskQueueFamilyBloomFilterOutcomeRejected) + return false, temporal.NewApplicationError( + fmt.Sprintf("cannot add task queue %v since maximum number of task queues (%d) have been registered in deployment", args.GetTaskQueueName(), args.GetMaxTaskQueues()), + errMaxTaskQueuesInVersionType, + ) +} + +func (d *WorkflowRunner) recordTaskQueueFamilyBloomFilterOutcome(outcome string) { + d.metrics.WithTags(map[string]string{"outcome": outcome}). + Counter(metrics.WorkerDeploymentTaskQueueFamilyBloomFilterOutcome.Name()).Inc(1) +} + func (d *WorkflowRunner) validateDeleteDeployment() error { if len(d.State.Versions) > 0 { return serviceerror.NewFailedPrecondition("deployment has versions, can't be deleted") @@ -1625,7 +1671,7 @@ func (d *WorkflowRunner) syncVersion(ctx workflow.Context, targetVersion string, d.updateVersionSummary(sum) } else { //nolint:staticcheck // SA1019 - d.updateVersionSummary(versionStateToSummary(res.GetVersionState())) + d.updateVersionSummary(versionStateToSummary(res.GetVersionState(), false)) } } else if revisionNumber > 0 { // Activity failed meaning the synchronous part of the sync request failed. we need to untrack the revision number. diff --git a/service/worker/workerdeployment/workflow_test.go b/service/worker/workerdeployment/workflow_test.go index 7cd1cb45939..b7f5486d3ba 100644 --- a/service/worker/workerdeployment/workflow_test.go +++ b/service/worker/workerdeployment/workflow_test.go @@ -34,7 +34,7 @@ type WorkerDeploymentSuite struct { func TestWorkerDeploymentSuite(t *testing.T) { t.Parallel() - suite.Run(t, &WorkerDeploymentSuite{workflowVersion: VersionDataRevisionNumber}) + suite.Run(t, &WorkerDeploymentSuite{workflowVersion: TaskQueueFamilySummary}) } func (s *WorkerDeploymentSuite) SetupTest() { diff --git a/tests/versioning_test_env.go b/tests/versioning_test_env.go index 5649a910bbd..9a7b55dee4c 100644 --- a/tests/versioning_test_env.go +++ b/tests/versioning_test_env.go @@ -60,7 +60,7 @@ const ( versionStatusDraining = versionStatus(4) versionStatusDrained = versionStatus(5) - versioning3DeploymentWorkflowVersion = workerdeployment.VersionDataRevisionNumber + versioning3DeploymentWorkflowVersion = workerdeployment.TaskQueueFamilySummary ) var _ = testhooks.MatchingIgnoreRoutingConfigRevisionCheck diff --git a/tests/worker_deployment_test.go b/tests/worker_deployment_test.go index ff6dad424a5..4823ee8c25e 100644 --- a/tests/worker_deployment_test.go +++ b/tests/worker_deployment_test.go @@ -7,6 +7,8 @@ import ( "testing" "time" + "github.com/bits-and-blooms/bloom/v3" + "github.com/dgryski/go-farm" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" commonpb "go.temporal.io/api/common/v1" @@ -19,6 +21,7 @@ import ( "go.temporal.io/server/api/matchingservice/v1" "go.temporal.io/server/common" "go.temporal.io/server/common/dynamicconfig" + "go.temporal.io/server/common/metrics" "go.temporal.io/server/common/testing/parallelsuite" "go.temporal.io/server/common/testing/testhooks" "go.temporal.io/server/common/testing/testvars" @@ -317,6 +320,7 @@ func (s *WorkerDeploymentSuite) TestDeploymentVersionLimits() { func (s *WorkerDeploymentSuite) TestDeploymentVersionTaskQueueFamilyLimitAllowsNewType() { env := s.newTestEnv( + testcore.WithDynamicConfig(dynamicconfig.MatchingDeploymentWorkflowVersion, int(workerdeployment.TaskQueueFamilySummary)), testcore.WithDynamicConfig(dynamicconfig.MatchingMaxTaskQueuesInDeploymentVersion, 1), ) tv := env.Tv() @@ -324,6 +328,52 @@ func (s *WorkerDeploymentSuite) TestDeploymentVersionTaskQueueFamilyLimitAllowsN go s.pollFromDeployment(env, tv) s.ensureCreateVersionWithExpectedTaskQueues(env, tv, 1) + deploymentWorkflowID := workerdeployment.GenerateDeploymentWorkflowID(tv.DeploymentSeries()) + var taskQueueFamilySummary *deploymentspb.TaskQueueFamilySummary + s.Await(func(s *WorkerDeploymentSuite) { + queryResult, err := env.SdkClient().QueryWorkflow( + s.Context(), + deploymentWorkflowID, + "", + workerdeployment.QueryDescribeDeployment, + ) + if err != nil { + s.NoError(err) + return + } + + var queryResponse deploymentspb.QueryDescribeWorkerDeploymentResponse + if err := queryResult.Get(&queryResponse); err != nil { + s.NoError(err) + return + } + versionSummary := queryResponse.GetState().GetVersions()[tv.DeploymentVersionString()] + if versionSummary == nil { + s.NotNil(versionSummary) + return + } + taskQueueFamilySummary = versionSummary.GetTaskQueueFamilySummary() + if taskQueueFamilySummary == nil { + s.NotNil(taskQueueFamilySummary) + return + } + s.Equal(int32(1), taskQueueFamilySummary.GetCount()) + s.NotZero(taskQueueFamilySummary.GetBloomFilterSize()) + s.NotZero(taskQueueFamilySummary.GetBloomFilterHashCount()) + s.NotEmpty(taskQueueFamilySummary.GetBloomFilterWords()) + }, 10*time.Second, 200*time.Millisecond) + + metricCapture := env.StartNamespaceMetricCapture() + metricOutcomeCount := func(outcome string) int { + count := 0 + for _, recording := range metricCapture.Metric(metrics.WorkerDeploymentTaskQueueFamilyBloomFilterOutcome.Name()) { + if recording.Tags["outcome"] == outcome { + count++ + } + } + return count + } + go pollActivityFromDeployment(s.Context(), env.TestEnv, tv) s.Await(func(s *WorkerDeploymentSuite) { resp, err := env.FrontendClient().DescribeWorkerDeploymentVersion( @@ -342,13 +392,108 @@ func (s *WorkerDeploymentSuite) TestDeploymentVersionTaskQueueFamilyLimitAllowsN resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos(), ) }, 10*time.Second, 200*time.Millisecond) + s.Await(func(s *WorkerDeploymentSuite) { + s.Equal(1, metricOutcomeCount("accepted")) + }, 10*time.Second, 200*time.Millisecond) - secondTaskQueue := tv.WithTaskQueueNumber(2) - expectedError := fmt.Sprintf( - "cannot add task queue %v since maximum number of task queues (1) have been registered in deployment", - secondTaskQueue.TaskQueue().GetName(), + words := make([]uint64, len(taskQueueFamilySummary.GetBloomFilterWords())) + for index, word := range taskQueueFamilySummary.GetBloomFilterWords() { + words[index] = uint64(word) + } + filter := bloom.FromWithM( + words, + uint(taskQueueFamilySummary.GetBloomFilterSize()), + uint(taskQueueFamilySummary.GetBloomFilterHashCount()), + ) + var falsePositiveTaskQueue *testvars.TestVars + for taskQueueNumber := 2; taskQueueNumber < 1000; taskQueueNumber++ { + candidate := tv.WithTaskQueueNumber(taskQueueNumber) + if filter.TestString(candidate.TaskQueue().GetName()) { + falsePositiveTaskQueue = candidate + break + } + } + s.Require().NotNil(falsePositiveTaskQueue, "expected to find a Bloom filter false positive") + currentDeploymentRunID := func() string { + describeResponse, err := env.FrontendClient().DescribeWorkflowExecution( + s.Context(), + &workflowservice.DescribeWorkflowExecutionRequest{ + Namespace: env.Namespace().String(), + Execution: &commonpb.WorkflowExecution{WorkflowId: deploymentWorkflowID}, + }, + ) + s.Require().NoError(err) + runID := describeResponse.GetWorkflowExecutionInfo().GetExecution().GetRunId() + s.Require().NotEmpty(runID) + return runID + } + firstRelevantRunID := currentDeploymentRunID() + + expectedError := func(taskQueue *testvars.TestVars) string { + return fmt.Sprintf( + "cannot add task queue %v since maximum number of task queues (1) have been registered in deployment", + taskQueue.TaskQueue().GetName(), + ) + } + s.pollFromDeploymentExpectFail(env, falsePositiveTaskQueue, expectedError(falsePositiveTaskQueue)) + s.Await(func(s *WorkerDeploymentSuite) { + s.Equal(1, metricOutcomeCount("false_positive")) + }, 10*time.Second, 200*time.Millisecond) + + var rejectedTaskQueue *testvars.TestVars + for taskQueueNumber := 2; taskQueueNumber < 1000; taskQueueNumber++ { + candidate := tv.WithTaskQueueNumber(taskQueueNumber) + if !filter.TestString(candidate.TaskQueue().GetName()) { + rejectedTaskQueue = candidate + break + } + } + s.Require().NotNil(rejectedTaskQueue, "expected to find a task queue that is definitely absent from the Bloom filter") + + s.pollFromDeploymentExpectFail(env, rejectedTaskQueue, expectedError(rejectedTaskQueue)) + s.Await(func(s *WorkerDeploymentSuite) { + s.Equal(1, metricOutcomeCount("rejected")) + s.Len(metricCapture.Metric(metrics.WorkerDeploymentTaskQueueFamilyBloomFilterOutcome.Name()), 3) + }, 10*time.Second, 200*time.Millisecond) + + registerWorkerUpdateID := func(taskQueue *testvars.TestVars) string { + return fmt.Sprintf( + "%s%v-%v-%d", + workerdeployment.AutoCreateRequestIDPrefix, + farm.Fingerprint64([]byte(taskQueue.BuildID())), + farm.Fingerprint64([]byte(taskQueue.TaskQueue().GetName())), + enumspb.TASK_QUEUE_TYPE_WORKFLOW, + ) + } + acceptedUpdateIDs := make(map[string]bool) + lastRelevantRunID := currentDeploymentRunID() + for runID := firstRelevantRunID; ; { + var nextRunID string + for _, event := range env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{ + WorkflowId: deploymentWorkflowID, + RunId: runID, + }) { + if accepted := event.GetWorkflowExecutionUpdateAcceptedEventAttributes(); accepted != nil { + acceptedUpdateIDs[accepted.GetProtocolInstanceId()] = true + } + if continuedAsNew := event.GetWorkflowExecutionContinuedAsNewEventAttributes(); continuedAsNew != nil { + nextRunID = continuedAsNew.GetNewExecutionRunId() + } + } + if runID == lastRelevantRunID { + break + } + s.Require().NotEmpty(nextRunID, "expected Deployment workflow run %q to continue as new", runID) + runID = nextRunID + } + s.True( + acceptedUpdateIDs[registerWorkerUpdateID(falsePositiveTaskQueue)], + "expected the Bloom false-positive registration update to be accepted", + ) + s.False( + acceptedUpdateIDs[registerWorkerUpdateID(rejectedTaskQueue)], + "task queue registration update was accepted by the Deployment workflow", ) - s.pollFromDeploymentExpectFail(env, secondTaskQueue, expectedError) } func (s *WorkerDeploymentSuite) TestNamespaceDeploymentsLimit() {