263
}
264
return newInternalStartedTask(&startedTaskInfo{workflowTaskInfo: resp}), nil
266
>
resp, err := fwdr.client.PollActivityTaskQueue(ctx, &matchingservice.PollActivityTaskQueueRequest{
267
>
NamespaceId: fwdr.partition.TaskQueue().NamespaceId(),
268
>
PollerId: pollerID,
269
>
PollRequest: &workflowservice.PollActivityTaskQueueRequest{
270
>
TaskQueue: &taskqueuepb.TaskQueue{
271
>
Name: target.RpcName(),
272
>
Kind: fwdr.partition.Kind(),
273
>
},
274
>
Identity: identity,
275
>
TaskQueueMetadata: pollMetadata.taskQueueMetadata,
276
>
WorkerVersionCapabilities: pollMetadata.workerVersionCapabilities,
277
>
DeploymentOptions: pollMetadata.deploymentOptions,
278
>
WorkerInstanceKey: pollMetadata.workerInstanceKey,
279
>
WorkerControlTaskQueue: pollMetadata.workerControlTaskQueue,
280
>
},
281
>
ForwardedSource: fwdr.partition.RpcName(),
282
>
Conditions: pollMetadata.conditions,
283
>
})
284
>
if err != nil {
285
return nil, fwdr.handleErr(err)
287
return nil, errNoTasks
288
}
289
>
return newInternalStartedTask(&startedTaskInfo{activityTaskInfo: resp}), nil
forwarder.go
290
case enumspb.TASK_QUEUE_TYPE_NEXUS:
291
resp, err := fwdr.client.PollNexusTaskQueue(ctx, &matchingservice.PollNexusTaskQueueRequest{