Skip to content

Stop keeping a string copy of each consumed message while its handler runs (0.18.1) - #120

Open
yosiat wants to merge 1 commit into
masterfrom
fix/consumer-message-string-retention
Open

yosiat wants to merge 1 commit into
masterfrom
fix/consumer-message-string-retention

Conversation

@yosiat

@yosiat yosiat commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

Summary

_processMessage converted every delivery to a string up front:

const messageString = msg.content.toString();

The string was used by a debug log before the handler and by the error log in
the catch block. Because it is a local of an async function, V8 keeps it
alive for as long as the function is paused at await action.callback(...).
So every in-flight message held a full string copy of its payload for the
whole run of its handler, on top of msg.content itself. A payload with any
non-Latin-1 character is stored as a two-byte string, which doubles the copy.

With large messages, a high prefetch and slow handlers, these copies add up
to prefetch x payload size of heap that nothing uses.

Changes

  • Remove the debug log of the full payload. The pluggable logger has no way
    to check whether debug is enabled, so the copy was paid on every message.
  • The error log now builds the payload string inside the catch block and
    cuts it to the first 8192 bytes: msg.content.toString('utf8', 0, 8192).
  • Version bump to 0.18.1 (patch).

No public API change. The error log keeps the queue name and the error; only
params.message is cut.

How it was checked

A standalone script drives the real Consumer with a fake channel (no
broker). It subscribes a handler that never finishes and delivers 200
messages of 1 MB each. The payload is JSON padded with whitespace, so the
raw message is large while the parsed body is tiny: any heap growth comes
from a string copy of the payload. Heap measured after GC, Node 22:

payload before this change after this change
1 MB, ASCII only 200 MB 0 MB
1 MB, one non-Latin-1 char 400 MB 0 MB
Script
// node --expose-gc consumer-heap.js <path to src/modules/consumer.js>
const Consumer = require(process.argv[2]);

const MESSAGES = 200;
const payload = Buffer.from(`{"id":1}${' '.repeat(1_000_000)}`);
// Kept reachable, like a handler waiting on a socket or a timer.
const handlerDone = new Promise(() => {});

let deliver;
const channel = {
  addListener() {},
  removeListener() {},
  assertQueue: async () => ({}),
  consume: async (queue, onMessage) => {
    deliver = onMessage;
    return { consumerTag: 'tag' };
  },
};
const connection = {
  config: { prefetch: MESSAGES, timeout: 1000, requeue: true, consumerSuffix: '' },
  isClosed: false,
  getChannel: async () => channel,
};
const heapMb = () => {
  global.gc();
  return process.memoryUsage().heapUsed / 1e6;
};

(async () => {
  const consumer = new Consumer(connection);
  await consumer.subscribe('queue', () => handlerDone);
  const before = heapMb();
  for (let i = 0; i < MESSAGES; i += 1) {
    deliver({ content: payload, properties: { contentType: 'application/json', headers: {} }, fields: {} });
  }
  await new Promise((resolve) => setImmediate(resolve));
  console.log(`heap kept after GC: ${(heapMb() - before).toFixed(0)} MB`);
})();

_processMessage built `msg.content.toString()` up front and kept it as a
local for the debug log and the error log. V8 keeps every local of a
paused async function alive, so each in-flight message held a full
string copy of its payload for as long as the handler ran. With large
messages and slow handlers, this copy grows with prefetch and can
dominate the heap.

- Drop the debug log of the full payload. The logger contract has no
  way to ask whether debug is enabled, so it cost a copy every time.
- The error log builds the payload string inside the catch block and
  cuts it to the first 8192 bytes.
@yosiat
yosiat marked this pull request as ready for review October 3, 2026 14:05
@yosiat
yosiat requested review from ramhr and shamil as code owners October 3, 2026 14:05
Copilot AI balanced review requested due to automatic review settings October 3, 2026 14:05
Comment thread src/modules/consumer.js
async _processMessage(channel, subscription, msg) {
const { queue, callback } = subscription;
const messageString = msg.content.toString();
logger.debug({

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I can add support to logger interface to do check for debug enabled, but in reality we never look at debug logs here.
If you think it's crucial I'll do so.

Copilot AI 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.

Copilot review overview

🟡 Changes recommended

The core memory and truncation behavior lacks automated regression coverage.

Review effort: Balanced
Findings: 1 Medium severity

Open (1)
What changed in this PR

Reduces memory retained while consumer handlers run by deferring and limiting payload stringification.

Changes:

  • Removes full-payload debug logging.
  • Truncates error-log payloads to 8192 bytes.
  • Bumps the package version to 0.18.1.
File Description
src/​modules/​consumer.js Avoids persistent payload string copies during processing.
package.json Bumps the patch version.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread src/modules/consumer.js
params: { queue, message: messageString },
// Built here and cut, never kept as a local: V8 keeps every local of this async function
// alive for as long as the handler is awaited.
params: { queue, message: msg.content.toString('utf8', 0, ERROR_LOG_PAYLOAD_BYTES) },
@yosiat yosiat unassigned ramhr Oct 4, 2026
@shamil shamil assigned yosiat and unassigned shamil Oct 4, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Development

Successfully merging this pull request may close these issues.

4 participants