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 123 | 28x 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
}
}
|