diff --git a/internal/onvif/onvif.go b/internal/onvif/onvif.go index c305b706e..e587bee8b 100644 --- a/internal/onvif/onvif.go +++ b/internal/onvif/onvif.go @@ -29,6 +29,36 @@ func Init() { // ONVIF client autodiscovery api.HandleFunc("api/onvif", apiOnvif) + + var enableWSD bool + if api.Port > 0 { + var cfg struct { + Mod struct { + MaxConcurrentProbes int `yaml:"maxConcurrentProbes"` + SendHelloMessage bool `yaml:"sendHelloMessage"` + DisableWSDiscovery bool `yaml:"disableWSDiscovery"` + } `yaml:"onvif"` + } + cfg.Mod.MaxConcurrentProbes = maxConcurrentProbes + cfg.Mod.SendHelloMessage = sendHelloMessage + + // load config from YAML + app.LoadConfig(&cfg) + maxConcurrentProbes = cfg.Mod.MaxConcurrentProbes + sendHelloMessage = cfg.Mod.SendHelloMessage + + servicePort = strconv.Itoa(api.Port) + + if !cfg.Mod.DisableWSDiscovery && maxConcurrentProbes > 0 { + enableWSD = true + } + } + + if enableWSD { + go HandleProbe() + } else { + log.Info().Msg("WS-Discovery disabled") + } } var log zerolog.Logger diff --git a/internal/onvif/ws-discovery.go b/internal/onvif/ws-discovery.go new file mode 100644 index 000000000..880270eb7 --- /dev/null +++ b/internal/onvif/ws-discovery.go @@ -0,0 +1,316 @@ +package onvif + +import ( + "bytes" + "encoding/xml" + "net" + "net/netip" + "regexp" + "runtime/debug" + "strings" + "sync" + "time" + + "github.com/google/uuid" +) + +const ( + multicastAddr = "239.255.255.250:3702" // Standard WS-Discovery multicast address and port +) + +var ( + metadataVer = "1" // Metadata version + maxConcurrentProbes = 32 + sendHelloMessage = true + + serviceUUID = uuid.Must(uuid.NewRandom()).URN() // Unique service identifier, generated on each startup + + servicePort string +) + +func HandleProbe() { + // 1. Resolve the UDP multicast address. + addr, err := net.ResolveUDPAddr("udp", multicastAddr) + if err != nil { + log.Error().Err(err).Msg("failed to resolve multicast address") + return + } + + // 2. Start listening for multicast messages. + conn, err := net.ListenMulticastUDP("udp", nil, addr) + if err != nil { + log.Error().Err(err).Msg("failed to listen on multicast address") + return + } + defer conn.Close() + + err = conn.SetReadBuffer(2 << 20) + if err != nil { + log.Error().Err(err).Msg("failed to set read buffer") + } + + log.Info().Str("addr", multicastAddr).Msg("started listening for multicast probes") + + // 3. Process incoming messages. + pool := &sync.Pool{ + New: func() any { + return make([]byte, 2048) + }, + } + workers := make(chan struct{}, maxConcurrentProbes) + + if sendHelloMessage { + go sendHello() + } + for { + buffer := pool.Get().([]byte) + n, remoteAddr, err := conn.ReadFromUDPAddrPort(buffer) + if err != nil { + log.Error().Err(err).Msg("failed to read packet") + continue + } + + // Handle each message concurrently. + select { + case workers <- struct{}{}: + go func() { + defer func() { + <-workers + if r := recover(); r != nil { + log.Error(). + Interface("panic", r). + Bytes("stack", debug.Stack()). + Msg("panic while handling Probe request") + } + }() + + handleProbe(conn, remoteAddr, buffer, n, pool) + }() + default: + pool.Put(buffer) + log.Warn().Str("remoteAddr", remoteAddr.String()).Msg("too many pending Probe requests, dropping packet") + } + } +} + +func sendHello() { + addr, err := net.ResolveUDPAddr("udp4", multicastAddr) + if err != nil { + log.Error().Err(err).Msg("failed to resolve multicast address for Hello") + return + } + + conn, err := net.DialUDP("udp4", nil, addr) + if err != nil { + log.Error().Err(err).Msg("failed to dial multicast address for Hello") + return + } + defer conn.Close() + + localAddr, ok := conn.LocalAddr().(*net.UDPAddr) + if !ok || localAddr.IP == nil || localAddr.IP.IsUnspecified() { + log.Error().Msg("failed to resolve local address for Hello") + return + } + + msg := buildHelloMessage(localAddr.IP.String()) + if _, err = conn.Write([]byte(msg)); err != nil { + log.Error().Err(err).Msg("failed to send Hello") + return + } + + log.Info().Str("addr", localAddr.IP.String()).Msg("sent Hello") +} + +func buildHelloMessage(localIP string) string { + return ` + + + ` + uuid.New().URN() + ` + urn:schemas-xmlsoap-org:ws:2005:04:discovery + http://schemas.xmlsoap.org/ws/2005/04/discovery/Hello + + + + + ` + serviceUUID + ` + dn:NetworkVideoTransmitter + onvif://www.onvif.org/type/NetworkVideoTransmitter + http://` + net.JoinHostPort(localIP, servicePort) + `/onvif/device_service + ` + metadataVer + ` + + +` +} + +// handleProbe handles an incoming Probe request. +func handleProbe(conn *net.UDPConn, remoteAddr netip.AddrPort, buffer []byte, n int, pool *sync.Pool) { + defer pool.Put(buffer) + + data := buffer[:n] + + // 1. Parse XML and extract the Probe message ID. + messageID, probeTypes, ok := extractMessageIDAndTypes(data) + if !ok { + // Ignore invalid Probe messages. + return + } + + log.Info(). + Str("messageID", messageID). + Str("remoteAddr", remoteAddr.String()). + Str("probe", probeTypes).Msg("received Probe request") + + switch { + case strings.HasSuffix(probeTypes, ":Device"): + fallthrough + case strings.EqualFold(probeTypes, "Device"): + probeTypes = "Device" + case strings.HasSuffix(probeTypes, ":NetworkVideoTransmitter"): + fallthrough + case strings.EqualFold(probeTypes, "NetworkVideoTransmitter"): + probeTypes = "NetworkVideoTransmitter" + default: + log.Warn(). + Str("messageID", messageID). + Str("remoteAddr", remoteAddr.String()). + Str("probe", probeTypes). + Msg("unsupported Probe type, ignoring request") + return + } + + localIP, err := getOutboundIP(remoteAddr) + if err != nil { + log.Err(err). + Msg("failed to resolve outbound IP") + return + } + + log.Info().Str("addr", localIP.String()).Msg("Local Addr") + + // 2. Build the ProbeMatch response. + responseXML := buildProbeMatchResponse(messageID, probeTypes, localIP.String()) + + // 3. Unicast the response back to the client. + _, err = conn.WriteToUDPAddrPort([]byte(responseXML), remoteAddr) + if err != nil { + log.Error().Err(err).Msg("failed to send ProbeMatch") + } else { + log.Info().Str("messageID", messageID).Str("remoteAddr", remoteAddr.String()).Msg("sent ProbeMatch response") + } +} + +type Envelope struct { + Header struct { + MessageID string `xml:"http://schemas.xmlsoap.org/ws/2004/08/addressing MessageID"` + } `xml:"http://www.w3.org/2003/05/soap-envelope Header"` + Body struct { + Probe struct { + Types string `xml:"http://schemas.xmlsoap.org/ws/2005/04/discovery Types"` + Scopes string `xml:"http://schemas.xmlsoap.org/ws/2005/04/discovery Scopes"` + } `xml:"http://schemas.xmlsoap.org/ws/2005/04/discovery Probe"` + } `xml:"http://www.w3.org/2003/05/soap-envelope Body"` +} + +var ( + messageIDRegex = regexp.MustCompile(`<[^\s>]*MessageID[^<>]*>(.*)]*MessageID>`) +) + +// extractMessageIDAndTypes extracts the MessageID and Types from a Probe request. +func extractMessageIDAndTypes(xmlData []byte) (msgID, probeTypes string, ok bool) { + var envelope Envelope + + err := xml.Unmarshal(xmlData, &envelope) + if err != nil { + log.Err(err).Msg("xml.Unmarshal failed") + matches := messageIDRegex.FindStringSubmatch(string(xmlData)) + if len(matches) > 1 { + msgID = matches[1] + if bytes.Contains(xmlData, []byte("NetworkVideoTransmitter")) { + probeTypes = "NetworkVideoTransmitter" + } else { + probeTypes = "Device" + } + ok = true + return + } + + log.Error().Err(err).Bytes("xml", xmlData).Msg("failed to parse XML") + return "", "", false + } + + return envelope.Header.MessageID, envelope.Body.Probe.Types, true +} + +// buildProbeMatchResponse builds a ProbeMatch response that follows the WS-Discovery specification. +func buildProbeMatchResponse(relatesToID, probeType, localIP string) string { + if probeType == "Device" { + return ` + + + ` + serviceUUID + ` + ` + relatesToID + ` + http://schemas.xmlsoap.org/ws/2004/08/addressing/role/anonymous + http://schemas.xmlsoap.org/ws/2005/04/discovery/ProbeMatches + + + + + ` + serviceUUID + ` + tds:Device + onvif://www.onvif.org/type/Device + http://` + net.JoinHostPort(localIP, servicePort) + `/onvif/device_service + ` + metadataVer + ` + + + +` + } + + return ` + + + ` + serviceUUID + ` + ` + relatesToID + ` + http://schemas.xmlsoap.org/ws/2004/08/addressing/role/anonymous + http://schemas.xmlsoap.org/ws/2005/04/discovery/ProbeMatches + + + + + + ` + serviceUUID + ` + dn:NetworkVideoTransmitter + onvif://www.onvif.org/type/NetworkVideoTransmitter + http://` + net.JoinHostPort(localIP, servicePort) + `/onvif/device_service + ` + metadataVer + ` + + + +` +} + +func getOutboundIP(remote netip.AddrPort) (net.IP, error) { + // Create a temporary connection to determine the route. + d := net.Dialer{LocalAddr: nil, Timeout: 5 * time.Second} + conn, err := d.Dial("udp", remote.String()) + if err != nil { + return nil, err + } + defer conn.Close() + + localAddr := conn.LocalAddr().(*net.UDPAddr) + return localAddr.IP, nil +}