From ea0238abe92433cff4a6a5dea4735b5491f512f6 Mon Sep 17 00:00:00 2001 From: Josiah Clark Date: Wed, 15 Jul 2026 11:15:27 +1200 Subject: [PATCH 1/4] isapi: add AAC two-way talk with length-prefixed ADTS Hikvision AAC /audioData uses [u32be len][ADTS]. Accept AAC channels, pass sessionId on open, and transcode WebRTC PCMU/PCMA to AAC via ffmpeg. Co-authored-by: Cursor --- internal/isapi/README.md | 15 ++ pkg/isapi/backchannel.go | 69 --------- pkg/isapi/client.go | 320 ++++++++++++++++++++++++++++++++++++--- 3 files changed, 310 insertions(+), 94 deletions(-) delete mode 100644 pkg/isapi/backchannel.go diff --git a/internal/isapi/README.md b/internal/isapi/README.md index 63892f725..2b76a819e 100644 --- a/internal/isapi/README.md +++ b/internal/isapi/README.md @@ -12,3 +12,18 @@ streams: - rtsp://admin:password@192.168.1.123:554/Streaming/Channels/101 - isapi://admin:password@192.168.1.123:80/ ``` + +## Codecs + +| Camera TwoWayAudio | Support | +|--------------------|---------| +| G.711ulaw / G.711alaw | Raw PCMU/PCMA (classic path) | +| AAC | Length-prefixed ADTS over `/audioData` (`[u32be len][ADTS]`). WebRTC mic (PCMU/PCMA) is transcoded to AAC via `ffmpeg` | + +Requires `ffmpeg` on PATH when the camera is set to AAC and the browser sends G.711. + +## Notes + +- Session: `close` → brief settle → `open` → `PUT .../audioData?sessionId=...` +- AAC sample rate is read from `audioSamplingRate` (typically 16 kHz) +- Some firmware lists G.711 in capabilities but rejects `open` while AAC works — leave the camera on AAC diff --git a/pkg/isapi/backchannel.go b/pkg/isapi/backchannel.go deleted file mode 100644 index ade16255b..000000000 --- a/pkg/isapi/backchannel.go +++ /dev/null @@ -1,69 +0,0 @@ -package isapi - -import ( - "encoding/json" - - "github.com/AlexxIT/go2rtc/pkg/core" - "github.com/pion/rtp" -) - -func (c *Client) GetMedias() []*core.Media { - return c.medias -} - -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 { - if c.sender == nil { - c.sender = core.NewSender(media, track.Codec) - c.sender.Handler = func(packet *rtp.Packet) { - if c.conn == nil { - return - } - c.send += len(packet.Payload) - _, _ = c.conn.Write(packet.Payload) - } - } - - c.sender.HandleRTP(track) - return nil -} - -func (c *Client) Start() (err error) { - if err = c.Open(); err != nil { - return - } - return -} - -func (c *Client) Stop() (err error) { - if c.sender != nil { - c.sender.Close() - } - - if c.conn != nil { - _ = c.Close() - return c.conn.Close() - } - - return nil -} - -func (c *Client) MarshalJSON() ([]byte, error) { - info := &core.Connection{ - ID: core.ID(c), - FormatName: "isapi", - Protocol: "http", - Medias: c.medias, - Send: c.send, - } - if c.conn != nil { - info.RemoteAddr = c.conn.RemoteAddr().String() - } - if c.sender != nil { - info.Senders = []*core.Sender{c.sender} - } - return json.Marshal(info) -} diff --git a/pkg/isapi/client.go b/pkg/isapi/client.go index ba3e68874..155e62449 100644 --- a/pkg/isapi/client.go +++ b/pkg/isapi/client.go @@ -1,14 +1,23 @@ package isapi import ( + "encoding/binary" + "encoding/hex" + "encoding/json" "errors" "io" "net" "net/http" "net/url" + "os/exec" + "strconv" + "sync" + "time" + "github.com/AlexxIT/go2rtc/pkg/aac" "github.com/AlexxIT/go2rtc/pkg/core" "github.com/AlexxIT/go2rtc/pkg/tcp" + "github.com/pion/rtp" ) // Deprecated: should be rewritten to core.Connection @@ -19,13 +28,21 @@ type Client struct { channel string conn net.Conn + codecName string // core.CodecAAC / CodecPCMU / CodecPCMA + sampleRate uint32 // AAC talk sample rate (Hz), typically 16000 + sessionID string + medias []*core.Media sender *core.Sender send int + + mu sync.Mutex + ffCmd *exec.Cmd + ffIn io.WriteCloser + ffDone chan struct{} } func Dial(rawURL string) (*Client, error) { - // check if url is valid url u, err := url.Parse(rawURL) if err != nil { return nil, err @@ -59,6 +76,7 @@ func (c *Client) Dial() (err error) { } b, err := io.ReadAll(res.Body) + tcp.Close(res) if err != nil { return err } @@ -68,11 +86,26 @@ func (c *Client) Dial() (err error) { codec := core.Between(xml, ``, `<`) switch codec { case "G.711ulaw": - codec = core.CodecPCMU + c.codecName = core.CodecPCMU + c.sampleRate = 8000 case "G.711alaw": - codec = core.CodecPCMA + c.codecName = core.CodecPCMA + c.sampleRate = 8000 + case "AAC": + c.codecName = core.CodecAAC + c.sampleRate = 16000 + if s := core.Between(xml, ``, `<`); s != "" { + if v, err := strconv.Atoi(s); err == nil && v > 0 { + // Camera XML uses kHz (e.g. 16) not Hz. + if v < 1000 { + c.sampleRate = uint32(v) * 1000 + } else { + c.sampleRate = uint32(v) + } + } + } default: - return nil + return errors.New("isapi: unsupported two-way codec: " + codec) } c.channel = core.Between(xml, ``, `<`) @@ -80,23 +113,38 @@ func (c *Client) Dial() (err error) { media := &core.Media{ Kind: core.KindAudio, Direction: core.DirectionSendonly, - Codecs: []*core.Codec{ - {Name: codec, ClockRate: 8000}, - }, } - c.medias = append(c.medias, media) + if c.codecName == core.CodecAAC { + conf := aac.EncodeConfig(aac.TypeAACLC, c.sampleRate, 1, false) + media.Codecs = []*core.Codec{ + { + Name: core.CodecAAC, + ClockRate: c.sampleRate, + Channels: 1, + FmtpLine: aac.FMTP + hex.EncodeToString(conf), + }, + // WebRTC mic offers Opus/PCMU/PCMA. Match G.711 and transcode to AAC. + {Name: core.CodecPCMU, ClockRate: 8000}, + {Name: core.CodecPCMA, ClockRate: 8000}, + } + } else { + media.Codecs = []*core.Codec{{ + Name: c.codecName, + ClockRate: c.sampleRate, + }} + } + + c.medias = append(c.medias, media) return nil } func (c *Client) Open() (err error) { - // Hikvision ISAPI may not accept a new open request if the previous one was not closed (e.g. - // using the test button on-camera or via curl command) but a close request can be sent even if - // the audio is already closed. So, we send a close request first and then open it again. Seems - // janky but it works. + // Hikvision may reject open if a previous session was not closed. if err = c.Close(); err != nil { return err } + time.Sleep(300 * time.Millisecond) link := c.url + "/ISAPI/System/TwoWayAudio/channels/" + c.channel req, err := http.NewRequest("PUT", link+"/open", nil) @@ -109,10 +157,22 @@ func (c *Client) Open() (err error) { return } + b, _ := io.ReadAll(res.Body) tcp.Close(res) + if res.StatusCode != http.StatusOK { + return errors.New("isapi: open: " + res.Status) + } + + c.sessionID = core.Between(string(b), ``, `<`) + + audioData := link + "/audioData" + if c.sessionID != "" { + audioData += "?sessionId=" + url.QueryEscape(c.sessionID) + } + ctx, pconn := tcp.WithConn() - req, err = http.NewRequestWithContext(ctx, "PUT", link+"/audioData", nil) + req, err = http.NewRequestWithContext(ctx, "PUT", audioData, nil) if err != nil { return err } @@ -127,12 +187,10 @@ func (c *Client) Open() (err error) { c.conn = *pconn - // just block until c.conn closed - b := make([]byte, 1) - _, _ = c.conn.Read(b) + buf := make([]byte, 1) + _, _ = c.conn.Read(buf) tcp.Close(res) - return nil } @@ -149,16 +207,228 @@ func (c *Client) Close() (err error) { } tcp.Close(res) + return nil +} +func (c *Client) GetMedias() []*core.Media { + return c.medias +} + +func (c *Client) GetTrack(media *core.Media, codec *core.Codec) (*core.Receiver, error) { + return nil, core.ErrCantGetTrack +} + +func (c *Client) AddTrack(media *core.Media, codec *core.Codec, track *core.Receiver) error { + if c.sender != nil { + c.sender.HandleRTP(track) + return nil + } + + c.sender = core.NewSender(media, track.Codec) + + switch { + case c.codecName == core.CodecAAC && track.Codec.Name == core.CodecAAC: + c.sender.Handler = func(packet *rtp.Packet) { + c.writeADTSFrames(packet.Payload) + } + if track.Codec.IsRTP() { + c.sender.Handler = aac.RTPToADTS(codec, c.sender.Handler) + } else { + c.sender.Handler = aac.EncodeToADTS(codec, c.sender.Handler) + } + + case c.codecName == core.CodecAAC && (track.Codec.Name == core.CodecPCMU || track.Codec.Name == core.CodecPCMA): + srcCodec := track.Codec.Name + c.sender.Handler = func(packet *rtp.Packet) { + c.writePCMUToAAC(srcCodec, packet.Payload) + } + + default: + c.sender.Handler = func(packet *rtp.Packet) { + if c.conn == nil { + return + } + c.send += len(packet.Payload) + _, _ = c.conn.Write(packet.Payload) + } + } + + c.sender.HandleRTP(track) return nil } -//type XMLChannels struct { -// Channels []Channel `xml:"TwoWayAudioChannel"` -//} +func (c *Client) Start() (err error) { + return c.Open() +} + +func (c *Client) Stop() (err error) { + c.stopFFmpeg() + + if c.sender != nil { + c.sender.Close() + } + + if c.conn != nil { + _ = c.Close() + return c.conn.Close() + } + + return nil +} + +func (c *Client) MarshalJSON() ([]byte, error) { + info := &core.Connection{ + ID: core.ID(c), + FormatName: "isapi", + Protocol: "http", + Medias: c.medias, + Send: c.send, + } + if c.conn != nil { + info.RemoteAddr = c.conn.RemoteAddr().String() + } + if c.sender != nil { + info.Senders = []*core.Sender{c.sender} + } + return json.Marshal(info) +} + +// writeADTSFrames writes Hikvision ISAPI AAC framing: [u32be len][ADTS]... +func (c *Client) writeADTSFrames(b []byte) { + if c.conn == nil || len(b) == 0 { + return + } + + for len(b) >= aac.ADTSHeaderSize { + if !aac.IsADTS(b) { + b = b[1:] + continue + } + size := int(aac.ReadADTSSize(b)) + if size < aac.ADTSHeaderSize || size > len(b) { + return + } + frame := b[:size] + b = b[size:] + + var hdr [4]byte + binary.BigEndian.PutUint32(hdr[:], uint32(size)) + if _, err := c.conn.Write(hdr[:]); err != nil { + return + } + if _, err := c.conn.Write(frame); err != nil { + return + } + c.send += 4 + size + } +} + +func (c *Client) writePCMUToAAC(codecName string, payload []byte) { + if len(payload) == 0 { + return + } + if err := c.ensureFFmpeg(codecName); err != nil { + return + } + c.mu.Lock() + in := c.ffIn + c.mu.Unlock() + if in == nil { + return + } + _, _ = in.Write(payload) +} + +func (c *Client) ensureFFmpeg(codecName string) error { + c.mu.Lock() + defer c.mu.Unlock() + if c.ffCmd != nil { + return nil + } -//type Channel struct { -// ID string `xml:"id"` -// Enabled string `xml:"enabled"` -// Codec string `xml:"audioCompressionType"` -//} + sampleFmt := "mulaw" + if codecName == core.CodecPCMA { + sampleFmt = "alaw" + } + + cmd := exec.Command( + "ffmpeg", + "-hide_banner", "-loglevel", "error", + "-f", sampleFmt, "-ar", "8000", "-ac", "1", "-i", "pipe:0", + "-c:a", "aac", "-profile:a", "aac_low", + "-ar", strconv.Itoa(int(c.sampleRate)), "-ac", "1", "-b:a", "64k", + "-f", "adts", "pipe:1", + ) + stdin, err := cmd.StdinPipe() + if err != nil { + return err + } + stdout, err := cmd.StdoutPipe() + if err != nil { + _ = stdin.Close() + return err + } + if err = cmd.Start(); err != nil { + _ = stdin.Close() + return err + } + + c.ffCmd = cmd + c.ffIn = stdin + c.ffDone = make(chan struct{}) + + go func() { + defer close(c.ffDone) + buf := make([]byte, 0, 4096) + tmp := make([]byte, 2048) + for { + n, err := stdout.Read(tmp) + if n > 0 { + buf = append(buf, tmp[:n]...) + for { + if len(buf) < aac.ADTSHeaderSize { + break + } + if !aac.IsADTS(buf) { + buf = buf[1:] + continue + } + size := int(aac.ReadADTSSize(buf)) + if size < aac.ADTSHeaderSize || size > len(buf) { + break + } + frame := append([]byte(nil), buf[:size]...) + buf = buf[size:] + c.writeADTSFrames(frame) + } + } + if err != nil { + return + } + } + }() + + return nil +} + +func (c *Client) stopFFmpeg() { + c.mu.Lock() + cmd := c.ffCmd + in := c.ffIn + done := c.ffDone + c.ffCmd = nil + c.ffIn = nil + c.ffDone = nil + c.mu.Unlock() + + if in != nil { + _ = in.Close() + } + if cmd != nil && cmd.Process != nil { + _ = cmd.Process.Kill() + _, _ = cmd.Process.Wait() + } + if done != nil { + <-done + } +} From cc5bb615b65b1a247d13b15ec8c76e7009310869 Mon Sep 17 00:00:00 2001 From: Josiah Clark Date: Wed, 15 Jul 2026 12:02:04 +1200 Subject: [PATCH 2/4] =?UTF-8?q?isapi:=20find=20Frigate=20ffmpeg=20path=20f?= =?UTF-8?q?or=20PCMU=E2=86=92AAC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Hardcoded `ffmpeg` often missing from go2rtc PATH in Frigate; probe common /usr/lib/ffmpeg/*/bin/ffmpeg locations. Co-authored-by: Cursor --- pkg/isapi/client.go | 43 ++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 42 insertions(+), 1 deletion(-) diff --git a/pkg/isapi/client.go b/pkg/isapi/client.go index 155e62449..b2ae70c58 100644 --- a/pkg/isapi/client.go +++ b/pkg/isapi/client.go @@ -9,6 +9,7 @@ import ( "net" "net/http" "net/url" + "os" "os/exec" "strconv" "sync" @@ -339,6 +340,32 @@ func (c *Client) writePCMUToAAC(codecName string, payload []byte) { _, _ = in.Write(payload) } +// findFFmpeg returns a usable ffmpeg binary. +// Frigate's go2rtc process often does NOT have `ffmpeg` on PATH; the ffmpeg +// module is configured with something like /usr/lib/ffmpeg/7.0/bin/ffmpeg. +func findFFmpeg() (string, error) { + candidates := []string{ + "ffmpeg", + "/usr/lib/ffmpeg/7.0/bin/ffmpeg", + "/usr/lib/ffmpeg/5.0/bin/ffmpeg", + "/usr/local/bin/ffmpeg", + "/usr/bin/ffmpeg", + } + for _, bin := range candidates { + path, err := exec.LookPath(bin) + if err == nil { + return path, nil + } + // Absolute paths: LookPath fails if not in PATH; check directly. + if len(bin) > 0 && bin[0] == '/' { + if st, err := os.Stat(bin); err == nil && !st.IsDir() { + return bin, nil + } + } + } + return "", errors.New("isapi: ffmpeg not found (needed for PCMU/PCMA → AAC)") +} + func (c *Client) ensureFFmpeg(codecName string) error { c.mu.Lock() defer c.mu.Unlock() @@ -346,13 +373,18 @@ func (c *Client) ensureFFmpeg(codecName string) error { return nil } + bin, err := findFFmpeg() + if err != nil { + return err + } + sampleFmt := "mulaw" if codecName == core.CodecPCMA { sampleFmt = "alaw" } cmd := exec.Command( - "ffmpeg", + bin, "-hide_banner", "-loglevel", "error", "-f", sampleFmt, "-ar", "8000", "-ac", "1", "-i", "pipe:0", "-c:a", "aac", "-profile:a", "aac_low", @@ -368,6 +400,11 @@ func (c *Client) ensureFFmpeg(codecName string) error { _ = stdin.Close() return err } + stderr, err := cmd.StderrPipe() + if err != nil { + _ = stdin.Close() + return err + } if err = cmd.Start(); err != nil { _ = stdin.Close() return err @@ -377,6 +414,10 @@ func (c *Client) ensureFFmpeg(codecName string) error { c.ffIn = stdin c.ffDone = make(chan struct{}) + go func() { + _, _ = io.Copy(io.Discard, stderr) + }() + go func() { defer close(c.ffDone) buf := make([]byte, 0, 4096) From e133542c0a170ef1fdd358c22d9a5dd1a2968c3d Mon Sep 17 00:00:00 2001 From: Josiah Clark Date: Wed, 15 Jul 2026 12:34:36 +1200 Subject: [PATCH 3/4] =?UTF-8?q?isapi:=20prefer=20Opus=E2=86=92AAC=20for=20?= =?UTF-8?q?talk=20with=20low-delay=20ffmpeg?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Match WebRTC Opus before PCMU/PCMA when the camera is AAC, wrap Opus in Ogg for ffmpeg, and use nobuffer/low_delay/flush. Keep G.711 passthrough unchanged and add [isapi] debug logs. Co-authored-by: Cursor --- internal/isapi/init.go | 7 +- pkg/isapi/client.go | 211 +++++++++++++++++++++++++++++++++++++---- pkg/isapi/log.go | 10 ++ 3 files changed, 206 insertions(+), 22 deletions(-) create mode 100644 pkg/isapi/log.go diff --git a/internal/isapi/init.go b/internal/isapi/init.go index 887a6748e..2bfc8f339 100644 --- a/internal/isapi/init.go +++ b/internal/isapi/init.go @@ -1,13 +1,16 @@ package isapi import ( + "github.com/AlexxIT/go2rtc/internal/app" "github.com/AlexxIT/go2rtc/internal/streams" "github.com/AlexxIT/go2rtc/pkg/core" - "github.com/AlexxIT/go2rtc/pkg/isapi" + pkg "github.com/AlexxIT/go2rtc/pkg/isapi" ) func Init() { + pkg.SetLogger(app.GetLogger("isapi")) + streams.HandleFunc("isapi", func(source string) (core.Producer, error) { - return isapi.Dial(source) + return pkg.Dial(source) }) } diff --git a/pkg/isapi/client.go b/pkg/isapi/client.go index b2ae70c58..8a9840b3f 100644 --- a/pkg/isapi/client.go +++ b/pkg/isapi/client.go @@ -1,6 +1,7 @@ package isapi import ( + "bufio" "encoding/binary" "encoding/hex" "encoding/json" @@ -13,12 +14,14 @@ import ( "os/exec" "strconv" "sync" + "sync/atomic" "time" "github.com/AlexxIT/go2rtc/pkg/aac" "github.com/AlexxIT/go2rtc/pkg/core" "github.com/AlexxIT/go2rtc/pkg/tcp" "github.com/pion/rtp" + "github.com/pion/webrtc/v4/pkg/media/oggwriter" ) // Deprecated: should be rewritten to core.Connection @@ -37,10 +40,15 @@ type Client struct { sender *core.Sender send int - mu sync.Mutex - ffCmd *exec.Cmd - ffIn io.WriteCloser - ffDone chan struct{} + mu sync.Mutex + ffCmd *exec.Cmd + ffIn io.WriteCloser + ffOgg *oggwriter.OggWriter + ffSrc string // matched source codec for current ffmpeg session + ffDone chan struct{} + inPkts atomic.Uint64 + inBytes atomic.Uint64 + outFrm atomic.Uint64 } func Dial(rawURL string) (*Client, error) { @@ -119,13 +127,16 @@ func (c *Client) Dial() (err error) { if c.codecName == core.CodecAAC { conf := aac.EncodeConfig(aac.TypeAACLC, c.sampleRate, 1, false) media.Codecs = []*core.Codec{ + // Prefer Opus (WebRTC mic) → ffmpeg → AAC. Avoid PCMU mush when possible. + {Name: core.CodecOpus, ClockRate: 48000, Channels: 2}, + {Name: core.CodecOpus, ClockRate: 48000}, { Name: core.CodecAAC, ClockRate: c.sampleRate, Channels: 1, FmtpLine: aac.FMTP + hex.EncodeToString(conf), }, - // WebRTC mic offers Opus/PCMU/PCMA. Match G.711 and transcode to AAC. + // Fallback: same G.711 match as stock, then transcode to AAC. {Name: core.CodecPCMU, ClockRate: 8000}, {Name: core.CodecPCMA, ClockRate: 8000}, } @@ -137,6 +148,12 @@ func (c *Client) Dial() (err error) { } c.medias = append(c.medias, media) + Log.Info(). + Str("host", c.url). + Str("cam_codec", c.codecName). + Uint32("sample_rate", c.sampleRate). + Str("channel", c.channel). + Msg("[isapi] dial two-way") return nil } @@ -192,6 +209,10 @@ func (c *Client) Open() (err error) { _, _ = c.conn.Read(buf) tcp.Close(res) + Log.Info(). + Str("session", c.sessionID). + Str("cam_codec", c.codecName). + Msg("[isapi] open audioData OK") return nil } @@ -226,10 +247,18 @@ func (c *Client) AddTrack(media *core.Media, codec *core.Codec, track *core.Rece } c.sender = core.NewSender(media, track.Codec) + src := track.Codec.Name + Log.Info(). + Str("src", src). + Uint32("src_rate", track.Codec.ClockRate). + Uint8("src_ch", track.Codec.Channels). + Str("cam", c.codecName). + Msg("[isapi] AddTrack") switch { case c.codecName == core.CodecAAC && track.Codec.Name == core.CodecAAC: c.sender.Handler = func(packet *rtp.Packet) { + c.noteIn(packet) c.writeADTSFrames(packet.Payload) } if track.Codec.IsRTP() { @@ -238,19 +267,31 @@ func (c *Client) AddTrack(media *core.Media, codec *core.Codec, track *core.Rece c.sender.Handler = aac.EncodeToADTS(codec, c.sender.Handler) } + case c.codecName == core.CodecAAC && track.Codec.Name == core.CodecOpus: + c.sender.Handler = func(packet *rtp.Packet) { + c.noteIn(packet) + c.writeOpusToAAC(packet) + } + case c.codecName == core.CodecAAC && (track.Codec.Name == core.CodecPCMU || track.Codec.Name == core.CodecPCMA): srcCodec := track.Codec.Name c.sender.Handler = func(packet *rtp.Packet) { + c.noteIn(packet) c.writePCMUToAAC(srcCodec, packet.Payload) } default: + // G.711 cam: raw bytes straight through (no ffmpeg). + Log.Info().Str("path", "raw").Str("codec", src).Msg("[isapi] G.711 passthrough") c.sender.Handler = func(packet *rtp.Packet) { if c.conn == nil { return } + c.noteIn(packet) c.send += len(packet.Payload) - _, _ = c.conn.Write(packet.Payload) + if _, err := c.conn.Write(packet.Payload); err != nil { + Log.Debug().Err(err).Msg("[isapi] G.711 write") + } } } @@ -258,11 +299,41 @@ func (c *Client) AddTrack(media *core.Media, codec *core.Codec, track *core.Rece return nil } +func (c *Client) noteIn(packet *rtp.Packet) { + n := c.inPkts.Add(1) + c.inBytes.Add(uint64(len(packet.Payload))) + if n == 1 || n%50 == 0 { + Log.Debug(). + Uint64("in_pkts", n). + Uint64("in_bytes", c.inBytes.Load()). + Uint64("out_frames", c.outFrm.Load()). + Int("sent_bytes", c.send). + Str("ff_src", c.ffSrcSnapshot()). + Msg("[isapi] talk stats") + } +} + +func (c *Client) ffSrcSnapshot() string { + c.mu.Lock() + defer c.mu.Unlock() + return c.ffSrc +} + func (c *Client) Start() (err error) { return c.Open() } func (c *Client) Stop() (err error) { + inPkts := c.inPkts.Load() + outFrm := c.outFrm.Load() + Log.Info(). + Uint64("in_pkts", inPkts). + Uint64("in_bytes", c.inBytes.Load()). + Uint64("out_frames", outFrm). + Int("sent_bytes", c.send). + Str("ff_src", c.ffSrcSnapshot()). + Msg("[isapi] stop") + c.stopFFmpeg() if c.sender != nil { @@ -315,12 +386,33 @@ func (c *Client) writeADTSFrames(b []byte) { var hdr [4]byte binary.BigEndian.PutUint32(hdr[:], uint32(size)) if _, err := c.conn.Write(hdr[:]); err != nil { + Log.Debug().Err(err).Msg("[isapi] AAC len write") return } if _, err := c.conn.Write(frame); err != nil { + Log.Debug().Err(err).Msg("[isapi] AAC frame write") return } c.send += 4 + size + c.outFrm.Add(1) + } +} + +func (c *Client) writeOpusToAAC(packet *rtp.Packet) { + if packet == nil || len(packet.Payload) == 0 { + return + } + if err := c.ensureFFmpeg(core.CodecOpus, packet); err != nil { + return + } + c.mu.Lock() + ogg := c.ffOgg + c.mu.Unlock() + if ogg == nil { + return + } + if err := ogg.WriteRTP(packet); err != nil { + Log.Debug().Err(err).Msg("[isapi] opus ogg write") } } @@ -328,7 +420,7 @@ func (c *Client) writePCMUToAAC(codecName string, payload []byte) { if len(payload) == 0 { return } - if err := c.ensureFFmpeg(codecName); err != nil { + if err := c.ensureFFmpeg(codecName, nil); err != nil { return } c.mu.Lock() @@ -337,7 +429,9 @@ func (c *Client) writePCMUToAAC(codecName string, payload []byte) { if in == nil { return } - _, _ = in.Write(payload) + if _, err := in.Write(payload); err != nil { + Log.Debug().Err(err).Msg("[isapi] g711→aac write") + } } // findFFmpeg returns a usable ffmpeg binary. @@ -363,10 +457,10 @@ func findFFmpeg() (string, error) { } } } - return "", errors.New("isapi: ffmpeg not found (needed for PCMU/PCMA → AAC)") + return "", errors.New("isapi: ffmpeg not found (needed for Opus/PCMU/PCMA → AAC)") } -func (c *Client) ensureFFmpeg(codecName string) error { +func (c *Client) ensureFFmpeg(codecName string, firstOpus *rtp.Packet) error { c.mu.Lock() defer c.mu.Unlock() if c.ffCmd != nil { @@ -375,22 +469,44 @@ func (c *Client) ensureFFmpeg(codecName string) error { bin, err := findFFmpeg() if err != nil { + Log.Error().Err(err).Str("src", codecName).Msg("[isapi] ffmpeg missing") return err } - sampleFmt := "mulaw" - if codecName == core.CodecPCMA { - sampleFmt = "alaw" + args := []string{ + "-hide_banner", "-loglevel", "warning", + "-fflags", "nobuffer", + "-flags", "low_delay", + "-probesize", "32", + "-analyzeduration", "0", + } + + switch codecName { + case core.CodecOpus: + args = append(args, + "-f", "ogg", "-i", "pipe:0", + ) + case core.CodecPCMU: + args = append(args, + "-f", "mulaw", "-ar", "8000", "-ac", "1", "-i", "pipe:0", + ) + case core.CodecPCMA: + args = append(args, + "-f", "alaw", "-ar", "8000", "-ac", "1", "-i", "pipe:0", + ) + default: + return errors.New("isapi: unsupported ffmpeg input: " + codecName) } - cmd := exec.Command( - bin, - "-hide_banner", "-loglevel", "error", - "-f", sampleFmt, "-ar", "8000", "-ac", "1", "-i", "pipe:0", + args = append(args, "-c:a", "aac", "-profile:a", "aac_low", "-ar", strconv.Itoa(int(c.sampleRate)), "-ac", "1", "-b:a", "64k", - "-f", "adts", "pipe:1", + "-f", "adts", + "-flush_packets", "1", + "pipe:1", ) + + cmd := exec.Command(bin, args...) stdin, err := cmd.StdinPipe() if err != nil { return err @@ -407,15 +523,58 @@ func (c *Client) ensureFFmpeg(codecName string) error { } if err = cmd.Start(); err != nil { _ = stdin.Close() + Log.Error().Err(err).Str("bin", bin).Msg("[isapi] ffmpeg start") return err } c.ffCmd = cmd c.ffIn = stdin + c.ffSrc = codecName c.ffDone = make(chan struct{}) + if codecName == core.CodecOpus { + _ = firstOpus + ch := uint16(2) + rate := uint32(48000) + if c.sender != nil && c.sender.Codec != nil { + if c.sender.Codec.ClockRate > 0 { + rate = c.sender.Codec.ClockRate + } + if c.sender.Codec.Channels > 0 { + ch = uint16(c.sender.Codec.Channels) + } + } + ogg, err := oggwriter.NewWith(stdin, rate, ch) + if err != nil { + _ = stdin.Close() + _ = cmd.Process.Kill() + Log.Error().Err(err).Msg("[isapi] oggwriter") + c.ffCmd = nil + c.ffIn = nil + c.ffDone = nil + return err + } + c.ffOgg = ogg + Log.Info(). + Str("bin", bin). + Str("src", "OPUS"). + Uint32("opus_rate", rate). + Uint16("opus_ch", ch). + Uint32("aac_rate", c.sampleRate). + Msg("[isapi] ffmpeg Opus→AAC started") + } else { + Log.Info(). + Str("bin", bin). + Str("src", codecName). + Uint32("aac_rate", c.sampleRate). + Msg("[isapi] ffmpeg G.711→AAC started") + } + go func() { - _, _ = io.Copy(io.Discard, stderr) + sc := bufio.NewScanner(stderr) + for sc.Scan() { + Log.Warn().Str("ffmpeg", sc.Text()).Msg("[isapi] ffmpeg stderr") + } }() go func() { @@ -444,6 +603,9 @@ func (c *Client) ensureFFmpeg(codecName string) error { } } if err != nil { + if !errors.Is(err, io.EOF) { + Log.Debug().Err(err).Msg("[isapi] ffmpeg stdout") + } return } } @@ -456,13 +618,19 @@ func (c *Client) stopFFmpeg() { c.mu.Lock() cmd := c.ffCmd in := c.ffIn + ogg := c.ffOgg done := c.ffDone + src := c.ffSrc c.ffCmd = nil c.ffIn = nil + c.ffOgg = nil c.ffDone = nil + c.ffSrc = "" c.mu.Unlock() - if in != nil { + if ogg != nil { + _ = ogg.Close() + } else if in != nil { _ = in.Close() } if cmd != nil && cmd.Process != nil { @@ -472,4 +640,7 @@ func (c *Client) stopFFmpeg() { if done != nil { <-done } + if src != "" { + Log.Debug().Str("src", src).Msg("[isapi] ffmpeg stopped") + } } diff --git a/pkg/isapi/log.go b/pkg/isapi/log.go new file mode 100644 index 000000000..320d395c0 --- /dev/null +++ b/pkg/isapi/log.go @@ -0,0 +1,10 @@ +package isapi + +import "github.com/rs/zerolog" + +// Log is set from internal/isapi.Init. Defaults to no-op so Dial stays usable in tests. +var Log = zerolog.Nop() + +func SetLogger(l zerolog.Logger) { + Log = l +} From 064454353ddb0be3b41f065258be94b914a434e8 Mon Sep 17 00:00:00 2001 From: Josiah Clark Date: Wed, 15 Jul 2026 12:55:06 +1200 Subject: [PATCH 4/4] isapi: use matched codec for talk path, not track ANY WebRTC often hands track.Codec=ANY; path selection must use the negotiated Opus/PCMU/AAC codec so AAC cams don't fall through to raw G.711 write. Co-authored-by: Cursor --- pkg/isapi/client.go | 32 +++++++++++++++++++++----------- 1 file changed, 21 insertions(+), 11 deletions(-) diff --git a/pkg/isapi/client.go b/pkg/isapi/client.go index 8a9840b3f..1f4af5c4d 100644 --- a/pkg/isapi/client.go +++ b/pkg/isapi/client.go @@ -246,39 +246,49 @@ func (c *Client) AddTrack(media *core.Media, codec *core.Codec, track *core.Rece return nil } - c.sender = core.NewSender(media, track.Codec) - src := track.Codec.Name + // Prefer matched producer codec. WebRTC often passes track.Codec=ANY; + // path must follow the negotiated codec (Opus/PCMU/…), not ANY. + pathCodec := codec + if pathCodec == nil || pathCodec.Name == "" || pathCodec.Name == core.CodecAny || pathCodec.Name == core.CodecAll { + pathCodec = track.Codec + } + c.sender = core.NewSender(media, pathCodec) + src := pathCodec.Name Log.Info(). Str("src", src). - Uint32("src_rate", track.Codec.ClockRate). - Uint8("src_ch", track.Codec.Channels). + Str("track", track.Codec.Name). + Uint32("src_rate", pathCodec.ClockRate). + Uint8("src_ch", pathCodec.Channels). Str("cam", c.codecName). Msg("[isapi] AddTrack") switch { - case c.codecName == core.CodecAAC && track.Codec.Name == core.CodecAAC: + case c.codecName == core.CodecAAC && src == core.CodecAAC: c.sender.Handler = func(packet *rtp.Packet) { c.noteIn(packet) c.writeADTSFrames(packet.Payload) } - if track.Codec.IsRTP() { - c.sender.Handler = aac.RTPToADTS(codec, c.sender.Handler) + if track.Codec.IsRTP() || pathCodec.IsRTP() { + c.sender.Handler = aac.RTPToADTS(pathCodec, c.sender.Handler) } else { - c.sender.Handler = aac.EncodeToADTS(codec, c.sender.Handler) + c.sender.Handler = aac.EncodeToADTS(pathCodec, c.sender.Handler) } + Log.Info().Str("path", "AAC→AAC").Msg("[isapi] talk path") - case c.codecName == core.CodecAAC && track.Codec.Name == core.CodecOpus: + case c.codecName == core.CodecAAC && src == core.CodecOpus: c.sender.Handler = func(packet *rtp.Packet) { c.noteIn(packet) c.writeOpusToAAC(packet) } + Log.Info().Str("path", "Opus→AAC").Msg("[isapi] talk path") - case c.codecName == core.CodecAAC && (track.Codec.Name == core.CodecPCMU || track.Codec.Name == core.CodecPCMA): - srcCodec := track.Codec.Name + case c.codecName == core.CodecAAC && (src == core.CodecPCMU || src == core.CodecPCMA): + srcCodec := src c.sender.Handler = func(packet *rtp.Packet) { c.noteIn(packet) c.writePCMUToAAC(srcCodec, packet.Payload) } + Log.Info().Str("path", "G.711→AAC").Str("codec", src).Msg("[isapi] talk path") default: // G.711 cam: raw bytes straight through (no ffmpeg).