2199
ctx context.Context,
2200
req *matchingservice.SyncDeploymentUserDataRequest,
2202
>
taskQueueFamily, err := tqid.NewTaskQueueFamily(req.NamespaceId, req.GetTaskQueue())
2203
>
applyUpdatesToRoutingConfig := false
2204
>
2205
>
if err != nil {
2206
return nil, err
2207
}
2209
return nil, errMissingDeploymentVersion
2210
}
2211
2212
>
tqMgr, _, err := e.getTaskQueuePartitionManager(ctx, taskQueueFamily.TaskQueue(enumspb.TASK_QUEUE_TYPE_WORKFLOW).RootPartition(), true, loadCauseOtherWrite)
matching_engine.go
2213
>
if err != nil {
2214
return nil, err
2215
}
2216
2217
>
updateOptions := UserDataUpdateOptions{Source: "SyncDeploymentUserData"}
matching_engine.go
2218
>
2219
>
version, err := tqMgr.GetUserDataManager().UpdateUserData(ctx, updateOptions, func(data *persistencespb.TaskQueueUserData) (*persistencespb.TaskQueueUserData, bool, error) {
2220
>
clk := data.GetClock()
2221
>
if clk == nil {
2222
>
clk = hlc.Zero(e.clusterMeta.GetClusterID())
2223
>
}
2224
>
now := hlc.Next(clk, e.timeSource)
2225
>
// clone the whole thing so we can just mutate
2226
>
data = common.CloneProto(data)
2227
>
2228
>
// fill in enough structure so that we can set/append the new deployment data
2229
>
if data == nil {
2230
data = &persistencespb.TaskQueueUserData{}
2231
}
2233
data.PerType = make(map[int32]*persistencespb.TaskQueueTypeUserData)
2234
}
2235
2237
>
for _, t := range req.TaskQueueTypes {
2238
>
if data.PerType[int32(t)] == nil {
2239
data.PerType[int32(t)] = &persistencespb.TaskQueueTypeUserData{}
2240
}
2242
data.PerType[int32(t)].DeploymentData = &persistencespb.DeploymentData{}
2243
}
2244
2245
// set/append the new data
2247
>
2248
>
//nolint:staticcheck // SA1019
2249
>
if vd := req.GetUpdateVersionData(); vd != nil {
2250
// [cleanup-public-preview-versioning]
2251
if vd.GetVersion() == nil { // unversioned ramp