{"version":3,"file":"instrumentQueue.js","sources":["../../../../src/instrumentations/worker/instrumentQueue.ts"],"sourcesContent":["import type { ExportedHandler, MessageBatch } from '@cloudflare/workers-types';\nimport type { env as cloudflareEnv, WorkerEntrypoint } from 'cloudflare:workers';\nimport {\n  captureException,\n  SEMANTIC_ATTRIBUTE_SENTRY_OP,\n  SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN,\n  SEMANTIC_ATTRIBUTE_SENTRY_SOURCE,\n  startSpan,\n  withIsolationScope,\n} from '@sentry/core';\nimport type { CloudflareOptions } from '../../client';\nimport { flushAndDispose } from '../../flush';\nimport { ensureInstrumented } from '../../instrument';\nimport { getFinalOptions } from '../../options';\nimport { addCloudResourceContext } from '../../scope-utils';\nimport { init } from '../../sdk';\nimport { instrumentContext } from '../../utils/instrumentContext';\nimport { instrumentEnv } from './instrumentEnv';\n\n/**\n * Core queue handler logic - wraps execution with Sentry instrumentation.\n */\nfunction wrapQueueHandler(\n  batch: MessageBatch,\n  options: CloudflareOptions,\n  context: ExecutionContext,\n  fn: () => unknown,\n): unknown {\n  return withIsolationScope(isolationScope => {\n    const waitUntil = context.waitUntil.bind(context);\n\n    const client = init({ ...options, ctx: context });\n    isolationScope.setClient(client);\n\n    addCloudResourceContext(isolationScope);\n\n    return startSpan(\n      {\n        op: 'faas.queue',\n        name: `process ${batch.queue}`,\n        attributes: {\n          'faas.trigger': 'pubsub',\n          'messaging.destination.name': batch.queue,\n          'messaging.system': 'cloudflare',\n          'messaging.operation.type': 'process',\n          'messaging.operation.name': 'process',\n          'messaging.batch.message_count': batch.messages.length,\n          'messaging.message.retry.count': batch.messages.reduce((acc, message) => acc + message.attempts - 1, 0),\n          [SEMANTIC_ATTRIBUTE_SENTRY_OP]: 'queue.process',\n          [SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN]: 'auto.faas.cloudflare.queue',\n          [SEMANTIC_ATTRIBUTE_SENTRY_SOURCE]: 'task',\n        },\n      },\n      async () => {\n        try {\n          return await fn();\n        } catch (e) {\n          captureException(e, { mechanism: { handled: false, type: 'auto.faas.cloudflare.queue' } });\n          throw e;\n        } finally {\n          waitUntil(flushAndDispose(client));\n        }\n      },\n    );\n  });\n}\n\n/**\n * Instruments a queue handler for ExportedHandler (env/ctx come from args).\n */\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport function instrumentExportedHandlerQueue<T extends ExportedHandler<any, any, any>>(\n  handler: T,\n  optionsCallback: (env: typeof cloudflareEnv) => CloudflareOptions | undefined,\n): void {\n  if (!('queue' in handler) || typeof handler.queue !== 'function') {\n    return;\n  }\n\n  handler.queue = ensureInstrumented(\n    handler.queue,\n    original =>\n      new Proxy(original, {\n        apply(target, thisArg, args: Parameters<NonNullable<T['queue']>>) {\n          const [batch, env, ctx] = args;\n          const context = instrumentContext(ctx);\n          const options = getFinalOptions(optionsCallback(env), env);\n          args[1] = instrumentEnv(env, options);\n          args[2] = context;\n\n          return wrapQueueHandler(batch, options, context, () => target.apply(thisArg, args));\n        },\n      }),\n  );\n}\n\n/**\n * Instruments a queue method for WorkerEntrypoint (options/context already available).\n */\nexport function instrumentWorkerEntrypointQueue<T extends WorkerEntrypoint>(\n  instance: T,\n  options: CloudflareOptions,\n  context: ExecutionContext,\n): void {\n  if (!instance.queue) {\n    return;\n  }\n\n  const original = instance.queue.bind(instance);\n  instance.queue = new Proxy(original, {\n    apply(target, thisArg, args: [MessageBatch]) {\n      const [batch] = args;\n\n      return wrapQueueHandler(batch, options, context, () => Reflect.apply(target, thisArg, args));\n    },\n  });\n}\n"],"names":[],"mappings":";;;;;;;;;AAmBA;AACA;AACA;AACA,SAAS,gBAAgB;AACzB,EAAE,KAAK;AACP,EAAE,OAAO;AACT,EAAE,OAAO;AACT,EAAE,EAAE;AACJ,EAAW;AACX,EAAE,OAAO,kBAAkB,CAAC,cAAA,IAAkB;AAC9C,IAAI,MAAM,SAAA,GAAY,OAAO,CAAC,SAAS,CAAC,IAAI,CAAC,OAAO,CAAC;;AAErD,IAAI,MAAM,MAAA,GAAS,IAAI,CAAC,EAAE,GAAG,OAAO,EAAE,GAAG,EAAE,OAAA,EAAS,CAAC;AACrD,IAAI,cAAc,CAAC,SAAS,CAAC,MAAM,CAAC;;AAEpC,IAAI,uBAAuB,CAAC,cAAc,CAAC;;AAE3C,IAAI,OAAO,SAAS;AACpB,MAAM;AACN,QAAQ,EAAE,EAAE,YAAY;AACxB,QAAQ,IAAI,EAAE,CAAC,QAAQ,EAAE,KAAK,CAAC,KAAK,CAAC,CAAA;AACA,QAAA,UAAA,EAAA;AACA,UAAA,cAAA,EAAA,QAAA;AACA,UAAA,4BAAA,EAAA,KAAA,CAAA,KAAA;AACA,UAAA,kBAAA,EAAA,YAAA;AACA,UAAA,0BAAA,EAAA,SAAA;AACA,UAAA,0BAAA,EAAA,SAAA;AACA,UAAA,+BAAA,EAAA,KAAA,CAAA,QAAA,CAAA,MAAA;AACA,UAAA,+BAAA,EAAA,KAAA,CAAA,QAAA,CAAA,MAAA,CAAA,CAAA,GAAA,EAAA,OAAA,KAAA,GAAA,GAAA,OAAA,CAAA,QAAA,GAAA,CAAA,EAAA,CAAA,CAAA;AACA,UAAA,CAAA,4BAAA,GAAA,eAAA;AACA,UAAA,CAAA,gCAAA,GAAA,4BAAA;AACA,UAAA,CAAA,gCAAA,GAAA,MAAA;AACA,SAAA;AACA,OAAA;AACA,MAAA,YAAA;AACA,QAAA,IAAA;AACA,UAAA,OAAA,MAAA,EAAA,EAAA;AACA,QAAA,CAAA,CAAA,OAAA,CAAA,EAAA;AACA,UAAA,gBAAA,CAAA,CAAA,EAAA,EAAA,SAAA,EAAA,EAAA,OAAA,EAAA,KAAA,EAAA,IAAA,EAAA,4BAAA,EAAA,EAAA,CAAA;AACA,UAAA,MAAA,CAAA;AACA,QAAA,CAAA,SAAA;AACA,UAAA,SAAA,CAAA,eAAA,CAAA,MAAA,CAAA,CAAA;AACA,QAAA;AACA,MAAA,CAAA;AACA,KAAA;AACA,EAAA,CAAA,CAAA;AACA;;AAEA;AACA;AACA;AACA;AACA,SAAA,8BAAA;AACA,EAAA,OAAA;AACA,EAAA,eAAA;AACA,EAAA;AACA,EAAA,IAAA,EAAA,OAAA,IAAA,OAAA,CAAA,IAAA,OAAA,OAAA,CAAA,KAAA,KAAA,UAAA,EAAA;AACA,IAAA;AACA,EAAA;;AAEA,EAAA,OAAA,CAAA,KAAA,GAAA,kBAAA;AACA,IAAA,OAAA,CAAA,KAAA;AACA,IAAA,QAAA;AACA,MAAA,IAAA,KAAA,CAAA,QAAA,EAAA;AACA,QAAA,KAAA,CAAA,MAAA,EAAA,OAAA,EAAA,IAAA,EAAA;AACA,UAAA,MAAA,CAAA,KAAA,EAAA,GAAA,EAAA,GAAA,CAAA,GAAA,IAAA;AACA,UAAA,MAAA,OAAA,GAAA,iBAAA,CAAA,GAAA,CAAA;AACA,UAAA,MAAA,OAAA,GAAA,eAAA,CAAA,eAAA,CAAA,GAAA,CAAA,EAAA,GAAA,CAAA;AACA,UAAA,IAAA,CAAA,CAAA,CAAA,GAAA,aAAA,CAAA,GAAA,EAAA,OAAA,CAAA;AACA,UAAA,IAAA,CAAA,CAAA,CAAA,GAAA,OAAA;;AAEA,UAAA,OAAA,gBAAA,CAAA,KAAA,EAAA,OAAA,EAAA,OAAA,EAAA,MAAA,MAAA,CAAA,KAAA,CAAA,OAAA,EAAA,IAAA,CAAA,CAAA;AACA,QAAA,CAAA;AACA,OAAA,CAAA;AACA,GAAA;AACA;;AAEA;AACA;AACA;AACA,SAAA,+BAAA;AACA,EAAA,QAAA;AACA,EAAA,OAAA;AACA,EAAA,OAAA;AACA,EAAA;AACA,EAAA,IAAA,CAAA,QAAA,CAAA,KAAA,EAAA;AACA,IAAA;AACA,EAAA;;AAEA,EAAA,MAAA,QAAA,GAAA,QAAA,CAAA,KAAA,CAAA,IAAA,CAAA,QAAA,CAAA;AACA,EAAA,QAAA,CAAA,KAAA,GAAA,IAAA,KAAA,CAAA,QAAA,EAAA;AACA,IAAA,KAAA,CAAA,MAAA,EAAA,OAAA,EAAA,IAAA,EAAA;AACA,MAAA,MAAA,CAAA,KAAA,CAAA,GAAA,IAAA;;AAEA,MAAA,OAAA,gBAAA,CAAA,KAAA,EAAA,OAAA,EAAA,OAAA,EAAA,MAAA,OAAA,CAAA,KAAA,CAAA,MAAA,EAAA,OAAA,EAAA,IAAA,CAAA,CAAA;AACA,IAAA,CAAA;AACA,GAAA,CAAA;AACA;;;;"}