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
13 changes: 1 addition & 12 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ By package, bottom-up along the dependency stack:
| 10.9 | REQUEST_UPDATE | 0x02 | DONE | A REQUEST_UPDATE opening a request stream closes the session with PROTOCOL_VIOLATION (`ErrUnexpectedRequestUpdate`). |
| 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. |
| 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. |
| 10.13 | FETCH | 0x16 | DONE | Standalone, the only kind in draft-20. From the cache, a Location is non-existent only on a signal: a Prior Group or Object ID Gap, a Group's or the Track's end, or an upstream's FETCH. Other uncached Locations are FETCHed from a fetch-capable upstream in one span, within FILL_TIMEOUT, or else marked End of Unknown (or Timed-Out) Range. |
| 10.14 | FETCH_OK | 0x18 | DONE | An End Location before the FETCH's Start closes the session. A Start relative to the Largest Object is compared through End ≤ Largest; an End of {0,0} is let through, as it cannot be told apart from "no content yet". |
| 10.15 | TRACK_STATUS | 0x0D | DONE | Reply via REQUEST_OK, then FIN; any follow-up from the requester closes the session. |
Expand Down Expand Up @@ -527,17 +527,6 @@ Session layer:
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.
- `Publication`'s automatic PUBLISH_DONE UPDATE_FAILED is sent while its subgroup
streams are open, `WriteObject` still succeeds after `Done`, and a subgroup
opened concurrently with `Done` is missing from the Stream Count (§10.12).
- `Publication`'s REQUEST_UPDATE_OK carries LARGEST_OBJECT only for Objects it
wrote itself, not the one its SUBSCRIBE_OK or PUBLISH reported (§10.2.17,
§10.9.1).
- A rejected request sends STOP_SENDING with INTERNAL_ERROR (§3.3.4 SHOULD use a
relevant code).
- Mandatory Track Property enforcement is off unless configured (§2.5.1).
- SETUP options are sorted unstably, so with more than 12 the Token order on the
wire can differ from the order `heldSetupAliases` replays (§10.3.1.4).

Relay:

Expand Down
58 changes: 46 additions & 12 deletions pkg/moqt/session/datastream_out.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,9 +86,13 @@ type OutgoingSubgroupStream struct {
encHavePrev bool

// Set by [Publication.OpenSubgroup]: onObject is told each written
// object's Location, and paused reports a Forward State of 0 (§11.4.3).
// object's Location, paused reports a Forward State of 0 (§11.4.3), ended
// reports that the publication ended (Done), and onEnd is told when the
// stream is FINished or reset.
onObject func(group, object uint64)
paused func() bool
ended func() bool
onEnd func()
}

// WithDeliveryTimeouts returns a shallow copy of s configured with the §8
Expand Down Expand Up @@ -168,9 +172,13 @@ func (s *OutgoingSubgroupStream) WriteObjectReceivedAt(
}
s.sawFirstObject = true

// §10.12: Done reset the stream before PUBLISH_DONE.
if s.ended != nil && s.ended() {
return ErrPublicationEnded
}
// §5.1: no Objects while the Forward State is 0; §11.4.3: reset.
if s.paused != nil && s.paused() {
s.dst.CancelWrite(uint64(moqt.StreamResetCancelled))
s.Cancel(moqt.StreamResetCancelled)
return ErrForwardPaused
}
if err := s.checkObjectTimeout(receivedAt); err != nil {
Expand Down Expand Up @@ -247,7 +255,7 @@ func (s *OutgoingSubgroupStream) checkObjectTimeout(receivedAt time.Time) error
}
elapsed := time.Since(receivedAt)
if elapsed > s.objectTimeout {
s.dst.CancelWrite(uint64(moqt.StreamResetDeliveryTimeout))
s.Cancel(moqt.StreamResetDeliveryTimeout)
return fmt.Errorf("%w (elapsed %s, limit %s)",
ErrDeliveryTimeout, elapsed, s.objectTimeout)
}
Expand All @@ -262,6 +270,9 @@ func (s *OutgoingSubgroupStream) checkObjectTimeout(receivedAt time.Time) error
// stream if the peer has not acknowledged all data within the timeout (§8).
func (s *OutgoingSubgroupStream) Close() error {
err := s.dst.Close()
if s.onEnd != nil {
s.onEnd()
}
tracked, ok := s.dst.(DeliveryTrackingSendStream)
if s.subgroupTimeout > 0 && ok {
finished := tracked.Finished()
Expand All @@ -285,6 +296,9 @@ func (s *OutgoingSubgroupStream) Close() error {
// Cancel resets the stream with the given application code (§3.3.4).
func (s *OutgoingSubgroupStream) Cancel(code moqt.StreamResetCode) {
s.dst.CancelWrite(uint64(code))
if s.onEnd != nil {
s.onEnd()
}
}

// SetSendPriority forwards the composite §7.2 scheduling key to the underlying
Expand Down Expand Up @@ -367,25 +381,45 @@ func (s *Session) OpenSubgroupContext(
ctx context.Context,
h message.SubgroupHeader,
) (*OutgoingSubgroupStream, error) {
dst, err := s.conn.OpenUniStream()
sg, reset, err := s.openSubgroup(ctx, h, false)
if err != nil {
return nil, err
}
if reset {
// ctx was cancelled just after the header write went through, and
// the stream is reset.
return nil, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", ctx.Err())
}
return sg, nil
}

// openSubgroup opens a subgroup stream and writes its header, which
// cancelling ctx interrupts by resetting the stream, first marking what was
// written reliable when markReliable is set (§11.4.3). It returns the stream
// whenever the whole header was written, since the peer can then attribute
// the stream to its track, with reset reporting that ctx reset it just after.
func (s *Session) openSubgroup(
ctx context.Context,
h message.SubgroupHeader,
markReliable bool,
) (sg *OutgoingSubgroupStream, reset bool, err error) {
dst, err := s.conn.OpenUniStream()
if err != nil {
return nil, false, err
}
stop := context.AfterFunc(ctx, func() {
if r, ok := dst.(ReliableResetStream); ok && markReliable {
r.SetReliableBoundary()
}
dst.CancelWrite(uint64(moqt.StreamResetCancelled))
})
if err := message.WriteSubgroupHeader(dst, h); err != nil {
stop()
dst.CancelWrite(uint64(moqt.StreamResetInternalError))
if ctx.Err() != nil {
return nil, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", ctx.Err())
return nil, false, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", ctx.Err())
}
return nil, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", err)
}
if !stop() {
// The AfterFunc already ran: ctx was cancelled while (or just
// after) the header write went through — the stream is reset.
return nil, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", ctx.Err())
return nil, false, fmt.Errorf("moqt/session: write SUBGROUP_HEADER: %w", err)
}
return &OutgoingSubgroupStream{header: h, dst: dst}, nil
return &OutgoingSubgroupStream{header: h, dst: dst}, !stop(), nil
}
8 changes: 8 additions & 0 deletions pkg/moqt/session/export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,3 +37,11 @@ func OpenRequestForTest(s *Session, first message.Message) (Stream, error) {
func WithSetupOptionForTest(kv wire.KVPair) Option {
return func(c *config) { c.setupOptions = append(c.setupOptions, kv) }
}

// OpenSubgroupsForTest reports how many subgroups p tracks as open, which Done
// would reset.
func (p *Publication) OpenSubgroupsForTest() int {
p.subMu.Lock()
defer p.subMu.Unlock()
return len(p.open)
}
11 changes: 4 additions & 7 deletions pkg/moqt/session/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -166,13 +166,10 @@ func WithGrease() Option {
// set, it returns *ErrUnsupportedMandatoryTrackProperty; [Request.AcceptPublish]
// refuses such a PUBLISH with UNSUPPORTED_EXTENSION.
//
// If this option is never called, enforcement is disabled and all properties
// pass through. Leave it unset only when the application checks the
// properties itself (§2.5.1).
//
// End subscribers that interpret track data should call this option to opt
// in to enforcement. Pass an empty (non-nil) map to reject all mandatory
// properties, or populate the map with the types you support.
// If this option is never called, or types is empty or nil, no Mandatory Track
// Property is known and every one is refused: an endpoint that does not
// understand one "MUST NOT process or forward that track" (§2.5.1). List the
// types this endpoint understands to accept them.
func WithKnownMandatoryTrackProperties(types map[message.PropertyType]struct{}) Option {
return func(c *config) {
c.knownMandatoryTrackProperties = types
Expand Down
Loading
Loading