Skip to content

fix(server-utils): Record amqplib consumer info before the first delivery and instrument the callback API - #25047

Open
s1gr1d wants to merge 6 commits into
developfrom
sig/amqplib-fix-consumer-info
Open

s1gr1d wants to merge 6 commits into
developfrom
sig/amqplib-fix-consumer-info

Conversation

@s1gr1d

@s1gr1d s1gr1d commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

When a worker subscribes to a queue that already has messages, amqplib can decode BasicConsumeOk and the first message in one socket read. We recorded each consumer's queue and noAck flag when consume() resolved. That happens too late:

one socket read, decoded synchronously
│
├─ BasicConsumeOk ─► registerConsumer(tag)         ◄─ now: info recorded here
├─ BasicDeliver   ─► dispatchMessage(tag)          ◄─ lookup missed here
│
└─ decode ends ─► microtasks ─► consume().then()   ◄─ before: info recorded here

The first message's span got the routing key as its name, and noAck spans ended as errors at channel close. Now we record the info in registerConsumer, which amqplib always calls before dispatch.

This PR also instruments the callback API (lib/callback_model.js). It has the same classes and signatures as the promise API.

Follow-up on this PR (noticed the bug while working on this): #23652

@s1gr1d
s1gr1d requested a review from a team as a code owner October 5, 2026 11:29
@s1gr1d
s1gr1d requested review from JPeer264, isaacs and mydea and removed request for a team October 5, 2026 11:29
@s1gr1d

s1gr1d commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

bugbot run

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/server-utils/src/integrations/amqplib.ts Outdated
@github-actions

github-actions Bot commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

size-limit report 📦

Path Size % Change Change
@sentry/browser 29.6 kB - -
@sentry/browser - with treeshaking flags 27.75 kB - -
@sentry/browser - with treeshaking flags tracing without tracing 27.65 kB - -
@sentry/browser (incl. Tracing) 51.52 kB - -
@sentry/browser (incl. Tracing + Span Streaming) 51.52 kB - -
@sentry/browser (incl. Tracing, Profiling) 54.5 kB - -
@sentry/browser (incl. Tracing, Replay) 91.23 kB - -
@sentry/browser (incl. Tracing, Replay) - with treeshaking flags 80.18 kB - -
@sentry/browser (incl. Tracing, Replay with Canvas) 95.93 kB - -
@sentry/browser (incl. Tracing, Replay, Feedback) 108.89 kB - -
@sentry/browser (incl. Feedback) 47.12 kB - -
@sentry/browser (incl. sendFeedback) 34.65 kB - -
@sentry/browser (incl. FeedbackAsync) 39.76 kB - -
@sentry/browser (incl. Metrics) 30.61 kB - -
@sentry/browser (incl. Logs) 30.89 kB - -
@sentry/browser (incl. Metrics & Logs) 31.55 kB - -
@sentry/react 31.43 kB - -
@sentry/react (incl. Tracing) 53.84 kB - -
@sentry/vue 37.55 kB - -
@sentry/vue (incl. Tracing) 54.4 kB - -
@sentry/svelte 29.63 kB - -
@sentry/remix (Remix 3 client bundle) 56.54 kB - -
CDN Bundle 31.33 kB - -
CDN Bundle (incl. Tracing) 52.07 kB - -
CDN Bundle (incl. Logs, Metrics) 33.56 kB - -
CDN Bundle (incl. Tracing, Logs, Metrics) 54.03 kB - -
CDN Bundle (incl. Replay, Logs, Metrics) 74.38 kB - -
CDN Bundle (incl. Tracing, Replay) 89.74 kB - -
CDN Bundle (incl. Tracing, Replay, Logs, Metrics) 91.69 kB - -
CDN Bundle (incl. Tracing, Replay, Feedback) 95.9 kB - -
CDN Bundle (incl. Tracing, Replay, Feedback, Logs, Metrics) 97.87 kB - -
CDN Bundle - uncompressed 92.46 kB - -
CDN Bundle (incl. Tracing) - uncompressed 154.77 kB - -
CDN Bundle (incl. Logs, Metrics) - uncompressed 99.04 kB - -
CDN Bundle (incl. Tracing, Logs, Metrics) - uncompressed 160.72 kB - -
CDN Bundle (incl. Replay, Logs, Metrics) - uncompressed 228.98 kB - -
CDN Bundle (incl. Tracing, Replay) - uncompressed 274.89 kB - -
CDN Bundle (incl. Tracing, Replay, Logs, Metrics) - uncompressed 280.83 kB - -
CDN Bundle (incl. Tracing, Replay, Feedback) - uncompressed 288.59 kB - -
CDN Bundle (incl. Tracing, Replay, Feedback, Logs, Metrics) - uncompressed 294.52 kB - -
@sentry/nextjs (client) 56.19 kB - -
@sentry/sveltekit (client) 51.9 kB - -
@sentry/core/server 40.52 kB - -
@sentry/core/browser 13.51 kB - -
@sentry/node 145.16 kB +0.11% +158 B 🔺
@sentry/node/import (ESM hook with diagnostics-channel injection) 83.25 kB +0.03% +24 B 🔺
@sentry/node - without tracing 93.49 kB +0.04% +37 B 🔺
@sentry/node - without channel injection 123.31 kB +0.12% +145 B 🔺
@sentry/aws-serverless 101.72 kB +0.03% +29 B 🔺
@sentry/cloudflare (withSentry) - minified 208.95 kB - -
@sentry/cloudflare (withSentry) 517.52 kB - -

View base workflow run

@isaacs isaacs left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'd recommend adding the error handler, but this looks good. Even if all the suggestions get put off to follow-up PRs, it's a net improvement.

const pending = consumerChannel[CHANNEL_PENDING_CONSUMERS] ?? [];
// A failed `consume` leaves its entry behind, so pair by callback, not by position.
const index = pending.findIndex(entry => entry.callback === callback);
const entry = pending[index];

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A failed consume on a live channel can poison a later consume that reuses the same callback.

(Also relevant to line 227 above)

Pairing by callback avoids the problem only when the next consume uses a different callback. Sharing one handler across several queues is a common pattern.

So the sequence is:

  1. channel.consume(badName, handler) rejects, but the channel stays open. This happens for client-side encode errors, for example a queue name or consumerTag over 255 bytes (defs.js throws TypeError: Field 'queue' is the wrong type). Server-side errors (404, 403) close the channel, so they are harmless here.
  2. channel.consume('jobs', handler, { noAck: true }) succeeds.
  3. registerConsumer finds the stale entry first (findIndex returns the oldest match) and records { queue: badName, noAck: false } for the jobs tag.

Every message on jobs then gets the wrong static span name, and its spans stay open until the 1-minute timeout or channel close, then end as errors.

That's the same bug being fixed here, but narrower.

Suggestion: drop the entry when the consume call fails. The tracing-channel context object is the same object for start and error, so keep a reference to the entry on it:

channel.start.subscribe(message => {
  // ...
  const entry = { callback, info: { noAck: !!options?.noAck, queue } };
  consumerChannel[CHANNEL_PENDING_CONSUMERS]?.push(entry);
  (data as AmqpConsumeContext)._sentryPendingConsumer = entry;
});

channel.error.subscribe(message => {
  const data = message as AmqpConsumeContext;
  const pending = data.self?.[CHANNEL_PENDING_CONSUMERS];
  const index = pending?.indexOf(data._sentryPendingConsumer!) ?? -1;
  if (index >= 0) pending!.splice(index, 1);
});

// `consume` knows the queue and `noAck`, and `registerConsumer` knows the tag. Together they tell the
// dispatch hook how to name the consumer span and when to end it.
{
channelName: 'consume',

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The callback API is still affected by the bug, right? Only lib/channel_model.js's Channel.consume is hooked, and lib/callback_model.js has its own Channel.consume (it also passes callback unchanged to registerConsumer).

Dispatch runs through the shared BaseChannel, so callback-API consumers get consumer spans but never get consumer info. Their spans always use the routing key as the static name, and noAck spans always leak until close and end as errors.

This PR makes the fix cheap: subscribeConsume uses only start, so one more config entry for lib/callback_model.js with kind: 'Sync' would be enough, and the new registerConsumer hook already covers both APIs.

Could be a follow-up PR (with an integration scenario using amqplib/callback_api), or include it here since the hook is shared now.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I did it in this PR and added a test.

The classes and signatures from callback_model.js match channel_model.js, so the files share one flatMap and also the same channels.

@s1gr1d s1gr1d changed the title fix(server-utils): Record amqplib consumer info before the first delivery fix(server-utils): Record amqplib consumer info before the first delivery and instrument the callback API Oct 6, 2026
Comment thread packages/server-utils/src/integrations/amqplib.ts
@s1gr1d

s1gr1d commented Oct 6, 2026

Copy link
Copy Markdown
Member Author

bugbot run

@cursor cursor Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cursor Bugbot has reviewed your changes and found 2 potential issues.

Fix All in Cursor

❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.

Want reviews to match your repository better? Bugbot Learning can learn team-specific rules from PR activity. A team admin can enable Learning in the Cursor dashboard.

Reviewed by Cursor Bugbot for commit 2b5a531. Configure here.

Comment thread packages/server-utils/src/orchestrion/config/amqplib.ts
Comment thread dev-packages/node-integration-tests/suites/tracing/amqplib/test.ts

@JPeer264 JPeer264 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM + what isaac mentioned

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants