rtc

package module
v0.0.0-...-f458fbb Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 30, 2026 License: BSD-3-Clause Imports: 45 Imported by: 0

README

Official Go WebRTC client for Stream Video

Go Reference License

Join a Stream Video call from Go as a real participant: consume remote tracks and publish local ones. Built for bots, ingress, and Vision Agents. For creating calls, tokens, and moderation, use getstream-go.

Install

go get github.com/GetStream/getstream-go-webrtc

Then import as rtc "github.com/GetStream/getstream-go-webrtc".

Quickstart

Join a call, subscribe to everyone else, and drain each remote track (an unread track stalls):

client, err := rtc.NewRTCClient("apiKey", "apiSecret",
	rtc.WithUser(rtc.User{ID: "agent", Name: "Agent"}),
	getstream.WithTimeout(10*time.Second),
)
if err != nil {
	log.Fatal(err)
}
defer client.Close()

call := client.Call("default", "my-call")
join, err := call.Join(ctx, rtc.WithOnTrack(rtc.SubscriberFunc(func(t rtc.OnTrackReceived) {
	go func() {
		for {
			if _, _, err := t.Track.ReadRTP(); err != nil {
				return
			}
		}
	}()
})))
if err != nil {
	log.Fatal(err)
}
defer call.Leave("done")

var subs []*signal_rpc.TrackSubscriptionDetails
for _, p := range join.GetCallState().GetParticipants() {
	for _, trackType := range p.GetPublishedTracks() {
		subs = append(subs, &signal_rpc.TrackSubscriptionDetails{
			UserId: p.GetUserId(), SessionId: p.GetSessionId(), TrackType: trackType,
		})
	}
}
_ = call.SubscribeToTracks(ctx, subs...)

On an end-user device, use NewClient with a token your backend minted. Never ship the API secret.

Audio

Read a remote track as 16 kHz PCM, or publish PCM as Opus. Runnable version: examples/audio.

reader, err := audiortc.NewTrackReader(t.Track, audiortc.ReaderConfig{
	Opus: opus.Config{SampleRate: 16000, Channels: 1},
})
if err != nil {
	return
}
defer reader.Close()
for pcm, err := range reader.Frames() {
	if err != nil {
		return
	}
	stt.Send(pcm.Bytes())
}
writer, err := audiortc.NewTrackWriter(audiortc.WriterConfig{})
if err != nil {
	return err
}
track, err := audiortc.NewAudioTrack(&sfu_models.TrackInfo{
	TrackId: uuid.NewString(), TrackType: sfu_models.TrackType_TRACK_TYPE_AUDIO,
}, writer)
if err != nil {
	return err
}
if _, err := call.AddTrack(track.TrackInfo(), track); err != nil {
	return err
}
return writer.Write(audio.FromInt16(samples, 24000, 1))

Examples

export STREAM_API_KEY=... STREAM_API_SECRET=... STREAM_CALL_ID=my-call
go run ./examples/join     # subscribe and print RTP counts
go run ./examples/audio    # publish a tone or WAV, optionally record

See examples/README.md.

What is Stream?

Stream allows developers to rapidly deploy scalable feeds, chat messaging and video with an industry leading 99.999% uptime SLA guarantee. All calls run on Stream's network of edge servers around the world, ensuring optimal latency and reliability. See the Video documentation.

Free for makers

Stream is free for most side and hobby projects. To qualify, your project/company needs to have < 5 team members and < $10k in monthly revenue. Makers get $100 in monthly credit for video for free. For more details, see the Maker Account.

Contributing

See CONTRIBUTING.md.

License

BSD 3-Clause. See LICENSE.

Documentation

Index

Constants

View Source
const (
	// HeaderCloudFrontPop is the CloudFront response header carrying the edge
	// point of presence that served the request. Its first three characters are
	// the airport code Stream uses as a location hint.
	HeaderCloudFrontPop = "X-Amz-Cf-Pop"

	// FallbackLocationName is the hint used when discovery fails.
	FallbackLocationName = "IAD"

	// DefaultLocationHintURL is the endpoint probed to discover the nearest
	// CloudFront POP. Override it with WithLocationHintURL.
	DefaultLocationHintURL = "https://hint.stream-io-video.com/"
)

Variables

View Source
var ErrNoSender = errors.New("sender not available")

Functions

func AwaitCallCoordinatorEvent

func AwaitCallCoordinatorEvent[T coordinator.CallEvents](call *Call, match func(T) bool) *event.EventAwaiter[T, models.WebsocketEvent]

AwaitCallCoordinatorEvent returns a coordinator event awaiter that blocks until a coordinator event that matches the specified predicate occurs within a call context.

func AwaitCallSignalEvent

func AwaitCallSignalEvent[T signal.Events](call *Call, match func(T) bool) *event.EventAwaiter[T, *sfu_events.SfuEvent]

AwaitCallSignalEvent returns a signal event awaiter that blocks until a signal event that matches the specified predicate occurs within a call context.

func AwaitCoordinatorEvent

func AwaitCoordinatorEvent[T coordinator.Events](client *Client, match func(T) bool) *event.EventAwaiter[T, models.WebsocketEvent]

AwaitCoordinatorEvent returns an event awaiter that blocks until an event that matches the specified predicate occurs.

func GetCodecPreferencesByMimeType

func GetCodecPreferencesByMimeType(mimeType string) []webrtc.RTPCodecParameters

func HandleCallCoordinatorEvent

func HandleCallCoordinatorEvent[T coordinator.CallEvents](call *Call, onEvent func(T)) func()

HandleCallCoordinatorEvent registers a handler for coordinator events specific to a call. It returns a function to unregister the event handler.

func HandleCallEvent

func HandleCallEvent[T CallEvents](client *Call, onEvent func(T)) func()

HandleCallEvent registers a handler for a call event from either family, routing coordinator events to the coordinator event store and SFU signal events to the signal one. It returns a function to unregister the handler.

The routing switch is generated by internal/cmd/genwsevent from the same sources as the CallEvents constraint, so every type the constraint admits has a case; see events_dispatch.gen.go. A miss can only mean the generated files are stale, which TestCallEventDispatchIsExhaustive catches. At runtime it degrades to a logged error and a no-op handler rather than the panic it used to be.

func HandleCallSignalEvent

func HandleCallSignalEvent[T signal.Events](client *Call, onEvent func(T)) func()

HandleCallSignalEvent registers a handler for signal events related to a call. It returns a function to unregister the event handler.

func HandleCoordinatorEvent

func HandleCoordinatorEvent[T coordinator.Events](client *Client, onEvent func(T)) func()

HandleCoordinatorEvent registers a handler for coordinator events. It returns a function to unregister the event handler.

func SDKBuildVersion

func SDKBuildVersion() string

SDKBuildVersion reports this SDK's module version, as recorded in the binary's build info. It falls back to "0.0.0" when the build has none.

func SafeInt64ToInt32

func SafeInt64ToInt32(i int64) int32

SafeInt64ToInt32 saturates at both ends for the same reason.

func SafeUint64ToUint32

func SafeUint64ToUint32(i uint64) uint32

SafeUint64ToUint32 saturates rather than wrapping, so a counter that has overflowed the narrower stats field reports the maximum instead of a small number that would look like a reset.

func WebRTCBuildVersion

func WebRTCBuildVersion() string

WebRTCBuildVersion reports the version of the pion/webrtc module this binary was built against, so the SFU's telemetry attributes behaviour to the right stack instead of a hardcoded placeholder.

Types

type AudioOutboundRtpProvider

type AudioOutboundRtpProvider interface {
	GetAudioOutboundRtpStats() *webrtc.OutboundRTPStreamStats
	SetAudioOutboundRtpStats(stats *webrtc.OutboundRTPStreamStats)
}

type AudioPlayoutStatsProvider

type AudioPlayoutStatsProvider interface {
	GetAudioPlayoutStats() *webrtc.AudioPlayoutStats
}

type AudioSourceStatsProvider

type AudioSourceStatsProvider interface {
	GetAudioSourceStats() *webrtc.AudioSourceStats
	SetAudioSourceStats(stats *webrtc.AudioSourceStats)
	AudioOutboundRtpProvider
}

type Call

type Call struct {
	Type   string
	Id     string
	UserID string

	GetCred GetCredentialsFunc

	SessionID atomicx.AtomicValue[string]

	// Tracer for API Calls
	Tracing atomic.Pointer[rtcstats.TraceBuffer]
	// contains filtered or unexported fields
}

func (*Call) AddSimulcastTracks

func (c *Call) AddSimulcastTracks(trackInfo *sfu_models.TrackInfo, tracks ...webrtc.TrackLocal) (*webrtc.RTPTransceiver, error)

func (*Call) AddTrack

func (c *Call) AddTrack(trackInfo *sfu_models.TrackInfo, track webrtc.TrackLocal) (*webrtc.RTPTransceiver, error)

func (*Call) CID

func (c *Call) CID() string

func (*Call) Client

func (c *Call) Client() *signal.Client

func (*Call) GetState

func (c *Call) GetState() *CallState

func (*Call) GetStatsReportingMetrics

func (c *Call) GetStatsReportingMetrics() *CallStatsReportingMetrics

func (*Call) Join

func (c *Call) Join(ctx context.Context, opts ...JoinOption) (*sfu_events.JoinResponse, error)

Join performs the coordinator join, opens the SFU websocket, sends the SFU JoinRequest and brings up the publisher and subscriber peer connections.

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.

func (*Call) Leave

func (c *Call) Leave(reason string) error

func (*Call) OnAudioLevelChanged

func (c *Call) OnAudioLevelChanged(changed *sfu_events.SfuEvent_AudioLevelChanged)

func (*Call) OnCallEnded

func (c *Call) OnCallEnded(eventError *sfu_events.SfuEvent_CallEnded)

func (*Call) OnCallGrantsUpdated

func (c *Call) OnCallGrantsUpdated(updated *sfu_events.SfuEvent_CallGrantsUpdated)

func (*Call) OnChangePublishOptions

func (c *Call) OnChangePublishOptions(options *sfu_events.SfuEvent_ChangePublishOptions)

OnChangePublishOptions records the SFU's new encoding configuration and hands it to the application.

Swift reconciles this by cloning tracks onto a transceiver per publish option, so it can publish the same source under two codecs at once. That needs the SDK to own the encoder, which this one does not, so it reports the change and leaves the response to the application. Dual-codec publishing is a follow-up.

func (*Call) OnChangePublishQuality

func (c *Call) OnChangePublishQuality(quality *sfu_events.SfuEvent_ChangePublishQuality)

OnChangePublishQuality applies the SFU's bandwidth decision to the published simulcast layers.

The SFU tracks what its subscribers actually need and how much bandwidth this publisher has, and turns layers off when nobody is watching them or the uplink cannot carry them. Ignoring this event -- which the SDK used to do -- means a bot publishing three simulcast layers keeps pushing all three into a congested uplink, degrading every layer instead of the one nobody wants.

func (*Call) OnConnectionQualityChanged

func (c *Call) OnConnectionQualityChanged(changed *sfu_events.SfuEvent_ConnectionQualityChanged)

func (*Call) OnDominantSpeakerChanged

func (c *Call) OnDominantSpeakerChanged(changed *sfu_events.SfuEvent_DominantSpeakerChanged)

func (*Call) OnError

func (c *Call) OnError(eventError *sfu_events.SfuEvent_Error)

func (*Call) OnGoAway

func (c *Call) OnGoAway(away *sfu_events.SfuEvent_GoAway)

func (*Call) OnHealthCheckResponse

func (c *Call) OnHealthCheckResponse(response *sfu_events.SfuEvent_HealthCheckResponse)

func (*Call) OnIceRestart

func (c *Call) OnIceRestart(restart *sfu_events.SfuEvent_IceRestart)

func (*Call) OnIceTrickle

func (c *Call) OnIceTrickle(trickle *sfu_events.SfuEvent_IceTrickle)

func (*Call) OnInboundStateChanged

func (c *Call) OnInboundStateChanged(handler func([]InboundTrackState))

OnInboundStateChanged registers a callback fired when the SFU pauses or resumes tracks this client subscribes to. Participant state has already been updated when it runs. The handler runs on the signalling read loop, so it must not block.

func (*Call) OnInboundStateNotification

func (c *Call) OnInboundStateNotification(notification *sfu_events.SfuEvent_InboundStateNotification)

OnInboundStateNotification records which subscribed tracks the SFU has paused.

The SFU pauses a track when nobody is rendering it, saving the bandwidth of sending media that would be thrown away. It only does so for clients that advertised CLIENT_CAPABILITY_SUBSCRIBER_VIDEO_PAUSE, which this SDK does, so without handling this a paused track is indistinguishable from a dead one.

func (*Call) OnJoinResponse

func (c *Call) OnJoinResponse(response *sfu_events.SfuEvent_JoinResponse)

func (*Call) OnParticipantJoined

func (c *Call) OnParticipantJoined(joined *sfu_events.SfuEvent_ParticipantJoined)

func (*Call) OnParticipantLeft

func (c *Call) OnParticipantLeft(left *sfu_events.SfuEvent_ParticipantLeft)

func (*Call) OnPinsUpdated

func (c *Call) OnPinsUpdated(updated *sfu_events.SfuEvent_PinsUpdated)

func (*Call) OnPublishOptionsChanged

func (c *Call) OnPublishOptionsChanged(handler func([]*sfu_models.PublishOption))

OnPublishOptionsChanged registers a callback fired whenever the SFU sends new publish options, either in the join response or through a mid-call ChangePublishOptions event. The handler runs on the signalling read loop, so it must not block.

func (*Call) OnPublishQualityChanged

func (c *Call) OnPublishQualityChanged(handler func([]PublishQualityTarget))

OnPublishQualityChanged registers a callback fired when the SFU changes the quality it wants published, typically because its bandwidth estimate moved.

The layers' active flags have already been applied by the time the callback runs; the targets are passed so an application that controls its encoder can match the requested bitrate, frame rate, and resolution. The handler runs on the signalling read loop, so it must not block.

func (*Call) OnPublisherAnswer

func (c *Call) OnPublisherAnswer(answer *sfu_events.SfuEvent_PublisherAnswer)

func (*Call) OnPublisherNegotiationFailed

func (c *Call) OnPublisherNegotiationFailed(handler func(err *pc.NegotiationError))

func (*Call) OnSubscriberNegotiationFailed

func (c *Call) OnSubscriberNegotiationFailed(handler func(err *pc.NegotiationError))

func (*Call) OnSubscriberOffer

func (c *Call) OnSubscriberOffer(offer *sfu_events.SfuEvent_SubscriberOffer)

OnSubscriberOffer applies an SFU offer to the subscriber peer connection.

The subscriber can legitimately be absent: offers race with the teardown a REJOIN or a Leave performs, and this runs on the signalling read loop, so a panic here would take the process down with it.

func (*Call) OnTrackPublished

func (c *Call) OnTrackPublished(published *sfu_events.SfuEvent_TrackPublished)

func (*Call) OnTrackUnpublished

func (c *Call) OnTrackUnpublished(e *sfu_events.SfuEvent_TrackUnpublished)

func (*Call) OnUnretryableError

func (c *Call) OnUnretryableError(handler func(error))

func (*Call) PublishOptions

func (c *Call) PublishOptions() []*sfu_models.PublishOption

PublishOptions returns the encoding configuration the SFU wants this client to publish with: codec, target bitrate, frame rate, and how many spatial and temporal layers to produce.

This SDK does not own an encoder -- the application supplies already-encoded samples -- so the options are advisory. Read them at join time, and register OnPublishOptionsChanged to be told when the SFU changes its mind mid-call.

func (*Call) PublisherPC

func (c *Call) PublisherPC() *webrtc.PeerConnection

func (*Call) RawHandler

func (c *Call) RawHandler(event *sfu_events.SfuEvent)

func (*Call) RefreshState

func (c *Call) RefreshState(ctx context.Context) error

RefreshState re-runs the coordinator join to pick up the current call state (members, capabilities, settings) and credentials. Join must have succeeded first.

func (*Call) SendSubscriberRTCP

func (c *Call) SendSubscriberRTCP(pkts []rtcp.Packet) error

func (*Call) SetCredentials

func (c *Call) SetCredentials(cred models.Credentials)

func (*Call) SubscribeToTracks

func (c *Call) SubscribeToTracks(ctx context.Context, trackDetails ...*signal_rpc.TrackSubscriptionDetails) error

func (*Call) SubscriberPC

func (c *Call) SubscriberPC() *webrtc.PeerConnection

type CallConnectionState

type CallConnectionState string
const (
	// CallConnectionStateConnected indicates that the client has
	// successfully connected to the SFU and the connection is healthy.
	CallConnectionStateConnected CallConnectionState = "CONNECTED"

	// CallConnectionStateConnecting indicates that the client has not
	// left the call, but currently does not have a connection to the
	// SFU. In this state, signal server RPCs may not work because the
	// client may not even have valid credentials to the SFU
	CallConnectionStateConnecting CallConnectionState = "CONNECTING"

	// CallConnectionStateMigrating indicates the client is in the process of
	// migrating to another SFU
	CallConnectionStateMigrating CallConnectionState = "MIGRATING"

	// CallConnectionStateDisconnected is a terminal state indicating that
	// the client left the call (gracefully or otherwise)
	CallConnectionStateDisconnected CallConnectionState = "DISCONNECTED"
)

type CallEvents

type CallEvents interface {
	coordinator.CallEvents | signal.Events
}

CallEvents combines coordinator and signal event interfaces, representing events that can occur within a call.

type CallState

type CallState struct {
	models.CallResponse
	EdgeName        string                  `json:"edge_name"`
	Url             string                  `json:"url"`
	Token           string                  `json:"token"`
	WebsocketUrl    string                  `json:"websocket_url"`
	BlockedUsers    []models.UserResponse   `json:"blocked_users"`
	Members         []models.MemberResponse `json:"members"`
	OwnCapabilities []models.OwnCapability  `json:"own_capabilities"`
	StatsOptions    models.StatsOptions     `json:"stats_options"`
	Membership      *models.MemberResponse  `json:"membership"`
	JoinCallRequest *models.JoinCallRequest
}

type CallStatsReportingMetrics

type CallStatsReportingMetrics struct {
	LastSuccessfulReportAt time.Time
	SuccessfulReports      uint64
	FailedReports          uint64
	SkippedReports         uint64
	// contains filtered or unexported fields
}

CallStatsReportingMetrics counts the outcome of each stats reporting tick. The counters are written from the stats worker goroutine and read by callers of GetStatsReportingMetrics, so they are guarded; read a consistent snapshot through that method rather than touching the fields of a live instance.

func NewCallStatsReportingMetrics

func NewCallStatsReportingMetrics() *CallStatsReportingMetrics

type Client

type Client struct {
	coordinator.CoordinatorClientInterface
	User         User
	UserID       string
	ConnectionID atomicx.AtomicValue[string]
	OwnUser      atomic.Pointer[models.OwnUserResponse]
	Tracing      atomic.Pointer[rtcstats.TraceBuffer]
	// contains filtered or unexported fields
}

func NewClient

func NewClient(apiKey string, user User, token TokenProvider, opts ...Option) (*Client, error)

NewClient is the primary, client-side constructor. It mints no tokens of its own: token supplies the JWT for user, exactly as the JS SDK's StreamVideoClient takes a token or token provider.

The returned Client is connected: unless WithoutCoordinatorWS is passed, the coordinator websocket is up and delivering events.

func NewRTCClient

func NewRTCClient(apiKey, apiSecret string, opts ...ClientOption) (*Client, error)

NewRTCClient is the server-side constructor, for bots, tests and ingress/egress-style workloads. It builds a getstream-go Stream from the API secret, mints a token for the WithUser user (User{ID: "agent", Name: "Agent"} by default) and returns a connected Client. The Stream stays reachable through Server.

opts may mix this package's Option values with getstream.ClientOption values, which configure the Stream.

Never ship an API secret to an end-user device; use NewClient there.

func (*Client) AddAudioSourceStatsProviders

func (c *Client) AddAudioSourceStatsProviders(providers ...AudioSourceStatsProvider)

func (*Client) AddVideoSourceStatsProviders

func (c *Client) AddVideoSourceStatsProviders(providers ...VideoSourceStatsProvider)

func (*Client) Call

func (c *Client) Call(callType, id string) *Call

Call constructs a handle for the given call. It performs no I/O; call (*Call).Join to actually join.

func (*Client) ClientDetails

func (c *Client) ClientDetails() ClientDetails

func (*Client) GetAudioSourceStatsProviders

func (c *Client) GetAudioSourceStatsProviders() []AudioSourceStatsProvider

func (*Client) GetVideoSourceStatsProviders

func (c *Client) GetVideoSourceStatsProviders() []VideoSourceStatsProvider

func (*Client) Server

func (c *Client) Server() *getstream.Stream

Server returns the embedded server-side SDK, or nil when the Client was built without an API secret.

func (*Client) StatsReportingInterval

func (c *Client) StatsReportingInterval() time.Duration

type ClientDetails

type ClientDetails struct {
	SDKVersion SDKVersion
	OSName     string
	// SDKType overrides the SDK reported to the SFU. It defaults to
	// SDK_TYPE_GO, for which the SFU substitutes a fixed set of decode codecs
	// instead of the advertised ones, so tests that need their own capabilities
	// respected have to present a different SDK.
	SDKType     sfu_models.SdkType
	BrowserName string
}

func (ClientDetails) SDKVersionString

func (c ClientDetails) SDKVersionString() string

type ClientOption

type ClientOption = any

ClientOption is an Option or a getstream.ClientOption. NewRTCClient rejects anything else.

type CloudFrontDiscovery

type CloudFrontDiscovery struct {
	// contains filtered or unexported fields
}

func NewCloudFrontDiscovery

func NewCloudFrontDiscovery(url string, maxRetries int, client HTTPClient, log logger.ILogger) *CloudFrontDiscovery

func (*CloudFrontDiscovery) Discover

func (c *CloudFrontDiscovery) Discover(ctx context.Context) string

type Events

type Events = coordinator.Events

Events is an alias for coordinator.Events, representing a set of standard events handled by the coordinator.

type GetCredentialsFunc

type GetCredentialsFunc func(forceReload bool, excludeSFU string) (models.Credentials, error)

type HTTPClient

type HTTPClient interface {
	Do(req *http.Request) (*http.Response, error)
}

type ICERecoveryConfig

type ICERecoveryConfig struct {
	DisconnectedRestartDelay   time.Duration
	MaxRestartsWithoutRecovery int
}

ICERecoveryConfig tunes iceRecovery. The zero value is not useful; use defaultICERecoveryConfig.

type InboundTrackState

type InboundTrackState struct {
	ParticipantID ParticipantID
	TrackType     sfu_models.TrackType
	Paused        bool
}

InboundTrackState is one entry of an InboundStateNotification: whether the SFU is currently forwarding a given participant's track.

type JoinOption

type JoinOption func(*joinOptions)

JoinOption configures a single (*Call).Join.

func WithClientCapabilities

func WithClientCapabilities(caps ...sfu_models.ClientCapability) JoinOption

WithClientCapabilities adds SFU behaviours this client opts into, on top of the defaults. Notably CLIENT_CAPABILITY_COORDINATOR_STATS, which lets the SFU forward this client's stats to the coordinator.

func WithCodecsFromMediaEngine

func WithCodecsFromMediaEngine() JoinOption

WithCodecsFromMediaEngine populates the join request publisher and subscriber SDP fields from the configured media engines.

func WithE2EE

func WithE2EE() JoinOption

WithE2EE marks the join as end-to-end encrypted. Key exchange is the application's responsibility; this only tells the coordinator.

func WithExternalRTCP

func WithExternalRTCP() JoinOption

func WithLocation

func WithLocation(location string) JoinOption

WithLocation pins the join to a location hint (an airport code such as "AMS"), skipping discovery.

func WithMembersLimit

func WithMembersLimit(limit int32) JoinOption

WithMembersLimit caps how many members the coordinator returns in the call response.

func WithMigratingFrom

func WithMigratingFrom(sfuID string) JoinOption

WithMigratingFrom asks the coordinator for an edge other than the named one.

func WithNotify

func WithNotify() JoinOption

WithNotify sends a notification to the call members on join.

func WithOnTrack

func WithOnTrack(onTrack Subscriber) JoinOption

func WithPreferredPublishOptions

func WithPreferredPublishOptions(opts ...*sfu_models.PublishOption) JoinOption

WithPreferredPublishOptions asks the SFU to let this client publish with the given codecs, bitrates, and layer counts.

It is a request, not a setting: the SFU weighs it against what the call's subscribers can decode and answers in the join response's publish options, which PublishOptions returns. Without this the SFU picks entirely on its own, which for a bot with a fixed encoder may not be what it can actually produce.

func WithPublisherPeerConfiguration

func WithPublisherPeerConfiguration(conf pc.PeerConfig) JoinOption

func WithRing

func WithRing() JoinOption

WithRing rings the call members on join.

func WithSessionID

func WithSessionID(sessionID string) JoinOption

func WithSubscriberPeerConfiguration

func WithSubscriberPeerConfiguration(conf pc.PeerConfig) JoinOption

func WithoutCreate

func WithoutCreate() JoinOption

WithoutCreate joins only an existing call, failing if it does not exist. Joining creates the call by default.

type LocationDiscovery

type LocationDiscovery interface {
	Discover(_ context.Context) string
}

type OnTrackReceived

type OnTrackReceived struct {
	ParticipantID ParticipantID
	TrackType     sfu_models.TrackType

	Participant *Participant
	Track       *webrtc.TrackRemote
	RTPReceiver *webrtc.RTPReceiver
	WriteRTCP   func([]rtcp.Packet) error
}

type Option

type Option func(*options)

func WithAudioPlayoutStatsProviders

func WithAudioPlayoutStatsProviders(providers ...AudioPlayoutStatsProvider) Option

func WithAudioSourceStatsProviders

func WithAudioSourceStatsProviders(providers ...AudioSourceStatsProvider) Option

func WithClientDetails

func WithClientDetails(clientDetails ClientDetails) Option

func WithConnectTimeout

func WithConnectTimeout(d time.Duration) Option

func WithCoordinatorOptions

func WithCoordinatorOptions(coordinatorOptions ...coordinator.Option) Option

func WithDisconnectionTimeout

func WithDisconnectionTimeout(d time.Duration) Option

WithDisconnectionTimeout abandons a call that has been disconnected for this long without recovering, reporting the failure through the unretryable error handler. Zero, the default, retries for as long as the call is open.

It is the one bound that spans a whole reconnect burst: the per-strategy caps and the rejoin rate limit each govern one kind of attempt, so without this a call can keep cycling through them indefinitely.

func WithHealthCheck

func WithHealthCheck(interval, timeout time.Duration) Option

WithHealthCheck overrides how often the SDK pings the SFU and how long it waits for a reply before treating the connection as dead and reconnecting. Non-positive values keep the defaults (5s and 15s).

func WithICERecovery

func WithICERecovery(cfg ICERecoveryConfig) Option

WithICERecovery tunes how a peer connection failure is handled: how long a disconnected connection is given to recover on its own before ICE is restarted, and how many restarts are attempted before the call gives up and reconnects from scratch. Zero fields keep their defaults.

func WithLocationHintURL

func WithLocationHintURL(url string) Option

WithLocationHintURL overrides the endpoint probed to find the nearest edge. It is only consulted when location discovery is enabled.

func WithLogger

func WithLogger(l logger.ILogger) Option

func WithReconnectConfig

func WithReconnectConfig(cfg ReconnectConfig) Option

WithReconnectConfig tunes the reconnect loop: backoff bounds, how many fast reconnects are attempted, the rejoin rate limit, and the overall budget after which a call that cannot reconnect is abandoned. Zero fields keep their defaults.

func WithSource

func WithSource(source Source) Option

func WithStatsInterval

func WithStatsInterval(interval time.Duration) Option

func WithUser

func WithUser(user User) Option

WithUser sets the user NewRTCClient connects as. NewClient ignores it.

func WithVideoSourceStatsProviders

func WithVideoSourceStatsProviders(providers ...VideoSourceStatsProvider) Option

func WithoutCoordinatorWS

func WithoutCoordinatorWS() Option

func WithoutLocationDiscovery

func WithoutLocationDiscovery() Option

WithoutLocationDiscovery disables the automatic CloudFront-based location lookup.

type Participant

type Participant struct {
	Call *Call
	ParticipantID
	PublishedTracks   *skipset.SkipSet[sfu_models.TrackType]
	JoinedAt          time.Time
	TrackLookupPrefix string
	ConnectionQuality atomicx.AtomicValue[sfu_models.ConnectionQuality]
	IsSpeaking        atomic.Bool
	IsDominantSpeaker atomic.Bool
	AudioLevel        atomicx.AtomicValue[float32]
	Name              string
	Image             string
	Custom            map[string]any
	// todo make roles thread safe
	Roles []string
	// contains filtered or unexported fields
}

func NewParticipant

func NewParticipant(call *Call, p *sfu_models.Participant) *Participant

func (*Participant) IsTrackPaused

func (p *Participant) IsTrackPaused(trackType sfu_models.TrackType) bool

IsTrackPaused reports whether the SFU has paused forwarding the given track type for this participant.

A paused track produces no media even though it is still subscribed, so without this an application cannot tell "the SFU is deliberately not sending this" from "the track is broken".

func (*Participant) PausedTracks

func (p *Participant) PausedTracks() []sfu_models.TrackType

PausedTracks returns the track types the SFU is currently not forwarding for this participant.

func (*Participant) SetTrackPaused

func (p *Participant) SetTrackPaused(trackType sfu_models.TrackType, paused bool)

SetTrackPaused records whether the SFU has paused forwarding the given track type for this participant.

type ParticipantID

type ParticipantID struct {
	UserID    UserID    `json:"user_id"`
	SessionID SessionID `json:"session_id"`
}

ParticipantID identifies a participant: the user plus the session they joined with. It is comparable so it can be used as a map key.

func (ParticipantID) Equal

func (p ParticipantID) Equal(other ParticipantID) bool

func (ParticipantID) IsZero

func (p ParticipantID) IsZero() bool

type ParticipantStore

type ParticipantStore struct {
	Total     atomic.Int32
	Anonymous atomic.Int32
	// contains filtered or unexported fields
}

func NewParticipantStore

func NewParticipantStore(call *Call, callState *sfu_models.CallState) *ParticipantStore

func (*ParticipantStore) GetByID

func (s *ParticipantStore) GetByID(id ParticipantID) *Participant

func (*ParticipantStore) GetByTrackPrefix

func (s *ParticipantStore) GetByTrackPrefix(prefix string) *Participant

func (*ParticipantStore) LookupParticipantByTrack

func (s *ParticipantStore) LookupParticipantByTrack(streamID string) (*Participant, sfu_models.TrackType)

func (*ParticipantStore) Remove

func (s *ParticipantStore) Remove(p *Participant)

func (*ParticipantStore) Set

func (s *ParticipantStore) Set(p *Participant)

type PublishQualityTarget

type PublishQualityTarget struct {
	// RID identifies the simulcast layer: "f" (full), "h" (half), "q" (quarter).
	// Empty for a single-layer track.
	RID string
	// TrackType is the track these layers belong to, e.g. video or screenshare.
	TrackType sfu_models.TrackType
	// Active is whether the SFU wants this layer at all. Already applied.
	Active bool
	// MaxBitrate is the target for this layer, in bits per second.
	MaxBitrate int32
	// MaxFramerate is the target frame rate, in frames per second. Zero means
	// the SFU did not express a preference.
	MaxFramerate uint32
	// ScaleResolutionDownBy is how much to downscale the source resolution for
	// this layer: 1 for full size, 2 for half, 4 for a quarter.
	ScaleResolutionDownBy float32
	// ScalabilityMode is the SVC mode for codecs that support it, e.g. "L3T3_KEY".
	ScalabilityMode string
	// DegradationPreference says what to sacrifice when the target cannot be
	// met: frame rate, resolution, or a balance of the two.
	DegradationPreference sfu_models.DegradationPreference
	// Codec is the codec the SFU expects for this layer, when it specified one.
	Codec *sfu_models.Codec
}

PublishQualityTarget is what the SFU's bandwidth estimator wants a single simulcast layer to be encoded at.

The SDK acts on Active by itself, because whether a layer's packets go on the wire is something it controls. The rest describes the encoder, which in this SDK belongs to the application: it supplies already-encoded samples, and pion exposes no way to reconfigure an encoder from the RTP sender. So those fields are reported rather than applied, and an application that can retune its encoder should do so from OnPublishQualityChanged.

type QueryCallsResponse

type QueryCallsResponse struct {
	Calls []*Call
	Next  *string
	Prev  *string
}

type ReconnectConfig

type ReconnectConfig struct {
	// InitialBackoff is the first wait after a failed reconnect attempt; it
	// doubles up to MaxBackoff.
	InitialBackoff time.Duration
	MaxBackoff     time.Duration

	// MaxFastAttempts bounds consecutive fast reconnects before escalating.
	MaxFastAttempts int

	// FallbackFastReconnectDeadline is used when the SFU's JoinResponse does not
	// carry one.
	FallbackFastReconnectDeadline time.Duration

	// RejoinLimit and RejoinWindow rate-limit rejoins and migrations. A call
	// that exceeds them is abandoned through the unretryable error handler.
	RejoinLimit  int
	RejoinWindow time.Duration

	// DisconnectionTimeout abandons a call that has been disconnected for this
	// long without recovering. Zero means no limit, matching the previous
	// behaviour of retrying forever.
	DisconnectionTimeout time.Duration

	// MigrationCompleteTimeout is how long a migration waits for the SFU it is
	// leaving to confirm the handover before failing and escalating.
	MigrationCompleteTimeout time.Duration
}

ReconnectConfig bounds how aggressively a dropped call is reconnected. Zero fields fall back to defaults; use defaultReconnectConfig for those.

type SDKVersion

type SDKVersion struct {
	Major string
	Minor string
	Patch string
}

type SessionID

type SessionID string

SessionID identifies one participant session inside a call. The same user joining twice gets two session IDs.

type Source

type Source int32
const (
	SourceWebRTC Source = iota
	SourceRTMP
	SourceSIP
	SourceSRT
	SourceRTSP
)

type Subscriber

type Subscriber interface {
	OnTrack(track OnTrackReceived)
}

type SubscriberFunc

type SubscriberFunc func(track OnTrackReceived)

func (SubscriberFunc) OnTrack

func (s SubscriberFunc) OnTrack(track OnTrackReceived)

type TokenProvider

type TokenProvider = coordinator.TokenProvider

TokenProvider returns a Stream JWT for the given user ID. It is called once when the Client is constructed.

func StaticToken

func StaticToken(token string) TokenProvider

StaticToken returns a TokenProvider that always yields token.

type TrackDetails

type TrackDetails struct {
	Info        *sfu_models.TrackInfo
	Tracks      []webrtc.TrackLocal
	Transceiver *webrtc.RTPTransceiver
}

type User

type User struct {
	ID   string
	Name string
	Type UserType
}

User is the identity the Client connects as.

type UserID

type UserID string

UserID identifies a Stream user.

type UserType

type UserType string

UserType distinguishes the three ways a user can be authenticated against Stream. It mirrors the JS SDK's user types.

const (
	// UserTypeAuthenticated is a regular user identified by a JWT minted with
	// the app's secret.
	UserTypeAuthenticated UserType = "authenticated"
	// UserTypeGuest is a temporary user created on the fly by the coordinator.
	UserTypeGuest UserType = "guest"
	// UserTypeAnonymous is an unidentified viewer, only useful for watching
	// livestreams.
	UserTypeAnonymous UserType = "anonymous"
)

type VideoOutboundRtpProvider

type VideoOutboundRtpProvider interface {
	GetVideoOutboundRtpStats() *webrtc.OutboundRTPStreamStats
	SetVideoOutboundRtpStats(stats *webrtc.OutboundRTPStreamStats)
}

type VideoSourceStatsProvider

type VideoSourceStatsProvider interface {
	GetVideoSourceStats() *webrtc.VideoSourceStats
	SetVideoSourceStats(stats *webrtc.VideoSourceStats)
	VideoOutboundRtpProvider
}

Directories

Path Synopsis
Package audio provides a PCM audio buffer and the conversions that voice pipelines need: sample-rate and channel conversion, chunking for VAD and STT, WAV and G.711 encoding, and streaming adapters.
Package audio provides a PCM audio buffer and the conversions that voice pipelines need: sample-rate and channel conversion, chunking for VAD and STT, WAV and G.711 encoding, and streaming adapters.
opus
Package opus encodes and decodes Opus using a pure-Go codec.
Package opus encodes and decodes Opus using a pure-Go codec.
rtc
Package rtc connects the audio package to WebRTC tracks.
Package rtc connects the audio package to WebRTC tracks.
Package coordinator talks to the Stream Video coordinator: the REST call that joins a call and hands back SFU credentials, and the websocket that delivers call events.
Package coordinator talks to the Stream Video coordinator: the REST call that joins a call and hands back SFU credentials, and the websocket that delivers call events.
models
Package models provides primitives to interact with the openapi HTTP API.
Package models provides primitives to interact with the openapi HTTP API.
Package event provides a small pub/sub store for intercepting events and waiting for a specific one to arrive.
Package event provides a small pub/sub store for intercepting events and waiting for a specific one to arrive.
examples
audio command
Command audio joins a Stream call, publishes a WAV file as real Opus, and records every remote participant to their own WAV file.
Command audio joins a Stream call, publishes a WAV file as real Opus, and records every remote participant to their own WAV file.
join command
Command join joins a Stream call as a subscriber and reports how many RTP packets it receives per remote track.
Command join joins a Stream call as a subscriber and reports how many RTP packets it receives per remote track.
publish command
Command publish joins a Stream call and publishes a synthetic Opus audio track.
Command publish joins a Stream call and publishes a synthetic Opus audio track.
internal
atomicx
Package atomicx provides small generic helpers over sync/atomic.
Package atomicx provides small generic helpers over sync/atomic.
cmd/genwsevent command
Command genwsevent regenerates the event plumbing that neither oapi-codegen nor protoc-gen-go produces.
Command genwsevent regenerates the event plumbing that neither oapi-codegen nor protoc-gen-go produces.
ratelimit
Package ratelimit provides a sliding-window counter used to bound how often an expensive recovery action may be retried.
Package ratelimit provides a sliding-window counter used to bound how often an expensive recovery action may be retried.
readmecheck
Package readmecheck compiles the code snippets from the README so they cannot drift from the API.
Package readmecheck compiles the code snippets from the README so they cannot drift from the API.
red
Package red implements RFC 2198 redundant audio for Opus: each outgoing packet repeats the frames sent just before it, so a receiver can rebuild a lost packet from the one that follows.
Package red implements RFC 2198 redundant audio for Opus: each outgoing packet repeats the frames sent just before it, so a receiver can rebuild a lost packet from the one that follows.
rtretry
Package rtretry provides an http.RoundTripper that retries requests which are safe to replay: idempotent methods, plus requests explicitly marked with NewRetryableRequestWithContext.
Package rtretry provides an http.RoundTripper that retries requests which are safe to replay: idempotent methods, plus requests explicitly marked with NewRetryableRequestWithContext.
sdputil
Package sdputil holds small helpers for reading values out of parsed SDP.
Package sdputil holds small helpers for reading values out of parsed SDP.
testutil
Package testutil holds fixtures shared by tests across this module: locally signed tokens, synthetic media providers, canned track layouts and FakeSFU, a loopback stand-in for the SFU signalling websocket.
Package testutil holds fixtures shared by tests across this module: locally signed tokens, synthetic media providers, canned track layouts and FakeSFU, a loopback stand-in for the SFU signalling websocket.
xerr
Package xerr provides the small set of error helpers this module needs.
Package xerr provides the small set of error helpers this module needs.
xrand
Package xrand generates random identifiers used for WebRTC track and session IDs.
Package xrand generates random identifiers used for WebRTC track and session IDs.
Package logger defines the logging interface used throughout this module.
Package logger defines the logging interface used throughout this module.
Package rtcstats collects the WebRTC trace events the SDK reports to the SFU.
Package rtcstats collects the WebRTC trace events the SDK reports to the SFU.
Package websocket wraps a gobwas/ws connection in a typed, codec-driven read/write API.
Package websocket wraps a gobwas/ws connection in a typed, codec-driven read/write API.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL