Skip to content

Commit 9eba04c

Browse files
okdaichiclaude
andauthored
fix: end JS subscriptions on TrackReader close (#447)
* fix: end JS subscriptions on TrackReader close * fix: align JS track close with Go and reject late groups * docs: note JS TrackReader close fix * fix(moq-web): untrack a GroupReader whose stream ends with any read error TrackReader tracks the GroupReaders it hands out until they end, so close() can cancel the active ones. A GroupReader reported its end only at EOF, so a group stream reset by the publisher, or failing another read, stayed tracked until the subscription closed: on a long-lived subscription the set grew with every reset group. readFrame now reports the end on any stream read error. An oversized length or a sink error leaves the stream open, so it stays tracked until the caller cancels it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: OkutaniDaichi0106 <132345842+OkutaniDaichi0106@users.noreply.github.com> Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 570f5af commit 9eba04c

9 files changed

Lines changed: 262 additions & 13 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
2121
- **webtransport-go:** bumped to `v0.13.0-okdaichi.2`. The session context now ends with the close error as its cause (okdaichi/webtransport-go#9). A dialed session's QUIC connection now stays open until its close capsule can arrive, instead of closing at once and losing it (#11). The fork also syncs with upstream v0.13.0, and quic-go goes to v0.63.0.
2222
- **`moqt.Cause`** converts that `*webtransport.SessionError` into the same `SessionError` it gives for native QUIC. The WebTransport wrapper passes the session's context through unchanged.
2323
- **On the server,** a WebTransport session still takes its values from the upgrade request. It now ends with the transport, or with the connection's `ConnContext` context. Before, the request's context could end first, with no cause, and hide the close reason. Cancelling `ConnContext`'s context still ends the session, as since v0.20.1.
24+
- **@qumo/moq: closing a `TrackReader` now ends its subscription.** A normal
25+
`close()` sends FIN on the subscribe stream and cancels queued and active
26+
group streams. `closeWithError()` cancels them with the appropriate error
27+
codes, and late group streams are cancelled instead of left unread.
2428

2529
## [v0.22.0] - 2026-10-04
2630

‎moq-web/src/group_stream.ts‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -134,22 +134,40 @@ export class GroupReader {
134134
#reader: ReceiveStream;
135135
readonly context: Context;
136136
#cancelFunc: CancelCauseFunc;
137+
#onDone?: () => void;
137138
// Previous frame's timestamp on this stream, used to resolve the next
138139
// frame's delta-encoded timestamp.
139140
#prevTimestamp: number = 0;
140141
/** Timestamp of the most recently read frame, in timescale units. */
141142
lastTimestamp: number = 0;
142143

143-
constructor(trackCtx: Context, reader: ReceiveStream, group: GroupMessage) {
144+
constructor(
145+
trackCtx: Context,
146+
reader: ReceiveStream,
147+
group: GroupMessage,
148+
onDone?: () => void,
149+
) {
144150
this.sequence = group.sequence;
145151
this.#reader = reader;
152+
this.#onDone = onDone;
146153
[this.context, this.#cancelFunc] = withCancelCause(trackCtx);
147154

148155
trackCtx.done().then(() => {
149156
this.cancel(GroupErrorCode.PublishAborted);
150157
});
151158
}
152159

160+
/**
161+
* Records that the stream itself ended, at EOF or with a reset or another
162+
* read error, so the reader's owner stops tracking it, and returns err.
163+
* An oversized length or a sink error leaves the stream open: the caller
164+
* still owns it and cancels it.
165+
*/
166+
#streamEnded(err: Error): Error {
167+
this.#onDone?.();
168+
return err;
169+
}
170+
153171
/**
154172
* Read a single frame into a sink.
155173
* @param sink - A {@link ByteSink}, {@link ByteSinkFunc}, or callback receiving the raw bytes.
@@ -170,14 +188,14 @@ export class GroupReader {
170188
// propagate it verbatim so callers can detect a normal end‑of‑stream
171189
// and break out of their read loops (see `frames()` below).
172190
if (errTs) {
173-
return errTs;
191+
return this.#streamEnded(errTs);
174192
}
175193
const ts = this.#prevTimestamp + zigzagDecode(delta);
176194

177195
// Read length prefix as varint
178196
const [len, , err1] = await readVarint(this.#reader);
179197
if (err1) {
180-
return err1;
198+
return this.#streamEnded(err1);
181199
}
182200

183201
// Reject an oversized length before allocating — guards against an
@@ -196,7 +214,7 @@ export class GroupReader {
196214
const dst = sink.reserve(len);
197215
const [, err2] = await readFull(this.#reader, dst);
198216
if (err2) {
199-
return err2;
217+
return this.#streamEnded(err2);
200218
}
201219
sink.timestamp = ts;
202220
this.#prevTimestamp = ts;
@@ -210,7 +228,7 @@ export class GroupReader {
210228
// Read the frame data
211229
const [, err2] = await readFull(this.#reader, buf);
212230
if (err2) {
213-
return err2;
231+
return this.#streamEnded(err2);
214232
}
215233

216234
// Write to sink (handle both ByteSink and ByteSinkFunc)
@@ -243,6 +261,7 @@ export class GroupReader {
243261
false,
244262
);
245263
this.#cancelFunc(reason);
264+
this.#onDone?.();
246265
await this.#reader.cancel(code);
247266
}
248267

‎moq-web/src/group_stream_test.ts‎

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -402,3 +402,51 @@ Deno.test("GroupReader", async (t) => {
402402
}
403403
});
404404
});
405+
406+
Deno.test("GroupReader.readFrame tells its owner when the stream ends", async (t) => {
407+
const reader = (read: ReceiveStream["read"]): ReceiveStream => ({
408+
read,
409+
cancel: async (_code: number) => {},
410+
closed: () => new Promise<void>(() => {}),
411+
});
412+
const cases = [
413+
{ name: "at EOF", read: async () => [0, new EOFError()] as [number, Error] },
414+
{ name: "on a reset", read: async () => [0, new Error("stream reset")] as [number, Error] },
415+
];
416+
for (const c of cases) {
417+
await t.step(c.name, async () => {
418+
const onDone = spy(() => {});
419+
const gr = new GroupReader(
420+
background(),
421+
reader(c.read),
422+
new GroupMessage({ sequence: 1 }),
423+
onDone,
424+
);
425+
426+
const err = await gr.readFrame(new Frame(new ArrayBuffer(8)));
427+
428+
assertInstanceOf(err, Error);
429+
assertEquals(onDone.calls.length, 1);
430+
});
431+
}
432+
433+
await t.step("not for an oversized length: the stream is still open", async () => {
434+
const buf = new Buffer(new ArrayBuffer(0));
435+
await writeVarint(buf, 0); // timestamp delta
436+
await writeVarint(buf, MAX_FRAME_SIZE + 1);
437+
const src = new Buffer(new ArrayBuffer(0));
438+
src.write(buf.bytes());
439+
const onDone = spy(() => {});
440+
const gr = new GroupReader(
441+
background(),
442+
reader(src.read.bind(src)),
443+
new GroupMessage({ sequence: 1 }),
444+
onDone,
445+
);
446+
447+
const err = await gr.readFrame(new Frame(new ArrayBuffer(8)));
448+
449+
assertInstanceOf(err, Error);
450+
assertEquals(onDone.calls.length, 0);
451+
});
452+
});

‎moq-web/src/internal/queue.ts‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,22 @@ export class Queue<T> {
2121
}
2222
}
2323

24+
async tryEnqueue(item: T): Promise<boolean> {
25+
await this.#mutex.lock();
26+
try {
27+
if (this.#closed) return false;
28+
this.#items.push(item);
29+
if (this.#pending) {
30+
const [resolve] = this.#pending;
31+
this.#pending = undefined;
32+
resolve();
33+
}
34+
return true;
35+
} finally {
36+
this.#mutex.unlock();
37+
}
38+
}
39+
2440
async dequeue(): Promise<T | undefined> {
2541
while (true) {
2642
await this.#mutex.lock();
@@ -65,6 +81,15 @@ export class Queue<T> {
6581
}
6682
}
6783

84+
async drain(): Promise<T[]> {
85+
await this.#mutex.lock();
86+
try {
87+
return this.#items.splice(0);
88+
} finally {
89+
this.#mutex.unlock();
90+
}
91+
}
92+
6893
close(): void {
6994
if (this.#closed) {
7095
return;

‎moq-web/src/internal/queue_test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -115,6 +115,13 @@ Deno.test("internal/queue - basic enqueue/dequeue and close behavior", async (t)
115115
assertEquals(v, 1);
116116
});
117117

118+
await t.step("tryEnqueue rejects items after close", async () => {
119+
const q = new Queue<number>();
120+
q.close();
121+
assertEquals(await q.tryEnqueue(1), false);
122+
assertEquals(await q.dequeue(), undefined);
123+
});
124+
118125
await t.step("dequeue after multiple enqueues and closes", async () => {
119126
const q = new Queue<number>();
120127
await q.enqueue(1);

‎moq-web/src/session.ts‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -709,12 +709,13 @@ export class Session {
709709

710710
const queue = this.#queues.get(req.subscribeId);
711711
if (!queue) {
712-
// No enqueue function yet.
713-
// This can happen if the subscribe call is not completed yet.
712+
await reader.cancel(GroupErrorCode.SubscribeCanceled);
714713
return;
715714
}
716715
try {
717-
await queue.enqueue([reader, req]);
716+
if (!await queue.tryEnqueue([reader, req])) {
717+
await reader.cancel(GroupErrorCode.SubscribeCanceled);
718+
}
718719
} catch (e) {
719720
console.error(
720721
`moq: failed to enqueue group for subscribe ID ${req.subscribeId}:`,

‎moq-web/src/subscribe_stream.ts‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -192,12 +192,19 @@ export class SendSubscribeStream {
192192
}
193193

194194
async closeWithError(code: SubscribeErrorCode): Promise<void> {
195+
if (this.context.err()) return;
195196
const err = new WebTransportStreamError({
196197
source: "stream",
197198
streamErrorCode: code,
198199
}, false);
199-
await this.#stream.writable.cancel(code);
200200
this.#cancelFunc(err);
201+
await this.#stream.writable.cancel(code);
202+
}
203+
204+
async close(): Promise<void> {
205+
if (this.context.err()) return;
206+
this.#cancelFunc(undefined);
207+
await this.#stream.writable.close();
201208
}
202209
}
203210

‎moq-web/src/track_reader.ts‎

Lines changed: 34 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import { GroupMessage } from "./internal/message/mod.ts";
77
import type { BroadcastPath } from "./broadcast_path.ts";
88
import { Queue } from "./internal/queue.ts";
99
import type { SubscribeID } from "./alias.ts";
10+
import { GroupErrorCode } from "./error.ts";
1011

1112
/**
1213
* Subscriber-side handle for reading groups from a subscribed track.
@@ -22,6 +23,8 @@ export class TrackReader {
2223
#subscribeStream: SendSubscribeStream;
2324
#queue: Queue<[ReceiveStream, GroupMessage]>;
2425
#onCloseFunc: () => void;
26+
#groups = new Set<GroupReader>();
27+
#closing = false;
2528

2629
constructor(
2730
broadcastPath: BroadcastPath,
@@ -45,6 +48,7 @@ export class TrackReader {
4548
async acceptGroup(
4649
signal: Promise<void>,
4750
): Promise<[GroupReader, undefined] | [undefined, Error]> {
51+
if (this.#closing) return [undefined, new ContextCancelledError()];
4852
// Check if context is already cancelled
4953
const err = this.context.err();
5054
if (err) {
@@ -69,13 +73,22 @@ export class TrackReader {
6973
return [undefined, dequeued];
7074
}
7175
if (dequeued === undefined) {
72-
// This is
73-
throw new Error("dequeue returned undefined");
76+
return [undefined, new ContextCancelledError()];
7477
}
7578

7679
const [reader, msg] = dequeued;
80+
if (this.#closing) {
81+
await reader.cancel(GroupErrorCode.SubscribeCanceled);
82+
return [undefined, new ContextCancelledError()];
83+
}
7784

78-
const group = new GroupReader(this.context, reader, msg);
85+
const group = new GroupReader(
86+
this.context,
87+
reader,
88+
msg,
89+
() => this.#groups.delete(group),
90+
);
91+
this.#groups.add(group);
7992

8093
return [group, undefined];
8194
}
@@ -100,12 +113,29 @@ export class TrackReader {
100113
}
101114

102115
async closeWithError(code: number): Promise<void> {
103-
await this.#subscribeStream.closeWithError(code);
116+
if (this.#closing) return;
117+
this.#closing = true;
104118
this.#onCloseFunc();
119+
await this.#cancelGroups(code);
120+
await this.#subscribeStream.closeWithError(code);
105121
}
106122

107123
async close(): Promise<void> {
124+
if (this.#closing) return;
125+
this.#closing = true;
108126
this.#onCloseFunc();
127+
await this.#cancelGroups(GroupErrorCode.SubscribeCanceled);
128+
await this.#subscribeStream.close();
129+
}
130+
131+
async #cancelGroups(queuedCode: number): Promise<void> {
132+
this.#queue.close();
133+
const pending = await this.#queue.drain();
134+
await Promise.all([
135+
...pending.map(([reader]) => reader.cancel(queuedCode)),
136+
...Array.from(this.#groups, (group) => group.cancel(GroupErrorCode.SubscribeCanceled)),
137+
]);
138+
this.#groups.clear();
109139
}
110140

111141
get subscribeId(): SubscribeID {

0 commit comments

Comments
 (0)