diff --git a/doc/concept/standard.md b/doc/concept/standard.md index f96904d2ef..de72754f61 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -31,6 +31,11 @@ This is an implementation limit, not a limit in the IETF draft. A larger declared block stops its subgroup stream with `MALFORMED_TRACK` before reading the block; other groups and the session stay open. +Rust and JavaScript read an incoming padding stream (draft-18 and later) to +the end and discard it, without sending `STOP_SENDING`. A unidirectional stream type the negotiated draft does not +define, or a `SUBGROUP_HEADER` type it marks invalid, closes the session with +`PROTOCOL_VIOLATION`, as the draft requires. + An IETF publisher declares the track's default priority in `SUBSCRIBE_OK` or `PUBLISH` when that draft carries track properties. Groups without a priority flag inherit it. If the property is absent, the IETF wire default of 128 maps diff --git a/js/net/src/ietf/connection.test.ts b/js/net/src/ietf/connection.test.ts new file mode 100644 index 0000000000..8b65bc6db1 --- /dev/null +++ b/js/net/src/ietf/connection.test.ts @@ -0,0 +1,78 @@ +import { expect, spyOn, test } from "bun:test"; +import { SessionCode } from "../error.ts"; +import { createMockTransportPair } from "../mock.ts"; +import { Stream, Writer } from "../stream.ts"; +import { Connection } from "./connection.ts"; +import { ALPN, type IetfVersion, Version } from "./version.ts"; + +const PADDING = 0x132b3e28n; + +/** A server-side session over `version`, and a client uni stream that sends only `type`. */ +async function openUni( + version: IetfVersion, + alpn: string, + type: bigint, +): Promise<{ pair: ReturnType; connection: Connection; writer: Writer }> { + const pair = createMockTransportPair(alpn); + const control = await Stream.open(pair.server, { version }); + const connection = new Connection({ + url: new URL("https://example.com"), + quic: pair.server, + control, + maxRequestId: 100n, + version, + client: false, + }); + + const writer = new Writer(await pair.client.createUnidirectionalStream(), version); + await writer.u62(type); + + return { pair, connection, writer }; +} + +/** Padding is read to the end and dropped: no STOP_SENDING, and the session stays up. */ +test("a padding stream is discarded", async () => { + const { pair, connection, writer } = await openUni(Version.DRAFT_19, ALPN.DRAFT_19, PADDING); + let closed = false; + void pair.server.closed.then(() => { + closed = true; + }); + + try { + await writer.write(new Uint8Array(4096)); + writer.close(); + + // A STOP_SENDING would reject this; a clean close means every byte was read. + await writer.closed; + + // A session close would come from the same handler that read the stream. + await Bun.sleep(0); + expect(closed).toBe(false); + } finally { + connection.close(); + } +}); + +/** An unknown or invalid stream type MUST close the session with PROTOCOL_VIOLATION. */ +test("an unknown uni stream type closes the session", async () => { + for (const [version, alpn, type] of [ + [Version.DRAFT_19, ALPN.DRAFT_19, 0n], + // Past 2^53, so it cannot be read as a number. + [Version.DRAFT_19, ALPN.DRAFT_19, 2n ** 53n], + // Padding arrived in draft-18. + [Version.DRAFT_17, ALPN.DRAFT_17, PADDING], + // A SUBGROUP_HEADER with the reserved SUBGROUP_ID_MODE (0b11). + [Version.DRAFT_19, ALPN.DRAFT_19, 0x56n], + ] as const) { + const logged = spyOn(console, "error").mockImplementation(() => void 0); + const { pair, connection } = await openUni(version, alpn, type); + + try { + const info = await pair.client.closed; + expect(info.closeCode).toBe(SessionCode.ProtocolViolation); + } finally { + logged.mockRestore(); + connection.close(); + } + } +}); diff --git a/js/net/src/ietf/connection.ts b/js/net/src/ietf/connection.ts index d4eb993892..ece52824b7 100644 --- a/js/net/src/ietf/connection.ts +++ b/js/net/src/ietf/connection.ts @@ -3,14 +3,14 @@ import type * as announce from "../announced.ts"; import type { Established } from "../connection/established.ts"; import { type Probe, type Stats, transportStats } from "../connection/stats.ts"; import { type Transport, transportOf } from "../connection/transport.ts"; -import { error, fromClose, ProtocolViolation, StreamCode, StreamError } from "../error.ts"; +import { error, fromClose, ProtocolViolation, SessionCode, StreamCode, StreamError } from "../error.ts"; import type { Consumer as OriginConsumer } from "../origin.ts"; import type * as Path from "../path.ts"; import { type Reader, Readers, type Stream } from "../stream.ts"; import { registerWire } from "../wire.ts"; import { ControlStreamAdapter, NativeSession, type Session } from "./adapter.ts"; import * as Cluster from "./cluster.ts"; -import { Fetch } from "./fetch.ts"; +import { Fetch, FetchHeader } from "./fetch.ts"; import { GoAway } from "./goaway.ts"; import { Group } from "./object.ts"; import { Publish } from "./publish.ts"; @@ -22,6 +22,9 @@ import { Subscriber } from "./subscriber.ts"; import { TrackStatusRequest } from "./track.ts"; import { type IetfVersion, Version, versionName } from "./version.ts"; +// The PADDING stream type (draft-18+): bytes a peer sends to probe for bandwidth. +const PADDING = 0x132b3e28n; + /** * Represents a connection to a MoQ server using moq-transport protocol. * @@ -158,17 +161,29 @@ export class Connection implements Established { * Closes the connection. */ close() { + this.#close(); + } + + // Close with the session code the peer should see, a clean close by default. + #close(info?: WebTransportCloseInfo) { if (this.#closed) return; this.#closed = true; - this.#session.close(); - + // Before the session, whose own close would send a clean code first. try { - this.#quic.close(); + this.#quic.close(info); } catch { // ignore } + + this.#session.close(); + } + + // The peer broke the protocol, so losing the stream is not enough: nothing stops it + // repeating the violation on the next one. + #violated(err: ProtocolViolation) { + this.#close({ closeCode: SessionCode.ProtocolViolation, reason: err.message }); } async #run(): Promise { @@ -199,10 +214,7 @@ export class Connection implements Established { void this.#runBidi(stream).catch((err: unknown) => { console.error("error processing bidi stream", err); stream.abort(new Error("bidi stream error")); - - // The peer broke the protocol, so losing the stream is not enough: nothing - // stops it repeating the violation on the next one. - if (err instanceof ProtocolViolation) this.close(); + if (err instanceof ProtocolViolation) this.#violated(err); }); } } @@ -313,13 +325,37 @@ export class Connection implements Established { .catch((err: unknown) => { console.error("error processing object stream", err); stream.stop(err); + + // An unknown or invalid stream type MUST close the session, not just the stream. + if (err instanceof ProtocolViolation) this.#violated(err); }); } } async #runUni(stream: Reader) { - const header = await Group.decode(stream, this.#session.version); - await this.#subscriber.handleGroup(header, stream); + const version = this.#session.version; + // Full width, so an unknown type past 2^53 is still classified rather than thrown. + const type = await stream.u62(); + + // SUBGROUP_HEADER types match 0b0XX1XXXX; Group.decode validates the bits per draft. + if (type <= 0xffn && (type & 0x90n) === 0x10n) { + const header = await Group.decode(stream, version, Number(type)); + await this.#subscriber.handleGroup(header, stream); + return; + } + + // The receiver MUST discard padding. We read it to the end rather than cancel, + // so a peer probing for bandwidth gets the throughput it is measuring. + if (type === PADDING && version >= Version.DRAFT_18) { + await stream.discard(); + return; + } + + // We never FETCH, so a fetch response answers nothing of ours. + if (type === BigInt(FetchHeader.type)) throw new Error("unexpected fetch stream"); + + // Anything else is unknown, and a second SETUP is a violation too. + throw new ProtocolViolation(`unknown uni stream type: 0x${type.toString(16)}`); } /** diff --git a/js/net/src/ietf/object.ts b/js/net/src/ietf/object.ts index 9d6a63b883..769c933252 100644 --- a/js/net/src/ietf/object.ts +++ b/js/net/src/ietf/object.ts @@ -1,4 +1,4 @@ -import { StreamCode, Stream as StreamError } from "../error.ts"; +import { ProtocolViolation, StreamCode, Stream as StreamError } from "../error.ts"; import { asIetf, type Cursor, type Reader, Writer } from "../stream.ts"; import { Timescale, Timestamp } from "../time.ts"; import { type IetfVersion, Version } from "./version.ts"; @@ -213,8 +213,9 @@ export class Group { } } - static async decode(r: Reader, version: IetfVersion): Promise { - const raw = await r.u53(); + /** Decode a SUBGROUP_HEADER. Pass `type` when the caller already read it to classify the stream. */ + static async decode(r: Reader, version: IetfVersion, type?: number): Promise { + const raw = type ?? (await r.u53()); // Strip the draft-18 FIRST_OBJECT bit before the range check, but keep the value: // it is the only signal that a subgroup starts partway through. Drafts that predate // it carry no such signal, so they are taken at their word. @@ -222,18 +223,17 @@ export class Group { const firstObject = legacy || (raw & FIRST_OBJECT_BIT) !== 0; const id = legacy ? raw : raw & ~FIRST_OBJECT_BIT; - let hasPriority: boolean; - let baseId: number; - if (id >= 0x10 && id <= 0x1f) { - hasPriority = true; - baseId = id; - } else if (id >= 0x30 && id <= 0x3f) { - hasPriority = false; - baseId = id - (0x30 - 0x10); - } else { - throw new Error(`Unsupported group type: ${id}`); + // An invalid type MUST close the session, including the reserved SUBGROUP_ID_MODE + // 0b11 (draft-21 section 11.3.1). + const known = (id >= 0x10 && id <= 0x1f) || (id >= 0x30 && id <= 0x3f); + if (!known || (id & 0x06) === 0x06) { + throw new ProtocolViolation(`Unsupported group type: ${raw}`); } + // 0x30-0x3F omit the priority and inherit it from the control message. + const hasPriority = id < 0x30; + const baseId = hasPriority ? id : id - (0x30 - 0x10); + const flags: GroupFlags = { hasExtensions: (baseId & 0x01) !== 0, hasSubgroupObject: (baseId & 0x02) !== 0, diff --git a/js/net/src/stream.ts b/js/net/src/stream.ts index c2847f395a..e4c8401c5a 100644 --- a/js/net/src/stream.ts +++ b/js/net/src/stream.ts @@ -407,6 +407,15 @@ export class Reader { return this.#slice(this.#buffer.byteLength); } + // Reads to the end of the stream, dropping every byte instead of buffering it. + async discard(): Promise { + this.#buffer = new Uint8Array(); + do { + this.#chunks = []; + this.#chunked = 0; + } while (await this.#fill()); + } + async string(): Promise { return this.decode(STRING); } diff --git a/quest/m0/README.md b/quest/m0/README.md index 1f4d926607..0dd8c24bdb 100644 --- a/quest/m0/README.md +++ b/quest/m0/README.md @@ -41,7 +41,6 @@ Published API or wire breaks still land on dev; each quest's Plan says so. ## Required - [IETF FIN semantics](/quest/m0/ietf-fin-not-cancel.md) - a request stream FIN stops updates without cancelling, and REQUEST_UPDATE on a subscribe is parsed -- [IETF stream types](/quest/m0/ietf-uni-stream-types.md) - padding streams are discarded stream-only and an unknown uni type closes the session, per draft-21 - [Request caps](/quest/m0/request-caps.md) - lite message sizes, IETF request IDs, and per-session announces and subscriptions are bounded - [quest check everywhere](/quest/m0/quest-check-everywhere.md) - `quest check` guards `main`, `dev`, and the line branches on push and PR, not only PRs into `main` - [noq reassembly cap](/quest/m0/noq-reassembly-cap.md) - noq carries quinn's stream reassembly cap and the connection receive window is finite by default diff --git a/quest/m0/ietf-uni-stream-types.md b/quest/m0/ietf-uni-stream-types.md deleted file mode 100644 index 1dfd4e0104..0000000000 --- a/quest/m0/ietf-uni-stream-types.md +++ /dev/null @@ -1,45 +0,0 @@ -# [S] Validate IETF unidirectional stream types - -## Goal - -Accept valid padding streams and close the session for genuinely unknown -stream types according to the negotiated moq-transport draft. Behavior only, -so it ships on main. - -## Plan - -Moved from m1 to m0 in the 2026-09-30 audit: a padding stream answered with -INTERNAL_ERROR is legal input mishandled, which m0 fixes before Seattle -interop on 2026-10-12. - -`run_unis` in `rs/moq-net/src/ietf/session.rs` routes every non-SETUP uni -stream to `run_uni_group`, which rejects padding and unknown types alike while -leaving the session alive. That stream-only rejection reaches the wire as -INTERNAL_ERROR, because nothing registers a code for it: the handler maps a -session-scoped error to `StreamError::Internal` and aborts the reader. - -draft-21 settles what each stream type means: a stream whose type the -endpoint does not recognize MUST close the session, and a padding stream -(type 0x132B3E28) MUST be discarded, which an endpoint may do by cancelling -it. The tree negotiates drafts 14 through 22 (`ietf::Version` in -`rs/moq-net/src/ietf/version.rs`), so apply that split to every supported draft and check the earlier ones for -the padding type value and whether draining is required. - -- Classify stream types before spawning a group handler. Handle PADDING per - draft, draining it where required and otherwise cancelling it with a - stream-only code, and propagate a genuinely unknown type to the session - driver as a protocol violation. Keep ordinary group failures scoped to their - streams. -- The test `unknown_uni_type_does_not_claim_the_session_closed` in - `session.rs` asserts the current behavior, an INTERNAL_ERROR stop and - no session close, and flips: an unknown type now closes the session and - stops nothing on its own. -- Add a padding test asserting a stream-only cancel with no session close, and - keep `a_group_for_a_retired_alias_is_stopped_with_cancelled`, which - pins that a dropped group never closes the session. - -Consult [draft-21 section 11.5](https://www.ietf.org/archive/id/draft-ietf-moq-transport-21.html) -for the wording, and [draft-19 section 3.4 and section 11.5.1](https://www.ietf.org/archive/id/draft-ietf-moq-transport-19.html) -plus [draft-20 section 11.5.1](https://www.ietf.org/archive/id/draft-ietf-moq-transport-20.html) -and [draft-22](https://www.ietf.org/archive/id/draft-ietf-moq-transport-22.html) -for the other drafts the tree negotiates. diff --git a/rs/moq-net/src/ietf/group.rs b/rs/moq-net/src/ietf/group.rs index 0196865a1a..e3798baf00 100644 --- a/rs/moq-net/src/ietf/group.rs +++ b/rs/moq-net/src/ietf/group.rs @@ -702,9 +702,10 @@ mod tests { assert!(flags.has_end); } - /// Regression: a publisher-emitted Draft18 GroupHeader byte must satisfy the - /// subscriber's uni-stream classifier mask `(byte & 0x90) == 0x10`. Otherwise - /// the uni stream is dropped as UnexpectedStream and the data plane stalls. + /// Regression: a publisher-emitted Draft18 GroupHeader byte must match the + /// SUBGROUP_HEADER form `(byte & 0x90) == 0x10` and decode as flags, which is how + /// the uni-stream classifier recognizes it. Otherwise the uni stream is an unknown + /// type, which closes the session. #[test] fn test_draft18_group_header_passes_stream_classifier() { let header = GroupHeader { @@ -719,10 +720,14 @@ mod tests { header.encode(&mut buf, Version::Draft18).unwrap(); let type_byte = buf[0] as u64; - // The check in session.rs::run_uni_group. assert_eq!( type_byte & 0x90, 0x10, + "draft-18 SUBGROUP_HEADER type 0x{type_byte:02x} is outside the SUBGROUP_HEADER form", + ); + // The check in session.rs::UniType::classify. + assert!( + GroupFlags::decode(type_byte, Version::Draft18).is_ok(), "draft-18 SUBGROUP_HEADER type 0x{type_byte:02x} not recognized by uni-stream classifier", ); } diff --git a/rs/moq-net/src/ietf/session.rs b/rs/moq-net/src/ietf/session.rs index 439c259b4e..f1c966669c 100644 --- a/rs/moq-net/src/ietf/session.rs +++ b/rs/moq-net/src/ietf/session.rs @@ -613,10 +613,53 @@ async fn run_setup( Ok(()) } +/// The PADDING stream type (draft-18+): bytes a peer sends to probe for bandwidth. +const PADDING: u64 = 0x132B3E28; + +/// What a unidirectional stream's type names on the negotiated draft. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum UniType { + /// The peer's SETUP, which then carries its GOAWAY (draft-17+). + Setup, + /// A SUBGROUP_HEADER, carrying one subgroup of a subscription. + Subgroup, + /// A FETCH_HEADER, carrying a fetch response. + Fetch, + /// Data to discard (draft-18+). + Padding, +} + +impl UniType { + /// `None` is a type the draft does not define, which MUST close the session + /// (draft-21 section 6.4.1, and its equivalent in every draft we negotiate). + fn classify(kind: u64, version: Version) -> Option { + // Draft-14-17 use SUBGROUP_HEADER types 0x10-0x1D and 0x30-0x3D; draft-18 adds + // 0x40 (FIRST_OBJECT), also covering 0x50-0x5D and 0x70-0x7D. A reserved + // SUBGROUP_ID_MODE (0b11) or a bit the draft lacks is invalid, which MUST close the + // session too (draft-21 section 11.3.1). + if ietf::GroupFlags::decode(kind, version).is_ok() { + return Some(Self::Subgroup); + } + + match kind { + FetchHeader::TYPE => Some(Self::Fetch), + setup::SETUP_V17 => match version { + // SETUP rides the bidi control stream. + Version::Draft14 | Version::Draft15 | Version::Draft16 => None, + _ => Some(Self::Setup), + }, + PADDING => match version { + Version::Draft14 | Version::Draft15 | Version::Draft16 | Version::Draft17 => None, + _ => Some(Self::Padding), + }, + _ => None, + } + } +} + /// Accept incoming uni streams and dispatch each to a handler. /// -/// For v17, this also handles the SETUP stream (0x2F00) and GOAWAY. -/// For v14-16, all uni streams are group data. +/// For v17+, this also handles the SETUP stream (0x2F00) and GOAWAY. async fn run_unis( mut session: S, subscriber: Subscriber, @@ -679,69 +722,84 @@ where Err(err) => return Err(err), }; - // v17+: SETUP arrives on a uni stream, then becomes the GOAWAY channel. - // We accept it in the background without blocking; the one thing that does - // need it (the MoQ Cluster negotiation) waits on `peer_setup` instead, so a - // slow SETUP delays announcements rather than the whole session. - if kind == setup::SETUP_V17 { - // Exactly one SETUP per endpoint. A second would let a peer restate its - // declared identity mid-session, silently re-attributing every route - // already built from the first. - if std::mem::replace(&mut seen_setup, true) { - return Err(Error::ProtocolViolation); - } + let Some(ty) = UniType::classify(kind, version) else { + tracing::warn!(kind, "unknown uni stream type"); + return Err(Error::UnexpectedStream); + }; - let peer_setup = peer_setup.clone(); - let mut session = session.clone(); - let goaway = goaway.clone(); - tasks.push(async move { - // The negotiation gates the announce and dispatch loops, so a SETUP we - // cannot read must end the session rather than leave them parked on a - // slot nothing will ever fill. - let msg = match reader.decode::().await { - Ok(msg) => msg, - Err(err) => { - tracing::warn!(%err, "setup decode error"); - session.close(SessionError::ProtocolViolation.to_code(), "invalid setup"); - return; - } - }; + match ty { + // SETUP then becomes the GOAWAY channel. We accept it in the background + // without blocking; the one thing that does need it (the MoQ Cluster + // negotiation) waits on `peer_setup` instead, so a slow SETUP delays + // announcements rather than the whole session. + UniType::Setup => { + // Exactly one SETUP per endpoint. A second would let a peer restate its + // declared identity mid-session, silently re-attributing every route + // already built from the first. + if std::mem::replace(&mut seen_setup, true) { + return Err(Error::ProtocolViolation); + } - if let Some(peer_setup) = peer_setup { - let peer = match decode_peer_setup(msg.parameters, version) { - Ok(peer) => peer, + let peer_setup = peer_setup.clone(); + let mut session = session.clone(); + let goaway = goaway.clone(); + tasks.push(async move { + // The negotiation gates the announce and dispatch loops, so a SETUP we + // cannot read must end the session rather than leave them parked on a + // slot nothing will ever fill. + let msg = match reader.decode::().await { + Ok(msg) => msg, Err(err) => { - tracing::warn!(%err, "setup parameter decode error"); - session.close(SessionError::ProtocolViolation.to_code(), "invalid setup parameters"); + tracing::warn!(%err, "setup decode error"); + session.close(SessionError::ProtocolViolation.to_code(), "invalid setup"); return; } }; - peer_setup.set(peer); - } - - // Monitor for GOAWAY after setup completes. - if let Err(err) = run_goaway(reader.with_version(version), version, goaway).await { - tracing::warn!(%err, "goaway error"); - } - }); - continue; - } + if let Some(peer_setup) = peer_setup { + let peer = match decode_peer_setup(msg.parameters, version) { + Ok(peer) => peer, + Err(err) => { + tracing::warn!(%err, "setup parameter decode error"); + session.close(SessionError::ProtocolViolation.to_code(), "invalid setup parameters"); + return; + } + }; + peer_setup.set(peer); + } - // Poll one child handler for each group stream. - let mut sub = subscriber.clone(); - tasks.push(async move { - let mut reader = reader.with_version(version); - if let Err(err) = run_uni_group(&mut sub, &mut reader).await { - tracing::debug!(%err, "uni stream error"); - // This handler stops only the stream, so it cannot claim the session closed. - let reset = match StreamError::from(&err) { - StreamError::Session(_) => StreamError::Internal, - reset => reset, - }; - reader.abort(reset); + // Monitor for GOAWAY after setup completes. + if let Err(err) = run_goaway(reader.with_version(version), version, goaway).await { + tracing::warn!(%err, "goaway error"); + } + }); } - }); + UniType::Subgroup => { + let mut sub = subscriber.clone(); + tasks.push(async move { + let mut reader = reader.with_version(version); + let res = sub.recv_group(&mut reader).await; + stop_on_error(&mut reader, res); + }); + } + // A fill fetch stream carries the head of the group a draft-20 subscription + // joined part way through. One answering no fill of ours is refused inside. + UniType::Fetch => { + let mut sub = subscriber.clone(); + tasks.push(async move { + let mut reader = reader.with_version(version); + let res = sub.recv_fill(&mut reader).await; + stop_on_error(&mut reader, res); + }); + } + // The receiver MUST discard padding. We read it to the end rather than cancel, + // so a peer probing for bandwidth gets the throughput it is measuring. + UniType::Padding => { + tasks.push(async move { + while let Ok(Some(_)) = std::future::poll_fn(|cx| reader.poll_read_chunk(cx, usize::MAX)).await {} + }); + } + } } } @@ -762,30 +820,19 @@ where } } -async fn run_uni_group( - subscriber: &mut Subscriber, - stream: &mut Reader, -) -> Result<(), Error> -where - S: crate::transport::poll::Boxable, -{ - let kind: u64 = stream.decode_peek().await?; - - // SUBGROUP_HEADER type bytes match the form 0b0XX1XXXX (spec ยง11.4.2): - // draft-14-17 use 0x10-0x1D and 0x30-0x3D, draft-18 adds 0x40 (FIRST_OBJECT) - // extending the form to also cover 0x50-0x5D and 0x70-0x7D. Per-version and - // per-bit validation (e.g., FIRST_OBJECT must be 0 on draft-17) is done in - // `GroupFlags::decode`. - if kind <= 0xff && (kind & 0x90) == 0x10 { - return subscriber.recv_group(stream).await; - } +/// Stop a data stream whose handler failed. The handler owns only the stream, so it +/// cannot claim the session closed. +fn stop_on_error(reader: &mut Reader, res: Result<(), Error>) { + let Err(err) = res else { + return; + }; - match kind { - // A fill fetch stream carries the head of the group a draft-20 subscription joined - // part way through. One answering no fill of ours is refused inside. - FetchHeader::TYPE => subscriber.recv_fill(stream).await, - _ => Err(Error::UnexpectedStream), - } + tracing::debug!(%err, "uni stream error"); + let reset = match StreamError::from(&err) { + StreamError::Session(_) => StreamError::Internal, + reset => reset, + }; + reader.abort(reset); } /// Accept incoming bidi streams and dispatch to the correct handler based on message type. @@ -1325,9 +1372,13 @@ mod tests { writes.clone() } - async fn dispatch_uni(payload: Vec, retired_alias: Option) -> crate::lite::test_transport::Log { - const VERSION: Version = Version::Draft19; - + /// Run the uni dispatch loop over one incoming stream. Returns what reached the wire, + /// and the loop's result if that stream ended it. + async fn dispatch_uni( + version: Version, + payload: Vec, + retired_alias: Option, + ) -> (crate::lite::test_transport::Log, Option>) { let origin = crate::origin::Config::new(crate::Hop::new(1).unwrap()).produce(); // The peer opens one uni stream and then goes quiet, so // the loop is still running when the assertion is taken. @@ -1345,7 +1396,7 @@ mod tests { peer_setup.clone(), crate::Hop::new(1).unwrap(), None, - VERSION, + version, tasks, Default::default(), ); @@ -1357,11 +1408,11 @@ mod tests { let (_goaway, goaway) = crate::goaway::Handle::new(false); // The peer's SETUP has not arrived yet, which is the state a group stream racing // ahead of it lands in. - let mut unis = std::pin::pin!(run_unis(session, subscriber, Some(peer_setup), false, VERSION, goaway)); + let mut unis = std::pin::pin!(run_unis(session, subscriber, Some(peer_setup), false, version, goaway)); for _ in 0..100 { if let std::task::Poll::Ready(result) = futures::poll!(unis.as_mut()) { - panic!("the dispatch loop ended over one rejected stream: {result:?}"); + return (log, Some(result)); } if !log.stops().is_empty() { break; @@ -1369,36 +1420,87 @@ mod tests { tokio::time::sleep(std::time::Duration::from_millis(1)).await; } - log + (log, None) } /// A late group must reach the dispatch loop and stop with CANCELLED. #[tokio::test(start_paused = true)] async fn a_group_for_a_retired_alias_is_stopped_with_cancelled() { - let log = dispatch_uni(subgroup_header(Version::Draft19, 7, 0).await, Some(7)).await; + let (log, result) = + dispatch_uni(Version::Draft19, subgroup_header(Version::Draft19, 7, 0).await, Some(7)).await; assert_eq!( log.stops(), vec![crate::ietf::error::CANCELLED], "the group stream must be stopped with the cancelled code", ); + assert!(result.is_none(), "one dropped group ended the session: {result:?}"); assert_eq!(log.closes(), vec![], "one dropped group may not close the session"); } /// A non-zero subgroup is refused on its own stream, never by closing the session. #[tokio::test(start_paused = true)] async fn a_non_zero_subgroup_is_stopped_without_closing_the_session() { - let log = dispatch_uni(subgroup_header(Version::Draft19, 7, 1).await, None).await; + let (log, result) = dispatch_uni(Version::Draft19, subgroup_header(Version::Draft19, 7, 1).await, None).await; assert_eq!(log.stops(), vec![crate::ietf::error::INTERNAL_ERROR]); + assert!(result.is_none(), "one refused subgroup ended the session: {result:?}"); assert_eq!(log.closes(), vec![], "one refused subgroup may not close the session"); } + /// A stream type encoded for `version`, followed by a few bytes of body. + fn uni_stream(version: Version, kind: u64) -> Vec { + let mut buf = Vec::new(); + kind.encode(&mut buf, version).unwrap(); + buf.extend_from_slice(&[0; 4]); + buf + } + + /// Padding is read and dropped: no STOP_SENDING, and the session stays up. #[tokio::test(start_paused = true)] - async fn unknown_uni_type_does_not_claim_the_session_closed() { - let log = dispatch_uni(vec![0], None).await; - assert_eq!(log.stops(), vec![crate::ietf::error::INTERNAL_ERROR]); - assert!(log.closes().is_empty()); + async fn a_padding_stream_is_discarded() { + for version in [ + Version::Draft18, + Version::Draft19, + Version::Draft20, + Version::Draft21, + Version::Draft22, + ] { + let (log, result) = dispatch_uni(version, uni_stream(version, PADDING), None).await; + + assert!( + log.stops().is_empty(), + "{version:?}: padding was stopped: {:?}", + log.stops() + ); + assert!(result.is_none(), "{version:?}: padding ended the session: {result:?}"); + assert_eq!(log.closes(), vec![], "{version:?}"); + } + } + + /// An unknown or invalid stream type MUST close the session, so it stops nothing on its own: + /// the session close takes the stream with it. + #[tokio::test(start_paused = true)] + async fn an_unknown_uni_type_closes_the_session() { + for (version, kind) in [ + (Version::Draft19, 0), + // Padding and uni SETUP arrived in later drafts, so earlier ones do not know them. + (Version::Draft17, PADDING), + (Version::Draft16, setup::SETUP_V17), + // SUBGROUP_HEADER types with the reserved SUBGROUP_ID_MODE (0b11). + (Version::Draft19, 0x56), + (Version::Draft14, 0x16), + // FIRST_OBJECT (0x40) arrived in draft-18. + (Version::Draft17, 0x50), + ] { + let (log, result) = dispatch_uni(version, uni_stream(version, kind), None).await; + + let Some(Err(err)) = result else { + panic!("{version:?}: type {kind:#x} did not end the session: {result:?}"); + }; + assert_eq!(SessionError::from(&err), SessionError::ProtocolViolation, "{version:?}"); + assert!(log.stops().is_empty(), "{version:?}: stopped {:?}", log.stops()); + } } /// A peer's advertisement of `room/host`, then two namespace-keyed withdrawals of it.