Skip to content

Commit db6a1a6

Browse files
committed
test: cover writable-to-web batching and context
Add regression coverage for byte reuse, large chunks, mixed writev batches, and errors reported after early acceptance of copied writes. Check ALS context across queued writes and writev completion with async hooks and legacy context propagation, including init-hook boundaries. Signed-off-by: seungwoo <zoozoo1302@gmail.com>
1 parent aa9aca5 commit db6a1a6

2 files changed

Lines changed: 386 additions & 0 deletions

File tree

Lines changed: 151 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,151 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const { AsyncLocalStorage, createHook } = require('async_hooks');
6+
const { Writable } = require('stream');
7+
8+
const mode = process.argv[2];
9+
const hook = mode === 'hooks' ? createHook({ init() {} }).enable() : undefined;
10+
11+
async function checkContext(origin) {
12+
const storage = new AsyncLocalStorage();
13+
const callbacks = [];
14+
const consumed = [];
15+
let started = Promise.withResolvers();
16+
17+
function hold(chunks, callback) {
18+
assert.strictEqual(storage.getStore(), origin);
19+
callbacks.push(() => {
20+
for (const chunk of chunks) consumed.push(chunk[0]);
21+
callback();
22+
});
23+
started.resolve();
24+
}
25+
26+
const writable = new Writable({
27+
highWaterMark: 8,
28+
write: common.mustCallAtLeast((chunk, encoding, callback) => {
29+
hold([chunk], callback);
30+
}),
31+
writev: common.mustCallAtLeast((chunks, callback) => {
32+
hold(chunks.map(({ chunk }) => chunk), callback);
33+
}),
34+
});
35+
const writer = storage.run(origin, () => Writable.toWeb(writable).getWriter());
36+
const writes = [];
37+
for (let value = 1; value <= 16; value++) {
38+
const chunk = Buffer.from([value]);
39+
writes.push(storage.run(`writer-${value}`, common.mustCall(() =>
40+
writer.write(chunk).then(common.mustCall(() => {
41+
assert.strictEqual(storage.getStore(), `writer-${value}`);
42+
chunk.fill(255);
43+
})))));
44+
}
45+
46+
let index = 0;
47+
while (consumed.length < 16) {
48+
if (callbacks.length === 0) await started.promise;
49+
started = Promise.withResolvers();
50+
const callback = callbacks.shift();
51+
storage.run(`callback-${index++}`, common.mustCall(() => {
52+
callback();
53+
assert.strictEqual(storage.getStore(), `callback-${index - 1}`);
54+
}));
55+
await new Promise(setImmediate);
56+
}
57+
await Promise.all(writes);
58+
await writer.close();
59+
assert.deepStrictEqual(consumed, [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16]);
60+
storage.disable();
61+
}
62+
63+
async function checkIndependentBatches(withString) {
64+
const storage = new AsyncLocalStorage();
65+
const callbacks = [];
66+
const contexts = [];
67+
const writable = new Writable({
68+
write: common.mustCall((chunk, encoding, callback) => {
69+
contexts.push(storage.getStore());
70+
if (chunk[0] === 1) setImmediate(callback);
71+
else callbacks.push(callback);
72+
}, withString ? 5 : 4),
73+
});
74+
const writer = Writable.toWeb(writable).getWriter();
75+
await writer.write(Buffer.from([1]));
76+
await storage.run('batch-A', () => writer.write(Buffer.from([2])));
77+
callbacks.shift()();
78+
79+
if (withString === 'native') {
80+
storage.run('batch-B', () => writable.write('3'));
81+
} else if (withString) {
82+
await storage.run('batch-B', () => writer.write('3'));
83+
}
84+
await storage.run('batch-B', () => writer.write(Buffer.from([3])));
85+
await storage.run('batch-B', () => writer.write(Buffer.from([4])));
86+
if (withString) storage.run('batch-B', () => callbacks.shift()());
87+
storage.run('completion-B', common.mustCall(() => {
88+
callbacks.shift()();
89+
assert.strictEqual(storage.getStore(), 'completion-B');
90+
}));
91+
callbacks.shift()();
92+
await writer.close();
93+
assert.deepStrictEqual(contexts, withString ?
94+
[undefined, 'batch-A', 'batch-B', 'batch-B', 'batch-B'] :
95+
[undefined, 'batch-A', 'batch-B', 'batch-B']);
96+
storage.disable();
97+
}
98+
99+
async function checkResourceInitHook(hookFirst) {
100+
let storage;
101+
const initHook = createHook({
102+
init(id, type) {
103+
if (type === 'WEBSTREAM_WRITABLE_WRITE') storage.enterWith('hook-origin');
104+
},
105+
});
106+
if (hookFirst) initHook.enable();
107+
storage = new AsyncLocalStorage();
108+
const contexts = [];
109+
const writable = new Writable({
110+
write: common.mustCall((chunk, encoding, callback) => {
111+
contexts.push(storage.getStore());
112+
setImmediate(callback);
113+
}, 2),
114+
});
115+
const writer = Writable.toWeb(writable).getWriter();
116+
await writer.write(Buffer.from([1]));
117+
if (!hookFirst) initHook.enable();
118+
try {
119+
const write = storage.run('stream-origin', () => writer.write(Buffer.from([2])));
120+
assert.strictEqual(storage.getStore(), undefined);
121+
await write;
122+
await writer.close();
123+
assert.deepStrictEqual(contexts, [undefined, 'stream-origin']);
124+
} finally {
125+
initHook.disable();
126+
storage.disable();
127+
}
128+
}
129+
130+
async function main() {
131+
// node:test enables async hooks and would hide the default frame path.
132+
for (const origin of [undefined, 'stream-origin']) {
133+
await checkContext(origin);
134+
}
135+
for (const withString of [false, true, 'native']) await checkIndependentBatches(withString);
136+
for (const hookFirst of [false, true]) await checkResourceInitHook(hookFirst);
137+
hook?.disable();
138+
139+
if (mode === undefined) {
140+
for (const args of [
141+
[__filename, 'hooks'],
142+
['--no-async-context-frame', __filename, 'legacy'],
143+
]) {
144+
const { code, signal, stderr } = await common.spawnPromisified(process.execPath, args);
145+
assert.strictEqual(code, 0, stderr);
146+
assert.strictEqual(signal, null);
147+
}
148+
}
149+
}
150+
151+
main().then(common.mustCall());
Lines changed: 235 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,235 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const { Writable } = require('stream');
6+
7+
async function checkReusableBatches() {
8+
const consumed = [];
9+
let batches = 0;
10+
let drains = 0;
11+
const writable = new Writable({
12+
highWaterMark: 64 * 1024,
13+
write(chunk, encoding, callback) {
14+
setImmediate(() => {
15+
consumed.push(Buffer.from(chunk));
16+
callback();
17+
});
18+
},
19+
writev(chunks, callback) {
20+
batches++;
21+
setImmediate(() => {
22+
consumed.push(...chunks.map(({ chunk }) => Buffer.from(chunk)));
23+
callback();
24+
});
25+
},
26+
});
27+
writable.on('error', common.mustNotCall());
28+
writable.on('drain', () => { drains++; });
29+
const writer = Writable.toWeb(writable).getWriter();
30+
const input = Buffer.alloc(1024);
31+
for (let i = 1; i <= 128; i++) {
32+
input.fill(i);
33+
await writer.write(input);
34+
if (i === 1) assert.strictEqual(consumed.length, 1);
35+
input.fill(255);
36+
}
37+
await writer.close();
38+
39+
assert(batches > 0, 'small awaited writes should reach native writev');
40+
assert(drains > 0, 'buffered small writes should observe native backpressure');
41+
assert.strictEqual(consumed.length, 128);
42+
for (let i = 1; i <= 128; i++) {
43+
assert.deepStrictEqual(consumed[i - 1], Buffer.alloc(1024, i));
44+
}
45+
}
46+
47+
async function checkLateWriteError() {
48+
const error = new Error('late native write failure');
49+
const started = Promise.withResolvers();
50+
let finishWrite;
51+
const writable = new Writable({
52+
write: common.mustCall((chunk, encoding, callback) => {
53+
if (chunk[0] === 1) {
54+
setImmediate(callback);
55+
} else {
56+
finishWrite = callback;
57+
started.resolve();
58+
}
59+
}, 2),
60+
});
61+
const writer = Writable.toWeb(writable).getWriter();
62+
await writer.write(Buffer.from([1]));
63+
64+
const closed = assert.rejects(writer.closed, (actual) => actual === error);
65+
let fulfilled = false;
66+
const smallWrite = writer.write(Buffer.from([2])).then(() => {
67+
fulfilled = true;
68+
return true;
69+
}, (actual) => {
70+
assert.strictEqual(actual, error);
71+
return false;
72+
});
73+
await started.promise;
74+
await new Promise(setImmediate);
75+
const fulfilledBeforeError = fulfilled;
76+
const close = assert.rejects(writer.close(), (actual) => actual === error);
77+
finishWrite(error);
78+
const accepted = await smallWrite;
79+
await Promise.all([closed, close]);
80+
81+
assert.strictEqual(fulfilledBeforeError, true);
82+
assert.strictEqual(accepted, true);
83+
}
84+
85+
async function checkLargeWrite(size) {
86+
const started = Promise.withResolvers();
87+
const input = new Uint8Array(new ArrayBuffer(size + 8), 4, size);
88+
input.fill(3);
89+
let consumed;
90+
let finishWrite;
91+
const writable = new Writable({
92+
highWaterMark: 4,
93+
write: common.mustCall((chunk, encoding, callback) => {
94+
if (chunk[0] === 1) {
95+
setImmediate(callback);
96+
} else {
97+
assert.strictEqual(chunk.buffer, input.buffer);
98+
assert.strictEqual(chunk.byteOffset, input.byteOffset);
99+
finishWrite = () => {
100+
consumed = Buffer.from(chunk);
101+
callback();
102+
};
103+
started.resolve();
104+
}
105+
}, 2),
106+
});
107+
writable.on('error', common.mustNotCall());
108+
const writer = Writable.toWeb(writable).getWriter();
109+
await writer.write(Buffer.from([1]));
110+
111+
let fulfilled = false;
112+
const pending = writer.write(input).then(() => { fulfilled = true; });
113+
await started.promise;
114+
await new Promise(setImmediate);
115+
const fulfilledBeforeCompletion = fulfilled;
116+
finishWrite();
117+
await pending;
118+
input.fill(255);
119+
await writer.close();
120+
121+
assert.strictEqual(fulfilledBeforeCompletion, false);
122+
assert.deepStrictEqual(consumed, Buffer.alloc(size, 3));
123+
}
124+
125+
async function checkLargeWriteAfterCopiedBatch() {
126+
const consumed = [];
127+
const batchStarted = Promise.withResolvers();
128+
const input = Buffer.alloc(1);
129+
const large = Buffer.alloc(4, 5);
130+
let finishInitial;
131+
let finishBatch;
132+
const writable = new Writable({
133+
highWaterMark: 4,
134+
write: common.mustCall((chunk, encoding, callback) => {
135+
const consume = () => {
136+
consumed.push(Buffer.from(chunk));
137+
callback();
138+
};
139+
if (chunk[0] === 1) setImmediate(consume);
140+
else finishInitial = consume;
141+
}, 2),
142+
writev: common.mustCall((chunks, callback) => {
143+
assert.strictEqual(chunks.length, 3);
144+
assert.strictEqual(chunks[2].chunk.buffer, large.buffer);
145+
finishBatch = () => {
146+
consumed.push(...chunks.map(({ chunk }) => Buffer.from(chunk)));
147+
callback();
148+
};
149+
batchStarted.resolve();
150+
}),
151+
});
152+
writable.on('error', common.mustNotCall());
153+
const writer = Writable.toWeb(writable).getWriter();
154+
await writer.write(Buffer.from([1]));
155+
for (const value of [2, 3, 4]) {
156+
input[0] = value;
157+
await writer.write(input);
158+
input[0] = 255;
159+
}
160+
161+
let fulfilled = false;
162+
const pending = writer.write(large).then(() => { fulfilled = true; });
163+
await new Promise(setImmediate);
164+
assert.strictEqual(fulfilled, false);
165+
finishInitial();
166+
await batchStarted.promise;
167+
await new Promise(setImmediate);
168+
const fulfilledBeforeBatchCompletion = fulfilled;
169+
finishBatch();
170+
await pending;
171+
large.fill(255);
172+
await writer.close();
173+
174+
assert.strictEqual(fulfilledBeforeBatchCompletion, false);
175+
assert.deepStrictEqual(consumed, [
176+
Buffer.from([1]), Buffer.from([2]), Buffer.from([3]),
177+
Buffer.from([4]), Buffer.from([5, 5, 5, 5]),
178+
]);
179+
}
180+
181+
async function checkQueuedWritevError() {
182+
const error = new Error('queued native writev failure');
183+
const batchStarted = Promise.withResolvers();
184+
let finishInitial;
185+
let finishBatch;
186+
const writable = new Writable({
187+
highWaterMark: 4,
188+
write: common.mustCall((chunk, encoding, callback) => {
189+
if (chunk[0] === 1) setImmediate(callback);
190+
else finishInitial = callback;
191+
}, 2),
192+
writev: common.mustCall((chunks, callback) => {
193+
assert.deepStrictEqual(chunks.map(({ chunk }) => Array.from(chunk)), [[3], [4]]);
194+
finishBatch = callback;
195+
batchStarted.resolve();
196+
}),
197+
});
198+
const writer = Writable.toWeb(writable).getWriter();
199+
await writer.write(Buffer.from([1]));
200+
for (const value of [2, 3, 4]) await writer.write(Buffer.from([value]));
201+
202+
const closed = assert.rejects(writer.closed, (actual) => actual === error);
203+
const close = assert.rejects(writer.close(), (actual) => actual === error);
204+
finishInitial();
205+
await batchStarted.promise;
206+
setImmediate(finishBatch, error);
207+
await Promise.all([closed, close]);
208+
}
209+
210+
async function checkSyncWriteErrorAfterAsyncWrite() {
211+
const error = new Error('synchronous native write failure');
212+
const writable = new Writable({
213+
write: common.mustCall((chunk, encoding, callback) => {
214+
if (chunk[0] === 1) setImmediate(callback);
215+
else callback(error);
216+
}, 2),
217+
});
218+
const writer = Writable.toWeb(writable).getWriter();
219+
await writer.write(Buffer.from([1]));
220+
await Promise.all([
221+
assert.rejects(writer.write(Buffer.from([2])), (actual) => actual === error),
222+
assert.rejects(writer.closed, (actual) => actual === error),
223+
]);
224+
}
225+
226+
async function main() {
227+
await checkReusableBatches();
228+
await checkLateWriteError();
229+
for (const size of [4, 5]) await checkLargeWrite(size);
230+
await checkLargeWriteAfterCopiedBatch();
231+
await checkQueuedWritevError();
232+
await checkSyncWriteErrorAfterAsyncWrite();
233+
}
234+
235+
main().then(common.mustCall());

0 commit comments

Comments
 (0)