MCPcopy Create free account
hub / github.com/poundifdef/smoothmq / Dequeue

Method Dequeue

queue/sqlite/sqlite.go:390–483  ·  view source on GitHub ↗
(tenantId int64, queueName string, numToDequeue int, requeueIn int)

Source from the content-addressed store, hash-verified

388}
389
390func (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,

Callers

nothing calls this directly

Calls 3

getQueueMethod · 0.95
calculateRateLimitMethod · 0.95
ToModelMethod · 0.80

Tested by

no test coverage detected