| 271 | } |
| 272 | |
| 273 | func (q *SQLiteQueue) Enqueue(tenantId int64, queueName string, message string, kv map[string]string, delay int) (int64, error) { |
| 274 | messageSnow := q.snow.Generate() |
| 275 | messageId := messageSnow.Int64() |
| 276 | |
| 277 | queue, err := q.getQueue(tenantId, queueName) |
| 278 | if err != nil { |
| 279 | return 0, err |
| 280 | } |
| 281 | |
| 282 | now := time.Now().UTC().Unix() |
| 283 | deliverAt := now + int64(delay) |
| 284 | |
| 285 | newMessage := &Message{ |
| 286 | ID: messageId, |
| 287 | TenantID: tenantId, |
| 288 | QueueID: queue.ID, |
| 289 | DeliverAt: deliverAt, |
| 290 | DeliveredAt: 0, |
| 291 | MaxTries: queue.MaxRetries, |
| 292 | Message: message, |
| 293 | |
| 294 | KV: make([]KV, 0), |
| 295 | } |
| 296 | |
| 297 | for k, v := range kv { |
| 298 | newKv := KV{ |
| 299 | TenantID: tenantId, |
| 300 | MessageID: messageId, |
| 301 | QueueID: queue.ID, |
| 302 | K: k, |
| 303 | V: v, |
| 304 | } |
| 305 | newMessage.KV = append(newMessage.KV, newKv) |
| 306 | } |
| 307 | |
| 308 | q.Mu.Lock() |
| 309 | defer q.Mu.Unlock() |
| 310 | |
| 311 | if err := q.DBG.Create(newMessage).Error; err != nil { |
| 312 | return 0, err |
| 313 | } |
| 314 | |
| 315 | log.Debug().Int64("message_id", messageId).Msg("Enqueued message") |
| 316 | |
| 317 | return messageId, nil |
| 318 | } |
| 319 | |
| 320 | // Calculate how many messages to allow the user to dequeue based on queue's rate limit |
| 321 | func (q *SQLiteQueue) calculateRateLimit(queue *Queue, now int64, numToDequeue int) (int, error) { |