139
func (r *reschedulerImpl) Reschedule(
140
namespaceID string,
142
>
r.Lock()
143
>
defer r.Unlock()
144
>
145
>
now := r.timeSource.Now()
146
>
updatedRescheduleTime := false
147
>
for key, pq := range r.pqMap {
148
>
if key.NamespaceID != namespaceID {
149
continue
150
}
151
153
>
// set reschedule time for all tasks in this pq to be now
154
>
items := make([]rescheduledExecuable, 0, pq.Len())
155
>
for !pq.IsEmpty() {
156
>
rescheduled := pq.Remove()
157
>
// scheduled queue pre-fetches tasks,
158
>
// so we need to make sure the reschedule time is not before the task scheduled time
159
>
rescheduled.rescheduleTime = util.MaxTime(
160
>
rescheduled.executable.GetKey().FireTime.Add(common.ScheduledTaskMinPrecision),
161
>
now,
162
>
)
163
>
items = append(items, rescheduled)
164
>
}
165
>
r.pqMap[key] = r.newPriorityQueue(items)
166
}
167
168
// then update timer gate to trigger the actual reschedule
170
>
r.timerGate.Update(now)
171
>
}
172
}
173