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/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/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..1f4af5c4d 100644 --- a/pkg/isapi/client.go +++ b/pkg/isapi/client.go @@ -1,14 +1,27 @@ package isapi import ( + "bufio" + "encoding/binary" + "encoding/hex" + "encoding/json" "errors" "io" "net" "net/http" "net/url" - + "os" + "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 @@ -19,13 +32,26 @@ 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 + 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) { - // check if url is valid url u, err := url.Parse(rawURL) if err != nil { return nil, err @@ -59,6 +85,7 @@ func (c *Client) Dial() (err error) { } b, err := io.ReadAll(res.Body) + tcp.Close(res) if err != nil { return err } @@ -68,11 +95,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 +122,47 @@ 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{ + // 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), + }, + // Fallback: same G.711 match as stock, then 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) + 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 } 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 +175,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 +205,14 @@ 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) - + Log.Info(). + Str("session", c.sessionID). + Str("cam_codec", c.codecName). + Msg("[isapi] open audioData OK") return nil } @@ -149,16 +229,428 @@ 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 + } + + // 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). + 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 && src == core.CodecAAC: + c.sender.Handler = func(packet *rtp.Packet) { + c.noteIn(packet) + c.writeADTSFrames(packet.Payload) + } + if track.Codec.IsRTP() || pathCodec.IsRTP() { + c.sender.Handler = aac.RTPToADTS(pathCodec, c.sender.Handler) + } else { + c.sender.Handler = aac.EncodeToADTS(pathCodec, c.sender.Handler) + } + Log.Info().Str("path", "AAC→AAC").Msg("[isapi] talk path") + + 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 && (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). + 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) + if _, err := c.conn.Write(packet.Payload); err != nil { + Log.Debug().Err(err).Msg("[isapi] G.711 write") + } + } + } + + c.sender.HandleRTP(track) + 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 { + c.sender.Close() + } + + if c.conn != nil { + _ = c.Close() + return c.conn.Close() + } return nil } -//type XMLChannels struct { -// Channels []Channel `xml:"TwoWayAudioChannel"` -//} +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 { + 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") + } +} + +func (c *Client) writePCMUToAAC(codecName string, payload []byte) { + if len(payload) == 0 { + return + } + if err := c.ensureFFmpeg(codecName, nil); err != nil { + return + } + c.mu.Lock() + in := c.ffIn + c.mu.Unlock() + if in == nil { + return + } + if _, err := in.Write(payload); err != nil { + Log.Debug().Err(err).Msg("[isapi] g711→aac write") + } +} + +// 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 Opus/PCMU/PCMA → AAC)") +} + +func (c *Client) ensureFFmpeg(codecName string, firstOpus *rtp.Packet) error { + c.mu.Lock() + defer c.mu.Unlock() + if c.ffCmd != nil { + return nil + } + + bin, err := findFFmpeg() + if err != nil { + Log.Error().Err(err).Str("src", codecName).Msg("[isapi] ffmpeg missing") + return err + } + + 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) + } + + args = append(args, + "-c:a", "aac", "-profile:a", "aac_low", + "-ar", strconv.Itoa(int(c.sampleRate)), "-ac", "1", "-b:a", "64k", + "-f", "adts", + "-flush_packets", "1", + "pipe:1", + ) + + cmd := exec.Command(bin, args...) + stdin, err := cmd.StdinPipe() + if err != nil { + return err + } + stdout, err := cmd.StdoutPipe() + if err != nil { + _ = stdin.Close() + return err + } + stderr, err := cmd.StderrPipe() + if err != nil { + _ = stdin.Close() + return err + } + 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() { + sc := bufio.NewScanner(stderr) + for sc.Scan() { + Log.Warn().Str("ffmpeg", sc.Text()).Msg("[isapi] ffmpeg stderr") + } + }() + + 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 { + if !errors.Is(err, io.EOF) { + Log.Debug().Err(err).Msg("[isapi] ffmpeg stdout") + } + return + } + } + }() + + return nil +} -//type Channel struct { -// ID string `xml:"id"` -// Enabled string `xml:"enabled"` -// Codec string `xml:"audioCompressionType"` -//} +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 ogg != nil { + _ = ogg.Close() + } else if in != nil { + _ = in.Close() + } + if cmd != nil && cmd.Process != nil { + _ = cmd.Process.Kill() + _, _ = cmd.Process.Wait() + } + 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 +}