File size: 5,154 Bytes
3144483 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 | import { Writable, type Readable } from "node:stream";
import { finished } from "node:stream/promises";
import { toErrorObject } from "../infra/errors.js";
import { createWindowsOutputDecoder } from "../infra/windows-encoding.js";
import { createDeferredCore } from "../shared/deferred.js";
export function onDecodedOutput(
stream: Readable,
listener: (chunk: string) => void,
onRaw?: (chunk: Buffer) => void,
): () => void {
const decoder = createWindowsOutputDecoder();
const emit = (text: string) => {
if (text) {
listener(text);
}
};
let flushed = false;
const flush = () => {
if (flushed) {
return;
}
flushed = true;
emit(decoder.flush());
};
const onData = (chunk: Buffer | string) => {
onRaw?.(typeof chunk === "string" ? Buffer.from(chunk) : chunk);
emit(decoder.decode(chunk));
};
stream.on("data", onData);
stream.once("end", flush);
stream.once("close", flush);
return () => {
// A queued close may still invoke its copied listener, so suppress flush before detaching.
flushed = true;
stream.off("data", onData);
stream.off("end", flush);
stream.off("close", flush);
};
}
/** One consumer holds native pipe backpressure until each decoded chunk settles. */
export function createAwaitedDecodedOutput(stream: Readable, onFailure: (error: unknown) => void) {
type Consumer = (chunk: string) => void | Promise<void>;
const subscription = createDeferredCore<Consumer | undefined>();
let subscribed = false;
let closed = false;
let decoder: ReturnType<typeof createWindowsOutputDecoder> | undefined;
let inFlight = Promise.resolve();
const deliver = (read: () => string, callback: (error?: Error | null) => void) => {
inFlight = (async () => {
const listener = await subscription.promise;
if (closed || !listener) {
return;
}
const text = read();
if (text) {
await listener(text);
}
})();
void inFlight.then(
() => callback(),
(error: unknown) => callback(toErrorObject(error, "Process stdout consumption failed")),
);
};
const sink = new Writable({
write(chunk: Buffer | string, _encoding, callback) {
deliver(() => (decoder ??= createWindowsOutputDecoder()).decode(chunk), callback);
},
final(callback) {
deliver(() => decoder?.flush() ?? "", callback);
},
});
const onSourceError = (error: Error) => sink.destroy(error);
const onSourceClose = () => {
if (stream.readableEnded) {
sink.end();
} else {
sink.destroy(new Error("Process stdout closed before EOF"));
}
};
stream.once("error", onSourceError);
stream.once("close", onSourceClose);
const done = (async () => {
try {
await finished(sink, { cleanup: true });
} catch (error) {
if (!closed) {
closed = true;
try {
onFailure(error);
} catch (stopError) {
throw new AggregateError(
[error, stopError],
"Process stdout consumption and stop failed",
{
cause: stopError,
},
);
} finally {
// Stop owns the failure; discard remaining bytes so the native pipe can reach EOF.
stream.unpipe(sink);
stream.resume();
}
}
throw error;
} finally {
stream.unpipe(sink);
subscription.resolve(undefined);
// Destroy may settle the stream before its accepted consumer has returned.
await inFlight.catch(() => undefined);
stream.off("error", onSourceError);
stream.off("close", onSourceClose);
}
})();
void done.catch(() => undefined);
// Node resumes child pipes on exit, so own writes before the caller subscribes.
stream.pipe(sink);
if (stream.errored) {
sink.destroy(stream.errored);
} else if (stream.destroyed) {
onSourceClose();
}
const consume = (listener: Consumer): Promise<void> => {
if (closed || subscribed) {
return Promise.reject(new Error("Process stdout consumption is already owned or closed"));
}
subscribed = true;
subscription.resolve(listener);
return done;
};
return {
consume,
drain: () => (subscribed || closed ? done : consume(() => {})),
close: () => {
closed = true;
subscription.resolve(undefined);
if (!sink.writableFinished) {
stream.unpipe(sink);
sink.destroy(new Error("Process stdout consumption closed"));
}
},
};
}
/** Output failure requests stop separately; both completion paths settle before returning. */
export async function joinProcessCompletionAndOutput<T>(
completion: Promise<T>,
output: Promise<void>,
): Promise<T> {
const [outcome, consumed] = await Promise.allSettled([completion, output]);
if (outcome.status === "rejected") {
if (consumed.status === "rejected") {
throw new AggregateError(
[outcome.reason, consumed.reason],
"Process and stdout consumption failed",
);
}
throw outcome.reason;
}
if (consumed.status === "rejected") {
throw consumed.reason;
}
return outcome.value;
}
|