1272
return &matchingservice.CancelOutstandingWorkerPollsResponse{}, nil
1273
}
1275
>
// TODO(dynamic partitioning): get real num read partitions from the partition manager.
1276
>
numPartitions := cfg.NumReadPartitions()
1277
>
1278
>
e.logger.Debug("Initiating fan-out for worker poll cancellation",
1279
>
tag.WorkflowNamespaceID(request.GetNamespaceId()),
1280
>
tag.WorkflowTaskQueueName(rootPartition.TaskQueue().Name()),
1281
>
tag.WorkflowTaskQueueType(request.GetTaskQueueType()),
1282
>
tag.NewStringTag("worker-instance-key", request.GetWorkerInstanceKey()),
1283
>
tag.NewInt32("partition-count", int32(numPartitions)),
1284
>
)
1285
>
1286
>
workers := []*matchingservice.CancelOutstandingWorkerPollsPartitionRequest_WorkerEntry{{
1287
>
WorkerInstanceKey: request.GetWorkerInstanceKey(),
1288
>
WorkerIdentity: request.GetWorkerIdentity(),
1289
>
}}
1290
>
1291
>
// Group partitions by destination host. When Route() is unavailable or fails, each
1292
>
// unroutable partition gets a synthetic key so it's sent as an individual RPC.
1293
>
routingClient, _ := e.matchingRawClient.(matching.RoutingClient) //nolint:revive // unchecked-type-assertion: nil is the desired zero value
1294
>
self := e.hostInfoProvider.HostInfo().Identity()
1295
>
tq := rootPartition.TaskQueue()
1296
>
partitionsByTarget := make(map[string][]*tqid.NormalPartition, numPartitions)
1297
>
1298
>
for i := range numPartitions {
1299
>
partition := tq.NormalPartition(i)
1300
>
target := ""
1301
>
if routingClient != nil {
1302
h, err := routingClient.Route(partition)
1303
if err != nil {