diff --git a/STATUS.md b/STATUS.md index 364ce8b9..a5abae7b 100644 --- a/STATUS.md +++ b/STATUS.md @@ -107,7 +107,7 @@ By package, bottom-up along the dependency stack: | § | Feature | Status | Notes | |---------|----------------------------------|--------|-------| -| 5.1 | Subscriptions | DONE | Subscribe/Publish/OK/Error state machine in `subscribe.go` and `publish.go`. A second response to this side's SUBSCRIBE or PUBLISH closes the session with PROTOCOL_VIOLATION: on a `Subscription`'s or `Session.Publish` `Publication`'s broker, a SUBSCRIBE_OK, or a REQUEST_OK / REQUEST_ERROR before any Update (after one it may answer an Update that gave up); in the relay, a REQUEST_OK / REQUEST_ERROR on a forwarded PUBLISH. | +| 5.1 | Subscriptions | DONE | Subscribe/Publish/OK/Error state machine in `subscribe.go` and `publish.go`. A second response to this side's SUBSCRIBE or PUBLISH closes the session with PROTOCOL_VIOLATION: a SUBSCRIBE_OK on a `Subscription`'s broker, and a REQUEST_OK / REQUEST_ERROR on any broker before this side sent an Update (after one it may answer an Update that gave up); see §10.9. | | 5.1.1 | Subscription state management | DONE | REQUEST_ERROR / STOP_SENDING / PUBLISH_DONE handling + cleanup. The relay resets a cancelled subscription's open subgroup and fill streams. | | 5.1.2 | Location filters | DONE | Every start/end form (unfiltered, Next Object, relative and absolute start, absolute range) + `Matches`. | | 5.1.3 | Fill semantics | PARTIAL | Fill fetch streams from FILL_PARAMETERS on SUBSCRIBE / REQUEST_UPDATE (`handler_fill.go`), and on SUBSCRIBE_TRACKS, one per forwarded PUBLISH's subscription, keyed to the PUBLISH's Request ID (§10.1). A fill inherits the subscription's Range Filters; the ones inside FILL_PARAMETERS override per type. A cancelled subscription's open fills are reset (§5.1.3.1). Not done: scheduling fills against their subscription (§7.2, see Limitations). | @@ -188,12 +188,12 @@ By package, bottom-up along the dependency stack: | 10.3.1.5| MOQT_IMPLEMENTATION | 0x07 | DONE | Advisory. | | 10.3.1.6| MAX_FILTER_RANGES | 0x06 | DONE | `WithMaxFilterRanges` advertises it; relay rejects over-limit/prohibited filters with INVALID_FILTER. The relay advertises `relay.DefaultMaxFilterRanges` (16) rather than inheriting the session default of 0, which would prohibit the Range Filters it fully implements; `relay.Config.MaxFilterRanges` overrides, negative to prohibit. | | 10.3.1.7| MAX_REQUEST_UPDATES | 0x08 | DONE | `WithMaxRequestUpdates` advertises the per-stream limit; enforced on inbound follow-ups via `RequestUpdateLimiter`, closing with `TOO_MANY_REQUEST_UPDATES` on overflow. | -| 10.4 | GOAWAY | 0x10 | DONE | Same encoding on control and request streams (draft-19 dropped the Request ID field); callback. As recipient the relay initiates no new SUBSCRIBE, FETCH or PUBLISH to the peer and leaves closing the session to the sender. On a request stream (`RequestBroker.Serve` and the relay's readers after the response; `session.RequestGoaways` for callers that read one themselves) a second GOAWAY, or one with a New Session URI received by a server, closes the session with PROTOCOL_VIOLATION; a single one is handed to the reader, and neither side migrates the request. | +| 10.4 | GOAWAY | 0x10 | DONE | Same encoding on control and request streams (draft-19 dropped the Request ID field); callback. As recipient the relay initiates no new SUBSCRIBE, FETCH or PUBLISH to the peer and leaves closing the session to the sender. On a request stream (`RequestBroker.Serve` and the relay's readers after the response; `session.RequestGoaways` for callers that read one themselves) a second GOAWAY, or one with a New Session URI received by a server, closes the session with PROTOCOL_VIOLATION; a single one is handed to the reader, and neither side migrates the request. One before the response is checked the same way, the request still awaits its response, and the GOAWAY is the stream's first follow-up; bar PUBLISH_NAMESPACE, SUBSCRIBE_NAMESPACE and SUBSCRIBE_TRACKS, whose response MUST come first (§6.1, §6.2). | | 10.5 | REQUEST_OK | 0x07 | DONE | Shared OK for PUBLISH/UPDATE/TRACK_STATUS/namespace reqs. Track Properties where they must be empty close the session on receipt and are refused on send (`ErrTrackPropertiesNotAllowed`). | | 10.6 | REQUEST_ERROR (+ Redirect) | 0x05 | DONE | Redirect required only when code==REDIRECT, and exposed as `RequestRejectedError.Redirect`; `Request.Reject` sends one. A Connect URI received by a server, or a Track Name for SUBSCRIBE_NAMESPACE / PUBLISH_NAMESPACE / SUBSCRIBE_TRACKS, closes the session (§10.6.1), and `Reject` refuses to send either; on a REQUEST_UPDATE's answer only the Connect URI is checked, as the reader does not know the request, and an update handler's REDIRECT is sent as INTERNAL_ERROR (§10.6.2 does not list REQUEST_UPDATE). The relay does not follow a Redirect: an upstream REDIRECT becomes INTERNAL_ERROR downstream. `Request.Reject` sends a Retry Interval; the relay invites a jittered ~1 s retry on EXCESSIVE_LOAD and passes an upstream SUBSCRIBE rejection on by meaning, Retry Interval kept. | | 10.7 | SUBSCRIBE | 0x03 | DONE | | | 10.8 | SUBSCRIBE_OK | 0x04 | DONE | Registers inbound track alias. | -| 10.9 | REQUEST_UPDATE | 0x02 | DONE | A REQUEST_UPDATE opening a request stream closes the session with PROTOCOL_VIOLATION (`ErrUnexpectedRequestUpdate`). | +| 10.9 | REQUEST_UPDATE | 0x02 | DONE | A REQUEST_UPDATE opening a request stream closes the session with PROTOCOL_VIOLATION (`ErrUnexpectedRequestUpdate`). A REQUEST_OK / REQUEST_ERROR on a request stream where this side sent no REQUEST_UPDATE answers nothing and closes the session too, read by a `RequestBroker` or by the relay: a deliberate choice, as the draft names no rule for it on a stream this side answered. | | 10.10 | PUBLISH_STATE_NOTIFY | 0x22 | DONE | Only the publisher may send it; enforced by brokers and the relay. | | 10.11 | PUBLISH | 0x1D | DONE | | | 10.12 | PUBLISH_DONE | 0x0B | DONE | Sent once every stream of the subscription has closed and no datagram send is in progress, with the exact Stream Count; written on its own goroutine, so subscribers do not wait on each other. When a track's last upstream ends, its PUBLISH_DONE code reaches subscribers if it is about the track (TRACK_ENDED, MALFORMED_TRACK); codes about the relay's own upstream subscription become INTERNAL_ERROR. Session `Publication.Done` resets the subgroups still open with CANCELLED, refuses later opens and writes (`ErrPublicationEnded`), and counts every subgroup opened, however the opens race it. | @@ -519,18 +519,6 @@ A second full review against draft-ietf-moq-transport-20 (2026-09-26, at `dbe571e`) found the gaps below. Each item names the rule it misses. Items already listed as Limitations above are not repeated here. -Session layer: - -- On a request stream the relay answered, a REQUEST_OK or REQUEST_ERROR from - the requester is ignored rather than closing the session, so a Connect URI in - one is not checked (§10.6.1). The draft defines no such message; only on a - forwarded PUBLISH, where it is a second response, does the relay close the - session (§5.1). -- A GOAWAY on a request stream before its initial response (other than to - SUBSCRIBE_NAMESPACE or SUBSCRIBE_TRACKS, where §10.19/§10.20 make it a - PROTOCOL_VIOLATION) is an error rather than a legal message checked - against §10.4. - Relay: - Objects age from when they were read whole rather than their beginning diff --git a/pkg/moqt/session/broker.go b/pkg/moqt/session/broker.go index e5eb06bd..0920dfc8 100644 --- a/pkg/moqt/session/broker.go +++ b/pkg/moqt/session/broker.go @@ -32,10 +32,11 @@ import ( // [RequestBroker.HandleUpdates], or declined with NOT_SUPPORTED when there // is none, since acknowledging an unapplied update would misstate the // request's state. -// - On the broker of a [Subscription], or of a [Publication] from -// [Session.Publish], a REQUEST_OK / REQUEST_ERROR before any Update, or a -// SUBSCRIBE_OK on a Subscription, is a second response to the request and -// closes the session (§5.1). +// - A REQUEST_OK / REQUEST_ERROR before this side sent any REQUEST_UPDATE +// answers nothing and closes the session: on a request this side sent it +// is a second response (§5.1, §5.2, §6.2), and on a stream this side +// answered only the requester sends REQUEST_UPDATE, bar a PUBLISH's +// subscriber (§10.9). So does a SUBSCRIBE_OK on a Subscription (§5.1). // - On a [NamespaceSubscription]'s broker, a NAMESPACE_DONE for a suffix no // NAMESPACE announced closes the session (§10.19). // - A second GOAWAY on the stream, or one with a New Session URI received @@ -248,7 +249,9 @@ var ErrRequestStreamClosed = errors.New("moqt/session: request stream closed") // NewRequestBroker builds a [RequestBroker] for an established request // stream. Typed request handles expose a Broker method that fills this in; // use this constructor for accept-side streams (a [Request] this endpoint -// accepted). +// accepted). On a request this side sent, read its response first: the +// broker reads any REQUEST_OK or REQUEST_ERROR as answering a REQUEST_UPDATE, +// and closes the session on one when none was sent. func (s *Session) NewRequestBroker(stream Stream) *RequestBroker { return &RequestBroker{stream: stream, sess: s} } @@ -513,10 +516,10 @@ func (b *RequestBroker) readFailed(ctx context.Context, err error) error { // // Responses route to Update waiters; a token cache fault closes the session // (§10.2.2); peer REQUEST_UPDATEs are answered as described on -// [RequestBroker], and a second response to this side's SUBSCRIBE or PUBLISH -// closes the session (§5.1). Every other message, including each -// REQUEST_UPDATE and any other unsolicited response, is passed to onMsg (nil -// means "discard"); return false from onMsg to stop serving. +// [RequestBroker], and a REQUEST_OK or REQUEST_ERROR before this side sent a +// REQUEST_UPDATE closes the session (§5.1, §10.9). Every other message, +// including each REQUEST_UPDATE, is passed to onMsg (nil means "discard"); +// return false from onMsg to stop serving. // // A read error resets the read side with INTERNAL_ERROR; a malformed follow-up // also closes the session with PROTOCOL_VIOLATION (§10). Serve returns nil on @@ -570,17 +573,17 @@ func (b *RequestBroker) Serve(ctx context.Context, onMsg func(message.Message) b case *message.RequestOK, *message.RequestError: // The request's own response was read before the broker // attached, so every REQUEST_OK here is a REQUEST_UPDATE_OK - // (§10.5), even one whose Update gave up. Before any Update, on - // this side's SUBSCRIBE or PUBLISH, it is a second response to - // the request, a second PUBLISH_OK among them (§5.1). - if b.answered() != 0 { - b.mu.Lock() - updated := b.updated - b.mu.Unlock() - if !updated { - return b.sess.closeProtocolViolation( - fmt.Errorf("moqt/session: %s after the response, before any REQUEST_UPDATE", m.Type())) - } + // (§10.5), even one whose Update gave up. Before any Update it + // answers nothing: on a request this side sent it is a second + // response, a second PUBLISH_OK among them (§5.1, §5.2, §6.2), + // and on a stream this side answered the peer, as requester, was + // sent no REQUEST_UPDATE to answer (§10.9). + b.mu.Lock() + updated := b.updated + b.mu.Unlock() + if !updated { + return b.sess.closeProtocolViolation( + fmt.Errorf("moqt/session: %s before any REQUEST_UPDATE", m.Type())) } if err := b.sess.checkRequestOKTrackProperties(nil, m); err != nil { return err diff --git a/pkg/moqt/session/early_goaway_test.go b/pkg/moqt/session/early_goaway_test.go new file mode 100644 index 00000000..ba16ca09 --- /dev/null +++ b/pkg/moqt/session/early_goaway_test.go @@ -0,0 +1,164 @@ +package session_test + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt" + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/moqt/wire" +) + +// A GOAWAY "MAY also be sent on a request stream to initiate migration of that +// individual request" (§10.4), before the request's response too, unless the +// response MUST come first (§6.1, §6.2). The requester keeps waiting for the +// response and reads the GOAWAY as the stream's first follow-up, as it would +// one sent after the response. + +// answerAfterGoaways accepts one request on sess, sends it goaways, then +// answers it with reply. +func answerAfterGoaways(t *testing.T, sess *session.Session, goaways []*message.Goaway, reply func(*session.Request)) { + t.Helper() + go func() { + r, err := sess.AcceptRequest(t.Context()) + if err != nil { + return + } + for _, g := range goaways { + if message.Marshal(r.Stream, g) != nil { + return + } + } + reply(r) + }() +} + +var migrate = &message.Goaway{NewSessionURI: []byte("https://relay.example/moq"), Timeout: 1000} + +// TestEarlyRequestGoawayThenResponse: the request succeeds, and the GOAWAY is +// the first follow-up its reader sees, through a broker or a direct read. +func TestEarlyRequestGoawayThenResponse(t *testing.T) { + t.Run("SUBSCRIBE", func(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{migrate}, func(r *session.Request) { + p, err := r.AcceptSubscribe(&message.SubscribeOK{TrackAlias: 1}) + if err == nil { + _ = p.Done(moqt.PublishDoneGoingAway, "") + } + }) + sub, err := cli.Subscribe(t.Context(), &message.Subscribe{Name: []byte("t")}) + if err != nil { + t.Fatalf("Subscribe after an early GOAWAY: %v", err) + } + var got []message.Message + if err := sub.Broker().Serve(t.Context(), func(m message.Message) bool { + got = append(got, m) + return true + }); err != nil { + t.Fatalf("Serve: %v", err) + } + if len(got) == 0 { + t.Fatal("Serve's callback never saw the GOAWAY") + } + if g, ok := got[0].(*message.Goaway); !ok || g.Timeout != migrate.Timeout { + t.Fatalf("first follow-up = %#v, want the early GOAWAY", got[0]) + } + requireStaysOpen(t, cli, 50*time.Millisecond) + }) + t.Run("FETCH", func(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{migrate}, func(r *session.Request) { + _, _ = r.AcceptFetch(&message.FetchOK{}) + }) + fr, err := cli.Fetch(t.Context(), &message.Fetch{Name: []byte("t")}) + if err != nil { + t.Fatalf("Fetch after an early GOAWAY: %v", err) + } + m, err := message.Parse(fr.Stream) + if err != nil { + t.Fatalf("read follow-up: %v", err) + } + if _, ok := m.(*message.Goaway); !ok { + t.Fatalf("first follow-up = %T, want *message.Goaway", m) + } + requireStaysOpen(t, cli, 50*time.Millisecond) + }) +} + +// TestEarlyRequestGoawayOnPublishNamespace: a PUBLISH_NAMESPACE's response +// MUST be "the first message on the bidi stream" (§6.2), so a GOAWAY ahead of +// it fails the request, leaving the session open as with any other unexpected +// first message. +func TestEarlyRequestGoawayOnPublishNamespace(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{migrate}, func(r *session.Request) { + _ = r.Reply(&message.RequestOK{}) + }) + if _, err := cli.PublishNamespace(t.Context(), &message.PublishNamespace{ + Namespace: wire.TrackNamespace{[]byte("ns")}, + }); err == nil { + t.Fatal("PublishNamespace succeeded after a GOAWAY ahead of its response") + } + requireStaysOpen(t, cli, 50*time.Millisecond) +} + +// TestEarlyRequestGoawayThenRejection: a REQUEST_ERROR after the GOAWAY is the +// request's answer. +func TestEarlyRequestGoawayThenRejection(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{migrate}, func(r *session.Request) { + _ = r.RejectError(moqt.RequestDoesNotExist, "gone") + }) + _, err := cli.Subscribe(t.Context(), &message.Subscribe{Name: []byte("t")}) + rej, ok := errors.AsType[*session.RequestRejectedError](err) + if !ok || rej.Code != moqt.RequestDoesNotExist { + t.Fatalf("Subscribe = %v, want its DOES_NOT_EXIST rejection", err) + } + requireStaysOpen(t, cli, 50*time.Millisecond) +} + +// TestEarlyRequestGoawayViolationCloses: the §10.4 checks apply from the +// first message: a second GOAWAY before the response, one after an early one, +// or a New Session URI to the server closes the session with +// PROTOCOL_VIOLATION. +func TestEarlyRequestGoawayViolationCloses(t *testing.T) { + t.Run("two before the response", func(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{{}, {}}, func(r *session.Request) { + _, _ = r.AcceptSubscribe(&message.SubscribeOK{TrackAlias: 1}) + }) + if _, err := cli.Subscribe(t.Context(), &message.Subscribe{Name: []byte("t")}); err == nil { + t.Fatal("Subscribe succeeded after two GOAWAYs") + } + requireClosedCode(t, cli, moqt.SessionProtocolViolation) + }) + t.Run("one before and one after", func(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, srv, []*message.Goaway{{}}, func(r *session.Request) { + if _, err := r.AcceptSubscribe(&message.SubscribeOK{TrackAlias: 1}); err == nil { + _ = message.Marshal(r.Stream, &message.Goaway{}) + } + }) + sub, err := cli.Subscribe(t.Context(), &message.Subscribe{Name: []byte("t")}) + if err != nil { + t.Fatalf("Subscribe: %v", err) + } + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second) + defer cancel() + _ = sub.Broker().Serve(ctx, nil) + requireClosedCode(t, cli, moqt.SessionProtocolViolation) + }) + t.Run("New Session URI to the server", func(t *testing.T) { + cli, srv := openPair(t) + answerAfterGoaways(t, cli, []*message.Goaway{migrate}, func(r *session.Request) { + _, _ = r.AcceptPublish() + }) + if _, err := srv.Publish(t.Context(), &message.Publish{Name: []byte("t"), TrackAlias: 1}); err == nil { + t.Fatal("Publish succeeded after a GOAWAY with a New Session URI to the server") + } + requireClosedCode(t, srv, moqt.SessionProtocolViolation) + }) +} diff --git a/pkg/moqt/session/request.go b/pkg/moqt/session/request.go index 271c87d1..682e8bee 100644 --- a/pkg/moqt/session/request.go +++ b/pkg/moqt/session/request.go @@ -1,6 +1,7 @@ package session import ( + "bytes" "context" "errors" "fmt" @@ -716,7 +717,10 @@ func (s *Session) readResponse(ctx context.Context, stream Stream) (message.Mess // response. An OK is handed to onOK, which then owns the stream; REQUEST_ERROR // (§10.6) becomes a *RequestRejectedError; anything else is an error, and for // SUBSCRIBE_NAMESPACE and SUBSCRIBE_TRACKS also closes the session (§10.19, -// §10.20). On either failure the stream is closed. +// §10.20). On either failure the stream is closed. One GOAWAY ahead of the +// response is checked, the response still awaited, and the GOAWAY put back as +// the stream's first follow-up (§10.4), bar where the response must come +// first. func awaitRequestResponse[OK message.Message, R any]( ctx context.Context, s *Session, @@ -729,6 +733,22 @@ func awaitRequestResponse[OK message.Message, R any]( return zero, err } resp, err := s.readResponse(ctx, stream) + // §10.4: "A GOAWAY MAY also be sent on a request stream to initiate + // migration of that individual request", before its response too, save + // where the response MUST be "the first message" (§6.1, §6.2). It is + // checked as any on the stream is, and left for the stream's next reader, + // as one after the response is; the request is still answered. + if g, isGoaway := resp.(*message.Goaway); isGoaway && err == nil && !responseFirst(m.Type()) { + var goaways RequestGoaways + if gerr := goaways.Received(s, g); gerr != nil { + return zero, gerr + } + resp, err = s.readResponse(ctx, stream) + if g2, again := resp.(*message.Goaway); again && err == nil { + return zero, goaways.Received(s, g2) // a second GOAWAY (§10.4) + } + stream = replayMessage(stream, g) + } if err != nil { _ = stream.Close() return zero, fmt.Errorf("moqt/session: read %s response: %w", m.Type(), err) @@ -763,6 +783,31 @@ func awaitRequestResponse[OK message.Message, R any]( return zero, err } +// responseFirst reports whether the response to a request of type t MUST be +// the first message on its stream: to SUBSCRIBE_NAMESPACE and SUBSCRIBE_TRACKS +// (§6.1), and to PUBLISH_NAMESPACE (§6.2). +func responseFirst(t message.Type) bool { + return t == message.TypePublishNamespace || t == message.TypeSubscribeNamespace || + t == message.TypeSubscribeTracks +} + +// replayStream is a request stream whose reads begin with a message already +// read off it. A handle's Stream is one after an early GOAWAY. +type replayStream struct { + Stream + + r io.Reader +} + +func (s *replayStream) Read(p []byte) (int, error) { return s.r.Read(p) } + +// replayMessage returns stream with m put back ahead of what it has left. +func replayMessage(stream Stream, m message.Message) Stream { + var buf bytes.Buffer + _ = message.Marshal(&buf, m) // a message just parsed re-encodes + return &replayStream{Stream: stream, r: io.MultiReader(&buf, stream)} +} + // UpdateRequest sends a REQUEST_UPDATE (§10.9) with a fresh Request ID (§10.1) // on an established request stream and awaits its REQUEST_OK or REQUEST_ERROR // (*RequestRejectedError). params carries only the fields to change. The diff --git a/pkg/moqt/session/second_response_test.go b/pkg/moqt/session/second_response_test.go index c040008e..fa687c11 100644 --- a/pkg/moqt/session/second_response_test.go +++ b/pkg/moqt/session/second_response_test.go @@ -8,6 +8,7 @@ import ( "github.com/floatdrop/moq-go/pkg/moqt" "github.com/floatdrop/moq-go/pkg/moqt/message" "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/moqt/wire" ) // TestSecondResponseCloses: a second response to a SUBSCRIBE or PUBLISH read @@ -51,29 +52,46 @@ func TestSecondResponseCloses(t *testing.T) { } } -// TestResponderStreamResponseKeepsSession: on a stream this side answered, -// the peer is the requester, so a REQUEST_OK or REQUEST_ERROR from it is no -// second response (§5.1) and reaches Serve's callback. -func TestResponderStreamResponseKeepsSession(t *testing.T) { - cli, srv := openPair(t) - sub, pub := subscribePair(t, cli, srv) - go func() { _ = message.Marshal(sub, &message.RequestError{ErrorCode: moqt.RequestInternalError}) }() - got := make(chan message.Message, 1) - go func() { - _ = pub.Broker().Serve(t.Context(), func(m message.Message) bool { - got <- m - return false +// TestResponderStreamStrayResponseCloses: on a stream this side answered, a +// REQUEST_OK or REQUEST_ERROR from the requester answers nothing when this +// side sent no REQUEST_UPDATE (§10.9), and closes the session with +// PROTOCOL_VIOLATION: on an accepted SUBSCRIBE, where only the requester may +// send one, and on an accepted PUBLISH that this side never updated. +func TestResponderStreamStrayResponseCloses(t *testing.T) { + for _, tc := range []struct { + name string + resp message.Message + // viaPublish: the client accepts the server's PUBLISH; else the + // server accepts the client's SUBSCRIBE. + viaPublish bool + }{ + {"REQUEST_OK on SUBSCRIBE", &message.RequestOK{}, false}, + {"REQUEST_ERROR on SUBSCRIBE", &message.RequestError{ErrorCode: moqt.RequestInternalError}, false}, + {"REQUEST_OK on PUBLISH", &message.RequestOK{}, true}, + {"REQUEST_ERROR on PUBLISH", &message.RequestError{ErrorCode: moqt.RequestInternalError}, true}, + } { + t.Run(tc.name, func(t *testing.T) { + cli, srv := openPair(t) + responder := srv + var ( + from session.Stream + broker *session.RequestBroker + ) + if tc.viaPublish { + responder = cli + h, pub := establishOnAlias(t, cli, srv, true, "t", 7) + from, broker = pub, h.(*session.IncomingPublication).Broker() + } else { + sub, pub := subscribePair(t, cli, srv) + from, broker = sub, pub.Broker() + } + go func() { _ = message.Marshal(from, tc.resp) }() + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second) + defer cancel() + _ = broker.Serve(ctx, nil) + requireClosedCode(t, responder, moqt.SessionProtocolViolation) }) - }() - select { - case m := <-got: - if _, ok := m.(*message.RequestError); !ok { - t.Fatalf("Serve's callback got %T, want *message.RequestError", m) - } - case <-time.After(2 * time.Second): - t.Fatal("the REQUEST_ERROR never reached Serve's callback") } - requireStaysOpen(t, srv, 50*time.Millisecond) } // TestSubscribeOKOnPublishKeepsSession: a SUBSCRIBE_OK is no response to a @@ -123,3 +141,54 @@ func TestLateUpdateAnswerKeepsSession(t *testing.T) { <-answered requireStaysOpen(t, cli, 100*time.Millisecond) } + +// TestRequesterStrayResponseCloses: on this side's FETCH or +// SUBSCRIBE_NAMESPACE, a REQUEST_OK before any REQUEST_UPDATE is a second +// response to the request (§5.2: "exactly one FETCH_OK or REQUEST_ERROR") and +// closes the session with PROTOCOL_VIOLATION. +func TestRequesterStrayResponseCloses(t *testing.T) { + for _, tc := range []struct { + name string + open func(t *testing.T, cli, srv *session.Session) *session.RequestBroker + }{ + {"FETCH", func(t *testing.T, cli, srv *session.Session) *session.RequestBroker { + go func() { + r, err := srv.AcceptRequest(t.Context()) + if err != nil { + return + } + if _, err := r.AcceptFetch(&message.FetchOK{}); err == nil { + _ = message.Marshal(r.Stream, &message.RequestOK{}) + } + }() + fr, err := cli.Fetch(t.Context(), &message.Fetch{Name: []byte("t")}) + must(t, err) + return fr.Broker() + }}, + {"SUBSCRIBE_NAMESPACE", func(t *testing.T, cli, srv *session.Session) *session.RequestBroker { + go func() { + r, err := srv.AcceptRequest(t.Context()) + if err != nil { + return + } + if _, err := r.AcceptSubscribeNamespace(); err == nil { + _ = message.Marshal(r.Stream, &message.RequestOK{}) + } + }() + ns, err := cli.SubscribeNamespace(t.Context(), &message.SubscribeNamespace{ + TrackNamespacePrefix: wire.TrackNamespace{[]byte("ns")}, + }) + must(t, err) + return ns.Broker() + }}, + } { + t.Run(tc.name, func(t *testing.T) { + cli, srv := openPair(t) + broker := tc.open(t, cli, srv) + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second) + defer cancel() + _ = broker.Serve(ctx, nil) + requireClosedCode(t, cli, moqt.SessionProtocolViolation) + }) + } +} diff --git a/pkg/relay/handler_forward.go b/pkg/relay/handler_forward.go index b95700c1..70f84ff2 100644 --- a/pkg/relay/handler_forward.go +++ b/pkg/relay/handler_forward.go @@ -153,7 +153,7 @@ func (h *sessionHandler) serveForwardedPublish( if p, ok := params.Find(message.ParamNewGroupRequest); ok { h.propagateNewGroupUpstream(ctx, fullName, p.Varint) } - h.readSubscribeUpdates(ctx, stream, sub, fullName, true) + h.readSubscribeUpdates(ctx, stream, sub, fullName) } // inflightSubscribe is a track's SUBSCRIBEs in flight on one session. diff --git a/pkg/relay/handler_subscribe.go b/pkg/relay/handler_subscribe.go index 902ba4c3..644f2a69 100644 --- a/pkg/relay/handler_subscribe.go +++ b/pkg/relay/handler_subscribe.go @@ -164,7 +164,7 @@ func (h *sessionHandler) handleSubscribe(ctx context.Context, req *session.Reque h.propagateForwardUpstream(ctx, fullName) } - h.readSubscribeUpdates(ctx, req.Stream, sub, fullName, false) + h.readSubscribeUpdates(ctx, req.Stream, sub, fullName) h.log.LogAttrs(ctx, slog.LevelDebug, "SUBSCRIBE stream ended", slog.String("name", string(msg.Name))) } @@ -173,31 +173,19 @@ func (h *sessionHandler) handleSubscribe(ctx context.Context, req *session.Reque // SUBSCRIBE's stream to [sessionHandler.handleSubscribeUpdate] until the // subscriber cancels, the stream turns undecodable (see [readRequestStream]) // or ctx ends. A subscriber FIN is not a cancellation (§3.3.2): the -// subscription lives on in [awaitRequestEnd]. forwarded marks a PUBLISH the -// relay forwarded, which the subscriber answered. +// subscription lives on in [awaitRequestEnd]. The stream is a SUBSCRIBE's or +// a forwarded PUBLISH's. func (h *sessionHandler) readSubscribeUpdates( ctx context.Context, stream session.Stream, sub *registry.DownstreamSub, fullName track.FullTrackName, - forwarded bool, ) { updates := h.sess.NewRequestUpdateLimiter() fin := readRequestStream(ctx, h.sess, stream, func(m message.Message) bool { if h.isPeerStateNotify(m) { return false } - switch m.(type) { - case *message.RequestOK, *message.RequestError: - // On a forwarded PUBLISH the subscriber answered already, and - // the relay sends no REQUEST_UPDATE here: a second PUBLISH_OK - // or REQUEST_ERROR (§5.1). - if forwarded { - _ = h.sess.Close(moqt.SessionProtocolViolation, - fmt.Sprintf("%s after the PUBLISH_OK", m.Type())) - return false - } - } if upd, ok := m.(*message.RequestUpdate); ok { // §10.2.1: parameters outside the scope of a subscriber's // update are session-fatal. diff --git a/pkg/relay/second_response_test.go b/pkg/relay/second_response_test.go index da2914b0..cc8e64ea 100644 --- a/pkg/relay/second_response_test.go +++ b/pkg/relay/second_response_test.go @@ -6,6 +6,7 @@ import ( "time" "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" "github.com/floatdrop/moq-go/pkg/relay" ) @@ -40,3 +41,64 @@ func TestRelay_SecondResponseCloses(t *testing.T) { }) } } + +// TestRelay_StrayResponseCloses: on a request stream the relay answered, the +// requester has nothing to answer, since only it sends REQUEST_UPDATE there +// (§10.9), so a REQUEST_OK or REQUEST_ERROR from it is a protocol violation +// and the relay closes the session. On a PUBLISH, the relay's subscriber side +// may send one, but has not. +func TestRelay_StrayResponseCloses(t *testing.T) { + t.Parallel() + for _, rc := range []struct { + request string + open func(t *testing.T) (*session.Session, session.Stream) + }{ + {"SUBSCRIBE", func(t *testing.T) (*session.Session, session.Stream) { + pubSess, _ := newCam1Publisher(t, nil) + subSess := dialAnotherClient(t, pubSess) + return subSess, subscribeCam1(t, subSess).Stream + }}, + {"FETCH", func(t *testing.T) (*session.Session, session.Stream) { + pubSess, subSess, alias := publishAndCache(t) + sendObjects(pubSess, alias, 0, 1) + waitRelayLargest(t, subSess, ns("video"), []byte("cam1"), 0, 0) + fr, err := subSess.Fetch(t.Context(), &message.Fetch{Namespace: ns("video"), Name: []byte("cam1")}) + if err != nil { + t.Fatalf("Fetch: %v", err) + } + return subSess, fr.Stream + }}, + {"PUBLISH_NAMESPACE", func(t *testing.T) (*session.Session, session.Stream) { + sess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + return sess, publishNS(t, sess, "video").Stream + }}, + {"SUBSCRIBE_NAMESPACE", func(t *testing.T) (*session.Session, session.Stream) { + sess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + s, _ := subscribeNS(t, sess, "video") + return sess, s.Stream + }}, + {"PUBLISH", func(t *testing.T) (*session.Session, session.Stream) { + sess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + return sess, publishVideoTrack(t, sess, "cam1", 7).Stream + }}, + {"SUBSCRIBE_TRACKS", func(t *testing.T) (*session.Session, session.Stream) { + sess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + return sess, subscribeTracks(t, sess, ns("video")) + }}, + } { + for _, resp := range []message.Message{&message.RequestOK{}, &message.RequestError{}} { + t.Run(resp.Type().String()+" on "+rc.request, func(t *testing.T) { + t.Parallel() + sess, stream := rc.open(t) + if err := message.Marshal(stream, resp); err != nil { + t.Fatalf("write %s: %v", resp.Type(), err) + } + requireSessionClosed(t, sess, "a stray "+resp.Type().String()) + }) + } + } +} diff --git a/pkg/relay/session_handler.go b/pkg/relay/session_handler.go index 8978c714..406ccb82 100644 --- a/pkg/relay/session_handler.go +++ b/pkg/relay/session_handler.go @@ -572,7 +572,8 @@ func ctxResetCode(ctx context.Context) moqt.StreamResetCode { // StreamResetSessionClosed to unblock the parse). A follow-up that cannot be // read — any non-EOF error — resets the read side with // StreamResetInternalError so the peer learns reads stopped; a malformed one -// also closes the session (§10). +// also closes the session (§10), as does a REQUEST_OK or REQUEST_ERROR, since +// the relay sends no REQUEST_UPDATE on the streams it reads this way (§10.9). // // It reports fin when the requester ended its side with a FIN, which is not a // cancellation (§3.3.2); see [awaitRequestEnd]. @@ -610,6 +611,16 @@ func readRequestStream( done <- false return } + // §10.9: the relay sends no REQUEST_UPDATE on these streams, so + // a REQUEST_OK or REQUEST_ERROR answers nothing: on a forwarded + // PUBLISH it is a second response (§5.1). + switch m.(type) { + case *message.RequestOK, *message.RequestError: + _ = sess.Close(moqt.SessionProtocolViolation, + fmt.Sprintf("%s with no REQUEST_UPDATE to answer", m.Type())) + done <- false + return + } if !onMsg(m) { done <- false return