diff --git a/README.md b/README.md index b15d57aba..a56004b88 100644 --- a/README.md +++ b/README.md @@ -199,6 +199,7 @@ A summary table of all modules and features can be found [here](internal/README. - [`tapo`](internal/tapo/README.md) - [TP-Link Tapo](https://www.tapo.com/) cameras with two-way audio support. - [`vigi`](internal/tapo/README.md#tp-link-vigi) - TP-Link Vigi cameras. - [`tuya`](internal/tuya/README.md) - [Tuya](https://www.tuya.com/) ecosystem cameras with two-way audio support. +- [`unifi`](internal/unifi/README.md) - UniFi Protect cameras with talkback support. - [`webtorrent`](internal/webtorrent/README.md) - Stream from another go2rtc via [WebTorrent](https://en.wikipedia.org/wiki/WebTorrent) protocol. - [`wyze`](internal/wyze/README.md) - [Wyze](https://wyze.com/) cameras using native P2P protocol - [`xiaomi`](internal/xiaomi/README.md) - [Xiaomi Mi Home](https://home.mi.com/) ecosystem cameras with two-way audio support. @@ -282,6 +283,7 @@ Supported for: [`rtsp`](internal/rtsp/README.md#two-way-audio), [`tapo`](internal/tapo/README.md), [`tuya`](internal/tuya/README.md), +[`unifi`](internal/unifi/README.md), [`webrtc`](internal/webrtc/README.md), [`wyze`](internal/wyze/README.md), [`xiaomi`](internal/xiaomi/README.md). diff --git a/internal/README.md b/internal/README.md index f3e9f3b34..66bf59f79 100644 --- a/internal/README.md +++ b/internal/README.md @@ -59,6 +59,7 @@ Some formats and protocols go2rtc supports exclusively. They have no equivalent | [`rtsp`] | `rtsp` | `rtsp` | yes | yes | yes | yes | | [`tapo`] | `mpegts` | `http` | yes | | | yes | | [`tuya`] | `srtp` | `webrtc` | yes | | | yes | +| [`unifi`] | `rtp` | `http` | | | | yes | | [`v4l2`] | `rawvideo` | `ioctl` | yes | | | | | [`webrtc`] | `srtp` | `webrtc` | yes | yes | yes | yes | | [`webtorrent`] | `srtp` | `webrtc` | yes | yes | | | @@ -104,6 +105,7 @@ Some formats and protocols go2rtc supports exclusively. They have no equivalent [`streams`]: streams/README.md [`tapo`]: tapo/README.md [`tuya`]: tuya/README.md +[`unifi`]: unifi/README.md [`v4l2`]: v4l2/README.md [`webrtc`]: webrtc/README.md [`webtorrent`]: webtorrent/README.md diff --git a/internal/app/README.md b/internal/app/README.md index 1fa99235a..30b6fa0a3 100644 --- a/internal/app/README.md +++ b/internal/app/README.md @@ -94,4 +94,4 @@ log: api: trace # module name: log level ``` -Modules: `api`, `streams`, `rtsp`, `webrtc`, `mp4`, `hls`, `mjpeg`, `hass`, `homekit`, `onvif`, `rtmp`, `webtorrent`, `wyoming`, `echo`, `exec`, `expr`, `ffmpeg`, `wyze`, `xiaomi`. +Modules: `api`, `streams`, `rtsp`, `webrtc`, `mp4`, `hls`, `mjpeg`, `hass`, `homekit`, `onvif`, `rtmp`, `webtorrent`, `wyoming`, `echo`, `exec`, `expr`, `ffmpeg`, `unifi`, `wyze`, `xiaomi`. diff --git a/internal/streams/producer.go b/internal/streams/producer.go index 09e2dcc58..15acdedaa 100644 --- a/internal/streams/producer.go +++ b/internal/streams/producer.go @@ -8,6 +8,7 @@ import ( "time" "github.com/AlexxIT/go2rtc/pkg/core" + "github.com/AlexxIT/go2rtc/pkg/creds" ) type state byte @@ -135,7 +136,7 @@ func (p *Producer) MarshalJSON() ([]byte, error) { if conn := p.conn; conn != nil { return json.Marshal(conn) } - info := map[string]string{"url": p.url} + info := map[string]string{"url": creds.SecretString(p.url)} return json.Marshal(info) } @@ -149,7 +150,7 @@ func (p *Producer) start() { return } - log.Debug().Msgf("[streams] start producer url=%s", p.url) + log.Debug().Msgf("[streams] start producer url=%s", creds.SecretString(p.url)) p.state = stateStart p.workerID++ @@ -167,7 +168,7 @@ func (p *Producer) worker(conn core.Producer, workerID int) { return } - log.Warn().Err(err).Str("url", p.url).Caller().Send() + log.Warn().Err(err).Str("url", creds.SecretString(p.url)).Caller().Send() } p.reconnect(workerID, 0) @@ -178,11 +179,11 @@ func (p *Producer) reconnect(workerID, retry int) { defer p.mu.Unlock() if p.workerID != workerID { - log.Trace().Msgf("[streams] stop reconnect url=%s", p.url) + log.Trace().Msgf("[streams] stop reconnect url=%s", creds.SecretString(p.url)) return } - log.Debug().Msgf("[streams] retry=%d to url=%s", retry, p.url) + log.Debug().Msgf("[streams] retry=%d to url=%s", retry, creds.SecretString(p.url)) conn, err := GetProducer(p.url) if err != nil { @@ -257,7 +258,7 @@ func (p *Producer) stop() { p.workerID++ } - log.Debug().Msgf("[streams] stop producer url=%s", p.url) + log.Debug().Msgf("[streams] stop producer url=%s", creds.SecretString(p.url)) if p.conn != nil { _ = p.conn.Stop() diff --git a/internal/unifi/README.md b/internal/unifi/README.md new file mode 100644 index 000000000..36f063d12 --- /dev/null +++ b/internal/unifi/README.md @@ -0,0 +1,24 @@ +# UniFi Protect + +UniFi Protect talkback is supported as a native go2rtc audio backchannel. + +## Configuration + +```yaml +streams: + gate: + - rtspx://192.168.1.1:7441/#backchannel=0 + - ffmpeg:gate#audio=opus + - unifi-talkback:https://192.168.1.1?camera_id=&api_key= +``` + +The `unifi-talkback:` source uses the UniFi Protect public API: + +- `GET /proxy/protect/integration/v1/cameras/{camera_id}` checks whether the camera has a speaker +- `POST /proxy/protect/integration/v1/cameras/{camera_id}/talkback-session` starts a talkback session + +Normal `video+audio` viewing does not activate talkback. go2rtc opens the UniFi talkback session only after a client sends microphone RTP into the audio backchannel. + +Push-to-talk clients should start sending microphone audio when talk begins and stop or disconnect that microphone send session when talk ends. go2rtc closes the FFmpeg output when the backchannel producer stops. + +Supported UniFi talkback session codecs are `opus` over RTP, `aac` over ADTS, and `vorbis` over Ogg. Any other codec is rejected. diff --git a/internal/unifi/unifi.go b/internal/unifi/unifi.go new file mode 100644 index 000000000..7b9fc0e18 --- /dev/null +++ b/internal/unifi/unifi.go @@ -0,0 +1,22 @@ +package unifi + +import ( + "github.com/AlexxIT/go2rtc/internal/app" + "github.com/AlexxIT/go2rtc/internal/streams" + "github.com/AlexxIT/go2rtc/pkg/core" + "github.com/AlexxIT/go2rtc/pkg/unifi" +) + +func Init() { + unifi.SetLogger(app.GetLogger("unifi")) + + streams.HandleFunc(unifi.SchemeTalkback, func(source string) (core.Producer, error) { + return unifi.DialTalkback(source) + }) + + for _, sources := range streams.GetAllSources() { + for _, source := range sources { + unifi.RegisterTalkbackSecrets(source) + } + } +} diff --git a/main.go b/main.go index 00c059e3e..6439f11db 100644 --- a/main.go +++ b/main.go @@ -41,6 +41,7 @@ import ( "github.com/AlexxIT/go2rtc/internal/streams" "github.com/AlexxIT/go2rtc/internal/tapo" "github.com/AlexxIT/go2rtc/internal/tuya" + "github.com/AlexxIT/go2rtc/internal/unifi" "github.com/AlexxIT/go2rtc/internal/v4l2" "github.com/AlexxIT/go2rtc/internal/webrtc" "github.com/AlexxIT/go2rtc/internal/webtorrent" @@ -105,6 +106,7 @@ func main() { {"roborock", roborock.Init}, {"tapo", tapo.Init}, {"tuya", tuya.Init}, + {"unifi", unifi.Init}, {"wyze", wyze.Init}, {"xiaomi", xiaomi.Init}, {"yandex", yandex.Init}, diff --git a/pkg/README.md b/pkg/README.md index 89c1aa698..3de903df3 100644 --- a/pkg/README.md +++ b/pkg/README.md @@ -42,6 +42,7 @@ Some formats and protocols go2rtc supports exclusively. They have no equivalent | Net (priv) | roborock | webrtc | | h264, opus | opus | `roborock:` | | Net (priv) | tapo | http | | h264, pcma | pcm_alaw | `tapo:` | | Net (priv) | tuya | webrtc | | | | `tuya:` | +| Net (priv) | unifi | http, rtp, udp | | | aac, opus, vorbis | `unifi-talkback:` | | Net (priv) | vigi | http | | | | `vigi:` | | Net (priv) | webtorrent | webrtc | TODO | TODO | TODO | `webtorrent:` | | Net (priv) | xiaomi* | cs2, tutk | | | | `xiaomi:` | diff --git a/pkg/creds/secrets.go b/pkg/creds/secrets.go index 95ab4828d..7f247ec05 100644 --- a/pkg/creds/secrets.go +++ b/pkg/creds/secrets.go @@ -35,8 +35,13 @@ func getReplacer() *strings.Replacer { defer secretsMu.Unlock() if secretsReplacer == nil { - oldnew := make([]string, 0, 2*len(secrets)) - for _, s := range secrets { + values := slices.Clone(secrets) + slices.SortFunc(values, func(a, b string) int { + return len(b) - len(a) + }) + + oldnew := make([]string, 0, 2*len(values)) + for _, s := range values { oldnew = append(oldnew, s, "***") } secretsReplacer = strings.NewReplacer(oldnew...) @@ -59,16 +64,26 @@ const ( func SecretString(s string) string { re := getReplacer() - s = userinfoRegexp.ReplaceAllString(s, `://***@`) + s = secretUserinfo(s) return re.Replace(s) } func SecretWrite(w io.Writer, s string) (n int, err error) { re := getReplacer() - s = userinfoRegexp.ReplaceAllString(s, `://***@`) + s = secretUserinfo(s) return re.WriteString(w, s) } +func secretUserinfo(s string) string { + return userinfoRegexp.ReplaceAllStringFunc(s, func(match string) string { + userinfo := match[3 : len(match)-1] + if strings.Contains(userinfo, ":") { + return `://***:***@` + } + return `://***@` + }) +} + func SecretWriter(w io.Writer) io.Writer { return &secretWriter{w} } diff --git a/pkg/creds/secrets_test.go b/pkg/creds/secrets_test.go index 83f1908a2..cf3f0c9d5 100644 --- a/pkg/creds/secrets_test.go +++ b/pkg/creds/secrets_test.go @@ -13,3 +13,11 @@ func TestString(t *testing.T) { s := SecretString("rtsp://admin:pa$$word@192.168.1.123/stream1") require.Equal(t, "rtsp://***:***@192.168.1.123/stream1", s) } + +func TestStringOverlappingSecrets(t *testing.T) { + AddSecret("rtp://127.0.0.1:6500") + AddSecret("rtp://127.0.0.1:6500/talkback?token=secret") + + s := SecretString(`{"url":"rtp://127.0.0.1:6500/talkback?token=secret"}`) + require.Equal(t, `{"url":"***"}`, s) +} diff --git a/pkg/unifi/ffmpeg.go b/pkg/unifi/ffmpeg.go new file mode 100644 index 000000000..9f4183152 --- /dev/null +++ b/pkg/unifi/ffmpeg.go @@ -0,0 +1,263 @@ +package unifi + +import ( + "bufio" + "fmt" + "io" + "net" + "strconv" + "strings" + "sync" + "time" + + "github.com/AlexxIT/go2rtc/pkg/core" + "github.com/AlexxIT/go2rtc/pkg/creds" + "github.com/AlexxIT/go2rtc/pkg/shell" +) + +type ffmpegCommand interface { + StdinPipe() (io.WriteCloser, error) + Start() error + Wait() error + Close() error +} + +type ffmpegStderr interface { + StderrPipe() (io.ReadCloser, error) +} + +type rtpWriter interface { + io.WriteCloser + SetWriteDeadline(time.Time) error +} + +type talkbackOutput struct { + cmd ffmpegCommand + rtp rtpWriter + closeErr error + close sync.Once +} + +var newFFmpegCommand = func(command string) ffmpegCommand { + return shell.NewCommand(command) +} + +var reserveRTPPort = func() (int, error) { + conn, err := net.ListenPacket("udp4", "127.0.0.1:0") + if err != nil { + return 0, err + } + defer conn.Close() + + addr := conn.LocalAddr().(*net.UDPAddr) + return addr.Port, nil +} + +var dialRTPPort = func(port int) (rtpWriter, error) { + addr, err := net.ResolveUDPAddr("udp4", net.JoinHostPort("127.0.0.1", strconv.Itoa(port))) + if err != nil { + return nil, err + } + + conn, err := net.ListenUDP("udp4", nil) + if err != nil { + return nil, err + } + + return &udpRTPWriter{UDPConn: conn, addr: addr}, nil +} + +type udpRTPWriter struct { + *net.UDPConn + addr *net.UDPAddr +} + +func (w *udpRTPWriter) Write(b []byte) (int, error) { + return w.WriteTo(b, w.addr) +} + +func startFFmpegRTP(inputCodec *core.Codec, session *TalkbackSession) (*talkbackOutput, error) { + command, err := buildFFmpegCommand(session.URL, session.Codec, session.SamplingRate) + if err != nil { + return nil, err + } + + port, err := reserveRTPPort() + if err != nil { + return nil, err + } + + log.Debug(). + Int("local_port", port). + Str("codec", session.Codec). + Int("sampling_rate", session.SamplingRate). + Str("url", creds.SecretString(session.URL)). + Msg("[unifi] start ffmpeg") + log.Trace(). + Str("cmd", creds.SecretString(command)). + Msg("[unifi] ffmpeg command") + + cmd := newFFmpegCommand(command) + + stdin, err := cmd.StdinPipe() + if err != nil { + return nil, err + } + + var stderr io.ReadCloser + if cmd, ok := cmd.(ffmpegStderr); ok { + stderr, _ = cmd.StderrPipe() + } + + if err = cmd.Start(); err != nil { + _ = stdin.Close() + return nil, err + } + + if stderr != nil { + go logFFmpegStderr(stderr) + } + + if _, err = io.WriteString(stdin, buildInputSDP(inputCodec, port)); err != nil { + _ = stdin.Close() + _ = cmd.Close() + return nil, err + } + _ = stdin.Close() + + rtp, err := dialRTPPort(port) + if err != nil { + _ = cmd.Close() + return nil, err + } + + return &talkbackOutput{cmd: cmd, rtp: rtp}, nil +} + +func logFFmpegStderr(stderr io.Reader) { + scanner := bufio.NewScanner(stderr) + for scanner.Scan() { + log.Debug(). + Str("line", creds.SecretString(scanner.Text())). + Msg("[unifi] ffmpeg") + } +} + +func buildFFmpegCommand(outputURL, codec string, samplingRate int) (string, error) { + audioConfig, err := ffmpegAudioCodec(codec) + if err != nil { + return "", err + } + + args := []string{ + "ffmpeg", + "-hide_banner", + "-loglevel", "error", + "-protocol_whitelist", "file,pipe,udp,rtp", + "-fflags", "nobuffer", + "-flags", "low_delay", + "-f", "sdp", + "-i", "pipe:0", + "-map", "0:a:0", + "-vn", + "-c:a", audioConfig.encoder, + } + args = append(args, audioConfig.options...) + args = append(args, + "-ar:a", strconv.Itoa(samplingRate), + "-ac:a", "1", + "-flush_packets", "1", + "-f", audioConfig.format, + outputURL, + ) + + return strings.Join(args, " "), nil +} + +type ffmpegAudioConfig struct { + encoder string + format string + options []string +} + +func ffmpegAudioCodec(codec string) (*ffmpegAudioConfig, error) { + switch strings.ToLower(strings.TrimSpace(codec)) { + case "aac": + return &ffmpegAudioConfig{encoder: "aac", format: "adts"}, nil + case "opus": + return &ffmpegAudioConfig{ + encoder: "libopus", + format: "rtp", + options: []string{"-application:a", "lowdelay"}, + }, nil + case "vorbis": + return &ffmpegAudioConfig{encoder: "libvorbis", format: "ogg"}, nil + default: + return nil, fmt.Errorf("unifi: unsupported talkback codec: %s", codec) + } +} + +func buildInputSDP(codec *core.Codec, port int) string { + payloadType, rtpmap := inputRTPMap(codec) + + return fmt.Sprintf( + "v=0\r\n"+ + "o=- 0 0 IN IP4 127.0.0.1\r\n"+ + "s=go2rtc-unifi-talkback\r\n"+ + "c=IN IP4 127.0.0.1\r\n"+ + "t=0 0\r\n"+ + "m=audio %d RTP/AVP %d\r\n"+ + "a=rtpmap:%d %s\r\n"+ + "a=recvonly\r\n", + port, payloadType, payloadType, rtpmap, + ) +} + +func inputRTPMap(codec *core.Codec) (uint8, string) { + switch codec.Name { + case core.CodecPCMU: + return 0, "PCMU/8000" + case core.CodecPCMA: + return 8, "PCMA/8000" + case core.CodecOpus: + return 111, "opus/48000/2" + } + + payloadType := codec.PayloadType + if payloadType == 0 { + payloadType = 111 + } + + clockRate := codec.ClockRate + if clockRate == 0 { + clockRate = 48000 + } + + codecName := strings.ToLower(codec.Name) + if codecName == "" { + codecName = "opus" + } + + rtpmap := fmt.Sprintf("%s/%d", codecName, clockRate) + if codec.Channels > 1 { + rtpmap += fmt.Sprintf("/%d", codec.Channels) + } + + return payloadType, rtpmap +} + +func (o *talkbackOutput) Close() error { + o.close.Do(func() { + log.Debug().Msg("[unifi] stop ffmpeg") + + if o.rtp != nil { + o.closeErr = o.rtp.Close() + } + if o.cmd != nil { + if closeErr := o.cmd.Close(); o.closeErr == nil { + o.closeErr = closeErr + } + } + }) + return o.closeErr +} diff --git a/pkg/unifi/log.go b/pkg/unifi/log.go new file mode 100644 index 000000000..73667a2b9 --- /dev/null +++ b/pkg/unifi/log.go @@ -0,0 +1,9 @@ +package unifi + +import "github.com/rs/zerolog" + +var log = zerolog.Nop() + +func SetLogger(logger zerolog.Logger) { + log = logger +} diff --git a/pkg/unifi/talkback.go b/pkg/unifi/talkback.go new file mode 100644 index 000000000..0f47ee819 --- /dev/null +++ b/pkg/unifi/talkback.go @@ -0,0 +1,452 @@ +package unifi + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "strings" + "sync" + "time" + + "github.com/AlexxIT/go2rtc/pkg/core" + "github.com/AlexxIT/go2rtc/pkg/creds" + "github.com/AlexxIT/go2rtc/pkg/tcp" + "github.com/pion/rtp" +) + +const SchemeTalkback = "unifi-talkback" + +var errStopped = errors.New("unifi: stopped") + +type Client struct { + core.Connection + + baseURL *url.URL + cameraID string + apiKey string + + ctx context.Context + cancel context.CancelFunc + do func(*http.Request) (*http.Response, error) + + mu sync.Mutex + stopped bool + active *talkbackOutput + activeOnce sync.Once + activeReady chan struct{} + + done chan struct{} + doneOnce sync.Once + err error +} + +type publicCamera struct { + FeatureFlags struct { + HasSpeaker bool `json:"hasSpeaker"` + } `json:"featureFlags"` +} + +type TalkbackSession struct { + URL string `json:"url"` + Codec string `json:"codec"` + SamplingRate int `json:"samplingRate"` + BitsPerSample int `json:"bitsPerSample,omitempty"` +} + +func DialTalkback(rawURL string) (*Client, error) { + innerURL := strings.TrimPrefix(rawURL, SchemeTalkback+":") + if innerURL == rawURL { + return nil, fmt.Errorf("unifi: unsupported scheme: %s", rawURL) + } + + u, err := url.Parse(innerURL) + if err != nil { + return nil, err + } + + query := u.Query() + cameraID := query.Get("camera_id") + apiKey := query.Get("api_key") + if cameraID == "" { + return nil, errors.New("unifi: camera_id required") + } + if apiKey == "" { + return nil, errors.New("unifi: api_key required") + } + + RegisterTalkbackSecrets(rawURL) + log.Debug(). + Str("url", creds.SecretString(rawURL)). + Str("camera_id", cameraID). + Msg("[unifi] dial talkback") + + u.RawQuery = "" + u.Fragment = "" + + ctx, cancel := context.WithCancel(context.Background()) + c := &Client{ + Connection: core.Connection{ + ID: core.NewID(), + FormatName: "rtp", + Protocol: "http+rtp", + URL: creds.SecretString(rawURL), + }, + baseURL: u, + cameraID: cameraID, + apiKey: apiKey, + ctx: ctx, + cancel: cancel, + do: tcp.Do, + activeReady: make(chan struct{}), + done: make(chan struct{}), + } + + camera, err := c.getCamera() + if err != nil { + cancel() + return nil, err + } + + log.Debug(). + Str("camera_id", cameraID). + Bool("has_speaker", camera.FeatureFlags.HasSpeaker). + Msg("[unifi] camera metadata") + + if camera.FeatureFlags.HasSpeaker { + c.Medias = []*core.Media{ + { + Kind: core.KindAudio, + Direction: core.DirectionSendonly, + Codecs: []*core.Codec{ + { + Name: core.CodecOpus, + ClockRate: 48000, + Channels: 2, + PayloadType: 111, + }, + { + Name: core.CodecPCMU, + ClockRate: 8000, + PayloadType: 0, + }, + { + Name: core.CodecPCMA, + ClockRate: 8000, + PayloadType: 8, + }, + }, + }, + } + } + + return c, nil +} + +func RegisterTalkbackSecrets(rawURL string) { + innerURL := strings.TrimPrefix(rawURL, SchemeTalkback+":") + if innerURL == rawURL { + return + } + + u, err := url.Parse(innerURL) + if err != nil { + return + } + + creds.AddSecret(u.Query().Get("api_key")) +} + +func (c *Client) GetTrack(media *core.Media, codec *core.Codec) (*core.Receiver, error) { + return nil, core.ErrCantGetTrack +} + +func (c *Client) AddTrack(media *core.Media, _ *core.Codec, track *core.Receiver) error { + log.Debug(). + Str("camera_id", c.cameraID). + Str("codec", track.Codec.String()). + Msg("[unifi] add microphone track") + + sender := core.NewSender(media, track.Codec) + sender.Handler = func(packet *rtp.Packet) { + output, err := c.ensureActive(track.Codec) + if err != nil { + return + } + + b, err := packet.Marshal() + if err != nil { + c.finish(err) + return + } + + _ = output.rtp.SetWriteDeadline(time.Now().Add(core.ConnDeadline)) + if n, err := output.rtp.Write(b); err == nil { + c.Send += n + } else { + log.Debug().Err(err).Msg("[unifi] write microphone RTP") + c.finish(err) + } + } + sender.HandleRTP(track) + c.Senders = append(c.Senders, sender) + + return nil +} + +func (c *Client) Start() error { + <-c.done + + c.mu.Lock() + defer c.mu.Unlock() + return c.err +} + +func (c *Client) Stop() error { + log.Debug().Str("camera_id", c.cameraID).Msg("[unifi] stop talkback") + + c.cancel() + + err := c.Connection.Stop() + + c.mu.Lock() + c.stopped = true + output := c.active + c.active = nil + c.mu.Unlock() + + if output != nil { + if closeErr := output.Close(); err == nil { + err = closeErr + } + } + + c.finish(nil) + return err +} + +func (c *Client) ensureActive(inputCodec *core.Codec) (*talkbackOutput, error) { + c.activeOnce.Do(func() { + c.mu.Lock() + stopped := c.stopped + c.mu.Unlock() + + var output *talkbackOutput + var err error + if stopped { + err = errStopped + } else { + log.Debug(). + Str("camera_id", c.cameraID). + Str("codec", inputCodec.String()). + Msg("[unifi] open talkback session") + output, err = c.openTalkback(inputCodec) + } + + c.mu.Lock() + stopped = c.stopped + if err == nil && stopped { + err = errStopped + } + if err == nil { + c.active = output + } + c.mu.Unlock() + + if err != nil { + if output != nil { + _ = output.Close() + } + if !stopped && !errors.Is(err, errStopped) { + log.Warn().Err(err).Msg("[unifi] open talkback session") + c.finish(err) + } + } + + close(c.activeReady) + }) + + <-c.activeReady + + c.mu.Lock() + defer c.mu.Unlock() + + if c.err != nil { + return nil, c.err + } + if c.stopped || c.active == nil { + return nil, errStopped + } + + return c.active, nil +} + +func (c *Client) openTalkback(inputCodec *core.Codec) (*talkbackOutput, error) { + session, err := c.createTalkbackSession() + if err != nil { + return nil, err + } + + registerSecretURL(session.URL) + + log.Debug(). + Str("camera_id", c.cameraID). + Str("codec", session.Codec). + Int("sampling_rate", session.SamplingRate). + Int("bits_per_sample", session.BitsPerSample). + Str("url", creds.SecretString(session.URL)). + Msg("[unifi] talkback session") + + if session.SamplingRate == 0 { + return nil, errors.New("unifi: talkback samplingRate required") + } + + output, err := startFFmpegRTP(inputCodec, session) + if err != nil { + return nil, err + } + + go c.waitOutput(output) + + return output, nil +} + +func (c *Client) waitOutput(output *talkbackOutput) { + err := output.cmd.Wait() + log.Debug().Err(err).Str("camera_id", c.cameraID).Msg("[unifi] ffmpeg stopped") + + c.mu.Lock() + stopped := c.stopped + if c.active == output { + c.active = nil + } + c.mu.Unlock() + + _ = output.Close() + + if !stopped { + c.finish(err) + } +} + +func (c *Client) getCamera() (*publicCamera, error) { + var camera publicCamera + if err := c.requestJSON(http.MethodGet, c.cameraPath(), nil, &camera); err != nil { + return nil, err + } + return &camera, nil +} + +func (c *Client) createTalkbackSession() (*TalkbackSession, error) { + var session TalkbackSession + if err := c.requestJSON(http.MethodPost, c.cameraPath()+"/talkback-session", nil, &session); err != nil { + return nil, err + } + return &session, nil +} + +func (c *Client) requestJSON(method, path string, body io.Reader, out any) error { + u := *c.baseURL + u.Path = path + u.RawQuery = "" + + ts := time.Now() + log.Debug(). + Str("method", method). + Str("path", path). + Msg("[unifi] protect api request") + + req, err := http.NewRequestWithContext(c.ctx, method, u.String(), body) + if err != nil { + return err + } + req.Header.Set("X-API-KEY", c.apiKey) + + res, err := c.do(req) + if err != nil { + return err + } + defer tcp.Close(res) + + b, err := io.ReadAll(res.Body) + if err != nil { + log.Debug(). + Err(err). + Str("method", method). + Str("path", path). + Int("status", res.StatusCode). + Stringer("duration", time.Since(ts)). + Msg("[unifi] protect api response") + return err + } + + if res.StatusCode < http.StatusOK || res.StatusCode >= http.StatusMultipleChoices { + logProtectResponse(method, path, res.StatusCode, time.Since(ts), b) + return fmt.Errorf("unifi: %s %s: %s", method, path, res.Status) + } + + if err = json.Unmarshal(b, out); err != nil { + logProtectResponse(method, path, res.StatusCode, time.Since(ts), b) + return err + } + + if session, ok := out.(*TalkbackSession); ok { + registerSecretURL(session.URL) + } + + logProtectResponse(method, path, res.StatusCode, time.Since(ts), b) + return nil +} + +func (c *Client) cameraPath() string { + return "/proxy/protect/integration/v1/cameras/" + url.PathEscape(c.cameraID) +} + +func logProtectResponse(method, path string, status int, duration time.Duration, body []byte) { + log.Debug(). + Str("method", method). + Str("path", path). + Int("status", status). + Stringer("duration", duration). + Msg("[unifi] protect api response") + + if log.Trace().Enabled() { + log.Trace(). + Str("method", method). + Str("path", path). + Str("body", creds.SecretString(string(body))). + Msg("[unifi] protect api response body") + } +} + +func registerSecretURL(value string) { + if value == "" { + return + } + + creds.AddSecret(value) + + jsonEscaped := strings.ReplaceAll(value, `/`, `\/`) + creds.AddSecret(jsonEscaped) + + htmlEscaped := strings.ReplaceAll(value, `&`, `\u0026`) + creds.AddSecret(htmlEscaped) + creds.AddSecret(strings.ReplaceAll(htmlEscaped, `/`, `\/`)) +} + +func (c *Client) finish(err error) { + if err != nil { + c.mu.Lock() + if c.err == nil { + c.err = err + } + c.mu.Unlock() + } + + c.doneOnce.Do(func() { + close(c.done) + }) +} diff --git a/pkg/unifi/talkback_test.go b/pkg/unifi/talkback_test.go new file mode 100644 index 000000000..462e04c8c --- /dev/null +++ b/pkg/unifi/talkback_test.go @@ -0,0 +1,427 @@ +package unifi + +import ( + "bytes" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/AlexxIT/go2rtc/pkg/core" + "github.com/AlexxIT/go2rtc/pkg/creds" + "github.com/pion/rtp" + "github.com/stretchr/testify/require" +) + +func TestDialTalkbackURLAndSecret(t *testing.T) { + apiKey := "unifi-talkback-secret-key" + server, _ := newProtectServer(t, true, TalkbackSession{}) + defer server.Close() + + rawURL := SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=" + apiKey + client, err := DialTalkback(rawURL) + require.NoError(t, err) + defer client.Stop() + + require.Equal(t, strings.ReplaceAll(rawURL, apiKey, "***"), creds.SecretString(rawURL)) + + b, err := json.Marshal(client) + require.NoError(t, err) + require.NotContains(t, string(b), apiKey) + require.Contains(t, string(b), "api_key=***") +} + +func TestRegisterTalkbackSecrets(t *testing.T) { + apiKey := "unifi-register-secret-key" + rawURL := SchemeTalkback + ":https://192.168.1.1?camera_id=camera1&api_key=" + apiKey + + RegisterTalkbackSecrets(rawURL) + + require.NotContains(t, creds.SecretString(rawURL), apiKey) + require.Contains(t, creds.SecretString(rawURL), "api_key=***") +} + +func TestDialTalkbackAdvertisesMediaWithSpeaker(t *testing.T) { + server, _ := newProtectServer(t, true, TalkbackSession{}) + defer server.Close() + + client, err := DialTalkback(SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=key-with-speaker") + require.NoError(t, err) + defer client.Stop() + + require.Len(t, client.Medias, 1) + require.Equal(t, core.KindAudio, client.Medias[0].Kind) + require.Equal(t, core.DirectionSendonly, client.Medias[0].Direction) + require.Equal(t, []*core.Codec{ + { + Name: core.CodecOpus, + ClockRate: 48000, + Channels: 2, + PayloadType: 111, + }, + { + Name: core.CodecPCMU, + ClockRate: 8000, + PayloadType: 0, + }, + { + Name: core.CodecPCMA, + ClockRate: 8000, + PayloadType: 8, + }, + }, client.Medias[0].Codecs) +} + +func TestDialTalkbackDoesNotAdvertiseMediaWithoutSpeaker(t *testing.T) { + server, _ := newProtectServer(t, false, TalkbackSession{}) + defer server.Close() + + client, err := DialTalkback(SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=key-no-speaker") + require.NoError(t, err) + defer client.Stop() + + require.Empty(t, client.Medias) +} + +func TestTalkbackNoSessionDuringPassiveProbe(t *testing.T) { + server, posts := newProtectServer(t, true, TalkbackSession{ + URL: "rtp://127.0.0.1:6000", + Codec: "opus", + SamplingRate: 24000, + }) + defer server.Close() + + client, err := DialTalkback(SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=key-passive") + require.NoError(t, err) + + media := client.Medias[0] + track := core.NewReceiver(&core.Media{ + Kind: core.KindAudio, + Direction: core.DirectionRecvonly, + }, media.Codecs[0]) + require.NoError(t, client.AddTrack(media, media.Codecs[0], track)) + + done := make(chan error, 1) + go func() { + done <- client.Start() + }() + + time.Sleep(50 * time.Millisecond) + require.Zero(t, posts.Load()) + + require.NoError(t, client.Stop()) + require.NoError(t, <-done) + require.Zero(t, posts.Load()) +} + +func TestTalkbackStartsOnFirstRTPAndStopsResources(t *testing.T) { + restore, ffmpeg := fakeFFmpeg(t) + defer restore() + + server, posts := newProtectServer(t, true, TalkbackSession{ + URL: "rtp://127.0.0.1:6500", + Codec: "opus", + SamplingRate: 24000, + }) + defer server.Close() + + client, err := DialTalkback(SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=key-active") + require.NoError(t, err) + + media := client.Medias[0] + track := core.NewReceiver(&core.Media{ + Kind: core.KindAudio, + Direction: core.DirectionRecvonly, + }, media.Codecs[0]) + require.NoError(t, client.AddTrack(media, media.Codecs[0], track)) + + done := make(chan error, 1) + go func() { + done <- client.Start() + }() + + time.Sleep(50 * time.Millisecond) + require.Zero(t, posts.Load()) + require.Empty(t, ffmpeg.Commands()) + + track.WriteRTP(&rtp.Packet{ + Header: rtp.Header{ + Version: 2, + PayloadType: 111, + SequenceNumber: 1, + Timestamp: 960, + }, + Payload: []byte{0x01, 0x02, 0x03}, + }) + + require.Eventually(t, func() bool { + return posts.Load() == 1 && len(ffmpeg.Commands()) == 1 && ffmpeg.writer.writeCount.Load() == 1 + }, time.Second, 10*time.Millisecond) + + cmd := ffmpeg.Commands()[0] + require.Contains(t, cmd.command, "-c:a libopus") + require.Contains(t, cmd.command, "-ar:a 24000") + require.Contains(t, cmd.command, "-ac:a 1") + require.Contains(t, cmd.command, "-f rtp rtp://127.0.0.1:6500") + require.Contains(t, cmd.stdin.String(), "m=audio 45000 RTP/AVP 111") + require.Contains(t, cmd.stdin.String(), "a=rtpmap:111 opus/48000/2") + + require.NoError(t, client.Stop()) + require.NoError(t, <-done) + + require.True(t, cmd.closed.Load()) + require.True(t, ffmpeg.writer.closed.Load()) +} + +func TestTalkbackRejectsUnsupportedCodec(t *testing.T) { + restore, ffmpeg := fakeFFmpeg(t) + defer restore() + + server, _ := newProtectServer(t, true, TalkbackSession{ + URL: "rtp://127.0.0.1:6500", + Codec: "flac", + SamplingRate: 16000, + }) + defer server.Close() + + client, err := DialTalkback(SchemeTalkback + ":" + server.URL + "?camera_id=camera1&api_key=key-unsupported") + require.NoError(t, err) + defer client.Stop() + + media := client.Medias[0] + track := core.NewReceiver(&core.Media{ + Kind: core.KindAudio, + Direction: core.DirectionRecvonly, + }, media.Codecs[0]) + require.NoError(t, client.AddTrack(media, media.Codecs[0], track)) + + done := make(chan error, 1) + go func() { + done <- client.Start() + }() + + track.WriteRTP(&rtp.Packet{ + Header: rtp.Header{ + Version: 2, + PayloadType: 111, + SequenceNumber: 1, + Timestamp: 960, + }, + Payload: []byte{0x01, 0x02, 0x03}, + }) + + require.Eventually(t, func() bool { + return len(done) == 1 + }, time.Second, 10*time.Millisecond) + + err = <-done + require.EqualError(t, err, "unifi: unsupported talkback codec: flac") + require.Empty(t, ffmpeg.Commands()) +} + +func TestBuildFFmpegCommandOpus(t *testing.T) { + command, err := buildFFmpegCommand("rtp://10.0.0.2:4444", "opus", 24000) + require.NoError(t, err) + + require.Contains(t, command, "ffmpeg -hide_banner -loglevel error") + require.Contains(t, command, "-protocol_whitelist file,pipe,udp,rtp") + require.Contains(t, command, "-c:a libopus") + require.Contains(t, command, "-application:a lowdelay") + require.Contains(t, command, "-ar:a 24000") + require.Contains(t, command, "-ac:a 1") + require.Contains(t, command, "-f rtp rtp://10.0.0.2:4444") +} + +func TestBuildFFmpegCommandAAC(t *testing.T) { + command, err := buildFFmpegCommand("udp://10.0.0.2:4444", "aac", 16000) + require.NoError(t, err) + + require.Contains(t, command, "-c:a aac") + require.NotContains(t, command, "-application:a lowdelay") + require.Contains(t, command, "-ar:a 16000") + require.Contains(t, command, "-ac:a 1") + require.Contains(t, command, "-f adts udp://10.0.0.2:4444") +} + +func TestBuildFFmpegCommandVorbis(t *testing.T) { + command, err := buildFFmpegCommand("udp://10.0.0.2:4444", "vorbis", 16000) + require.NoError(t, err) + + require.Contains(t, command, "-c:a libvorbis") + require.NotContains(t, command, "-application:a lowdelay") + require.Contains(t, command, "-ar:a 16000") + require.Contains(t, command, "-ac:a 1") + require.Contains(t, command, "-f ogg udp://10.0.0.2:4444") +} + +func TestBuildFFmpegCommandUnsupportedCodec(t *testing.T) { + command, err := buildFFmpegCommand("rtp://10.0.0.2:4444", "flac", 24000) + + require.EqualError(t, err, "unifi: unsupported talkback codec: flac") + require.Empty(t, command) +} + +func TestBuildInputSDPG711(t *testing.T) { + pcmu := buildInputSDP(&core.Codec{Name: core.CodecPCMU, ClockRate: 8000}, 45000) + require.Contains(t, pcmu, "m=audio 45000 RTP/AVP 0") + require.Contains(t, pcmu, "a=rtpmap:0 PCMU/8000") + + pcma := buildInputSDP(&core.Codec{Name: core.CodecPCMA, ClockRate: 8000, PayloadType: 8}, 45002) + require.Contains(t, pcma, "m=audio 45002 RTP/AVP 8") + require.Contains(t, pcma, "a=rtpmap:8 PCMA/8000") +} + +func TestProtectResponseLogRedactsSessionURL(t *testing.T) { + registerSecretURL("rtp://127.0.0.1:6500") + sessionURL := "rtp://127.0.0.1:6500/talkback?token=unifi-session-token&other=1" + registerSecretURL(sessionURL) + + body := strings.ReplaceAll(`{"url":"`+sessionURL+`"}`, "/", `\/`) + body = strings.ReplaceAll(body, "&", `\u0026`) + redacted := creds.SecretString(body) + + require.NotContains(t, redacted, sessionURL) + require.NotContains(t, redacted, "unifi-session-token") + require.Contains(t, redacted, "***") +} + +func newProtectServer(t *testing.T, hasSpeaker bool, session TalkbackSession) (*httptest.Server, *atomic.Int64) { + t.Helper() + + var posts atomic.Int64 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + require.NotEmpty(t, r.Header.Get("X-API-KEY")) + require.Empty(t, r.URL.RawQuery) + + switch { + case r.Method == http.MethodGet && r.URL.Path == "/proxy/protect/integration/v1/cameras/camera1": + _ = json.NewEncoder(w).Encode(map[string]any{ + "featureFlags": map[string]any{ + "hasSpeaker": hasSpeaker, + }, + }) + case r.Method == http.MethodPost && r.URL.Path == "/proxy/protect/integration/v1/cameras/camera1/talkback-session": + posts.Add(1) + _ = json.NewEncoder(w).Encode(session) + default: + http.NotFound(w, r) + } + })) + + return server, &posts +} + +func fakeFFmpeg(t *testing.T) (func(), *fakeFFmpegState) { + t.Helper() + + oldCommand := newFFmpegCommand + oldReserve := reserveRTPPort + oldDial := dialRTPPort + + state := &fakeFFmpegState{ + writer: &fakeRTPWriter{}, + } + + newFFmpegCommand = func(command string) ffmpegCommand { + cmd := &fakeCommand{ + command: command, + done: make(chan struct{}), + } + state.mu.Lock() + state.commands = append(state.commands, cmd) + state.mu.Unlock() + return cmd + } + reserveRTPPort = func() (int, error) { + return 45000, nil + } + dialRTPPort = func(port int) (rtpWriter, error) { + require.Equal(t, 45000, port) + return state.writer, nil + } + + return func() { + newFFmpegCommand = oldCommand + reserveRTPPort = oldReserve + dialRTPPort = oldDial + }, state +} + +type fakeFFmpegState struct { + mu sync.Mutex + commands []*fakeCommand + writer *fakeRTPWriter +} + +func (s *fakeFFmpegState) Commands() []*fakeCommand { + s.mu.Lock() + defer s.mu.Unlock() + + commands := make([]*fakeCommand, len(s.commands)) + copy(commands, s.commands) + return commands +} + +type fakeCommand struct { + command string + stdin bytes.Buffer + done chan struct{} + once sync.Once + started atomic.Bool + closed atomic.Bool +} + +func (c *fakeCommand) StdinPipe() (io.WriteCloser, error) { + return nopWriteCloser{&c.stdin}, nil +} + +func (c *fakeCommand) Start() error { + c.started.Store(true) + return nil +} + +func (c *fakeCommand) Wait() error { + <-c.done + return nil +} + +func (c *fakeCommand) Close() error { + c.closed.Store(true) + c.once.Do(func() { + close(c.done) + }) + return nil +} + +type nopWriteCloser struct { + io.Writer +} + +func (n nopWriteCloser) Close() error { + return nil +} + +type fakeRTPWriter struct { + closed atomic.Bool + writeCount atomic.Int64 + buf bytes.Buffer +} + +func (w *fakeRTPWriter) Write(b []byte) (int, error) { + w.writeCount.Add(1) + return w.buf.Write(b) +} + +func (w *fakeRTPWriter) Close() error { + w.closed.Store(true) + return nil +} + +func (w *fakeRTPWriter) SetWriteDeadline(time.Time) error { + return nil +} diff --git a/website/.vitepress/config.js b/website/.vitepress/config.js index 792f2e75e..be00ba82a 100644 --- a/website/.vitepress/config.js +++ b/website/.vitepress/config.js @@ -26,7 +26,7 @@ export default defineConfig({ // second line of Telegram card (black bold), autodetect from site description ['meta', { property: 'og:title', content: 'go2rtc - Ultimate camera streaming application' }], // third line of Telegram card, autodetect from site description - ['meta', { property: 'og:description', content: 'Support alsa, doorbird, dvrip, eseecloud, ffmpeg, gopro, hass, hls, homekit, mjpeg, mp4, mpegts, nest, onvif, ring, roborock, rtmp, rtsp, tapo, vigi, tuya, v4l2, webrtc, wyze, xiaomi.' }], + ['meta', { property: 'og:description', content: 'Support alsa, doorbird, dvrip, eseecloud, ffmpeg, gopro, hass, hls, homekit, mjpeg, mp4, mpegts, nest, onvif, ring, roborock, rtmp, rtsp, tapo, vigi, tuya, unifi, v4l2, webrtc, wyze, xiaomi.' }], ['meta', { property: 'og:url', content: 'https://go2rtc.org/' }], ['meta', { property: 'og:image', content: 'https://go2rtc.org/images/logo.png' }], // important for Telegram - the image will be at the bottom and large @@ -142,6 +142,7 @@ export default defineConfig({ {text: 'roborock', link: '/internal/roborock/'}, {text: 'tapo', link: '/internal/tapo/'}, {text: 'tuya', link: '/internal/tuya/'}, + {text: 'unifi', link: '/internal/unifi/'}, {text: 'v4l2', link: '/internal/v4l2/'}, {text: 'wyze', link: '/internal/wyze/'}, {text: 'xiaomi', link: '/internal/xiaomi/'}, diff --git a/www/schema.json b/www/schema.json index 27fee57d7..55f226f53 100644 --- a/www/schema.json +++ b/www/schema.json @@ -178,6 +178,7 @@ "roborock", "tapo", "tuya", + "unifi", "xiaomi", "yandex", "debug",