All files / src/core/bindings QueueBinding.ts

33.33% Statements 6/18
0% Branches 0/10
14.28% Functions 1/7
37.5% Lines 6/16

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 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 12328x 28x                                           28x                                             28x   28x   28x                                                                                                                                                
import { Binding } from '@/core/bindings/Binding'
import { lazySettings } from '@/utils/lazySettings'
import { QueueSettings } from '@/types/QueueSettings'
import type { CloudflareEnv } from '@/types/CloudflareEnv'
 
/**
 * Accessor for the Queues producer binding (`queues.producers` in
 * `wrangler.json`).
 *
 * @typeParam Body The message payload type this queue carries.
 *
 * @example
 * ```ts
 * const jobs = new QueueBinding<EmailJob>() // binding name and send defaults from the configuration layer
 * await jobs.send(ctx.env, { to: 'user@example.com', template: 'welcome' })
 * ```
 *
 * @class
 *
 * @author Bayu Dwiyan Satria
 * @version 1.0.0
 * @since 1.0.0
 */
export class QueueBinding<Body = unknown, E extends CloudflareEnv = CloudflareEnv> extends Binding<Queue<Body>, E> {
  /**
   * This accessor’s settings, resolved on first read.
   */
  private readonly settings: () => QueueSettings
 
  /**
   * Default delivery delay applied to sends, in seconds.
   */
  private readonly delaySeconds?: number
 
  /**
   * Default content type applied to sends.
   */
  private readonly contentType?: QueueContentType
 
  /**
   * Constructs a QueueBinding instance.
   *
   * @param settings Resolved settings. Defaults to `resolve('queue')`, so callers
   *   normally construct this with no arguments at all.
   */
  constructor(settings?: QueueSettings) {
    const settle = lazySettings<QueueSettings>('queue', settings)
 
    super(() => settle().binding)
 
    this.settings = settle
  }
 
  /**
   * Sends a single message.
   *
   * @param env The Worker environment.
   * @param message The message payload.
   * @param options Send options — overrides the configured defaults.
   * @returns The send acknowledgement.
   * @throws {@link MissingBindingError} When the queue binding is absent.
   */
  public async send(env: E, message: Body, options: QueueSendOptions = {}): Promise<QueueSendResponse> {
    return await this.resolve(env).send(message, this.withDefaults(options))
  }
 
  /**
   * Sends a batch of messages in a single call.
   *
   * @remarks
   * Batching is the cheaper path — one call bills as one write regardless of
   * how many messages it carries, up to Cloudflare's batch limits.
   *
   * @param env The Worker environment.
   * @param messages The message payloads.
   * @param options Batch send options — overrides the configured defaults.
   * @returns The batch send acknowledgement.
   */
  public async sendBatch(
    env: E,
    messages: Body[],
    options: QueueSendBatchOptions = {}
  ): Promise<QueueSendBatchResponse> {
    const requests: MessageSendRequest<Body>[] = messages.map(body => ({
      body,
      contentType: this.settings().contentType,
      delaySeconds: this.settings().delaySeconds
    }))
 
    return await this.resolve(env).sendBatch(requests, this.withDefaults(options))
  }
 
  /**
   * Reads the queue's current backlog metrics.
   *
   * @param env The Worker environment.
   * @returns The queue metrics.
   */
  public async metrics(env: E): Promise<QueueMetrics> {
    return await this.resolve(env).metrics()
  }
 
  /**
   * Applies the configured defaults to send options that do not set them.
   *
   * @param options The caller-supplied send options.
   * @returns Send options with defaults filled in.
   */
  private withDefaults<T extends QueueSendOptions | QueueSendBatchOptions>(options: T): T {
    const merged = { ...options } as QueueSendOptions & QueueSendBatchOptions
 
    if (this.settings().delaySeconds !== undefined && merged.delaySeconds === undefined) {
      merged.delaySeconds = this.settings().delaySeconds
    }
 
    if (this.settings().contentType !== undefined && (merged as QueueSendOptions).contentType === undefined) {
      ;(merged as QueueSendOptions).contentType = this.settings().contentType
    }
 
    return merged as T
  }
}