All files / src/core/services QueueService.ts

15.78% Statements 3/19
0% Branches 0/8
0% Functions 0/4
15.78% Lines 3/19

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 10428x                         28x                                     28x                                                                                                                                              
import { QueueBinding } from '@/core/bindings/QueueBinding'
 
import type { MessageQueue } from '@bayudwiyansatria/core'
import type { CloudflareEnv } from '@/types/CloudflareEnv'
 
/**
 * Queue producer binding.
 *
 * Every setting it needs — binding name, default delay, content type —
 * arrives resolved from the configuration layer, so nothing about them is
 * decided here. The binding itself is read from `env` per call, which is what
 * makes one module-scope instance safe to share.
 */
const queue = new QueueBinding<unknown>()
 
/**
 * Background-work capability, backed by Cloudflare Queues.
 *
 * @remarks
 * Enqueuing is deliberately forgiving: work that cannot be queued should not
 * fail the request that triggered it, so these methods report failure through
 * their return value rather than throwing. Consumers live in a separate
 * Worker — this side only produces.
 *
 * @typeParam T The message payload this queue carries.
 *
 * @class
 *
 * @author Bayu Dwiyan Satria
 * @version 1.0.0
 * @since 1.0.0
 */
export class QueueService<T = unknown, E extends CloudflareEnv = CloudflareEnv> implements MessageQueue<T, E> {
  /**
   * Enqueues a message.
   *
   * @param env The Worker environment.
   * @param message The payload.
   * @param delaySeconds How long before consumers can see it. Falls back to the configured delay.
   * @returns `true` when the message was accepted.
   */
  public async enqueue(env: E, message: T, delaySeconds?: number): Promise<boolean> {
    if (!queue.isBound(env)) {
      return false
    }
 
    try {
      await queue.send(env, message, delaySeconds ? { delaySeconds } : {})
 
      return true
    } catch (e) {
      console.error('[Queue Error]', e)
 
      return false
    }
  }
 
  /**
   * Enqueues several messages in one call.
   *
   * A batch bills as a single write regardless of how many messages it
   * carries, so prefer this over sending in a loop.
   *
   * @param env The Worker environment.
   * @param messages The payloads.
   * @returns `true` when the batch was accepted.
   */
  public async enqueueAll(env: E, messages: T[]): Promise<boolean> {
    if (!queue.isBound(env) || messages.length === 0) {
      return false
    }
 
    try {
      await queue.sendBatch(env, messages)
 
      return true
    } catch (e) {
      console.error('[Queue Error]', e)
 
      return false
    }
  }
 
  /**
   * Reads the queue's current backlog.
   *
   * @param env The Worker environment.
   * @returns The queue metrics.
   */
  public async backlog(env: E): Promise<QueueMetrics> {
    return await queue.metrics(env)
  }
 
  /**
   * Reports whether the queue is bound to this Worker.
   *
   * @param env The Worker environment.
   * @returns `true` when the binding is present.
   */
  public isAvailable(env: E): boolean {
    return queue.isBound(env)
  }
}