Skip to content

Commit 0e5c5dd

Browse files
committed
Allow stream sinks to keep streams open
Add StreamSinkOptions.closeStream with a default of true so callers can preserve caller-owned writable streams during sink disposal. Disposal still waits for pending and buffered writes and releases the writer lock. Document the option and cover blocking, non-blocking, and active-flush disposal behavior. Closes #203 Assisted-by: Codex:gpt-5.6-sol
1 parent d289a98 commit 0e5c5dd

5 files changed

Lines changed: 198 additions & 6 deletions

File tree

CHANGES.md

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,13 @@ Version 2.4.0
88

99
To be released.
1010

11+
### @logtape/logtape
12+
13+
- Added the `StreamSinkOptions.closeStream` option to dispose stream sinks
14+
without closing caller-owned streams. [[#203]]
15+
16+
[#203]: https://github.com/dahlia/logtape/issues/203
17+
1118

1219
Version 2.3.1
1320
-------------
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
- Added the `StreamSinkOptions.closeStream` option to dispose stream sinks
2+
without closing caller-owned streams. [[#203]]

docs/manual/sinks.md

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,26 @@ await configure({
184184
> Node.js stream. You can use [`Writable.toWeb()`] method to convert a Node.js
185185
> stream to a `WritableStream`.
186186
187+
By default, disposing a stream sink closes its `WritableStream`. For a
188+
caller-owned stream that needs to remain open, such as a Node.js standard
189+
stream, set `closeStream` to `false`:
190+
191+
~~~~ typescript twoslash
192+
// @noErrors: 2345
193+
import "@types/node";
194+
import { getStreamSink } from "@logtape/logtape";
195+
import stream from "node:stream";
196+
197+
const stderrSink = getStreamSink(
198+
stream.Writable.toWeb(process.stderr),
199+
{ closeStream: false },
200+
);
201+
~~~~
202+
203+
Disposing the sink still waits for pending writes, flushes buffered records
204+
when non-blocking mode is enabled, and releases the writer lock without closing
205+
the stream.
206+
187207
See also `getStreamSink()` function and `StreamSinkOptions` interface
188208
in the API reference for more details.
189209

packages/logtape/src/sink.test.ts

Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -109,13 +109,17 @@ interface ConsoleMock extends Console {
109109

110110
test("getStreamSink()", async () => {
111111
let buffer: string = "";
112+
let closed = false;
112113
const decoder = new TextDecoder();
113114
const sink = getStreamSink(
114115
new WritableStream({
115116
write(chunk: Uint8Array) {
116117
buffer += decoder.decode(chunk);
117118
return Promise.resolve();
118119
},
120+
close() {
121+
closed = true;
122+
},
119123
}),
120124
);
121125
sink(trace);
@@ -136,6 +140,40 @@ test("getStreamSink()", async () => {
136140
2023-11-14 22:13:20.000 +00:00 [FTL] my-app·junk: Hello, 123 & 456!
137141
`,
138142
);
143+
assert.strictEqual(closed, true);
144+
});
145+
146+
test("getStreamSink() with closeStream: false", async () => {
147+
// Arrange
148+
let buffer = "";
149+
let closed = false;
150+
const decoder = new TextDecoder();
151+
const stream = new WritableStream({
152+
write(chunk: Uint8Array) {
153+
buffer += decoder.decode(chunk);
154+
},
155+
close() {
156+
closed = true;
157+
},
158+
});
159+
const sink = getStreamSink(stream, { closeStream: false });
160+
161+
// Act
162+
sink(info);
163+
await sink[Symbol.asyncDispose]();
164+
165+
// Assert
166+
assert.strictEqual(
167+
buffer,
168+
"2023-11-14 22:13:20.000 +00:00 [INF] my-app·junk: " +
169+
"Hello, 123 & 456!\n",
170+
);
171+
assert.strictEqual(closed, false);
172+
const writer = stream.getWriter();
173+
await writer.write(new TextEncoder().encode("after disposal"));
174+
await writer.close();
175+
writer.releaseLock();
176+
assert.ok(buffer.endsWith("after disposal"));
139177
});
140178

141179
test("getStreamSink() with nonBlocking - simple boolean", async () => {
@@ -300,6 +338,90 @@ test("getStreamSink() with nonBlocking - flush on dispose", async () => {
300338
);
301339
});
302340

341+
test(
342+
"getStreamSink() with nonBlocking and closeStream: false",
343+
async () => {
344+
// Arrange
345+
let buffer = "";
346+
let closed = false;
347+
const decoder = new TextDecoder();
348+
const stream = new WritableStream({
349+
write(chunk: Uint8Array) {
350+
buffer += decoder.decode(chunk);
351+
},
352+
close() {
353+
closed = true;
354+
},
355+
});
356+
const sink = getStreamSink(stream, {
357+
closeStream: false,
358+
nonBlocking: {
359+
bufferSize: 100,
360+
flushInterval: 5000,
361+
},
362+
});
363+
364+
// Act
365+
sink(info);
366+
await sink[Symbol.asyncDispose]();
367+
368+
// Assert
369+
assert.strictEqual(
370+
buffer,
371+
"2023-11-14 22:13:20.000 +00:00 [INF] my-app·junk: " +
372+
"Hello, 123 & 456!\n",
373+
);
374+
assert.strictEqual(closed, false);
375+
const writer = stream.getWriter();
376+
await writer.close();
377+
writer.releaseLock();
378+
},
379+
);
380+
381+
test(
382+
"getStreamSink() with nonBlocking waits for an active flush on dispose",
383+
async () => {
384+
// Arrange
385+
let markWriteStarted: () => void = () => {};
386+
let finishWrite: () => void = () => {};
387+
const writeStarted = new Promise<void>((resolve) => {
388+
markWriteStarted = resolve;
389+
});
390+
const writeCanFinish = new Promise<void>((resolve) => {
391+
finishWrite = resolve;
392+
});
393+
const stream = new WritableStream({
394+
async write() {
395+
markWriteStarted();
396+
await writeCanFinish;
397+
},
398+
});
399+
const sink = getStreamSink(stream, {
400+
closeStream: false,
401+
nonBlocking: { bufferSize: 1 },
402+
});
403+
404+
// Act
405+
sink(info);
406+
await beforeDeadline(writeStarted, "The buffered write did not start");
407+
let disposed = false;
408+
const disposePromise = Promise.resolve(sink[Symbol.asyncDispose]()).then(
409+
() => {
410+
disposed = true;
411+
},
412+
);
413+
await delay(0);
414+
415+
// Assert
416+
assert.strictEqual(disposed, false);
417+
finishWrite();
418+
await beforeDeadline(disposePromise, "The stream sink was not disposed");
419+
const writer = stream.getWriter();
420+
await writer.close();
421+
writer.releaseLock();
422+
},
423+
);
424+
303425
test("getStreamSink() with nonBlocking - buffer overflow protection", async () => {
304426
let buffer: string = "";
305427
const decoder = new TextDecoder();
@@ -3292,3 +3414,22 @@ function recordWithLevel(level: LogLevel): LogRecord {
32923414
properties: {},
32933415
};
32943416
}
3417+
3418+
async function beforeDeadline<T>(
3419+
promise: Promise<T>,
3420+
message: string,
3421+
): Promise<T> {
3422+
let timeoutId: ReturnType<typeof globalThis.setTimeout> | undefined;
3423+
try {
3424+
return await Promise.race([
3425+
promise,
3426+
new Promise<never>((_, reject) => {
3427+
timeoutId = globalThis.setTimeout(() => {
3428+
reject(new Error(message));
3429+
}, 1000);
3430+
}),
3431+
]);
3432+
} finally {
3433+
if (timeoutId !== undefined) globalThis.clearTimeout(timeoutId);
3434+
}
3435+
}

packages/logtape/src/sink.ts

Lines changed: 28 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,16 @@ export interface StreamSinkOptions {
8989
*/
9090
encoder?: { encode(text: string): Uint8Array };
9191

92+
/**
93+
* Whether to close the stream when the sink is disposed. Set this to
94+
* `false` for caller-owned streams that need to remain open after disposal.
95+
* The sink still waits for pending writes and releases its writer lock.
96+
*
97+
* @default `true`
98+
* @since 2.4.0
99+
*/
100+
closeStream?: boolean;
101+
92102
/**
93103
* Enable non-blocking mode with optional buffer configuration.
94104
* When enabled, log records are buffered and flushed in the background.
@@ -156,6 +166,7 @@ export function getStreamSink(
156166
): Sink & AsyncDisposable {
157167
const formatter = options.formatter ?? defaultTextFormatter;
158168
const encoder = options.encoder ?? new TextEncoder();
169+
const closeStream = options.closeStream ?? true;
159170
const writer = stream.getWriter();
160171

161172
if (!options.nonBlocking) {
@@ -167,8 +178,12 @@ export function getStreamSink(
167178
.then(() => writer.write(bytes));
168179
};
169180
sink[Symbol.asyncDispose] = async () => {
170-
await lastPromise;
171-
await writer.close();
181+
try {
182+
await lastPromise;
183+
if (closeStream) await writer.close();
184+
} finally {
185+
writer.releaseLock();
186+
}
172187
};
173188
return markSinkAsImmediate(sink);
174189
}
@@ -240,11 +255,18 @@ export function getStreamSink(
240255
clearInterval(flushTimer);
241256
flushTimer = null;
242257
}
243-
await flush();
244258
try {
245-
await writer.close();
246-
} catch {
247-
// Writer might already be closed or errored
259+
await activeFlush;
260+
await flush();
261+
if (closeStream) {
262+
try {
263+
await writer.close();
264+
} catch {
265+
// Writer might already be closed or errored
266+
}
267+
}
268+
} finally {
269+
writer.releaseLock();
248270
}
249271
};
250272

0 commit comments

Comments
 (0)