@@ -6,6 +6,7 @@ import type { Reader, Stream } from "../stream.ts";
66import type { Track } from "../track.ts" ;
77import { error } from "../util/error.ts" ;
88import type { Session } from "./adapter.ts" ;
9+ import { TrackAliases } from "./aliases.ts" ;
910import { Frame , type Group as GroupMessage } from "./object.ts" ;
1011import { type Publish , PublishError } from "./publish.ts" ;
1112import { type PublishNamespace , PublishNamespaceError , PublishNamespaceOk } from "./publish_namespace.ts" ;
@@ -31,8 +32,8 @@ import { Version } from "./version.ts";
3132export class Subscriber {
3233 #session: Session ;
3334
34- // Our subscribed tracks — keyed by trackAlias for group routing
35- #subscribes = new Map < bigint , Track > ( ) ;
35+ // Publisher-chosen aliases used by incoming group streams.
36+ #aliases = new TrackAliases < Track > ( ) ;
3637
3738 // Any currently active announcements.
3839 #announced = new Set < Path . Valid > ( ) ;
@@ -219,20 +220,19 @@ export class Subscriber {
219220 await msg . encode ( stream . writer , version ) ;
220221 console . debug ( `subscribe written: id=${ requestId } broadcast=${ broadcast } track=${ request . track . name } ` ) ;
221222
222- // Pre-register with requestId so early group uni streams aren't dropped.
223- // The publisher typically uses requestId as the trackAlias.
224- this . #subscribes. set ( requestId , request . track ) ;
225-
226223 // Read response (SubscribeOk or error)
227224 const respTypeId = await stream . reader . u53 ( ) ;
228225 if ( respTypeId === SubscribeOk . id ) {
229226 const ok = await SubscribeOk . decode ( stream . reader , version ) ;
230- // Update registration to use the actual trackAlias from SubscribeOk
231- if ( ok . trackAlias !== requestId ) {
232- this . #subscribes. delete ( requestId ) ;
233- this . #subscribes. set ( ok . trackAlias , request . track ) ;
227+ try {
228+ this . #aliases. set ( ok . trackAlias , request . track ) ;
229+ } catch ( err ) {
230+ this . #session. close ( ) ;
231+ throw err ;
234232 }
235- console . debug ( `subscribe ok: id=${ requestId } broadcast=${ broadcast } track=${ request . track . name } ` ) ;
233+ console . debug (
234+ `subscribe ok: id=${ requestId } broadcast=${ broadcast } track=${ request . track . name } alias=${ ok . trackAlias } ` ,
235+ ) ;
236236
237237 try {
238238 // Wait for stream close (= PublishDone) or track close (= local unsubscribe)
@@ -259,12 +259,9 @@ export class Subscriber {
259259 `subscribe close: id=${ requestId } broadcast=${ broadcast } track=${ request . track . name } ` ,
260260 ) ;
261261 } finally {
262- this . #subscribes . delete ( ok . trackAlias ) ;
262+ this . #aliases . delete ( ok . trackAlias , request . track ) ;
263263 }
264264 } else {
265- // Clean up pre-registered entry on error
266- this . #subscribes. delete ( requestId ) ;
267-
268265 // Error response
269266 let reasonPhrase = "unknown error" ;
270267 try {
@@ -282,7 +279,6 @@ export class Subscriber {
282279 throw new Error ( `SUBSCRIBE error: ${ reasonPhrase } ` ) ;
283280 }
284281 } catch ( err ) {
285- this . #subscribes. delete ( requestId ) ;
286282 stream . abort ( error ( err ) ) ;
287283 throw err ;
288284 }
@@ -417,12 +413,8 @@ export class Subscriber {
417413 }
418414
419415 try {
420- // Look up by trackAlias directly
421- const track = this . #subscribes. get ( group . trackAlias ) ;
422- if ( ! track ) {
423- // Fallback: try treating trackAlias as requestId (for compat)
424- throw new Error ( `unknown track: trackAlias=${ group . trackAlias } ` ) ;
425- }
416+ // The control message establishing this alias can arrive after the data stream.
417+ const track = await this . #aliases. get ( group . trackAlias ) ;
426418
427419 track . writeGroup ( producer ) ;
428420
0 commit comments