558
ctx context.Context,
559
params addTaskParams,
561
>
defer pm.sendPartitionCountTrailer(ctx)
562
>
if err := pm.checkPartitionCounts(ctx, true); err != nil {
563
return "", false, err
564
}
566
pm.signalPartitionScaler()
567
}
568
570
>
directive := params.taskInfo.GetVersionDirective()
571
>
572
>
pm.autoEnableIfNeeded(ctx, params)
573
>
// spoolQueue will be nil iff task is forwarded.
574
>
reredirectTask:
575
>
spoolQueue, syncMatchQueue, _, taskDispatchRevisionNumber, targetVersion, err := pm.getPhysicalQueuesForAdd(ctx, directive, params.forwardInfo, params.taskInfo.GetRunId(), params.taskInfo.GetWorkflowId(), false)
576
>
if err != nil {
577
return "", false, err
578
}
579
580
>
syncMatchTask := newInternalTaskForSyncMatch(params.taskInfo, params.forwardInfo, taskDispatchRevisionNumber, targetVersion)
task_queue_partition_manager.go
581
>
pm.config.setDefaultPriority(syncMatchTask)
582
>
if spoolQueue != nil && spoolQueue.QueueKey().Version().BuildId() != syncMatchQueue.QueueKey().Version().BuildId() {
583
// Task is not forwarded and build ID is different on the two queues -> redirect rule is being applied.
584
// Set redirectInfo in the task as it will be needed if we have to forward the task.