143
return weight
144
}
145
>
channelWeightUpdateCh := make(chan struct{}, 1)
scheduler.go
146
>
fifoSchedulerOptions := &tasks.FIFOSchedulerOptions{
147
>
QueueSize: prioritySchedulerProcessorQueueSize,
148
>
WorkerCount: options.WorkerCount,
149
>
}
150
>
151
>
fifoScheduler := tasks.NewFIFOScheduler[Executable](
152
>
fifoSchedulerOptions,
153
>
logger,
154
>
)
155
>
156
>
// Wrap the FIFO scheduler with ExecutionAwareScheduler for sequential per-execution processing
157
>
executionAwareScheduler := tasks.NewExecutionAwareScheduler[Executable](
158
>
fifoScheduler,
159
>
options.ExecutionAwareSchedulerOptions,
160
>
executableQueueKeyFn,
161
>
logger,
162
>
metricsHandler,
163
>
timeSource,
164
>
)
165
>
166
>
scheduler = tasks.NewInterleavedWeightedRoundRobinScheduler(
167
>
tasks.InterleavedWeightedRoundRobinSchedulerOptions[Executable, TaskChannelKey]{
168
>
TaskChannelKeyFn: taskChannelKeyFn,
169
>
ChannelWeightFn: channelWeightFn,
170
>
ChannelWeightUpdateCh: channelWeightUpdateCh,
171
>
InactiveChannelDeletionDelay: options.InactiveNamespaceDeletionDelay,
172
>
},
173
>
executionAwareScheduler,
174
>
logger,
175
>
)
176
>
177
>
return &schedulerImpl{
178
>
Scheduler: scheduler,
179
>
namespaceRegistry: namespaceRegistry,
180
>
taskChannelKeyFn: taskChannelKeyFn,
181
>
channelWeightFn: channelWeightFn,
182
>
channelWeightUpdateCh: channelWeightUpdateCh,
183
>
executionAwareScheduler: executionAwareScheduler,
184
>
}
185
}
186