mirror of
https://github.com/kerberos-io/agent.git
synced 2026-09-04 17:08:34 +00:00
Compare commits
8 Commits
feature/up
...
improve/mq
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2bb389bf1b | ||
|
|
bfa3c3f995 | ||
|
|
03c1458fd8 | ||
|
|
4e8fd673a8 | ||
|
|
2a3f4baa8c | ||
|
|
203d7b5518 | ||
|
|
d0a7efff85 | ||
|
|
95ea92b9ce |
@@ -81,8 +81,8 @@ RUN apk update && apk upgrade --available && sync
|
||||
|
||||
RUN mkdir -p /go/src/github.com/kerberos-io/agent/machinery/www
|
||||
COPY ui /go/src/github.com/kerberos-io/agent/ui
|
||||
RUN cd /go/src/github.com/kerberos-io/agent/ui && rm -rf yarn.lock && yarn config set network-timeout 300000 && \
|
||||
yarn && yarn build
|
||||
RUN cd /go/src/github.com/kerberos-io/agent/ui && yarn config set network-timeout 300000 && \
|
||||
yarn --frozen-lockfile && yarn build
|
||||
|
||||
####################################
|
||||
# Let's create a /dist folder containing just the files necessary for runtime.
|
||||
|
||||
@@ -35,10 +35,15 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
|
||||
// This can happen when the file was already removed (e.g. cleanup, or an
|
||||
// earlier successful upload). Skip it so the watcher drops the marker
|
||||
// instead of retrying indefinitely.
|
||||
if _, err := os.Stat("data/recordings/" + fileName); err != nil {
|
||||
info, err := os.Stat("data/recordings/" + fileName)
|
||||
if err != nil {
|
||||
log.Log.Info("UploadKerberosVault: skipping " + fileName + ", file doesn't exist anymore")
|
||||
return false, false, nil
|
||||
}
|
||||
if info.Size() == 0 {
|
||||
log.Log.Warning("UploadKerberosVault: skipping " + fileName + ", recording is empty")
|
||||
return false, false, nil
|
||||
}
|
||||
|
||||
// timestamp_microseconds_instanceName_regionCoordinates_numberOfChanges_token
|
||||
// 1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4
|
||||
|
||||
@@ -188,6 +188,9 @@ func HandleLiveStreamHLS(configuration *models.Configuration, communication *mod
|
||||
publishHLSReady(configuration, mqttClient, hubKey, deviceId, sessionID)
|
||||
lastReadyAnnounce = time.Now().Unix()
|
||||
})
|
||||
session.SetOnFailure(func(sessionID, reason string) {
|
||||
publishHLSFailure(configuration, mqttClient, hubKey, deviceId, sessionID, reason)
|
||||
})
|
||||
log.Log.Info("cloud.HandleLiveStreamHLS(): prewarming live HLS session " + session.SessionID())
|
||||
}
|
||||
|
||||
@@ -247,6 +250,9 @@ func HandleLiveStreamHLS(configuration *models.Configuration, communication *mod
|
||||
publishHLSReady(configuration, mqttClient, hubKey, deviceId, sessionID)
|
||||
lastReadyAnnounce = time.Now().Unix()
|
||||
})
|
||||
session.SetOnFailure(func(sessionID, reason string) {
|
||||
publishHLSFailure(configuration, mqttClient, hubKey, deviceId, sessionID, reason)
|
||||
})
|
||||
log.Log.Info("cloud.HandleLiveStreamHLS(): started live HLS session " + session.SessionID())
|
||||
}
|
||||
|
||||
@@ -284,6 +290,28 @@ func publishHLSReady(configuration *models.Configuration, mqttClient mqtt.Client
|
||||
}
|
||||
}
|
||||
|
||||
func publishHLSFailure(configuration *models.Configuration, mqttClient mqtt.Client, hubKey, deviceId, sessionID, reason string) {
|
||||
valueMap := map[string]interface{}{
|
||||
"session": sessionID,
|
||||
"device": deviceId,
|
||||
"reason": reason,
|
||||
}
|
||||
message := models.Message{
|
||||
Payload: models.Payload{
|
||||
Action: "receive-hls-error",
|
||||
DeviceId: deviceId,
|
||||
Value: valueMap,
|
||||
},
|
||||
}
|
||||
payload, err := models.PackageMQTTMessage(configuration, message)
|
||||
if err == nil {
|
||||
mqttClient.Publish("kerberos/hub/"+hubKey, 0, false, payload)
|
||||
log.Log.Warning("cloud.HandleLiveStreamHLS(): announced live HLS startup failure " + reason + " for " + sessionID)
|
||||
} else {
|
||||
log.Log.Error("cloud.HandleLiveStreamHLS(): failed to package receive-hls-error message: " + err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
// hlsStreamSource bundles everything the live HLS producer needs to mux one of
|
||||
// the camera's streams: the packet cursor it reads from plus the encoded
|
||||
// parameter sets and dimensions used to build that stream's init segment.
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
@@ -67,6 +68,7 @@ const (
|
||||
// cannot be delivered within roughly its own duration is stale, so the upload
|
||||
// is abandoned (dropped) rather than allowed to back up the pipeline.
|
||||
defaultPublishTimeout = 4 * time.Second
|
||||
maxErrorResponseBytes = 4 << 10
|
||||
)
|
||||
|
||||
// PublisherConfig carries the hub endpoint and credentials needed to ship live
|
||||
@@ -214,6 +216,11 @@ func (p *Publisher) post(ctx context.Context, params postParams) error {
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
body, _ := io.ReadAll(io.LimitReader(resp.Body, maxErrorResponseBytes))
|
||||
detail := strings.Join(strings.Fields(string(body)), " ")
|
||||
if detail != "" {
|
||||
return fmt.Errorf("livehls: upload %s rejected: %s: %s", params.name, resp.Status, detail)
|
||||
}
|
||||
return fmt.Errorf("livehls: upload %s rejected: %s", params.name, resp.Status)
|
||||
}
|
||||
log.Log.Debug("livehls.Publisher.post(): shipped " + params.name + " for session " + params.sessionID)
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -144,12 +145,19 @@ func TestPublisherPublishSegmentSendsSequenceAndDuration(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestPublisherReturnsErrorOnNon2xx(t *testing.T) {
|
||||
srv, _, _ := newCapturingServer(t, http.StatusInternalServerError)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
_, _ = w.Write([]byte(`{"error":true,"data":"No user found with this public and private key."}`))
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
p := testPublisher(srv.URL)
|
||||
|
||||
err := p.PublishSegment(context.Background(), "s", video.LiveSegment{SequenceNumber: 1, Data: []byte("x")})
|
||||
if err == nil {
|
||||
t.Fatal("expected an error on 500 response")
|
||||
t.Fatal("expected an error on 400 response")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "400 Bad Request") || !strings.Contains(err.Error(), "No user found with this public and private key") {
|
||||
t.Fatalf("error = %q, want status and bounded Hub response", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -281,7 +289,13 @@ func TestSessionRetriesInitWhenFirstAttemptFails(t *testing.T) {
|
||||
})
|
||||
|
||||
var ready int
|
||||
var failures int
|
||||
var failureReason string
|
||||
sess.SetOnReady(func(string) { ready++ })
|
||||
sess.SetOnFailure(func(_ string, reason string) {
|
||||
failures++
|
||||
failureReason = reason
|
||||
})
|
||||
|
||||
for i := 0; i < 60; i++ {
|
||||
isKey := i%25 == 0
|
||||
@@ -304,6 +318,9 @@ func TestSessionRetriesInitWhenFirstAttemptFails(t *testing.T) {
|
||||
if ready != 1 {
|
||||
t.Errorf("OnReady fired %d times, want 1", ready)
|
||||
}
|
||||
if failures != 1 || failureReason != "init-upload-failed" {
|
||||
t.Errorf("OnFailure = %d/%q, want 1/init-upload-failed", failures, failureReason)
|
||||
}
|
||||
}
|
||||
|
||||
// liveTestSPSForSession is the known-good baseline SPS reused across tests.
|
||||
|
||||
@@ -46,9 +46,11 @@ type Session struct {
|
||||
// lastInitAt is when the init segment was last (re)uploaded. The init is
|
||||
// re-sent periodically so its short TTL in the hub live window never lapses
|
||||
// mid-session; see refreshInitIfStale.
|
||||
lastInitAt time.Time
|
||||
readyFired bool
|
||||
onReady func(sessionID string)
|
||||
lastInitAt time.Time
|
||||
readyFired bool
|
||||
onReady func(sessionID string)
|
||||
failureFired bool
|
||||
onFailure func(sessionID, reason string)
|
||||
|
||||
// uploadsActive gates whether the init and completed segments are shipped to
|
||||
// hub-api. It is true for the default on-demand path. The prewarm path starts
|
||||
@@ -150,6 +152,7 @@ func NewSession(publisher *Publisher, opts SessionOptions) *Session {
|
||||
defer cancel()
|
||||
if err := s.publisher.PublishSegment(ctx, s.id, segment); err != nil {
|
||||
log.Log.Warning("livehls.Session: " + err.Error())
|
||||
s.fireFailureOnce("segment-upload-failed")
|
||||
return nil
|
||||
}
|
||||
s.fireReadyOnce()
|
||||
@@ -181,6 +184,7 @@ func NewSession(publisher *Publisher, opts SessionOptions) *Session {
|
||||
defer cancel()
|
||||
if err := s.publisher.PublishPart(ctx, s.id, part); err != nil {
|
||||
log.Log.Warning("livehls.Session: " + err.Error())
|
||||
s.fireFailureOnce("part-upload-failed")
|
||||
return nil
|
||||
}
|
||||
s.fireReadyOnce()
|
||||
@@ -216,6 +220,15 @@ func (s *Session) SetOnReady(fn func(sessionID string)) {
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// SetOnFailure registers a one-shot callback for startup upload failures. The
|
||||
// reason is a fixed code rather than an HTTP response body, so credentials or
|
||||
// server details cannot leak through MQTT diagnostics.
|
||||
func (s *Session) SetOnFailure(fn func(sessionID, reason string)) {
|
||||
s.mu.Lock()
|
||||
s.onFailure = fn
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// prewarmMaxBufferedSegments is how many of the most recent completed segments
|
||||
// the prewarm path keeps in memory while idle and flushes to a viewer on arrival.
|
||||
// One segment keeps startup instant (the viewer immediately gets a playable
|
||||
@@ -269,6 +282,7 @@ func (s *Session) SetUploadsActive(active bool) bool {
|
||||
ctx, cancel := s.newContext()
|
||||
if err := s.publisher.PublishSegment(ctx, s.id, buffered[i]); err != nil {
|
||||
log.Log.Warning("livehls.Session: prewarm flush: " + err.Error())
|
||||
s.fireFailureOnce("segment-upload-failed")
|
||||
cancel()
|
||||
continue
|
||||
}
|
||||
@@ -285,6 +299,7 @@ func (s *Session) SetUploadsActive(active bool) bool {
|
||||
ctx, cancel := s.newContext()
|
||||
if err := s.publisher.PublishPart(ctx, s.id, bufferedParts[i]); err != nil {
|
||||
log.Log.Warning("livehls.Session: prewarm flush (part): " + err.Error())
|
||||
s.fireFailureOnce("part-upload-failed")
|
||||
cancel()
|
||||
continue
|
||||
}
|
||||
@@ -382,6 +397,7 @@ func (s *Session) publishInitIfNeeded() bool {
|
||||
defer cancel()
|
||||
if err := s.publisher.PublishInit(ctx, s.id, initBytes); err != nil {
|
||||
log.Log.Warning("livehls.Session: init upload failed, will retry: " + err.Error())
|
||||
s.fireFailureOnce("init-upload-failed")
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -445,6 +461,18 @@ func (s *Session) fireReadyOnce() {
|
||||
fn(s.id)
|
||||
}
|
||||
|
||||
func (s *Session) fireFailureOnce(reason string) {
|
||||
s.mu.Lock()
|
||||
if s.failureFired || s.readyFired || s.onFailure == nil {
|
||||
s.mu.Unlock()
|
||||
return
|
||||
}
|
||||
s.failureFired = true
|
||||
fn := s.onFailure
|
||||
s.mu.Unlock()
|
||||
fn(s.id, reason)
|
||||
}
|
||||
|
||||
// newSessionID returns a short, unique, URL-safe session identifier of the form
|
||||
// <unix-seconds>-<random-hex>.
|
||||
func newSessionID() string {
|
||||
|
||||
@@ -264,6 +264,36 @@ func testVault(uri string) models.KStorage {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadKerberosVaultSkipsEmptyRecording(t *testing.T) {
|
||||
fileName := "1787015373_3-654_office-camera17_0-0-0-0_-1_1960.mp4"
|
||||
withRecording(t, fileName, nil)
|
||||
|
||||
requestCount := 0
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
requestCount++
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
vault := testVault(server.URL)
|
||||
configuration := &models.Configuration{Config: models.Config{
|
||||
Key: "device-key",
|
||||
KStorage: &vault,
|
||||
KStorageSecondary: &models.KStorage{},
|
||||
}}
|
||||
|
||||
uploaded, configured, err := UploadKerberosVault(configuration, fileName)
|
||||
if err != nil {
|
||||
t.Fatalf("UploadKerberosVault() error = %v", err)
|
||||
}
|
||||
if uploaded || configured {
|
||||
t.Fatalf("UploadKerberosVault() uploaded/configured = %v/%v, want false/false", uploaded, configured)
|
||||
}
|
||||
if requestCount != 0 {
|
||||
t.Fatalf("Vault received %d requests, want 0", requestCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadVaultResumable_HappyPath(t *testing.T) {
|
||||
srv := newFakeTus()
|
||||
ts := httptest.NewServer(srv)
|
||||
|
||||
@@ -173,6 +173,7 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
|
||||
|
||||
// Get the video stream from the RTSP server.
|
||||
videoStream := videoStreams[0]
|
||||
log.Log.Info(fmt.Sprintf("components.Kerberos.RunAgent(): detected main video stream: codec=%s resolution=%dx%d fps=%.2f", videoStream.Name, videoStream.Width, videoStream.Height, videoStream.FPS))
|
||||
|
||||
// Get some information from the video stream.
|
||||
width := videoStream.Width
|
||||
@@ -234,6 +235,7 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
|
||||
|
||||
// Get the video stream from the RTSP server.
|
||||
videoSubStream := videoSubStreams[0]
|
||||
log.Log.Info(fmt.Sprintf("components.Kerberos.RunAgent(): detected sub video stream: codec=%s resolution=%dx%d fps=%.2f", videoSubStream.Name, videoSubStream.Width, videoSubStream.Height, videoSubStream.FPS))
|
||||
|
||||
width := videoSubStream.Width
|
||||
height := videoSubStream.Height
|
||||
|
||||
@@ -72,6 +72,23 @@ func ConfigureMQTT(configDirectory string, configuration *models.Configuration,
|
||||
if config.Offline == "true" {
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): not starting as running in Offline mode.")
|
||||
} else {
|
||||
hubKey := ""
|
||||
if config.Cloud == "s3" && config.S3 != nil && config.S3.Publickey != "" {
|
||||
hubKey = config.S3.Publickey
|
||||
} else if config.Cloud == "kstorage" && config.KStorage != nil && config.KStorage.CloudKey != "" {
|
||||
hubKey = config.KStorage.CloudKey
|
||||
}
|
||||
if config.HubKey != "" {
|
||||
hubKey = config.HubKey
|
||||
}
|
||||
if hubKey == "" {
|
||||
log.Log.Warning("routers.mqtt.main.ConfigureMQTT(): not starting without a Hub key")
|
||||
return nil
|
||||
}
|
||||
if config.Key == "" {
|
||||
log.Log.Warning("routers.mqtt.main.ConfigureMQTT(): not starting without an Agent key")
|
||||
return nil
|
||||
}
|
||||
|
||||
opts := mqtt.NewClientOptions()
|
||||
|
||||
@@ -121,42 +138,27 @@ func ConfigureMQTT(configDirectory string, configuration *models.Configuration,
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): MQTT session is online")
|
||||
})
|
||||
|
||||
hubKey := ""
|
||||
// This is the old way ;)
|
||||
if config.Cloud == "s3" && config.S3 != nil && config.S3.Publickey != "" {
|
||||
hubKey = config.S3.Publickey
|
||||
} else if config.Cloud == "kstorage" && config.KStorage != nil && config.KStorage.CloudKey != "" {
|
||||
hubKey = config.KStorage.CloudKey
|
||||
}
|
||||
// This is the new way ;)
|
||||
if config.HubKey != "" {
|
||||
hubKey = config.HubKey
|
||||
rand.Seed(time.Now().UnixNano())
|
||||
random := rand.Intn(100)
|
||||
mqttClientID := config.Key + strconv.Itoa(random) // this random int is to avoid conflicts.
|
||||
|
||||
// This is a worked-around.
|
||||
// current S3 (Kerberos Hub SAAS) is using a secured MQTT, where the client id,
|
||||
// should match the kerberos hub key.
|
||||
if config.Cloud == "s3" {
|
||||
mqttClientID = config.Key
|
||||
}
|
||||
|
||||
if hubKey != "" {
|
||||
opts.SetClientID(mqttClientID)
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): Set ClientID " + mqttClientID)
|
||||
rand.Seed(time.Now().UnixNano())
|
||||
|
||||
rand.Seed(time.Now().UnixNano())
|
||||
random := rand.Intn(100)
|
||||
mqttClientID := config.Key + strconv.Itoa(random) // this random int is to avoid conflicts.
|
||||
opts.OnConnect = func(c mqtt.Client) {
|
||||
// We managed to connect to the MQTT broker, hurray!
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): " + mqttClientID + " connected to " + mqttURL)
|
||||
|
||||
// This is a worked-around.
|
||||
// current S3 (Kerberos Hub SAAS) is using a secured MQTT, where the client id,
|
||||
// should match the kerberos hub key.
|
||||
if config.Cloud == "s3" {
|
||||
mqttClientID = config.Key
|
||||
}
|
||||
|
||||
opts.SetClientID(mqttClientID)
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): Set ClientID " + mqttClientID)
|
||||
rand.Seed(time.Now().UnixNano())
|
||||
|
||||
opts.OnConnect = func(c mqtt.Client) {
|
||||
// We managed to connect to the MQTT broker, hurray!
|
||||
log.Log.Info("routers.mqtt.main.ConfigureMQTT(): " + mqttClientID + " connected to " + mqttURL)
|
||||
|
||||
// Create a susbcription for listen and reply
|
||||
MQTTListenerHandler(c, hubKey, configDirectory, configuration, communication)
|
||||
}
|
||||
// Create a susbcription for listen and reply
|
||||
MQTTListenerHandler(c, hubKey, configDirectory, configuration, communication)
|
||||
}
|
||||
mqc := mqtt.NewClient(opts)
|
||||
if token := mqc.Connect(); token.WaitTimeout(30 * time.Second) {
|
||||
|
||||
@@ -7,6 +7,22 @@ import (
|
||||
"github.com/kerberos-io/agent/machinery/src/models"
|
||||
)
|
||||
|
||||
func TestConfigureMQTTRequiresHubKey(t *testing.T) {
|
||||
configuration := &models.Configuration{Config: models.Config{Key: "agent-key"}}
|
||||
|
||||
if client := ConfigureMQTT("", configuration, &models.Communication{}); client != nil {
|
||||
t.Fatal("ConfigureMQTT() returned a client without a Hub key")
|
||||
}
|
||||
}
|
||||
|
||||
func TestConfigureMQTTRequiresAgentKey(t *testing.T) {
|
||||
configuration := &models.Configuration{Config: models.Config{HubKey: "hub-key"}}
|
||||
|
||||
if client := ConfigureMQTT("", configuration, &models.Communication{}); client != nil {
|
||||
t.Fatal("ConfigureMQTT() returned a client without an Agent key")
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnqueueLatestAudioReplacesOldestFrameWhenFull(t *testing.T) {
|
||||
audioChannel := make(chan models.AudioDataPartial, 2)
|
||||
audioChannel <- models.AudioDataPartial{Timestamp: 1}
|
||||
|
||||
@@ -226,6 +226,31 @@ func CreateWebRTC(name string, stunServers []string, turnServers []string, turnS
|
||||
}
|
||||
}
|
||||
|
||||
func nonEmptyICEURLs(urls []string) []string {
|
||||
nonEmpty := make([]string, 0, len(urls))
|
||||
for _, uri := range urls {
|
||||
if uri = strings.TrimSpace(uri); uri != "" {
|
||||
nonEmpty = append(nonEmpty, uri)
|
||||
}
|
||||
}
|
||||
return nonEmpty
|
||||
}
|
||||
|
||||
func buildICEServers(w WebRTC) []pionWebRTC.ICEServer {
|
||||
iceServers := make([]pionWebRTC.ICEServer, 0, 2)
|
||||
if stunURLs := nonEmptyICEURLs(w.StunServers); len(stunURLs) > 0 {
|
||||
iceServers = append(iceServers, pionWebRTC.ICEServer{URLs: stunURLs})
|
||||
}
|
||||
if turnURLs := nonEmptyICEURLs(w.TurnServers); len(turnURLs) > 0 {
|
||||
iceServers = append(iceServers, pionWebRTC.ICEServer{
|
||||
URLs: turnURLs,
|
||||
Username: w.TurnServersUsername,
|
||||
Credential: w.TurnServersCredential,
|
||||
})
|
||||
}
|
||||
return iceServers
|
||||
}
|
||||
|
||||
func (w WebRTC) DecodeSessionDescription(data string) ([]byte, error) {
|
||||
sd, err := base64.StdEncoding.DecodeString(data)
|
||||
if err != nil {
|
||||
@@ -423,21 +448,28 @@ func InitializeWebRTCConnection(configuration *models.Configuration, communicati
|
||||
|
||||
peerConnection, err := api.NewPeerConnection(
|
||||
pionWebRTC.Configuration{
|
||||
ICEServers: []pionWebRTC.ICEServer{
|
||||
{
|
||||
URLs: w.StunServers,
|
||||
},
|
||||
{
|
||||
URLs: w.TurnServers,
|
||||
Username: w.TurnServersUsername,
|
||||
Credential: w.TurnServersCredential,
|
||||
},
|
||||
},
|
||||
ICEServers: buildICEServers(*w),
|
||||
ICETransportPolicy: policy,
|
||||
},
|
||||
)
|
||||
|
||||
if err == nil && peerConnection != nil {
|
||||
if err != nil {
|
||||
globalConnectionManager.CloseCandidateChannel(sessionKey)
|
||||
log.Log.Error("webrtc.main.InitializeWebRTCConnection(): failed to create peer connection: " + err.Error() +
|
||||
" (STUN configured: " + strconv.FormatBool(strings.TrimSpace(config.STUNURI) != "") +
|
||||
", TURN configured: " + strconv.FormatBool(strings.TrimSpace(config.TURNURI) != "") +
|
||||
", TURN username configured: " + strconv.FormatBool(config.TURNUsername != "") +
|
||||
", TURN credential configured: " + strconv.FormatBool(config.TURNPassword != "") +
|
||||
", ForceTurn: " + config.ForceTurn + ")")
|
||||
return
|
||||
}
|
||||
if peerConnection == nil {
|
||||
globalConnectionManager.CloseCandidateChannel(sessionKey)
|
||||
log.Log.Error("webrtc.main.InitializeWebRTCConnection(): failed to create peer connection: pion returned a nil connection")
|
||||
return
|
||||
}
|
||||
|
||||
{
|
||||
|
||||
// Create context for this connection
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
38
machinery/src/webrtc/main_test.go
Normal file
38
machinery/src/webrtc/main_test.go
Normal file
@@ -0,0 +1,38 @@
|
||||
package webrtc
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestBuildICEServersOmitsEmptyURLs(t *testing.T) {
|
||||
webRTC := CreateWebRTC("camera", []string{""}, []string{""}, "", "")
|
||||
|
||||
iceServers := buildICEServers(*webRTC)
|
||||
|
||||
if len(iceServers) != 0 {
|
||||
t.Fatalf("buildICEServers() returned %d servers, want 0", len(iceServers))
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildICEServersIncludesConfiguredURLs(t *testing.T) {
|
||||
webRTC := CreateWebRTC(
|
||||
"camera",
|
||||
[]string{"", " stun:turn-fra1.kerberos.io:3478 "},
|
||||
[]string{" turn:turn-fra1.kerberos.io:3478 "},
|
||||
"username",
|
||||
"credential",
|
||||
)
|
||||
|
||||
iceServers := buildICEServers(*webRTC)
|
||||
|
||||
if len(iceServers) != 2 {
|
||||
t.Fatalf("buildICEServers() returned %d servers, want 2", len(iceServers))
|
||||
}
|
||||
if got := iceServers[0].URLs[0]; got != "stun:turn-fra1.kerberos.io:3478" {
|
||||
t.Fatalf("STUN URL = %q, want trimmed URL", got)
|
||||
}
|
||||
if got := iceServers[1].URLs[0]; got != "turn:turn-fra1.kerberos.io:3478" {
|
||||
t.Fatalf("TURN URL = %q, want trimmed URL", got)
|
||||
}
|
||||
if iceServers[1].Username != "username" || iceServers[1].Credential != "credential" {
|
||||
t.Fatal("TURN credentials were not preserved")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user