Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions pkg/repository/scheduler_queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -286,13 +286,26 @@ func (d *sharedRepository) markQueueItemsProcessed(ctx context.Context, tenantId
taskIds := make([]int64, 0, len(r.Assigned))
taskInsertedAts := make([]pgtype.Timestamptz, 0, len(r.Assigned))
workerIds := make([]uuid.UUID, 0, len(r.Assigned))
seenTaskIds := make(map[int64]*sqlcv1.V1QueueItem)

var minTaskInsertedAt pgtype.Timestamptz

// if there are any idsToUnqueue that are not in the queuedItems, this means they were
// deleted from the v1_queue_items table, so we should not assign them
for id, assignedItem := range queueItemIdsToAssignedItem {
if _, ok := queuedItemsMap[id]; ok {
if qi, seen := seenTaskIds[assignedItem.QueueItem.TaskID]; seen {
d.l.Error().
Int64("task_id", assignedItem.QueueItem.TaskID).
Str("task_inserted_at", qi.TaskInsertedAt.Time.String()).
Int64("retry_count", int64(qi.RetryCount)).
Msg("duplicate task id seen when preparing queue items for `UpdateTasksToAssigned`")
}

seenTaskIds[assignedItem.QueueItem.TaskID] = assignedItem.QueueItem
// todo: add this back if we remove the error log above for deduping
// taskIdToAssignedItem[assignedItem.QueueItem.TaskID] = assignedItem

taskIds = append(taskIds, assignedItem.QueueItem.TaskID)
taskInsertedAts = append(taskInsertedAts, assignedItem.QueueItem.TaskInsertedAt)
workerIds = append(workerIds, assignedItem.WorkerId)
Comment thread
mrkaye97 marked this conversation as resolved.
Expand Down
Loading