Files
kongruity/backend/lib/streams.ts

65 lines
1.7 KiB
TypeScript

import { Transform, type TransformCallback } from 'node:stream';
/**
* Backpressure on readable side is limits how many batches are in flight.
*
* @param size - maximum items per emitted batch
*/
export const batch = <T>(size: number): Transform => {
if (!Number.isInteger(size) || size < 1) {
throw new TypeError('batch(size) requires a positive integer size');
}
let pending: T[] = [];
return new Transform({
objectMode: true,
transform(item: T, _encoding: BufferEncoding, callback: TransformCallback) {
pending.push(item);
if (pending.length < size) {
callback();
return;
}
const full = pending;
pending = [];
callback(null, full);
},
flush(callback: TransformCallback) {
if (pending.length === 0) {
callback();
return;
}
const remainder = pending;
pending = [];
callback(null, remainder);
},
});
};
export const jsonArray = (): Transform => {
let wroteFirst = false;
return new Transform({
writableObjectMode: true,
transform(item: unknown, _encoding: BufferEncoding, callback: TransformCallback) {
let serialized: string;
try {
serialized = JSON.stringify(item);
} catch (err) {
callback(err as Error);
return;
}
const prefix = wroteFirst ? ',' : '[';
wroteFirst = true;
callback(null, prefix + serialized);
},
flush(callback: TransformCallback) {
callback(null, wroteFirst ? ']' : '[]');
},
});
};