import { Transform } from 'node:stream'; /** * Groups an object-mode stream into fixed-size arrays. Backpressure on the * readable side is what limits how many batches are ever in flight. * * @param {number} size - maximum items per emitted batch * @returns {Transform} */ export const batch = (size) => { if (!Number.isInteger(size) || size < 1) { throw new TypeError('batch(size) requires a positive integer size'); } let pending = []; return new Transform({ objectMode: true, transform(item, _encoding, callback) { pending.push(item); if (pending.length < size) { callback(); return; } const full = pending; pending = []; callback(null, full); }, flush(callback) { if (pending.length === 0) { callback(); return; } const remainder = pending; pending = []; callback(null, remainder); }, }); }; /** * Serializes an object-mode stream into a JSON array, one element at a time, * so no complete copy of the payload is ever held in memory. * * @returns {Transform} */ export const jsonArray = () => { let wroteFirst = false; return new Transform({ writableObjectMode: true, transform(item, _encoding, callback) { let serialized; try { serialized = JSON.stringify(item); } catch (err) { callback(err); return; } const prefix = wroteFirst ? ',' : '['; wroteFirst = true; callback(null, prefix + serialized); }, flush(callback) { callback(null, wroteFirst ? ']' : '[]'); }, }); };