Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions benchmark/webstreams/encoding-streams.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
TextEncoderStream,
TextDecoderStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e5],
kind: ['encode', 'decode'],
len: [16, 1024],
});

async function main({ n, kind, len }) {
const encoded = new TextEncoder().encode('a'.repeat(len));
const decoded = 'a'.repeat(len);
let i = 0;
const rs = new ReadableStream({
pull(controller) {
if (i++ < n) {
controller.enqueue(kind === 'encode' ? decoded : encoded);
} else {
controller.close();
}
},
});
const ts = kind === 'encode' ?
new TextEncoderStream() :
new TextDecoderStream();

const reader = rs.pipeThrough(ts).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
29 changes: 29 additions & 0 deletions benchmark/webstreams/from.js
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
'use strict';
const common = require('../common.js');
const {
ReadableStream,
} = require('node:stream/web');

const bench = common.createBenchmark(main, {
n: [1e6],
kind: ['sync', 'async'],
});

async function main({ n, kind }) {
function* syncGen() {
for (let i = 0; i < n; i++) yield i;
}

async function* asyncGen() {
for (let i = 0; i < n; i++) yield i;
}

const reader = ReadableStream.from(
kind === 'sync' ? syncGen() : asyncGen()).getReader();
bench.start();
for (;;) {
const { done } = await reader.read();
if (done) break;
}
bench.end(n);
}
48 changes: 22 additions & 26 deletions lib/internal/webstreams/encoding.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,6 +4,7 @@ const {
ObjectDefineProperties,
String,
StringPrototypeCharCodeAt,
StringPrototypeSlice,
Uint8Array,
} = primordials;

Expand DownExpand Up@@ -31,6 +32,9 @@ const {
kEnumerableProperty,
} = require('internal/util');

// Shared per-chunk decode options; decode() only reads the flag.
const kDecodeStreamingOptions = { __proto__: null, stream: true };

/**
* @typedef {import('./readablestream').ReadableStream} ReadableStream
* @typedef {import('./writablestream').WritableStream} WritableStream
Expand All@@ -46,34 +50,26 @@ class TextEncoderStream {
this.#transform = new TransformStream({
transform: (chunk, controller) => {
// https://encoding.spec.whatwg.org/#encode-and-enqueue-a-chunk
// The only cross-chunk state is a trailing high surrogate;
// encode() replaces interior lone surrogates with U+FFFD exactly
// like the spec's per-code-unit walk.
chunk = String(chunk);
let finalChunk = '';
for (let i = 0; i < chunk.length; i++) {
const item = chunk[i];
const codeUnit = StringPrototypeCharCodeAt(item, 0);
if (this.#pendingHighSurrogate !== null) {
const highSurrogate = this.#pendingHighSurrogate;
this.#pendingHighSurrogate = null;
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += highSurrogate + item;
continue;
}
finalChunk += '\uFFFD';
}
if (0xD800 <= codeUnit && codeUnit <= 0xDBFF) {
this.#pendingHighSurrogate = item;
continue;
}
if (0xDC00 <= codeUnit && codeUnit <= 0xDFFF) {
finalChunk += '\uFFFD';
continue;
}
finalChunk += item;
if (chunk.length === 0)
return;
if (this.#pendingHighSurrogate !== null) {
chunk = this.#pendingHighSurrogate + chunk;
this.#pendingHighSurrogate = null;
}
if (finalChunk) {
const value = this.#handle.encode(finalChunk);
controller.enqueue(value);
const lastCodeUnit =
StringPrototypeCharCodeAt(chunk, chunk.length - 1);
if (0xD800 <= lastCodeUnit && lastCodeUnit <= 0xDBFF) {
this.#pendingHighSurrogate =
StringPrototypeSlice(chunk, -1);
chunk = StringPrototypeSlice(chunk, 0, -1);
if (chunk.length === 0)
return;
}
controller.enqueue(this.#handle.encode(chunk));
},
flush: (controller) => {
// https://encoding.spec.whatwg.org/#encode-and-flush
Expand DownExpand Up@@ -137,7 +133,7 @@ class TextDecoderStream {
if (chunk === undefined) {
throw new ERR_INVALID_ARG_TYPE('chunk', 'string', chunk);
}
const value = this.#handle.decode(chunk, { stream: true });
const value = this.#handle.decode(chunk, kDecodeStreamingOptions);
if (value)
controller.enqueue(value);
},
Expand Down
83 changes: 75 additions & 8 deletions lib/internal/webstreams/readablestream.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -112,6 +112,7 @@ const {
getNonWritablePropertyDescriptor,
isBrandCheck,
kEmptyQueue,
kParkedAlgorithmResult,
kResolvedPromise,
kState,
kType,
Expand DownExpand Up@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
}
// Per GetIteratorDirect, the next method is looked up once.
const nextMethod = iterator.next;
const startAlgorithm = nonOpCallback;

async function pullAlgorithm() {
const iterResult = await iterator.next();
// Callback-style pull: the reaction steps are reused across chunks and
// completion is delivered to the controller's cached pull reactions
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
// single slot carries a non-thenable next() result between steps.
let pendingIterResult;

function rejectPull(error) {
readableStreamDefaultControllerError(stream[kState].controller, error);
}

function processIterResult(iterResult) {
const controller = stream[kState].controller;
if (typeof iterResult !== 'object' || iterResult === null) {
throw new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object');
rejectPull(new ERR_INVALID_STATE.TypeError(
'The promise returned by the iterator.next() method must fulfill with an object'));
return;
}
if (iterResult.done) {
readableStreamDefaultControllerClose(stream[kState].controller);
} else {
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
try {
if (iterResult.done) {
readableStreamDefaultControllerClose(controller);
} else {
const value = iterResult.value;
if (value !== null &&
(typeof value === 'object' || typeof value === 'function')) {
// Adopted like `await iterResult.value`, keeping the observable
// .then lookup on plain objects.
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
return;
}
readableStreamDefaultControllerEnqueue(controller, value);
}
} catch (error) {
rejectPull(error);
return;
}
// pullFulfilled exists: the controller creates it before the pull.
controller[kState].pullFulfilled();
}

function enqueueValue(value) {
const controller = stream[kState].controller;
try {
readableStreamDefaultControllerEnqueue(controller, value);
} catch (error) {
rejectPull(error);
return;
}
controller[kState].pullFulfilled();
}

function processPendingIterResult() {
const iterResult = pendingIterResult;
pendingIterResult = undefined;
processIterResult(iterResult);
}

function pullAlgorithm() {
let nextResult;
try {
nextResult = FunctionPrototypeCall(nextMethod, iterator);
} catch (error) {
return PromiseReject(error);
}
if (nextResult !== null &&
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
// Mirrors `await iterator.next()`: processIterResult runs at the
// microtask position the await resumed.
PromisePrototypeThen(
PromiseResolve(nextResult), processIterResult, rejectPull);
return kParkedAlgorithmResult;
}
// A non-thenable next() result fails validation a microtask later.
pendingIterResult = nextResult;
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
return kParkedAlgorithmResult;
}

async function cancelAlgorithm(reason) {
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/webstreams/util.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -359,7 +359,7 @@ const kResolvedPromise = PromiseResolve();
// operation and takes responsibility for delivering the fulfilled (or
// rejected) continuation itself later, instead of settling a promise
// (see the transform stream source pull algorithm).
const kParkedAlgorithmResult = { __proto__: null };
const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult');

// Wires the (possibly non-thenable) result of an underlying algorithm
// callback to its fulfilled/rejected continuations. A non-thenable result
Expand Down
Loading