{
  "version": 3,
  "sources": ["../../../../src/workers/queues/broker.worker.ts", "../../../../../../node_modules/.pnpm/kleur@4.1.5/node_modules/kleur/colors.mjs", "../../../../src/workers/queues/constants.ts", "../../../../src/workers/queues/schemas.ts"],
  "sourcesContent": ["import assert from \"node:assert\";\nimport { Buffer } from \"node:buffer\";\nimport { bold, green, grey, red, reset, yellow } from \"kleur/colors\";\nimport {\n\tGET,\n\tHttpError,\n\tLogLevel,\n\tMiniflareDurableObject,\n\tPOST,\n\tSharedBindings,\n\tviewToBuffer,\n} from \"miniflare:shared\";\nimport { HEADER_QUEUE_NAME, QueueBindings } from \"./constants\";\nimport {\n\tQueueConsumersSchema,\n\tQueueContentTypeSchema,\n\tQueueMessageDelaySchema,\n\tQueueProducersSchema,\n\tQueuesBatchRequestSchema,\n} from \"./schemas\";\nimport type {\n\tQueueConsumer,\n\tQueueContentType,\n\tQueueIncomingMessage,\n\tQueueMessageDelay,\n\tQueueOutgoingMessage,\n\tQueueProducer,\n\tQueuesOutgoingBatchRequest,\n} from \"./schemas\";\nimport type { Colorize } from \"kleur/colors\";\nimport type {\n\tMiniflareDurableObjectCf,\n\tMiniflareDurableObjectEnv,\n\tRouteHandler,\n\tTimerHandle,\n} from \"miniflare:shared\";\n\nconst MAX_MESSAGE_SIZE_BYTES = 128 * 1000;\nconst MAX_MESSAGE_BATCH_COUNT = 100;\nconst MAX_MESSAGE_BATCH_SIZE = (256 + 32) * 1000;\n\nconst DEFAULT_BATCH_SIZE = 5;\nconst DEFAULT_BATCH_TIMEOUT = 1; // second\nconst DEFAULT_RETRIES = 2;\n\nconst exceptionQueueResponse: FetcherQueueResult = {\n\toutcome: \"exception\",\n\tretryBatch: { retry: false },\n\tackAll: false,\n\tretryMessages: [],\n\texplicitAcks: [],\n};\n\nclass PayloadTooLargeError extends HttpError {\n\tconstructor(message: string) {\n\t\tsuper(413, message);\n\t}\n}\n\nfunction validateMessageSize(headers: Headers) {\n\tconst size = headers.get(\"Content-Length\");\n\tif (size !== null && parseInt(size) > MAX_MESSAGE_SIZE_BYTES) {\n\t\tthrow new PayloadTooLargeError(\n\t\t\t`message length of ${size} bytes exceeds limit of ${MAX_MESSAGE_SIZE_BYTES}`\n\t\t);\n\t}\n}\n\nfunction validateContentType(headers: Headers): QueueContentType {\n\tconst format = headers.get(\"X-Msg-Fmt\") ?? undefined; // zod will throw if null\n\tconst result = QueueContentTypeSchema.safeParse(format);\n\tif (!result.success) {\n\t\tthrow new HttpError(\n\t\t\t400,\n\t\t\t`message content type ${format} is invalid; if specified, must be one of 'text', 'json', 'bytes', or 'v8'`\n\t\t);\n\t}\n\treturn result.data;\n}\n\nfunction validateMessageDelay(headers: Headers): QueueMessageDelay {\n\tconst format = headers.get(\"X-Msg-Delay-Secs\");\n\tif (!format) return undefined;\n\tconst result = QueueMessageDelaySchema.safeParse(Number(format));\n\tif (!result.success) {\n\t\tthrow new HttpError(\n\t\t\t400,\n\t\t\t`message delay ${format} is invalid: ${result.error}`\n\t\t);\n\t}\n\treturn result.data;\n}\n\nfunction validateBatchSize(headers: Headers) {\n\tconst count = headers.get(\"CF-Queue-Batch-Count\");\n\tif (count !== null && parseInt(count) > MAX_MESSAGE_BATCH_COUNT) {\n\t\tthrow new PayloadTooLargeError(\n\t\t\t`batch message count of ${count} exceeds limit of ${MAX_MESSAGE_BATCH_COUNT}`\n\t\t);\n\t}\n\tconst largestSize = headers.get(\"CF-Queue-Largest-Msg\");\n\tif (largestSize !== null && parseInt(largestSize) > MAX_MESSAGE_SIZE_BYTES) {\n\t\tthrow new PayloadTooLargeError(\n\t\t\t`message in batch has length ${largestSize} bytes which exceeds single message size limit of ${MAX_MESSAGE_SIZE_BYTES}`\n\t\t);\n\t}\n\tconst batchSize = headers.get(\"CF-Queue-Batch-Bytes\");\n\tif (batchSize !== null && parseInt(batchSize) > MAX_MESSAGE_BATCH_SIZE) {\n\t\tthrow new PayloadTooLargeError(\n\t\t\t`batch size of ${batchSize} bytes exceeds limit of ${MAX_MESSAGE_BATCH_SIZE}`\n\t\t);\n\t}\n}\n\ntype QueueBody =\n\t| { contentType: \"text\"; body: string }\n\t| { contentType: \"json\"; body: unknown }\n\t| { contentType: \"bytes\"; body: ArrayBuffer }\n\t| { contentType: \"v8\"; body: Buffer };\n\nfunction deserialise({ contentType, body }: QueueIncomingMessage): QueueBody {\n\tif (contentType === \"text\") {\n\t\treturn { contentType, body: body.toString() };\n\t} else if (contentType === \"json\") {\n\t\treturn { contentType, body: JSON.parse(body.toString()) };\n\t} else if (contentType === \"bytes\") {\n\t\treturn { contentType, body: viewToBuffer(body) };\n\t} else {\n\t\treturn { contentType, body };\n\t}\n}\n\nfunction serialise(msg: QueueMessage): QueueOutgoingMessage {\n\tlet body: Buffer;\n\tif (msg.body.contentType === \"text\") {\n\t\tbody = Buffer.from(msg.body.body);\n\t} else if (msg.body.contentType === \"json\") {\n\t\tbody = Buffer.from(JSON.stringify(msg.body.body));\n\t} else if (msg.body.contentType === \"bytes\") {\n\t\tbody = Buffer.from(msg.body.body);\n\t} else {\n\t\tbody = msg.body.body;\n\t}\n\treturn {\n\t\tid: msg.id,\n\t\ttimestamp: msg.timestamp.getTime(),\n\t\tcontentType: msg.body.contentType,\n\t\tbody: body.toString(\"base64\"),\n\t};\n}\n\nclass QueueMessage {\n\tstatic #encoder = new TextEncoder();\n\t/**\n\t * Message body size in bytes. This is an approximation of production behaviour and may not be the exact same.\n\t */\n\t#bytes = 0;\n\t#failedAttempts = 0;\n\n\tconstructor(\n\t\treadonly id: string,\n\t\treadonly timestamp: Date,\n\t\treadonly body: QueueBody\n\t) {\n\t\tthis.#bytes = QueueMessage.#byteLength(body);\n\t}\n\n\tstatic #byteLength(body: QueueBody): number {\n\t\tswitch (body.contentType) {\n\t\t\tcase \"text\":\n\t\t\t\treturn this.#encoder.encode(body.body).byteLength;\n\t\t\tcase \"json\":\n\t\t\t\treturn this.#encoder.encode(JSON.stringify(body.body)).byteLength;\n\t\t\tcase \"bytes\":\n\t\t\tcase \"v8\":\n\t\t\t\treturn body.body.byteLength;\n\t\t\tdefault:\n\t\t\t\tthrow new Error(`Unexpected queue message contentType received`);\n\t\t}\n\t}\n\n\tincrementFailedAttempts(): number {\n\t\treturn ++this.#failedAttempts;\n\t}\n\n\tget failedAttempts() {\n\t\treturn this.#failedAttempts;\n\t}\n\n\tget bytes() {\n\t\treturn this.#bytes;\n\t}\n}\n\nfunction formatQueueResponse(\n\tqueueName: string,\n\tacked: number,\n\ttotal: number,\n\ttime?: number\n) {\n\tlet colour: Colorize;\n\tif (acked === total) colour = green;\n\telse if (acked > 0) colour = yellow;\n\telse colour = red;\n\n\tlet message = `${bold(\"QUEUE\")} ${queueName} ${colour(`${acked}/${total}`)}`;\n\tif (time !== undefined) message += grey(` (${time}ms)`);\n\treturn reset(message);\n}\n\ninterface PendingFlush {\n\timmediate: boolean;\n\ttimeout: TimerHandle;\n}\n\ntype QueueBrokerObjectEnv = MiniflareDurableObjectEnv & {\n\t// Reference to own Durable Object namespace for sending to dead-letter queues\n\t[SharedBindings.DURABLE_OBJECT_NAMESPACE_OBJECT]: DurableObjectNamespace;\n\t[QueueBindings.MAYBE_JSON_QUEUE_PRODUCERS]?: unknown;\n\t[QueueBindings.MAYBE_JSON_QUEUE_CONSUMERS]?: unknown;\n\t[QueueBindings.MAYBE_SERVICE_QUEUE_PROXY]?: Fetcher;\n} & {\n\t[K in `${typeof QueueBindings.SERVICE_WORKER_PREFIX}${string}`]:\n\t\t| Fetcher\n\t\t| undefined; // Won't have a `Fetcher` for every possible `string`\n};\n\nexport class QueueBrokerObject extends MiniflareDurableObject<QueueBrokerObjectEnv> {\n\treadonly #producers: Record<string, QueueProducer | undefined>;\n\treadonly #consumers: Record<string, QueueConsumer | undefined>;\n\treadonly #messages: QueueMessage[] = [];\n\t#pendingFlush?: PendingFlush;\n\t#backlogBytes = 0;\n\n\tconstructor(state: DurableObjectState, env: QueueBrokerObjectEnv) {\n\t\tsuper(state, env);\n\n\t\tconst maybeProducers = env[QueueBindings.MAYBE_JSON_QUEUE_PRODUCERS];\n\t\tif (maybeProducers === undefined) this.#producers = {};\n\t\telse this.#producers = QueueProducersSchema.parse(maybeProducers);\n\n\t\tconst maybeConsumers = env[QueueBindings.MAYBE_JSON_QUEUE_CONSUMERS];\n\t\tif (maybeConsumers === undefined) this.#consumers = {};\n\t\telse this.#consumers = QueueConsumersSchema.parse(maybeConsumers);\n\t}\n\n\tget #maybeProducer() {\n\t\treturn Object.values(this.#producers).find(\n\t\t\t(p) => p?.queueName === this.name\n\t\t);\n\t}\n\n\tget #maybeConsumer() {\n\t\treturn this.#consumers[this.name];\n\t}\n\n\t#dispatchBatch(\n\t\tworkerName: string,\n\t\tbatch: QueueMessage[],\n\t\tmetadata?: MessageBatchMetadata\n\t) {\n\t\tconst bindingName =\n\t\t\t`${QueueBindings.SERVICE_WORKER_PREFIX}${workerName}` as const;\n\t\tconst maybeService = this.env[bindingName];\n\t\tassert(\n\t\t\tmaybeService !== undefined,\n\t\t\t`Expected ${bindingName} service binding`\n\t\t);\n\t\tconst messages = batch.map(({ id, timestamp, body, failedAttempts }) => {\n\t\t\tconst attempts = failedAttempts + 1;\n\t\t\tif (body.contentType === \"v8\") {\n\t\t\t\treturn { id, timestamp, serializedBody: body.body, attempts };\n\t\t\t} else {\n\t\t\t\treturn { id, timestamp, body: body.body, attempts };\n\t\t\t}\n\t\t});\n\n\t\treturn maybeService.queue(this.name, messages, metadata);\n\t}\n\n\t#flush = async () => {\n\t\tconst consumer = this.#maybeConsumer;\n\t\tassert(consumer !== undefined);\n\n\t\tconst batchSize = consumer.maxBatchSize ?? DEFAULT_BATCH_SIZE;\n\t\tconst maxAttempts = (consumer.maxRetries ?? DEFAULT_RETRIES) + 1;\n\t\tconst maxAttemptsS = maxAttempts === 1 ? \"\" : \"s\";\n\n\t\t// Extract and dispatch a batch\n\t\tconst batch = this.#messages.splice(0, batchSize);\n\t\tthis.#backlogBytes -= batch.reduce((total, msg) => total + msg.bytes, 0);\n\t\tconst metadata: MessageBatchMetadata = {\n\t\t\tmetrics: {\n\t\t\t\tbacklogCount: this.#messages.length,\n\t\t\t\tbacklogBytes: this.#backlogBytes,\n\t\t\t\toldestMessageTimestamp: this.#messages[0]?.timestamp,\n\t\t\t},\n\t\t};\n\t\tconst startTime = Date.now();\n\t\tlet endTime: number;\n\t\tlet response: FetcherQueueResult;\n\t\ttry {\n\t\t\tresponse = await this.#dispatchBatch(\n\t\t\t\tconsumer.workerName,\n\t\t\t\tbatch,\n\t\t\t\tmetadata\n\t\t\t);\n\t\t\tendTime = Date.now();\n\t\t} catch (e: any) {\n\t\t\tendTime = Date.now();\n\t\t\tawait this.logWithLevel(LogLevel.ERROR, String(e));\n\t\t\tresponse = exceptionQueueResponse;\n\t\t}\n\n\t\t// Get messages to retry. If dispatching the batch failed for any reason,\n\t\t// retry all messages.\n\t\tconst retryAll = response.retryBatch.retry || response.outcome !== \"ok\";\n\t\tconst retryMessages = new Map(\n\t\t\tresponse.retryMessages?.map((r) => [r.msgId, r.delaySeconds])\n\t\t);\n\t\tconst globalDelay =\n\t\t\tresponse.retryBatch.delaySeconds ?? consumer.retryDelay ?? 0;\n\n\t\tlet failedMessages = 0;\n\t\tconst toDeadLetterQueue: QueueMessage[] = [];\n\t\tfor (const message of batch) {\n\t\t\tif (retryAll || retryMessages.has(message.id)) {\n\t\t\t\tfailedMessages++;\n\t\t\t\tconst failedAttempts = message.incrementFailedAttempts();\n\t\t\t\tif (failedAttempts < maxAttempts) {\n\t\t\t\t\tawait this.logWithLevel(\n\t\t\t\t\t\tLogLevel.DEBUG,\n\t\t\t\t\t\t`Retrying message \"${message.id}\" on queue \"${this.name}\"...`\n\t\t\t\t\t);\n\n\t\t\t\t\tconst fn = () => {\n\t\t\t\t\t\tthis.#messages.push(message);\n\t\t\t\t\t\tthis.#backlogBytes += message.bytes;\n\t\t\t\t\t\tthis.#ensurePendingFlush();\n\t\t\t\t\t};\n\t\t\t\t\tconst delay = retryMessages.get(message.id) ?? globalDelay;\n\t\t\t\t\tthis.timers.setTimeout(fn, delay * 1000);\n\t\t\t\t} else if (consumer.deadLetterQueue !== undefined) {\n\t\t\t\t\tawait this.logWithLevel(\n\t\t\t\t\t\tLogLevel.WARN,\n\t\t\t\t\t\t`Moving message \"${message.id}\" on queue \"${this.name}\" to dead letter queue \"${consumer.deadLetterQueue}\" after ${maxAttempts} failed attempt${maxAttemptsS}...`\n\t\t\t\t\t);\n\t\t\t\t\ttoDeadLetterQueue.push(message);\n\t\t\t\t} else {\n\t\t\t\t\tawait this.logWithLevel(\n\t\t\t\t\t\tLogLevel.WARN,\n\t\t\t\t\t\t`Dropped message \"${message.id}\" on queue \"${this.name}\" after ${maxAttempts} failed attempt${maxAttemptsS}!`\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t}\n\t\t}\n\t\tconst acked = batch.length - failedMessages;\n\t\tawait this.logWithLevel(\n\t\t\tLogLevel.INFO,\n\t\t\tformatQueueResponse(this.name, acked, batch.length, endTime - startTime)\n\t\t);\n\n\t\t// Ensure we flush again if we still have messages.\n\t\tthis.#pendingFlush = undefined;\n\t\tif (this.#messages.length > 0) this.#ensurePendingFlush();\n\n\t\tif (toDeadLetterQueue.length > 0) {\n\t\t\t// If we have messages to move to a dead letter queue, do so\n\t\t\tconst name = consumer.deadLetterQueue;\n\t\t\tassert(name !== undefined);\n\t\t\tconst ns = this.env[SharedBindings.DURABLE_OBJECT_NAMESPACE_OBJECT];\n\t\t\tconst id = ns.idFromName(name);\n\t\t\tconst stub = ns.get(id);\n\t\t\tconst cf: MiniflareDurableObjectCf = { miniflare: { name } };\n\t\t\tconst batchRequest: QueuesOutgoingBatchRequest = {\n\t\t\t\tmessages: toDeadLetterQueue.map(serialise),\n\t\t\t};\n\t\t\tconst res = await stub.fetch(\"http://placeholder/batch\", {\n\t\t\t\tmethod: \"POST\",\n\t\t\t\tbody: JSON.stringify(batchRequest),\n\t\t\t\tcf: cf as Record<string, unknown>,\n\t\t\t});\n\t\t\tassert(res.ok);\n\t\t}\n\t};\n\n\t#ensurePendingFlush() {\n\t\tconst consumer = this.#maybeConsumer;\n\t\tassert(consumer !== undefined);\n\n\t\tconst batchSize = consumer.maxBatchSize ?? DEFAULT_BATCH_SIZE;\n\t\tconst batchTimeout = consumer.maxBatchTimeout ?? DEFAULT_BATCH_TIMEOUT;\n\t\tconst batchHasSpace = this.#messages.length < batchSize;\n\n\t\tif (this.#pendingFlush !== undefined) {\n\t\t\t// If we have a pending immediate flush, or a delayed flush we haven't\n\t\t\t// filled the batch for yet, just wait for it\n\t\t\tif (this.#pendingFlush.immediate || batchHasSpace) return;\n\t\t\t// Otherwise, the batch is full, so clear the existing timeout, and\n\t\t\t// register an immediate flush\n\t\t\tthis.timers.clearTimeout(this.#pendingFlush.timeout);\n\t\t\tthis.#pendingFlush = undefined;\n\t\t}\n\n\t\t// Register a new flush timeout with the appropriate delay\n\t\tconst delay = batchHasSpace ? batchTimeout * 1000 : 0;\n\t\tconst timeout = this.timers.setTimeout(this.#flush, delay);\n\t\tthis.#pendingFlush = { immediate: delay === 0, timeout };\n\t}\n\n\t#enqueue(messages: QueueIncomingMessage[], globalDelay = 0) {\n\t\tfor (const message of messages) {\n\t\t\tconst randomness = crypto.getRandomValues(new Uint8Array(16));\n\t\t\tconst id = message.id ?? Buffer.from(randomness).toString(\"hex\");\n\t\t\tconst timestamp = new Date(message.timestamp ?? this.timers.now());\n\t\t\tconst body = deserialise(message);\n\t\t\tconst msg = new QueueMessage(id, timestamp, body);\n\n\t\t\tconst fn = () => {\n\t\t\t\tthis.#messages.push(msg);\n\t\t\t\tthis.#backlogBytes += msg.bytes;\n\t\t\t\tthis.#ensurePendingFlush();\n\t\t\t};\n\n\t\t\tconst delay = message.delaySecs ?? globalDelay;\n\t\t\tthis.timers.setTimeout(fn, delay * 1000);\n\t\t}\n\t}\n\n\t#getMetricsResponseObject() {\n\t\treturn {\n\t\t\tbacklogCount: this.#messages.length,\n\t\t\tbacklogBytes: this.#backlogBytes,\n\t\t\toldestMessageTimestamp: this.#messages[0]?.timestamp.getTime() ?? 0,\n\t\t};\n\t}\n\n\t// A queue with no local consumer may have its consumer in another dev\n\t// session. When the dev registry is enabled, this broker has a service\n\t// binding to the dev-registry proxy's `ExternalQueueProxy` entrypoint, which\n\t// resolves a consumer process by queue name and relays the request to that\n\t// process's broker. Returns `null` when the message should be dropped\n\t// instead (no proxy binding, no consumer registered, or forwarding failed),\n\t// mirroring the local no-consumer behaviour so the producer's `send()`\n\t// still succeeds.\n\tasync #tryRemoteConsumer(\n\t\treq: Request<unknown, unknown>\n\t): Promise<Response | null> {\n\t\tconst proxy = this.env[QueueBindings.MAYBE_SERVICE_QUEUE_PROXY];\n\t\tif (proxy === undefined) return null;\n\n\t\tconst headers = new Headers(req.headers);\n\t\theaders.set(HEADER_QUEUE_NAME, this.name);\n\t\t// The consumer's broker can't see this process's producer options, so\n\t\t// apply the local producer's default delivery delay before forwarding.\n\t\tconst delay = this.#maybeProducer?.deliveryDelay;\n\t\tif (delay !== undefined && !headers.has(\"X-Msg-Delay-Secs\")) {\n\t\t\theaders.set(\"X-Msg-Delay-Secs\", String(delay));\n\t\t}\n\n\t\ttry {\n\t\t\t// Buffer the body: if the proxy responds without draining the request\n\t\t\t// (e.g. the 503 no-consumer case), a streamed body would hang the\n\t\t\t// producer's `send()` on the unconsumed pipe.\n\t\t\tconst response = await proxy.fetch(req.url, {\n\t\t\t\tmethod: req.method,\n\t\t\t\theaders,\n\t\t\t\tbody: await req.arrayBuffer(),\n\t\t\t});\n\t\t\t// A 503 means the proxy found no dev session consuming this queue.\n\t\t\tif (response.status === 503) {\n\t\t\t\tawait this.logWithLevel(\n\t\t\t\t\tLogLevel.DEBUG,\n\t\t\t\t\t`Dropped message on queue \"${this.name}\": ${await response.text()}`\n\t\t\t\t);\n\t\t\t\treturn null;\n\t\t\t}\n\t\t\treturn response;\n\t\t} catch (e) {\n\t\t\tawait this.logWithLevel(\n\t\t\t\tLogLevel.WARN,\n\t\t\t\t`Failed to forward message on queue \"${this.name}\" to a consumer in another dev session, dropping it: ${e instanceof Error ? e.message : String(e)}`\n\t\t\t);\n\t\t\treturn null;\n\t\t}\n\t}\n\n\t@POST(\"/message\")\n\tmessage: RouteHandler = async (req) => {\n\t\tconst messageResponse = {\n\t\t\tmetadata: {\n\t\t\t\tmetrics: this.#getMetricsResponseObject(),\n\t\t\t},\n\t\t};\n\n\t\t// If we don't have a local consumer, try a consumer in another dev\n\t\t// session, then drop the message\n\t\tconst consumer = this.#maybeConsumer;\n\t\tif (consumer === undefined) {\n\t\t\tconst forwarded = await this.#tryRemoteConsumer(req);\n\t\t\treturn forwarded ?? Response.json(messageResponse);\n\t\t}\n\n\t\tvalidateMessageSize(req.headers);\n\t\tconst contentType = validateContentType(req.headers);\n\t\tconst delay =\n\t\t\tvalidateMessageDelay(req.headers) ?? this.#maybeProducer?.deliveryDelay;\n\t\tconst body = Buffer.from(await req.arrayBuffer());\n\n\t\tthis.#enqueue(\n\t\t\t[{ contentType, delaySecs: delay, body }],\n\t\t\tthis.#maybeProducer?.deliveryDelay\n\t\t);\n\t\treturn Response.json(messageResponse);\n\t};\n\n\t@POST(\"/batch\")\n\tbatch: RouteHandler = async (req) => {\n\t\tconst batchResponse = {\n\t\t\tmetadata: {\n\t\t\t\tmetrics: this.#getMetricsResponseObject(),\n\t\t\t},\n\t\t};\n\n\t\t// If we don't have a local consumer, try a consumer in another dev\n\t\t// session, then drop the batch\n\t\tconst consumer = this.#maybeConsumer;\n\t\tif (consumer === undefined) {\n\t\t\tconst forwarded = await this.#tryRemoteConsumer(req);\n\t\t\treturn forwarded ?? Response.json(batchResponse);\n\t\t}\n\n\t\t// NOTE: this endpoint is also used when moving messages to the dead-letter\n\t\t// queue. In this case, size headers won't be added and this validation is\n\t\t// a no-op. This allows us to enqueue a maximum size batch with additional\n\t\t// ID and timestamp information.\n\t\tvalidateBatchSize(req.headers);\n\t\tconst delay =\n\t\t\tvalidateMessageDelay(req.headers) ?? this.#maybeProducer?.deliveryDelay;\n\t\tconst body = QueuesBatchRequestSchema.parse(await req.json());\n\n\t\tthis.#enqueue(body.messages, delay);\n\t\treturn Response.json(batchResponse);\n\t};\n\n\t@GET(\"/metrics\")\n\tmetrics: RouteHandler = async (_req) => {\n\t\treturn Response.json(this.#getMetricsResponseObject());\n\t};\n}\n", "let FORCE_COLOR, NODE_DISABLE_COLORS, NO_COLOR, TERM, isTTY=true;\nif (typeof process !== 'undefined') {\n\t({ FORCE_COLOR, NODE_DISABLE_COLORS, NO_COLOR, TERM } = process.env || {});\n\tisTTY = process.stdout && process.stdout.isTTY;\n}\n\nexport const $ = {\n\tenabled: !NODE_DISABLE_COLORS && NO_COLOR == null && TERM !== 'dumb' && (\n\t\tFORCE_COLOR != null && FORCE_COLOR !== '0' || isTTY\n\t)\n}\n\nfunction init(x, y) {\n\tlet rgx = new RegExp(`\\\\x1b\\\\[${y}m`, 'g');\n\tlet open = `\\x1b[${x}m`, close = `\\x1b[${y}m`;\n\n\treturn function (txt) {\n\t\tif (!$.enabled || txt == null) return txt;\n\t\treturn open + (!!~(''+txt).indexOf(close) ? txt.replace(rgx, close + open) : txt) + close;\n\t};\n}\n\n// modifiers\nexport const reset = init(0, 0);\nexport const bold = init(1, 22);\nexport const dim = init(2, 22);\nexport const italic = init(3, 23);\nexport const underline = init(4, 24);\nexport const inverse = init(7, 27);\nexport const hidden = init(8, 28);\nexport const strikethrough = init(9, 29);\n\n// colors\nexport const black = init(30, 39);\nexport const red = init(31, 39);\nexport const green = init(32, 39);\nexport const yellow = init(33, 39);\nexport const blue = init(34, 39);\nexport const magenta = init(35, 39);\nexport const cyan = init(36, 39);\nexport const white = init(37, 39);\nexport const gray = init(90, 39);\nexport const grey = init(90, 39);\n\n// background colors\nexport const bgBlack = init(40, 49);\nexport const bgRed = init(41, 49);\nexport const bgGreen = init(42, 49);\nexport const bgYellow = init(43, 49);\nexport const bgBlue = init(44, 49);\nexport const bgMagenta = init(45, 49);\nexport const bgCyan = init(46, 49);\nexport const bgWhite = init(47, 49);\n", "export const QueueBindings = {\n\tSERVICE_WORKER_PREFIX: \"MINIFLARE_WORKER_\",\n\tMAYBE_JSON_QUEUE_PRODUCERS: \"MINIFLARE_QUEUE_PRODUCERS\",\n\tMAYBE_JSON_QUEUE_CONSUMERS: \"MINIFLARE_QUEUE_CONSUMERS\",\n\t// Optional service binding to the dev-registry proxy's `ExternalQueueProxy`\n\t// entrypoint, present when the dev registry is enabled for queue brokers.\n\tMAYBE_SERVICE_QUEUE_PROXY: \"MINIFLARE_QUEUE_PROXY\",\n} as const;\n\n// Header carrying the queue name on requests the broker forwards to the\n// dev-registry proxy, which resolves the consumer's process from it.\nexport const HEADER_QUEUE_NAME = \"MF-Queue-Name\";\n\n// Prefix for the workerd service backing a single queue's broker. Note this\n// must match the queues plugin name (\"queues\"), which lives on the Node.js\n// side of the build boundary.\nexport const SERVICE_QUEUE_PREFIX = \"queues:queue\";\n\n// The workerd service name backing a single queue's broker. Producers in other\n// dev sessions resolve a consumer process's broker by this exact name through\n// the dev registry's debug port, so it must be derived in one place rather\n// than reconstructed at each call site.\nexport function getQueueServiceName(queueId: string): string {\n\treturn `${SERVICE_QUEUE_PREFIX}:${queueId}`;\n}\n", "import { Base64DataSchema, z } from \"miniflare:zod\";\n\nexport const QueueMessageDelaySchema = z\n\t.number()\n\t.int()\n\t.min(0)\n\t.max(86400)\n\t.optional();\n\nexport const QueueProducerOptionsSchema = /* @__PURE__ */ z.object({\n\t// https://developers.cloudflare.com/queues/platform/configuration/#producer\n\tqueueName: z.string(),\n\tdeliveryDelay: QueueMessageDelaySchema,\n});\n\nexport const QueueProducerSchema = /* @__PURE__ */ z.intersection(\n\tQueueProducerOptionsSchema,\n\tz.object({ workerName: z.string() })\n);\nexport type QueueProducer = z.infer<typeof QueueProducerSchema>;\nexport const QueueProducersSchema =\n\t/* @__PURE__ */ z.record(z.string(), QueueProducerSchema);\n\nexport const QueueConsumerOptionsSchema = /* @__PURE__ */ z.object({\n\t// https://developers.cloudflare.com/queues/platform/configuration/#consumer\n\t// https://developers.cloudflare.com/queues/platform/limits/\n\tmaxBatchSize: z.number().min(0).max(100).optional(),\n\tmaxBatchTimeout: z.number().min(0).max(60).optional(), // seconds\n\tmaxRetries: z.number().min(0).max(100).optional(),\n\tdeadLetterQueue: z.string().optional(),\n\tretryDelay: QueueMessageDelaySchema,\n});\nexport const QueueConsumerSchema = /* @__PURE__ */ z.intersection(\n\tQueueConsumerOptionsSchema,\n\tz.object({ workerName: z.string() })\n);\nexport type QueueConsumer = z.infer<typeof QueueConsumerSchema>;\n// Maps queue names to the Worker that wishes to consume it. Note each queue\n// can only be consumed by one Worker, but one Worker may consume multiple\n// queues. Support for multiple consumers of a single queue is not planned\n// anytime soon.\nexport const QueueConsumersSchema =\n\t/* @__PURE__ */ z.record(z.string(), QueueConsumerSchema);\n\nexport const QueueContentTypeSchema = /* @__PURE__ */ z\n\t.enum([\"text\", \"json\", \"bytes\", \"v8\"])\n\t.default(\"v8\");\nexport type QueueContentType = z.infer<typeof QueueContentTypeSchema>;\n\nexport type QueueMessageDelay = z.infer<typeof QueueMessageDelaySchema>;\n\nexport const QueueIncomingMessageSchema = /* @__PURE__ */ z.object({\n\tcontentType: QueueContentTypeSchema,\n\tdelaySecs: QueueMessageDelaySchema,\n\tbody: Base64DataSchema,\n\t// When enqueuing messages on dead-letter queues, we want to reuse the same ID\n\t// and timestamp\n\tid: z.string().optional(),\n\ttimestamp: z.number().optional(),\n});\nexport type QueueIncomingMessage = z.infer<typeof QueueIncomingMessageSchema>;\nexport type QueueOutgoingMessage = z.input<typeof QueueIncomingMessageSchema>;\n\nexport const QueuesBatchRequestSchema = /* @__PURE__ */ z.object({\n\tmessages: z.array(QueueIncomingMessageSchema),\n});\nexport type QueuesOutgoingBatchRequest = z.input<\n\ttypeof QueuesBatchRequestSchema\n>;\n"],
  "mappings": ";;;;;;;;;AAAA,OAAO,YAAY;AACnB,SAAS,UAAAA,eAAc;;;ACDvB,IAAI,aAAa,qBAAqB,UAAU,MAAM,QAAM;AACxD,OAAO,UAAY,QACrB,EAAE,aAAa,qBAAqB,UAAU,KAAK,IAAI,QAAQ,OAAO,CAAC,GACxE,QAAQ,QAAQ,UAAU,QAAQ,OAAO;AAGnC,IAAM,IAAI;AAAA,EAChB,SAAS,CAAC,uBAAuB,YAAY,QAAQ,SAAS,WAC7D,eAAe,QAAQ,gBAAgB,OAAO;AAEhD;AAEA,SAAS,KAAK,GAAG,GAAG;AACnB,MAAI,MAAM,IAAI,OAAO,WAAW,CAAC,KAAK,GAAG,GACrC,OAAO,QAAQ,CAAC,KAAK,QAAQ,QAAQ,CAAC;AAE1C,SAAO,SAAU,KAAK;AACrB,WAAI,CAAC,EAAE,WAAW,OAAO,OAAa,MAC/B,QAAU,EAAE,KAAG,KAAK,QAAQ,KAAK,IAAI,IAAI,QAAQ,KAAK,QAAQ,IAAI,IAAI,OAAO;AAAA,EACrF;AACD;AAGO,IAAM,QAAQ,KAAK,GAAG,CAAC,GACjB,OAAO,KAAK,GAAG,EAAE,GACjB,MAAM,KAAK,GAAG,EAAE,GAChB,SAAS,KAAK,GAAG,EAAE,GACnB,YAAY,KAAK,GAAG,EAAE,GACtB,UAAU,KAAK,GAAG,EAAE,GACpB,SAAS,KAAK,GAAG,EAAE,GACnB,gBAAgB,KAAK,GAAG,EAAE,GAG1B,QAAQ,KAAK,IAAI,EAAE,GACnB,MAAM,KAAK,IAAI,EAAE,GACjB,QAAQ,KAAK,IAAI,EAAE,GACnB,SAAS,KAAK,IAAI,EAAE,GACpB,OAAO,KAAK,IAAI,EAAE,GAClB,UAAU,KAAK,IAAI,EAAE,GACrB,OAAO,KAAK,IAAI,EAAE,GAClB,QAAQ,KAAK,IAAI,EAAE,GACnB,OAAO,KAAK,IAAI,EAAE,GAClB,OAAO,KAAK,IAAI,EAAE,GAGlB,UAAU,KAAK,IAAI,EAAE,GACrB,QAAQ,KAAK,IAAI,EAAE,GACnB,UAAU,KAAK,IAAI,EAAE,GACrB,WAAW,KAAK,IAAI,EAAE,GACtB,SAAS,KAAK,IAAI,EAAE,GACpB,YAAY,KAAK,IAAI,EAAE,GACvB,SAAS,KAAK,IAAI,EAAE,GACpB,UAAU,KAAK,IAAI,EAAE;;;ADjDlC;AAAA,EACC;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,OACM;;;AEXA,IAAM,gBAAgB;AAAA,EAC5B,uBAAuB;AAAA,EACvB,4BAA4B;AAAA,EAC5B,4BAA4B;AAAA;AAAA;AAAA,EAG5B,2BAA2B;AAC5B,GAIa,oBAAoB;;;ACXjC,SAAS,kBAAkB,SAAS;AAE7B,IAAM,0BAA0B,EACrC,OAAO,EACP,IAAI,EACJ,IAAI,CAAC,EACL,IAAI,KAAK,EACT,SAAS,GAEE,6BAA6C,kBAAE,OAAO;AAAA;AAAA,EAElE,WAAW,EAAE,OAAO;AAAA,EACpB,eAAe;AAChB,CAAC,GAEY,sBAAsC,kBAAE;AAAA,EACpD;AAAA,EACA,EAAE,OAAO,EAAE,YAAY,EAAE,OAAO,EAAE,CAAC;AACpC,GAEa,uBACI,kBAAE,OAAO,EAAE,OAAO,GAAG,mBAAmB,GAE5C,6BAA6C,kBAAE,OAAO;AAAA;AAAA;AAAA,EAGlE,cAAc,EAAE,OAAO,EAAE,IAAI,CAAC,EAAE,IAAI,GAAG,EAAE,SAAS;AAAA,EAClD,iBAAiB,EAAE,OAAO,EAAE,IAAI,CAAC,EAAE,IAAI,EAAE,EAAE,SAAS;AAAA;AAAA,EACpD,YAAY,EAAE,OAAO,EAAE,IAAI,CAAC,EAAE,IAAI,GAAG,EAAE,SAAS;AAAA,EAChD,iBAAiB,EAAE,OAAO,EAAE,SAAS;AAAA,EACrC,YAAY;AACb,CAAC,GACY,sBAAsC,kBAAE;AAAA,EACpD;AAAA,EACA,EAAE,OAAO,EAAE,YAAY,EAAE,OAAO,EAAE,CAAC;AACpC,GAMa,uBACI,kBAAE,OAAO,EAAE,OAAO,GAAG,mBAAmB,GAE5C,yBAAyC,kBACpD,KAAK,CAAC,QAAQ,QAAQ,SAAS,IAAI,CAAC,EACpC,QAAQ,IAAI,GAKD,6BAA6C,kBAAE,OAAO;AAAA,EAClE,aAAa;AAAA,EACb,WAAW;AAAA,EACX,MAAM;AAAA;AAAA;AAAA,EAGN,IAAI,EAAE,OAAO,EAAE,SAAS;AAAA,EACxB,WAAW,EAAE,OAAO,EAAE,SAAS;AAChC,CAAC,GAIY,2BAA2C,kBAAE,OAAO;AAAA,EAChE,UAAU,EAAE,MAAM,0BAA0B;AAC7C,CAAC;;;AH5BD,IAAM,yBAAyB,MAAM,KAC/B,0BAA0B,KAC1B,yBAA0B,MAAY,KAEtC,qBAAqB,GACrB,wBAAwB,GACxB,kBAAkB,GAElB,yBAA6C;AAAA,EAClD,SAAS;AAAA,EACT,YAAY,EAAE,OAAO,GAAM;AAAA,EAC3B,QAAQ;AAAA,EACR,eAAe,CAAC;AAAA,EAChB,cAAc,CAAC;AAChB,GAEM,uBAAN,cAAmC,UAAU;AAAA,EAC5C,YAAY,SAAiB;AAC5B,UAAM,KAAK,OAAO;AAAA,EACnB;AACD;AAEA,SAAS,oBAAoB,SAAkB;AAC9C,MAAM,OAAO,QAAQ,IAAI,gBAAgB;AACzC,MAAI,SAAS,QAAQ,SAAS,IAAI,IAAI;AACrC,UAAM,IAAI;AAAA,MACT,qBAAqB,IAAI,2BAA2B,sBAAsB;AAAA,IAC3E;AAEF;AAEA,SAAS,oBAAoB,SAAoC;AAChE,MAAM,SAAS,QAAQ,IAAI,WAAW,KAAK,QACrC,SAAS,uBAAuB,UAAU,MAAM;AACtD,MAAI,CAAC,OAAO;AACX,UAAM,IAAI;AAAA,MACT;AAAA,MACA,wBAAwB,MAAM;AAAA,IAC/B;AAED,SAAO,OAAO;AACf;AAEA,SAAS,qBAAqB,SAAqC;AAClE,MAAM,SAAS,QAAQ,IAAI,kBAAkB;AAC7C,MAAI,CAAC,OAAQ;AACb,MAAM,SAAS,wBAAwB,UAAU,OAAO,MAAM,CAAC;AAC/D,MAAI,CAAC,OAAO;AACX,UAAM,IAAI;AAAA,MACT;AAAA,MACA,iBAAiB,MAAM,gBAAgB,OAAO,KAAK;AAAA,IACpD;AAED,SAAO,OAAO;AACf;AAEA,SAAS,kBAAkB,SAAkB;AAC5C,MAAM,QAAQ,QAAQ,IAAI,sBAAsB;AAChD,MAAI,UAAU,QAAQ,SAAS,KAAK,IAAI;AACvC,UAAM,IAAI;AAAA,MACT,0BAA0B,KAAK,qBAAqB,uBAAuB;AAAA,IAC5E;AAED,MAAM,cAAc,QAAQ,IAAI,sBAAsB;AACtD,MAAI,gBAAgB,QAAQ,SAAS,WAAW,IAAI;AACnD,UAAM,IAAI;AAAA,MACT,+BAA+B,WAAW,qDAAqD,sBAAsB;AAAA,IACtH;AAED,MAAM,YAAY,QAAQ,IAAI,sBAAsB;AACpD,MAAI,cAAc,QAAQ,SAAS,SAAS,IAAI;AAC/C,UAAM,IAAI;AAAA,MACT,iBAAiB,SAAS,2BAA2B,sBAAsB;AAAA,IAC5E;AAEF;AAQA,SAAS,YAAY,EAAE,aAAa,KAAK,GAAoC;AAC5E,SAAI,gBAAgB,SACZ,EAAE,aAAa,MAAM,KAAK,SAAS,EAAE,IAClC,gBAAgB,SACnB,EAAE,aAAa,MAAM,KAAK,MAAM,KAAK,SAAS,CAAC,EAAE,IAC9C,gBAAgB,UACnB,EAAE,aAAa,MAAM,aAAa,IAAI,EAAE,IAExC,EAAE,aAAa,KAAK;AAE7B;AAEA,SAAS,UAAU,KAAyC;AAC3D,MAAI;AACJ,SAAI,IAAI,KAAK,gBAAgB,SAC5B,OAAOC,QAAO,KAAK,IAAI,KAAK,IAAI,IACtB,IAAI,KAAK,gBAAgB,SACnC,OAAOA,QAAO,KAAK,KAAK,UAAU,IAAI,KAAK,IAAI,CAAC,IACtC,IAAI,KAAK,gBAAgB,UACnC,OAAOA,QAAO,KAAK,IAAI,KAAK,IAAI,IAEhC,OAAO,IAAI,KAAK,MAEV;AAAA,IACN,IAAI,IAAI;AAAA,IACR,WAAW,IAAI,UAAU,QAAQ;AAAA,IACjC,aAAa,IAAI,KAAK;AAAA,IACtB,MAAM,KAAK,SAAS,QAAQ;AAAA,EAC7B;AACD;AAEA,IAAM,eAAN,MAAM,cAAa;AAAA,EAQlB,YACU,IACA,WACA,MACR;AAHQ;AACA;AACA;AAET,SAAK,SAAS,cAAa,YAAY,IAAI;AAAA,EAC5C;AAAA,EALU;AAAA,EACA;AAAA,EACA;AAAA,EAVV,OAAO,WAAW,IAAI,YAAY;AAAA;AAAA;AAAA;AAAA,EAIlC,SAAS;AAAA,EACT,kBAAkB;AAAA,EAUlB,OAAO,YAAY,MAAyB;AAC3C,YAAQ,KAAK,aAAa;AAAA,MACzB,KAAK;AACJ,eAAO,KAAK,SAAS,OAAO,KAAK,IAAI,EAAE;AAAA,MACxC,KAAK;AACJ,eAAO,KAAK,SAAS,OAAO,KAAK,UAAU,KAAK,IAAI,CAAC,EAAE;AAAA,MACxD,KAAK;AAAA,MACL,KAAK;AACJ,eAAO,KAAK,KAAK;AAAA,MAClB;AACC,cAAM,IAAI,MAAM,+CAA+C;AAAA,IACjE;AAAA,EACD;AAAA,EAEA,0BAAkC;AACjC,WAAO,EAAE,KAAK;AAAA,EACf;AAAA,EAEA,IAAI,iBAAiB;AACpB,WAAO,KAAK;AAAA,EACb;AAAA,EAEA,IAAI,QAAQ;AACX,WAAO,KAAK;AAAA,EACb;AACD;AAEA,SAAS,oBACR,WACA,OACA,OACA,MACC;AACD,MAAI;AACJ,EAAI,UAAU,QAAO,SAAS,QACrB,QAAQ,IAAG,SAAS,SACxB,SAAS;AAEd,MAAI,UAAU,GAAG,KAAK,OAAO,CAAC,IAAI,SAAS,IAAI,OAAO,GAAG,KAAK,IAAI,KAAK,EAAE,CAAC;AAC1E,SAAI,SAAS,WAAW,WAAW,KAAK,KAAK,IAAI,KAAK,IAC/C,MAAM,OAAO;AACrB;AAmBO,IAAM,oBAAN,cAAgC,uBAA6C;AAAA,EAC1E;AAAA,EACA;AAAA,EACA,YAA4B,CAAC;AAAA,EACtC;AAAA,EACA,gBAAgB;AAAA,EAEhB,YAAY,OAA2B,KAA2B;AACjE,UAAM,OAAO,GAAG;AAEhB,QAAM,iBAAiB,IAAI,cAAc,0BAA0B;AACnE,IAAI,mBAAmB,SAAW,KAAK,aAAa,CAAC,IAChD,KAAK,aAAa,qBAAqB,MAAM,cAAc;AAEhE,QAAM,iBAAiB,IAAI,cAAc,0BAA0B;AACnE,IAAI,mBAAmB,SAAW,KAAK,aAAa,CAAC,IAChD,KAAK,aAAa,qBAAqB,MAAM,cAAc;AAAA,EACjE;AAAA,EAEA,IAAI,iBAAiB;AACpB,WAAO,OAAO,OAAO,KAAK,UAAU,EAAE;AAAA,MACrC,CAAC,MAAM,GAAG,cAAc,KAAK;AAAA,IAC9B;AAAA,EACD;AAAA,EAEA,IAAI,iBAAiB;AACpB,WAAO,KAAK,WAAW,KAAK,IAAI;AAAA,EACjC;AAAA,EAEA,eACC,YACA,OACA,UACC;AACD,QAAM,cACL,GAAG,cAAc,qBAAqB,GAAG,UAAU,IAC9C,eAAe,KAAK,IAAI,WAAW;AACzC;AAAA,MACC,iBAAiB;AAAA,MACjB,YAAY,WAAW;AAAA,IACxB;AACA,QAAM,WAAW,MAAM,IAAI,CAAC,EAAE,IAAI,WAAW,MAAM,eAAe,MAAM;AACvE,UAAM,WAAW,iBAAiB;AAClC,aAAI,KAAK,gBAAgB,OACjB,EAAE,IAAI,WAAW,gBAAgB,KAAK,MAAM,SAAS,IAErD,EAAE,IAAI,WAAW,MAAM,KAAK,MAAM,SAAS;AAAA,IAEpD,CAAC;AAED,WAAO,aAAa,MAAM,KAAK,MAAM,UAAU,QAAQ;AAAA,EACxD;AAAA,EAEA,SAAS,YAAY;AACpB,QAAM,WAAW,KAAK;AACtB,WAAO,aAAa,MAAS;AAE7B,QAAM,YAAY,SAAS,gBAAgB,oBACrC,eAAe,SAAS,cAAc,mBAAmB,GACzD,eAAe,gBAAgB,IAAI,KAAK,KAGxC,QAAQ,KAAK,UAAU,OAAO,GAAG,SAAS;AAChD,SAAK,iBAAiB,MAAM,OAAO,CAAC,OAAO,QAAQ,QAAQ,IAAI,OAAO,CAAC;AACvE,QAAM,WAAiC;AAAA,MACtC,SAAS;AAAA,QACR,cAAc,KAAK,UAAU;AAAA,QAC7B,cAAc,KAAK;AAAA,QACnB,wBAAwB,KAAK,UAAU,CAAC,GAAG;AAAA,MAC5C;AAAA,IACD,GACM,YAAY,KAAK,IAAI,GACvB,SACA;AACJ,QAAI;AACH,iBAAW,MAAM,KAAK;AAAA,QACrB,SAAS;AAAA,QACT;AAAA,QACA;AAAA,MACD,GACA,UAAU,KAAK,IAAI;AAAA,IACpB,SAAS,GAAQ;AAChB,gBAAU,KAAK,IAAI,GACnB,MAAM,KAAK,aAAa,SAAS,OAAO,OAAO,CAAC,CAAC,GACjD,WAAW;AAAA,IACZ;AAIA,QAAM,WAAW,SAAS,WAAW,SAAS,SAAS,YAAY,MAC7D,gBAAgB,IAAI;AAAA,MACzB,SAAS,eAAe,IAAI,CAAC,MAAM,CAAC,EAAE,OAAO,EAAE,YAAY,CAAC;AAAA,IAC7D,GACM,cACL,SAAS,WAAW,gBAAgB,SAAS,cAAc,GAExD,iBAAiB,GACf,oBAAoC,CAAC;AAC3C,aAAW,WAAW;AACrB,UAAI,YAAY,cAAc,IAAI,QAAQ,EAAE;AAG3C,YAFA,kBACuB,QAAQ,wBAAwB,IAClC,aAAa;AACjC,gBAAM,KAAK;AAAA,YACV,SAAS;AAAA,YACT,qBAAqB,QAAQ,EAAE,eAAe,KAAK,IAAI;AAAA,UACxD;AAEA,cAAM,KAAK,MAAM;AAChB,iBAAK,UAAU,KAAK,OAAO,GAC3B,KAAK,iBAAiB,QAAQ,OAC9B,KAAK,oBAAoB;AAAA,UAC1B,GACM,QAAQ,cAAc,IAAI,QAAQ,EAAE,KAAK;AAC/C,eAAK,OAAO,WAAW,IAAI,QAAQ,GAAI;AAAA,QACxC,MAAO,CAAI,SAAS,oBAAoB,UACvC,MAAM,KAAK;AAAA,UACV,SAAS;AAAA,UACT,mBAAmB,QAAQ,EAAE,eAAe,KAAK,IAAI,2BAA2B,SAAS,eAAe,WAAW,WAAW,kBAAkB,YAAY;AAAA,QAC7J,GACA,kBAAkB,KAAK,OAAO,KAE9B,MAAM,KAAK;AAAA,UACV,SAAS;AAAA,UACT,oBAAoB,QAAQ,EAAE,eAAe,KAAK,IAAI,WAAW,WAAW,kBAAkB,YAAY;AAAA,QAC3G;AAIH,QAAM,QAAQ,MAAM,SAAS;AAU7B,QATA,MAAM,KAAK;AAAA,MACV,SAAS;AAAA,MACT,oBAAoB,KAAK,MAAM,OAAO,MAAM,QAAQ,UAAU,SAAS;AAAA,IACxE,GAGA,KAAK,gBAAgB,QACjB,KAAK,UAAU,SAAS,KAAG,KAAK,oBAAoB,GAEpD,kBAAkB,SAAS,GAAG;AAEjC,UAAM,OAAO,SAAS;AACtB,aAAO,SAAS,MAAS;AACzB,UAAM,KAAK,KAAK,IAAI,eAAe,+BAA+B,GAC5D,KAAK,GAAG,WAAW,IAAI,GACvB,OAAO,GAAG,IAAI,EAAE,GAChB,KAA+B,EAAE,WAAW,EAAE,KAAK,EAAE,GACrD,eAA2C;AAAA,QAChD,UAAU,kBAAkB,IAAI,SAAS;AAAA,MAC1C,GACM,MAAM,MAAM,KAAK,MAAM,4BAA4B;AAAA,QACxD,QAAQ;AAAA,QACR,MAAM,KAAK,UAAU,YAAY;AAAA,QACjC;AAAA,MACD,CAAC;AACD,aAAO,IAAI,EAAE;AAAA,IACd;AAAA,EACD;AAAA,EAEA,sBAAsB;AACrB,QAAM,WAAW,KAAK;AACtB,WAAO,aAAa,MAAS;AAE7B,QAAM,YAAY,SAAS,gBAAgB,oBACrC,eAAe,SAAS,mBAAmB,uBAC3C,gBAAgB,KAAK,UAAU,SAAS;AAE9C,QAAI,KAAK,kBAAkB,QAAW;AAGrC,UAAI,KAAK,cAAc,aAAa,cAAe;AAGnD,WAAK,OAAO,aAAa,KAAK,cAAc,OAAO,GACnD,KAAK,gBAAgB;AAAA,IACtB;AAGA,QAAM,QAAQ,gBAAgB,eAAe,MAAO,GAC9C,UAAU,KAAK,OAAO,WAAW,KAAK,QAAQ,KAAK;AACzD,SAAK,gBAAgB,EAAE,WAAW,UAAU,GAAG,QAAQ;AAAA,EACxD;AAAA,EAEA,SAAS,UAAkC,cAAc,GAAG;AAC3D,aAAW,WAAW,UAAU;AAC/B,UAAM,aAAa,OAAO,gBAAgB,IAAI,WAAW,EAAE,CAAC,GACtD,KAAK,QAAQ,MAAMA,QAAO,KAAK,UAAU,EAAE,SAAS,KAAK,GACzD,YAAY,IAAI,KAAK,QAAQ,aAAa,KAAK,OAAO,IAAI,CAAC,GAC3D,OAAO,YAAY,OAAO,GAC1B,MAAM,IAAI,aAAa,IAAI,WAAW,IAAI,GAE1C,KAAK,MAAM;AAChB,aAAK,UAAU,KAAK,GAAG,GACvB,KAAK,iBAAiB,IAAI,OAC1B,KAAK,oBAAoB;AAAA,MAC1B,GAEM,QAAQ,QAAQ,aAAa;AACnC,WAAK,OAAO,WAAW,IAAI,QAAQ,GAAI;AAAA,IACxC;AAAA,EACD;AAAA,EAEA,4BAA4B;AAC3B,WAAO;AAAA,MACN,cAAc,KAAK,UAAU;AAAA,MAC7B,cAAc,KAAK;AAAA,MACnB,wBAAwB,KAAK,UAAU,CAAC,GAAG,UAAU,QAAQ,KAAK;AAAA,IACnE;AAAA,EACD;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAUA,MAAM,mBACL,KAC2B;AAC3B,QAAM,QAAQ,KAAK,IAAI,cAAc,yBAAyB;AAC9D,QAAI,UAAU,OAAW,QAAO;AAEhC,QAAM,UAAU,IAAI,QAAQ,IAAI,OAAO;AACvC,YAAQ,IAAI,mBAAmB,KAAK,IAAI;AAGxC,QAAM,QAAQ,KAAK,gBAAgB;AACnC,IAAI,UAAU,UAAa,CAAC,QAAQ,IAAI,kBAAkB,KACzD,QAAQ,IAAI,oBAAoB,OAAO,KAAK,CAAC;AAG9C,QAAI;AAIH,UAAM,WAAW,MAAM,MAAM,MAAM,IAAI,KAAK;AAAA,QAC3C,QAAQ,IAAI;AAAA,QACZ;AAAA,QACA,MAAM,MAAM,IAAI,YAAY;AAAA,MAC7B,CAAC;AAED,aAAI,SAAS,WAAW,OACvB,MAAM,KAAK;AAAA,QACV,SAAS;AAAA,QACT,6BAA6B,KAAK,IAAI,MAAM,MAAM,SAAS,KAAK,CAAC;AAAA,MAClE,GACO,QAED;AAAA,IACR,SAAS,GAAG;AACX,mBAAM,KAAK;AAAA,QACV,SAAS;AAAA,QACT,uCAAuC,KAAK,IAAI,wDAAwD,aAAa,QAAQ,EAAE,UAAU,OAAO,CAAC,CAAC;AAAA,MACnJ,GACO;AAAA,IACR;AAAA,EACD;AAAA,EAGA,UAAwB,OAAO,QAAQ;AACtC,QAAM,kBAAkB;AAAA,MACvB,UAAU;AAAA,QACT,SAAS,KAAK,0BAA0B;AAAA,MACzC;AAAA,IACD;AAKA,QADiB,KAAK,mBACL;AAEhB,aADkB,MAAM,KAAK,mBAAmB,GAAG,KAC/B,SAAS,KAAK,eAAe;AAGlD,wBAAoB,IAAI,OAAO;AAC/B,QAAM,cAAc,oBAAoB,IAAI,OAAO,GAC7C,QACL,qBAAqB,IAAI,OAAO,KAAK,KAAK,gBAAgB,eACrD,OAAOA,QAAO,KAAK,MAAM,IAAI,YAAY,CAAC;AAEhD,gBAAK;AAAA,MACJ,CAAC,EAAE,aAAa,WAAW,OAAO,KAAK,CAAC;AAAA,MACxC,KAAK,gBAAgB;AAAA,IACtB,GACO,SAAS,KAAK,eAAe;AAAA,EACrC;AAAA,EAGA,QAAsB,OAAO,QAAQ;AACpC,QAAM,gBAAgB;AAAA,MACrB,UAAU;AAAA,QACT,SAAS,KAAK,0BAA0B;AAAA,MACzC;AAAA,IACD;AAKA,QADiB,KAAK,mBACL;AAEhB,aADkB,MAAM,KAAK,mBAAmB,GAAG,KAC/B,SAAS,KAAK,aAAa;AAOhD,sBAAkB,IAAI,OAAO;AAC7B,QAAM,QACL,qBAAqB,IAAI,OAAO,KAAK,KAAK,gBAAgB,eACrD,OAAO,yBAAyB,MAAM,MAAM,IAAI,KAAK,CAAC;AAE5D,gBAAK,SAAS,KAAK,UAAU,KAAK,GAC3B,SAAS,KAAK,aAAa;AAAA,EACnC;AAAA,EAGA,UAAwB,OAAO,SACvB,SAAS,KAAK,KAAK,0BAA0B,CAAC;AAEvD;AA7DC;AAAA,EADC,KAAK,UAAU;AAAA,GApQJ,kBAqQZ,0BA6BA;AAAA,EADC,KAAK,QAAQ;AAAA,GAjSF,kBAkSZ,wBA6BA;AAAA,EADC,IAAI,UAAU;AAAA,GA9TH,kBA+TZ;",
  "names": ["Buffer", "Buffer"]
}
