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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
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
24 changes: 21 additions & 3 deletions doc/api/quic.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
separate APIs for creating each kind:
[`session.createBidirectionalStream()`][] and
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
peer are delivered via the [`session.onstream`][] callback.
peer are delivered via the [`session.onstream`][] callback. When the
negotiated application protocol supports the stream-level callbacks (e.g.
HTTP/3) and an `onheaders` callback is configured, incoming streams can
instead be consumed entirely through it and registering `onstream` is
optional.

There are two ways to write data to a stream:

Expand DownExpand Up@@ -409,7 +413,9 @@ A typical client session progresses through these stages:

On the server side, call [`quic.listen()`][] with a callback. The callback
fires for each incoming session after the TLS handshake begins. Incoming
streams arrive via the [`session.onstream`][] callback.
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
sessions with an `onheaders` callback configured, directly through that
callback (see the [minimal HTTP/3 server][] example).

[`session.destroy()`][] is available for immediate teardown — all open streams
are destroyed and the session is closed without waiting for them to finish.
Expand DownExpand Up@@ -1110,6 +1116,15 @@ added: v23.8.0

The callback to invoke when a new stream is initiated by a remote peer. Read/write.

If no `onstream` callback is set and the stream has no other consumer, an
incoming stream is destroyed on arrival and a warning is emitted. An
`onheaders` callback counts as a consumer when the negotiated application
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
`onwanttrailers`) do not, since they are conditional or outbound-only and
would leave the stream unobservable. An HTTP/3 server that handles requests
entirely through `onheaders` does not need to set `onstream`.

### `session.ondatagram`

<!-- YAML
Expand DownExpand Up@@ -4008,7 +4023,9 @@ import { listen } from 'node:quic';
const encoder = new TextEncoder();

const endpoint = await listen((session) => {
// The session.onstream callback fires for each new client-initiated stream.
// The session.onstream callback fires for each new client-initiated
// stream. It is optional here: with `onheaders` configured below,
// request streams are consumed through that callback.
}, {
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
// ALPN defaults to 'h3'.
Expand DownExpand Up@@ -4642,5 +4659,6 @@ throughput issues caused by flow control.
[`stream.writer`]: #streamwriter
[`writer.fail()`]: #streamwriter
[`writer.fail(reason)`]: #streamwriter
[minimal HTTP/3 server]: #minimal-http3-server
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
[qvis]: https://qvis.quictools.info/
38 changes: 33 additions & 5 deletions lib/internal/quic/quic.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -4129,6 +4129,24 @@ class QuicSession {
this.#inner.verifyPeer = value;
}

/**
* True if an incoming stream has a consumer registered on this session:
* either an onstream callback, or - when the negotiated application
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
* the application layer will invoke (onheaders et al).
* @returns {boolean}
*/
#hasStreamConsumer() {
if (typeof this.#inner.onstream === 'function') return true;
// Only onheaders is guaranteed to fire for every incoming stream when
// the negotiated application supports stream callbacks (e.g. HTTP/3).
// Other stream callbacks are conditional (ontrailers, oninfo) or
// outbound-only (onwanttrailers) and do not expose the stream, so they
// do not count as a consumer.
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
return getQuicSessionState(this).streamCallbacksSupported === 1;
}

/**
* @param {object} handle
* @param {number} direction
Expand All@@ -4141,10 +4159,13 @@ class QuicSession {
// Set the default byte budget for received streams.
stream.budget = kDefaultBudget;

// A new stream was received. If we don't have an onstream callback, then
// there's nothing we can do about it. Destroy the stream in this case.
if (typeof inner.onstream !== 'function') {
process.emitWarning('A new stream was received but no onstream callback was provided');
// A new stream was received. If the session has no consumer for it -
// neither an onstream callback nor, on a session whose application
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
// there's nothing that could ever read it. Destroy the stream in this
// case rather than letting it hold flow control credit.
if (!this.#hasStreamConsumer()) {
process.emitWarning('A new stream was received but no stream consumer callback was provided');
stream.destroy();
return;
}
Expand DownExpand Up@@ -4175,7 +4196,14 @@ class QuicSession {
});
}

safeCallbackInvoke(inner.onstream, this, stream);
// Deliver the stream to the onstream consumer if one is registered.
// Reaching this point without one means #hasStreamConsumer accepted
// the stream on behalf of the application layer: the session-level
// stream callbacks were applied above and the application (e.g.
// HTTP/3) drives the stream, so there is nothing to invoke here.
if (typeof inner.onstream === 'function') {
safeCallbackInvoke(inner.onstream, this, stream);
}
}

[kRemoveStream](stream) {
Expand Down
15 changes: 15 additions & 0 deletions lib/internal/quic/state.js
Original file line numberDiff line numberDiff line change
Expand Up@@ -72,6 +72,7 @@ const {
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
IDX_STATE_SESSION_HEADERS_SUPPORTED,
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
IDX_STATE_SESSION_WRAPPED,
IDX_STATE_SESSION_APPLICATION_TYPE,
IDX_STATE_SESSION_NO_ERROR_CODE,
Expand DownExpand Up@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
Expand DownExpand Up@@ -493,6 +495,19 @@ class QuicSessionState {
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
}

/**
* Whether the negotiated application dispatches the session-level
* stream callbacks (onheaders et al) for incoming streams.
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
* @type {number}
*/
get streamCallbacksSupported() {
const handle = this.#handle;
if (handle === undefined) return undefined;
return DataViewPrototypeGetUint8(
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
}

/** @type {boolean} */
get isWrapped() {
const handle = this.#handle;
Expand Down
5 changes: 5 additions & 0 deletions src/quic/application.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
// do not support headers should return false (the default).
virtual bool SupportsHeaders() const { return false; }

// True if this application dispatches the session-level stream
// callbacks (onheaders et al) for incoming streams when they are
// registered on the session.
virtual bool SupportsStreamCallbacks() const { return false; }

// Initiates application-level graceful shutdown signaling (e.g.,
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
virtual void BeginShutdown() {}
Expand Down
6 changes: 6 additions & 0 deletions src/quic/defs.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
UNSUPPORTED,
};

enum class StreamCallbacksSupportState : uint8_t {
UNKNOWN,
SUPPORTED,
UNSUPPORTED,
};

enum class PathValidationResult : uint8_t {
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,
Expand Down
2 changes: 2 additions & 0 deletions src/quic/http3.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {

bool SupportsHeaders() const override { return true; }

bool SupportsStreamCallbacks() const override { return true; }

bool is_started() const override { return started_; }

bool Start() override {
Expand Down
5 changes: 5 additions & 0 deletions src/quic/session.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
V(WRAPPED, wrapped, uint8_t) \
V(APPLICATION_TYPE, application_type, uint8_t) \
V(NO_ERROR_CODE, no_error_code, error_code) \
Expand DownExpand Up@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
impl_->state()->headers_supported = static_cast<uint8_t>(
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
: HeadersSupportState::UNSUPPORTED);
impl_->state()->stream_callbacks_supported =
static_cast<uint8_t>(app->SupportsStreamCallbacks()
? StreamCallbacksSupportState::SUPPORTED
: StreamCallbacksSupportState::UNSUPPORTED);
// Surface the application's "no error" and "internal error" codes via
// session state so that JS-side code (e.g. the stream writer's fail()
// path) can resolve the right wire code for the negotiated ALPN
Expand Down
Loading
Loading