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
18 changes: 3 additions & 15 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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). |
Expand Down Expand Up @@ -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. |
Expand Down Expand Up @@ -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
Expand Down
43 changes: 23 additions & 20 deletions pkg/moqt/session/broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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}
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
164 changes: 164 additions & 0 deletions pkg/moqt/session/early_goaway_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
Loading
Loading