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
4 changes: 2 additions & 2 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,9 +44,9 @@ The module builds with `CGO_ENABLED=0`.
| Path | What it is |
| --- | --- |
| `.` (`rtc`) | `Client`, `Call`, publisher, subscriber, media engine, stats, event fan-out. `call.go` is the most important file. |
| `coordinator/` | Coordinator REST + event websocket. Exactly one endpoint is called: `POST /api/v2/video/call/{type}/{id}/join`. |
| `coordinator/` | Coordinator REST + event websocket. Joins call `POST /api/v2/video/call/{type}/{id}/fast_join` (candidate SFUs, `fastjoin.go`), or `.../join` on the legacy flow, for reconnects and migrations, and where `fast_join` is missing. |
| `coordinator/models/` | Generated from the public OpenAPI spec. Do not hand-edit. |
| `signal/` | The SFU signalling client: the protobuf websocket plus the twirp `SignalServer` RPCs. |
| `signal/` | The SFU signalling client: the protobuf websocket plus the twirp `SignalServer` and `FastJoinServer` RPCs. |
| `pc/` | `pc.Transport`, the peer-connection wrapper. |
| `track/` | Local tracks driven by a `SampleProvider`, simulcast layers, RED transcoding for audio. |
| `audio/` | `audio.PCM`, resampling, chunking, WAV, G.711. |
Expand Down
9 changes: 9 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -104,8 +104,17 @@ if _, err := call.AddTrack(track.TrackInfo(), track); err != nil {
return writer.Write(audio.FromInt16(samples, 24000, 1))
```

To publish from the start, pass the track to `Join` instead: `rtc.WithTrack(track.TrackInfo(), track)`.
Its offer is then built while the coordinator request is in flight and answered by the SFU's
join itself, where `AddTrack` after `Join` costs a renegotiation.

## Join latency

`Join` takes the fast flow: one coordinator request returns candidate SFUs, and one request
to the SFU joins, answers the publisher offer and returns the subscriber offer. Where the
deployment has no fast join, it joins the legacy way; `call.JoinFlow()` says which ran, and
`rtc.WithJoinFlow(rtc.JoinFlowLegacy)` asks for the legacy one.

A call records how its first join went: every network step, timer and local step, what it
waited for, and each step in round trips to its peer. `String()` draws it as a DAG with
the critical path marked; it marshals to JSON.
Expand Down
188 changes: 135 additions & 53 deletions call.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,10 @@ type joinOptions struct {
location string
migratingFrom string

flow JoinFlow
// tracks are published from the start: see WithTrack.
tracks []trackWithInfo

// used by integration tests to induce specific behaviour
beforeSubscriberSendAnswer func(*signal_rpc.SendAnswerRequest) error
}
Expand All @@ -76,6 +80,7 @@ func defaultJoinOptions() joinOptions {
},
subscriber: SubscriberFunc(func(OnTrackReceived) {}),
create: true,
flow: JoinFlowFast,
}
}

Expand Down Expand Up @@ -283,6 +288,10 @@ type Call struct {
coordinatorState atomic.Pointer[CallState]
// trace records the first join's steps; see JoinTrace.
trace joinTracer
// joinFlow is the JoinFlow the first join took.
joinFlow atomic.Value
// attached is closed once a fast join's websocket attach has finished, well or not.
attached atomic.Pointer[chan struct{}]

GetCred GetCredentialsFunc
cred atomic.Pointer[models.Credentials]
Expand Down Expand Up @@ -754,6 +763,23 @@ func (c *Call) Client() *signal.Client {
return c.peer.Load().client.Load()
}

// credentials are the SFU credentials in use, or none yet: a fast join builds its peer
// connections before the coordinator has said which SFU to use.
func (c *Call) credentials() models.Credentials {
if cred := c.cred.Load(); cred != nil {
return *cred
}
return models.Credentials{}
}

func iceServers(servers []models.ICEServerResponse) []webrtc.ICEServer {
var out []webrtc.ICEServer
for _, s := range servers {
out = append(out, webrtc.ICEServer{URLs: s.Urls, Username: s.Username, Credential: s.Password})
}
return out
}

// unifiedSessionID returns an ID that stays the same for the lifetime of this
// Call, across every reconnect and migration, unlike SessionID which a rejoin
// rotates. It is what lets server-side stats be stitched back together into one
Expand Down Expand Up @@ -805,11 +831,13 @@ func humanizeSdkType(sdkType sfu_models.SdkType) string {
return sdkType.String()
}

// Join performs the coordinator join, opens the SFU websocket, sends the SFU
// JoinRequest and brings up the publisher and subscriber peer connections.
// Join joins the call: the coordinator request, the SFU join and both peer
// connections, returning once the SFU has accepted the client. Media starts flowing
// shortly after, as ICE and DTLS complete.
//
// The coordinator half runs only on the first call; the reconnect and migration
// paths re-enter Join to rebuild the SFU side against fresh credentials.
// By default it takes the fast join (JoinFlowFast): see WithJoinFlow. The reconnect
// and migration paths re-enter Join to rebuild the SFU side against fresh
// credentials, always through the SFU websocket's JoinRequest.
func (c *Call) Join(ctx context.Context, opts ...JoinOption) (*sfu_events.JoinResponse, error) {
options := defaultJoinOptions()
for _, o := range opts {
Expand All @@ -821,25 +849,106 @@ func (c *Call) Join(ctx context.Context, opts ...JoinOption) (*sfu_events.JoinRe
c.cc.claimConnectTrace(rec)
}

if err := c.joinCoordinator(ctx, options); err != nil {
return nil, err
if !reconnecting && c.GetCred == nil && options.flow == JoinFlowFast {
resp, err := c.fastJoin(ctx, opts, options, rec)
if !errors.Is(err, errFastJoinUnavailable) {
return resp, err
}
c.logger.WithField("err", err).Warn("fast join unavailable, joining the legacy way")
}

resp, err := c.legacyJoin(ctx, opts, options, rec, reconnecting)
if err != nil || reconnecting {
return resp, err
}
c.joinFlow.Store(JoinFlowLegacy)
for _, t := range options.tracks {
if _, err := c.AddTrack(t.info, t.tracks[0]); err != nil {
return nil, xerr.Wrap(err)
}
}
return resp, nil
}

// setSessionID picks the session ID of a join: the one asked for, or else the current
// one, or else a new one.
func (c *Call) setSessionID(options joinOptions) {
if options.sessionID != "" {
// always overwrite if sessionID is provided
c.SessionID.Store(options.sessionID)
} else if c.SessionID.Load() == "" {
c.SessionID.Store(uuid.New().String())
}
}

// rememberJoinOptions keeps the first join's options, which the reconnects reuse.
func (c *Call) rememberJoinOptions(opts []JoinOption) {
if c.joinOptions == nil {
if opts == nil {
// nil vs empty slice, todo better approach
opts = []JoinOption{}
}
c.joinOptions = opts
}
}

// joinRequest is the SFU websocket's JoinRequest for the current session and
// credentials.
func (c *Call) joinRequest(options joinOptions, publisherSDP, subscriberSDP string) *sfu_events.JoinRequest {
return &sfu_events.JoinRequest{
Token: c.credentials().Token,
SessionId: c.SessionID.Load(),
PublisherSdp: publisherSDP,
ClientDetails: c.sfuClientDetails(),
ReconnectDetails: options.reconnectDetails,
Source: c.cc.source.toSfuParticipantSource(),
SubscriberSdp: subscriberSDP,
// SessionId rotates on every rejoin, so without this the server cannot
// tell that the sessions before and after an outage were the same call
// from the same client, and its stats are split across them.
UnifiedSessionId: c.unifiedSessionID(),
Capabilities: options.clientCapabilities(),
PreferredPublishOptions: options.preferredPublishOptions,
}
}

func (c *Call) sfuClientDetails() *sfu_models.ClientDetails {
return &sfu_models.ClientDetails{
Sdk: &sfu_models.Sdk{
Type: c.cc.clientDetails.sdkType(),
Major: c.cc.clientDetails.SDKVersion.Major,
Minor: c.cc.clientDetails.SDKVersion.Minor,
Patch: c.cc.clientDetails.SDKVersion.Patch,
},
Os: &sfu_models.OS{
Name: c.cc.clientDetails.OSName,
},
Browser: &sfu_models.Browser{
Name: c.cc.clientDetails.BrowserName,
},
}
}

// startTracing starts the call's stats trace buffers, when the call reports stats.
func (c *Call) startTracing() {
if !c.statsEnabled() {
return
}
attempt := int64(c.reconnectAttempt.Load()) - 1
c.Tracing.Store(rtcstats.NewCallTraceBuffer("", attempt, c.credentials().Server.EdgeName))
c.Client().Tracing.Store(c.Tracing.Load())
}

// legacyJoin is the join through the coordinator's join and the SFU websocket's
// JoinRequest, which every reconnect takes.
func (c *Call) legacyJoin(
ctx context.Context, opts []JoinOption, options joinOptions, rec *jointrace.Recorder, reconnecting bool,
) (*sfu_events.JoinResponse, error) {
if err := c.joinCoordinator(ctx, options); err != nil {
return nil, err
}
c.setSessionID(options)
c.rememberJoinOptions(opts)
c.externalRTCP = options.externalRTCP

publisherSDP := ""
Expand All @@ -856,44 +965,8 @@ func (c *Call) Join(ctx context.Context, opts ...JoinOption) (*sfu_events.JoinRe
}
}

req := &sfu_events.JoinRequest{
Token: c.cred.Load().Token,
SessionId: c.SessionID.Load(),
PublisherSdp: publisherSDP,
ClientDetails: &sfu_models.ClientDetails{
Sdk: &sfu_models.Sdk{
Type: c.cc.clientDetails.sdkType(),
Major: c.cc.clientDetails.SDKVersion.Major,
Minor: c.cc.clientDetails.SDKVersion.Minor,
Patch: c.cc.clientDetails.SDKVersion.Patch,
},
Os: &sfu_models.OS{
Name: c.cc.clientDetails.OSName,
},
Browser: &sfu_models.Browser{
Name: c.cc.clientDetails.BrowserName,
},
},
ReconnectDetails: options.reconnectDetails,
Source: c.cc.source.toSfuParticipantSource(),
SubscriberSdp: subscriberSDP,
// SessionId rotates on every rejoin, so without this the server cannot
// tell that the sessions before and after an outage were the same call
// from the same client, and its stats are split across them.
UnifiedSessionId: c.unifiedSessionID(),
Capabilities: options.clientCapabilities(),
PreferredPublishOptions: options.preferredPublishOptions,
}

attempt := int64(c.reconnectAttempt.Load()) - 1
sfuid := c.cred.Load().Server.EdgeName

if c.statsEnabled() {
// Update Tracer for active call
c.Tracing.Store(rtcstats.NewCallTraceBuffer("", attempt, sfuid))
// Update Tracer for signal client
c.Client().Tracing.Store(c.Tracing.Load())
}
req := c.joinRequest(options, publisherSDP, subscriberSDP)
c.startTracing()

pcsStart := time.Now()
if err := c.initPubAndSub(options); err != nil {
Expand Down Expand Up @@ -1146,6 +1219,13 @@ func (c *Call) restoreICE(ctx context.Context) error {
}

func (c *Call) Leave(reason string) error {
// A fast join returns before its websocket has attached; the leave goes on it.
if attached := c.attached.Load(); attached != nil {
select {
case <-*attached:
case <-time.After(fastAttachTimeout):
}
}
// Report rtcstats for the last time before leaving
if err := c.reportRtcStats(c.callCtx); err != nil {
c.logger.WithField("err", err).Error("failed to report rtcstats to the SFU")
Expand Down Expand Up @@ -1252,13 +1332,20 @@ func (c *Call) OnIceTrickle(trickle *sfu_events.SfuEvent_IceTrickle) {
transport.AddICECandidate(candidate)
}

func (c *Call) AddSimulcastTracks(trackInfo *sfu_models.TrackInfo, tracks ...webrtc.TrackLocal) (*webrtc.RTPTransceiver, error) {
// recordPublishedTrack adds a track to what the call publishes, which is what a
// reconnect restores, and returns its info with the mid it will be sent on.
func (c *Call) recordPublishedTrack(trackInfo *sfu_models.TrackInfo, tracks ...webrtc.TrackLocal) *sfu_models.TrackInfo {
trackInfo = proto.Clone(trackInfo).(*sfu_models.TrackInfo)
c.publishedTracksMu.Lock()
defer c.publishedTracksMu.Unlock()
// todo, this logic needs to change once we support remove track
trackInfo.Mid = strconv.Itoa(len(c.publishedTracks))
c.publishedTracks = append(c.publishedTracks, trackWithInfo{tracks: tracks, info: trackInfo})
c.publishedTracksMu.Unlock()
return trackInfo
}

func (c *Call) AddSimulcastTracks(trackInfo *sfu_models.TrackInfo, tracks ...webrtc.TrackLocal) (*webrtc.RTPTransceiver, error) {
trackInfo = c.recordPublishedTrack(trackInfo, tracks...)
peerPub := c.publisherPeer()
if peerPub == nil {
return nil, xerr.Wrap(fmt.Errorf("add simulcast tracks: %w", errNoPeerConnection))
Expand All @@ -1272,12 +1359,7 @@ func (c *Call) AddSimulcastTracks(trackInfo *sfu_models.TrackInfo, tracks ...web
}

func (c *Call) AddTrack(trackInfo *sfu_models.TrackInfo, track webrtc.TrackLocal) (*webrtc.RTPTransceiver, error) {
trackInfo = proto.Clone(trackInfo).(*sfu_models.TrackInfo)
c.publishedTracksMu.Lock()
// todo, this logic needs to change once we support remove track
trackInfo.Mid = strconv.Itoa(len(c.publishedTracks))
c.publishedTracks = append(c.publishedTracks, trackWithInfo{tracks: []webrtc.TrackLocal{track}, info: trackInfo})
c.publishedTracksMu.Unlock()
trackInfo = c.recordPublishedTrack(trackInfo, track)
peerPub := c.publisherPeer()
if peerPub == nil {
return nil, xerr.Wrap(fmt.Errorf("add track: %w", errNoPeerConnection))
Expand Down Expand Up @@ -1321,7 +1403,7 @@ func (c *Call) sendSubscriptions(ctx context.Context, trackDetails []*signal_rpc
return xerr.Wrap(err)
}
rec.Add(jointrace.Span{
Name: jointrace.SubSubscribe, After: []string{jointrace.SFUJoin},
Name: jointrace.SubSubscribe, After: afterFirst(rec, jointrace.SFUJoin, jointrace.SFUFastJoin),
Start: start, End: time.Now(), Kind: jointrace.KindNet, Peer: jointrace.PeerSFU,
})

Expand Down
Loading
Loading