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