{"version":3,"file":"instrumentQueueProducer.js","sources":["../../../../src/instrumentations/worker/instrumentQueueProducer.ts"],"sourcesContent":["import type { MessageSendRequest, Queue, QueueSendBatchOptions, QueueSendOptions } from '@cloudflare/workers-types';\nimport { SEMANTIC_ATTRIBUTE_SENTRY_OP, SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, startSpan } from '@sentry/core';\n\nconst ORIGIN = 'auto.faas.cloudflare.queue';\n\nfunction startPublishSpan(\n options: {\n bindingName: string;\n bodySize: number | undefined;\n messageCount?: number;\n },\n callback: () => T,\n): T {\n const { bindingName, bodySize, messageCount } = options;\n\n return startSpan(\n {\n op: 'queue.publish',\n name: `send ${bindingName}`,\n attributes: {\n 'messaging.system': 'cloudflare',\n 'messaging.destination.name': bindingName,\n 'messaging.operation.type': 'send',\n 'messaging.operation.name': 'send',\n ...(messageCount !== undefined && { 'messaging.batch.message_count': messageCount }),\n 'messaging.message.body.size': bodySize,\n [SEMANTIC_ATTRIBUTE_SENTRY_OP]: 'queue.publish',\n [SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN]: ORIGIN,\n },\n },\n callback,\n );\n}\n\nfunction getBodySize(body: unknown): number | undefined {\n if (body == null) {\n return undefined;\n }\n\n if (typeof body === 'string') {\n return new TextEncoder().encode(body).byteLength;\n }\n\n if (body instanceof ArrayBuffer) {\n return body.byteLength;\n }\n\n if (ArrayBuffer.isView(body)) {\n return body.byteLength;\n }\n\n try {\n return new TextEncoder().encode(JSON.stringify(body)).byteLength;\n } catch {\n return undefined;\n }\n}\n\n/**\n * Wraps a Queue producer binding to create `queue.publish` spans on\n * `send` and `sendBatch` calls.\n *\n * The queue's own name is not available on the binding object, so we use\n * the env binding key (e.g. `MY_QUEUE`) as `messaging.destination.name`.\n */\nexport function instrumentQueueProducer(queue: T, bindingName: string): T {\n return new Proxy(queue, {\n get(target, prop, receiver) {\n if (prop === 'send') {\n const original = Reflect.get(target, prop, receiver) as Queue['send'];\n\n return function (this: unknown, message: unknown, options?: QueueSendOptions): Promise {\n return startPublishSpan({ bindingName, bodySize: getBodySize(message) }, () =>\n Reflect.apply(original, target, [message, options]),\n );\n };\n }\n\n if (prop === 'sendBatch') {\n const original = Reflect.get(target, prop, receiver) as Queue['sendBatch'];\n return function (\n this: unknown,\n messages: Iterable,\n options?: QueueSendBatchOptions,\n ): Promise {\n const messageArray = Array.from(messages);\n const totalBodySize = messageArray.reduce((acc, m) => {\n const size = getBodySize(m.body);\n if (size === undefined) {\n return acc;\n }\n return (acc ?? 0) + size;\n }, undefined);\n\n return startPublishSpan({ bindingName, bodySize: totalBodySize, messageCount: messageArray.length }, () =>\n Reflect.apply(original, target, [messageArray, options]),\n );\n };\n }\n\n return Reflect.get(target, prop, receiver);\n },\n });\n}\n"],"names":[],"mappings":";;AAGA,MAAM,MAAA,GAAS,4BAAA;AAEf,SAAS,gBAAA,CACP,SAKA,QAAA,EACG;AACH,EAAA,MAAM,EAAE,WAAA,EAAa,QAAA,EAAU,YAAA,EAAa,GAAI,OAAA;AAEhD,EAAA,OAAO,SAAA;AAAA,IACL;AAAA,MACE,EAAA,EAAI,eAAA;AAAA,MACJ,IAAA,EAAM,QAAQ,WAAW,CAAA,CAAA;AAAA,MACzB,UAAA,EAAY;AAAA,QACV,kBAAA,EAAoB,YAAA;AAAA,QACpB,4BAAA,EAA8B,WAAA;AAAA,QAC9B,0BAAA,EAA4B,MAAA;AAAA,QAC5B,0BAAA,EAA4B,MAAA;AAAA,QAC5B,GAAI,YAAA,KAAiB,MAAA,IAAa,EAAE,iCAAiC,YAAA,EAAa;AAAA,QAClF,6BAAA,EAA+B,QAAA;AAAA,QAC/B,CAAC,4BAA4B,GAAG,eAAA;AAAA,QAChC,CAAC,gCAAgC,GAAG;AAAA;AACtC,KACF;AAAA,IACA;AAAA,GACF;AACF;AAEA,SAAS,YAAY,IAAA,EAAmC;AACtD,EAAA,IAAI,QAAQ,IAAA,EAAM;AAChB,IAAA,OAAO,MAAA;AAAA,EACT;AAEA,EAAA,IAAI,OAAO,SAAS,QAAA,EAAU;AAC5B,IAAA,OAAO,IAAI,WAAA,EAAY,CAAE,MAAA,CAAO,IAAI,CAAA,CAAE,UAAA;AAAA,EACxC;AAEA,EAAA,IAAI,gBAAgB,WAAA,EAAa;AAC/B,IAAA,OAAO,IAAA,CAAK,UAAA;AAAA,EACd;AAEA,EAAA,IAAI,WAAA,CAAY,MAAA,CAAO,IAAI,CAAA,EAAG;AAC5B,IAAA,OAAO,IAAA,CAAK,UAAA;AAAA,EACd;AAEA,EAAA,IAAI;AACF,IAAA,OAAO,IAAI,aAAY,CAAE,MAAA,CAAO,KAAK,SAAA,CAAU,IAAI,CAAC,CAAA,CAAE,UAAA;AAAA,EACxD,CAAA,CAAA,MAAQ;AACN,IAAA,OAAO,MAAA;AAAA,EACT;AACF;AASO,SAAS,uBAAA,CAAyC,OAAU,WAAA,EAAwB;AACzF,EAAA,OAAO,IAAI,MAAM,KAAA,EAAO;AAAA,IACtB,GAAA,CAAI,MAAA,EAAQ,IAAA,EAAM,QAAA,EAAU;AAC1B,MAAA,IAAI,SAAS,MAAA,EAAQ;AACnB,QAAA,MAAM,QAAA,GAAW,OAAA,CAAQ,GAAA,CAAI,MAAA,EAAQ,MAAM,QAAQ,CAAA;AAEnD,QAAA,OAAO,SAAyB,SAAkB,OAAA,EAA2C;AAC3F,UAAA,OAAO,gBAAA;AAAA,YAAiB,EAAE,WAAA,EAAa,QAAA,EAAU,WAAA,CAAY,OAAO,CAAA,EAAE;AAAA,YAAG,MACvE,QAAQ,KAAA,CAAM,QAAA,EAAU,QAAQ,CAAC,OAAA,EAAS,OAAO,CAAC;AAAA,WACpD;AAAA,QACF,CAAA;AAAA,MACF;AAEA,MAAA,IAAI,SAAS,WAAA,EAAa;AACxB,QAAA,MAAM,QAAA,GAAW,OAAA,CAAQ,GAAA,CAAI,MAAA,EAAQ,MAAM,QAAQ,CAAA;AACnD,QAAA,OAAO,SAEL,UACA,OAAA,EACe;AACf,UAAA,MAAM,YAAA,GAAe,KAAA,CAAM,IAAA,CAAK,QAAQ,CAAA;AACxC,UAAA,MAAM,aAAA,GAAgB,YAAA,CAAa,MAAA,CAA2B,CAAC,KAAK,CAAA,KAAM;AACxE,YAAA,MAAM,IAAA,GAAO,WAAA,CAAY,CAAA,CAAE,IAAI,CAAA;AAC/B,YAAA,IAAI,SAAS,MAAA,EAAW;AACtB,cAAA,OAAO,GAAA;AAAA,YACT;AACA,YAAA,OAAA,CAAQ,OAAO,CAAA,IAAK,IAAA;AAAA,UACtB,GAAG,MAAS,CAAA;AAEZ,UAAA,OAAO,gBAAA;AAAA,YAAiB,EAAE,WAAA,EAAa,QAAA,EAAU,aAAA,EAAe,YAAA,EAAc,aAAa,MAAA,EAAO;AAAA,YAAG,MACnG,QAAQ,KAAA,CAAM,QAAA,EAAU,QAAQ,CAAC,YAAA,EAAc,OAAO,CAAC;AAAA,WACzD;AAAA,QACF,CAAA;AAAA,MACF;AAEA,MAAA,OAAO,OAAA,CAAQ,GAAA,CAAI,MAAA,EAAQ,IAAA,EAAM,QAAQ,CAAA;AAAA,IAC3C;AAAA,GACD,CAAA;AACH;;;;"}