(tenantId int64, queueName string, numToDequeue int, requeueIn int)
| 388 | } |
| 389 | |
| 390 | func (q *SQLiteQueue) Dequeue(tenantId int64, queueName string, numToDequeue int, requeueIn int) ([]*models.Message, error) { |
| 391 | queue, err := q.getQueue(tenantId, queueName) |
| 392 | if err != nil { |
| 393 | return nil, err |
| 394 | } |
| 395 | |
| 396 | // Queue is "paused" |
| 397 | if queue.RateLimit == 0 { |
| 398 | return nil, nil |
| 399 | } |
| 400 | |
| 401 | visibilityTimeout := queue.VisibilityTimeout |
| 402 | if requeueIn > -1 { |
| 403 | visibilityTimeout = requeueIn |
| 404 | } |
| 405 | |
| 406 | now := time.Now().UTC().Unix() |
| 407 | |
| 408 | q.Mu.Lock() |
| 409 | defer q.Mu.Unlock() |
| 410 | |
| 411 | maxToDequeue, err := q.calculateRateLimit(queue, now, numToDequeue) |
| 412 | if err != nil { |
| 413 | return nil, err |
| 414 | } |
| 415 | |
| 416 | var messages []Message |
| 417 | |
| 418 | res := q.DBG.Preload("KV").Where( |
| 419 | "deliver_at <= ? AND delivered_at <= ? AND (tries < max_tries OR max_tries = -1) AND tenant_id = ? AND queue_id = ?", |
| 420 | now, now, tenantId, queue.ID). |
| 421 | Limit(maxToDequeue). |
| 422 | Find(&messages) |
| 423 | |
| 424 | if res.Error != nil { |
| 425 | return nil, err |
| 426 | } |
| 427 | |
| 428 | if len(messages) == 0 { |
| 429 | return nil, nil |
| 430 | } |
| 431 | |
| 432 | rc := make([]*models.Message, len(messages)) |
| 433 | |
| 434 | for i, message := range messages { |
| 435 | rc[i] = message.ToModel() |
| 436 | } |
| 437 | |
| 438 | messageIDs := make([]int64, len(rc)) |
| 439 | for i, message := range rc { |
| 440 | messageIDs[i] = message.ID |
| 441 | } |
| 442 | |
| 443 | err = q.DBG.Transaction(func(tx *gorm.DB) error { |
| 444 | res = tx.Model(&Message{}).Where("tenant_id = ? AND queue_id = ? AND id in ?", tenantId, queue.ID, messageIDs). |
| 445 | UpdateColumns(map[string]any{ |
| 446 | "tries": gorm.Expr("tries+1"), |
| 447 | "delivered_at": now, |
nothing calls this directly
no test coverage detected