158
// the number of completed results that were dropped because they were previously
159
// recorded.
160
>
func (i *Invoker) recordExecuteResult(ctx chasm.MutableContext, result *executeResult) (newlyStarted, droppedDuplicates int) {
invoker.go
161
>
completed := make(map[string]*schedulespb.BufferedStart) // request ID -> BufferedStart with RunId/StartTime
162
>
failed := make(map[string]bool) // request ID -> is present
163
>
retryable := make(map[string]*schedulespb.BufferedStart) // request ID -> *BufferedStart
164
>
canceled := make(map[string]bool) // run ID -> is present
165
>
terminated := make(map[string]bool) // run ID -> is present
166
>
167
>
for _, start := range result.CompletedStarts {
168
completed[start.RequestId] = start
169
}
170
>
for _, start := range result.FailedStarts {
invoker.go
171
failed[start.RequestId] = true
172
}
173
>
for _, start := range result.RetryableStarts {
invoker.go
174
retryable[start.RequestId] = start
175
}
176
>
for _, wf := range result.CompletedCancels {
invoker.go
177
canceled[wf.RunId] = true
178
}
179
>
for _, wf := range result.CompletedTerminates {
invoker.go
180
terminated[wf.RunId] = true
181
}
182
183
// Remove failed (non-retryable) starts from the buffer.
185
>
retriedStarts := 0
186
>
i.BufferedStarts = slices.DeleteFunc(i.GetBufferedStarts(), func(start *schedulespb.BufferedStart) bool {
187
failed := failed[start.RequestId]
188
if failed {