From c836cef28dc6234ed21cc9ec4328de0a2dd4ef0e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?C=C3=A9dric=20Verstraeten?= Date: Wed, 15 Jul 2026 23:17:23 +0200 Subject: [PATCH] Add advanced discovery and stream verification Introduces a new ONVIF/network discovery pipeline that combines WS-Discovery, subnet-aware host/port scanning, banner fingerprinting, and MAC vendor enrichment to identify likely cameras. Adds API and CLI support for discovery options (`/api/camera/discover`, `-subnet`), plus a richer discovered-device response model. Also adds MQTT `verify-stream` handling to probe RTSP streams and return codec/resolution/fps, and persists detected stream FPS into config for main/sub streams. --- machinery/main.go | 11 +- machinery/src/capture/gortsplib.go | 17 + machinery/src/models/api_response.go | 21 ++ machinery/src/models/mqtt.go | 9 + machinery/src/onvif/discovery.go | 496 ++++++++++++++++++++++++++ machinery/src/onvif/fingerprint.go | 267 ++++++++++++++ machinery/src/onvif/main.go | 71 +++- machinery/src/routers/http/methods.go | 36 ++ machinery/src/routers/http/routes.go | 1 + machinery/src/routers/mqtt/main.go | 97 +++++ 10 files changed, 1012 insertions(+), 14 deletions(-) create mode 100644 machinery/src/onvif/discovery.go create mode 100644 machinery/src/onvif/fingerprint.go diff --git a/machinery/main.go b/machinery/main.go index b82b1b9..643332f 100644 --- a/machinery/main.go +++ b/machinery/main.go @@ -5,6 +5,7 @@ import ( "flag" "fmt" "os" + "strings" "time" "github.com/kerberos-io/agent/machinery/src/capture" @@ -76,12 +77,14 @@ func main() { var name string var port string var timeout string + var subnet string flag.StringVar(&action, "action", "version", "Tell us what you want do 'run' or 'version'") flag.StringVar(&configDirectory, "config", ".", "Where is the configuration stored") flag.StringVar(&name, "name", "agent", "Provide a name for the agent") flag.StringVar(&port, "port", "80", "On which port should the agent run") flag.StringVar(&timeout, "timeout", "2000", "Number of milliseconds to wait for the ONVIF discovery to complete") + flag.StringVar(&subnet, "subnet", "", "Optional subnet(s) to scan for discovery, e.g. '192.168.1.0/24' (comma-separated). Defaults to the local interfaces.") flag.Parse() // Specify the level of loggin: "info", "warning", "debug", "error" or "fatal." @@ -112,7 +115,13 @@ func main() { log.Log.Fatal("main.Main(): could not parse timeout: " + err.Error()) return } - onvif.Discover(timeout) + var subnets []string + for _, part := range strings.Split(subnet, ",") { + if trimmed := strings.TrimSpace(part); trimmed != "" { + subnets = append(subnets, trimmed) + } + } + onvif.Discover(timeout, subnets...) } case "decrypt": { diff --git a/machinery/src/capture/gortsplib.go b/machinery/src/capture/gortsplib.go index d41e954..ac57838 100644 --- a/machinery/src/capture/gortsplib.go +++ b/machinery/src/capture/gortsplib.go @@ -864,6 +864,7 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets // Get FPS using enhanced method fps := g.getEnhancedFPS(&sps, g.VideoH264Index) g.Streams[g.VideoH264Index].FPS = fps + g.persistStreamFPS(configuration, streamType, fps) log.Log.Debug(fmt.Sprintf("capture.golibrtsp.Start(%s): Final FPS=%.2f", streamType, fps)) g.VideoH264Forma.SPS = nalu if streamType == "main" && len(nalu) > 0 { @@ -1061,6 +1062,7 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets } if ptsFPS := ft.update(pts); ptsFPS > 0 && ptsFPS <= 120 { g.Streams[g.VideoH265Index].FPS = ptsFPS + g.persistStreamFPS(configuration, streamType, ptsFPS) } } @@ -1537,6 +1539,21 @@ func (g *Golibrtsp) initFPSCalculation() { } // Get enhanced FPS information from SPS with fallback to PTS-based calculation. +// persistStreamFPS stores the computed frame rate into the shared config so it +// is reported to the hub/UI (mirrors how width/height are persisted). The value +// is rounded to 2 decimals with trailing zeros trimmed (e.g. "25", "29.97"). +func (g *Golibrtsp) persistStreamFPS(configuration *models.Configuration, streamType string, fps float64) { + if fps <= 0 { + return + } + fpsStr := strconv.FormatFloat(float64(int(fps*100+0.5))/100, 'f', -1, 64) + if streamType == "main" { + configuration.Config.Capture.IPCamera.FPS = fpsStr + } else if streamType == "sub" { + configuration.Config.Capture.IPCamera.SubFPS = fpsStr + } +} + // The PTS-based FPS is computed per completed frame via fpsTracker.update(), // so by the time this is called we already have a good estimate. func (g *Golibrtsp) getEnhancedFPS(sps *h264.SPS, streamIndex int8) float64 { diff --git a/machinery/src/models/api_response.go b/machinery/src/models/api_response.go index 1b4fab2..885364b 100644 --- a/machinery/src/models/api_response.go +++ b/machinery/src/models/api_response.go @@ -19,6 +19,27 @@ type CameraStreams struct { SubRTSP string `json:"sub_rtsp"` } +// DiscoveredDevice describes a device found on the local network during a +// discovery scan (fing/wifiman-style). It combines ONVIF WS-Discovery results +// with an active port scan and MAC/vendor lookup so cameras can be +// auto-detected and pre-filled in the configuration UI. +type DiscoveredDevice struct { + IP string `json:"ip" bson:"ip"` + Hostname string `json:"hostname,omitempty" bson:"hostname"` + MAC string `json:"mac,omitempty" bson:"mac"` + Vendor string `json:"vendor,omitempty" bson:"vendor"` + Manufacturer string `json:"manufacturer,omitempty" bson:"manufacturer"` + Model string `json:"model,omitempty" bson:"model"` + Type string `json:"type,omitempty" bson:"type"` + Server string `json:"server,omitempty" bson:"server"` + OpenPorts []int `json:"open_ports,omitempty" bson:"open_ports"` + Services []string `json:"services,omitempty" bson:"services"` + ONVIF bool `json:"onvif" bson:"onvif"` + ONVIFXAddr string `json:"onvif_xaddr,omitempty" bson:"onvif_xaddr"` + RTSPURL string `json:"rtsp_url,omitempty" bson:"rtsp_url"` + IsCamera bool `json:"is_camera" bson:"is_camera"` +} + type OnvifPanTilt struct { OnvifCredentials OnvifCredentials `json:"onvif_credentials,omitempty" bson:"onvif_credentials"` Pan float64 `json:"pan,omitempty" bson:"pan"` diff --git a/machinery/src/models/mqtt.go b/machinery/src/models/mqtt.go index f64710e..d1802a3 100644 --- a/machinery/src/models/mqtt.go +++ b/machinery/src/models/mqtt.go @@ -175,6 +175,15 @@ type RequestConfigPayload struct { Timestamp int64 `json:"timestamp"` // timestamp of the preset request. } +// We received a verify-stream request: probe the given (or configured) RTSP +// stream and report whether it can be connected/decoded, along with the +// discovered codec/resolution/fps. Responds with action "verify-stream-result". +type VerifyStreamPayload struct { + Timestamp int64 `json:"timestamp"` // timestamp of the verify request. + Stream string `json:"stream"` // "main" or "sub". + RTSP string `json:"rtsp"` // optional RTSP url to verify; falls back to the configured one. +} + // We received a update config request, we'll update the current config and send a confirmation back. type UpdateConfigPayload struct { Timestamp int64 `json:"timestamp"` // timestamp of the preset request. diff --git a/machinery/src/onvif/discovery.go b/machinery/src/onvif/discovery.go new file mode 100644 index 0000000..c3ccfa0 --- /dev/null +++ b/machinery/src/onvif/discovery.go @@ -0,0 +1,496 @@ +package onvif + +import ( + "bufio" + "context" + "net" + "net/url" + "os" + "sort" + "strconv" + "strings" + "sync" + "time" + + onvifc "github.com/cedricve/go-onvif" + "github.com/kerberos-io/agent/machinery/src/log" + "github.com/kerberos-io/agent/machinery/src/models" +) + +// scanPort describes a TCP port we probe while scanning the local network, +// together with a human readable service name. +type scanPort struct { + Port int + Service string + // rtsp marks RTSP ports we can fingerprint via an OPTIONS request. + rtsp bool + // http marks HTTP ports we can fingerprint via a banner grab. + http bool + // camera marks ports that strongly hint the device is an IP camera or NVR + // (RTSP, dedicated ONVIF ports and well-known DVR/NVR control ports). + camera bool +} + +// commonCameraPorts is the list of TCP ports we probe on every host. These are +// the ports most commonly exposed by IP cameras (RTSP, HTTP(S) and ONVIF). +var commonCameraPorts = []scanPort{ + {Port: 554, Service: "RTSP", rtsp: true, camera: true}, + {Port: 8554, Service: "RTSP (alt)", rtsp: true, camera: true}, + {Port: 80, Service: "HTTP", http: true}, + {Port: 8080, Service: "HTTP (alt)", http: true}, + {Port: 8000, Service: "ONVIF", http: true, camera: true}, + {Port: 8899, Service: "ONVIF (alt)", camera: true}, + {Port: 443, Service: "HTTPS"}, + {Port: 37777, Service: "Dahua", camera: true}, + {Port: 34567, Service: "XMeye/Sofia", camera: true}, +} + +// ouiVendors maps the first three octets (OUI) of a MAC address, upper-cased and +// without separators, to a known camera/NVR vendor. This lets us flag likely +// cameras the same way tools such as Fing or WiFiman do, even when a device does +// not answer to ONVIF WS-Discovery. +var ouiVendors = map[string]string{ + "BCAD01": "Hikvision", "C056E3": "Hikvision", "4CBD8F": "Hikvision", + "44A642": "Hikvision", "E0509B": "Hikvision", "ACB927": "Hikvision", + "18800C": "Hikvision", "C40BCB": "Hikvision", + "3CEF8C": "Dahua", "90020A": "Dahua", "E0509B00": "Dahua", + "08ED02": "Dahua", "3CE376": "Dahua", "38AF29": "Dahua", "E45D51": "Dahua", + "00408C": "Axis", "AABBCC": "Axis", "B8A44F": "Axis", "ACCC8E": "Axis", + "E82725": "Bosch", "000CAB": "Bosch", + "001B9E": "Hanwha", "0009D2": "Hanwha", "E44CC7": "Hanwha", + "EC7196": "Reolink", "9CA3BA": "Reolink", + "3C33F1": "Amcrest", "9C8ECD": "Amcrest", + "000FFC": "Vivotek", "0002D1": "Vivotek", + "001C27": "Mobotix", "0003C5": "Mobotix", + "00126A": "Ubiquiti", "FCECDA": "Ubiquiti", "744401": "Ubiquiti", + "F0234B": "Foscam", "00626E": "Foscam", + "C09424": "TP-Link", "50C7BF": "TP-Link", +} + +// DiscoverDevices performs an advanced, Fing/WiFiman-style scan of the local +// network. It combines: +// +// 1. ONVIF WS-Discovery (multicast probe), and +// 2. an active TCP port scan of every host on the local IPv4 subnets for the +// ports typically exposed by IP cameras, and +// 3. MAC address + vendor (OUI) resolution from the local ARP table, and +// 4. best-effort reverse-DNS hostname lookup. +// +// The results are merged per IP address so a single device is reported once +// with all the information we could gather. Devices are flagged as cameras when +// they answer to ONVIF, expose an RTSP port, or have a MAC that belongs to a +// known camera vendor. +// +// Optional subnets (CIDR notation, e.g. "192.168.1.0/24") override the +// automatically detected local subnets. This is useful when the agent runs in a +// container/devcontainer whose interfaces are not on the same range as the +// cameras, but the target range is still routable from the host network. +func DiscoverDevices(timeout time.Duration, subnets ...string) []models.DiscoveredDevice { + devicesByIP := make(map[string]*models.DiscoveredDevice) + var mutex sync.Mutex + + // upsert returns the (possibly newly created) device entry for an IP in a + // concurrency-safe way. + upsert := func(ip string) *models.DiscoveredDevice { + mutex.Lock() + defer mutex.Unlock() + device, ok := devicesByIP[ip] + if !ok { + device = &models.DiscoveredDevice{IP: ip} + devicesByIP[ip] = device + } + return device + } + + // 1) ONVIF WS-Discovery. This is quick and reliable for ONVIF cameras. + onvifDevices, err := onvifc.StartDiscovery(timeout) + if err != nil { + log.Log.Error("onvif.DiscoverDevices(): WS-Discovery failed: " + err.Error()) + } else { + for _, onvifDevice := range onvifDevices { + ip := hostFromXAddr(onvifDevice.XAddr) + if ip == "" { + continue + } + device := upsert(ip) + device.ONVIF = true + device.ONVIFXAddr = onvifDevice.XAddr + device.IsCamera = true + if hostname, hostErr := onvifDevice.GetHostname(); hostErr == nil && hostname.Name != "" { + device.Hostname = hostname.Name + } + } + } + + // 2) Active port scan across the requested (or auto-detected) IPv4 subnets. + var targets []string + if len(subnets) > 0 { + targets = targetsFromSubnets(subnets) + } else { + targets = localScanTargets() + } + log.Log.Info("onvif.DiscoverDevices(): scanning " + strconv.Itoa(len(targets)) + " hosts on the local network(s)") + + // Bound the amount of concurrent dials so we do not exhaust file + // descriptors on constrained devices (e.g. Raspberry Pi). + semaphore := make(chan struct{}, 128) + dialTimeout := perHostTimeout(timeout) + var waitGroup sync.WaitGroup + + for _, ip := range targets { + waitGroup.Add(1) + semaphore <- struct{}{} + go func(ip string) { + defer waitGroup.Done() + defer func() { <-semaphore }() + + openPorts, services, isCamera := scanHost(ip, dialTimeout) + if len(openPorts) == 0 { + return + } + + // Fingerprint the host (RTSP/HTTP banner grab) to determine its + // manufacturer, model and type without any credentials. + fingerprint := fingerprintHost(ip, openPorts, dialTimeout) + + device := upsert(ip) + mutex.Lock() + device.OpenPorts = mergeSortedInts(device.OpenPorts, openPorts) + device.Services = mergeUniqueStrings(device.Services, services) + if isCamera || fingerprint.IsCamera { + device.IsCamera = true + } + if fingerprint.Manufacturer != "" { + device.Manufacturer = fingerprint.Manufacturer + } + if fingerprint.Model != "" { + device.Model = fingerprint.Model + } + if fingerprint.Type != "" { + device.Type = fingerprint.Type + } + if fingerprint.Server != "" { + device.Server = fingerprint.Server + } + for _, port := range openPorts { + if port == 554 || port == 8554 { + device.RTSPURL = "rtsp://" + ip + ":" + strconv.Itoa(port) + "/" + break + } + } + mutex.Unlock() + }(ip) + } + waitGroup.Wait() + + // 3) Enrich with MAC address / vendor from the ARP table and hostnames. + arpTable := readARPTable() + results := make([]models.DiscoveredDevice, 0, len(devicesByIP)) + for ip, device := range devicesByIP { + if mac, ok := arpTable[ip]; ok { + device.MAC = mac + if vendor := vendorFromMAC(mac); vendor != "" { + device.Vendor = vendor + device.IsCamera = true + } + } + // Fall back to the MAC vendor for the manufacturer, and make sure a + // camera always carries a device type. + if device.Manufacturer == "" && device.Vendor != "" { + device.Manufacturer = device.Vendor + } + if device.IsCamera && device.Type == "" { + device.Type = "IP Camera" + } + if device.Hostname == "" { + device.Hostname = reverseDNS(ip, dialTimeout) + } + results = append(results, *device) + } + + // Cameras first, then by IP, for a stable and useful ordering. + sort.Slice(results, func(i, j int) bool { + if results[i].IsCamera != results[j].IsCamera { + return results[i].IsCamera + } + return ipLess(results[i].IP, results[j].IP) + }) + + return results +} + +// scanHost probes the common camera ports on a single host and reports the open +// ports, their service names, and whether the host looks like a camera. +func scanHost(ip string, dialTimeout time.Duration) (openPorts []int, services []string, isCamera bool) { + for _, candidate := range commonCameraPorts { + address := net.JoinHostPort(ip, strconv.Itoa(candidate.Port)) + conn, err := net.DialTimeout("tcp", address, dialTimeout) + if err != nil { + continue + } + conn.Close() + openPorts = append(openPorts, candidate.Port) + services = append(services, candidate.Service) + if candidate.camera { + isCamera = true + } + } + return openPorts, services, isCamera +} + +// targetsFromSubnets expands one or more explicit CIDR ranges (e.g. +// "192.168.1.0/24") into a de-duplicated list of host addresses. Invalid or +// oversized ranges (mask < /22) are skipped so scans stay bounded. +func targetsFromSubnets(subnets []string) []string { + seen := make(map[string]struct{}) + var targets []string + + for _, subnet := range subnets { + subnet = strings.TrimSpace(subnet) + if subnet == "" { + continue + } + // Allow passing a bare host address (e.g. "192.168.1.50") too. + if !strings.Contains(subnet, "/") { + if net.ParseIP(subnet).To4() != nil { + if _, exists := seen[subnet]; !exists { + seen[subnet] = struct{}{} + targets = append(targets, subnet) + } + } else { + log.Log.Error("onvif.targetsFromSubnets(): invalid address '" + subnet + "'") + } + continue + } + _, ipNet, err := net.ParseCIDR(subnet) + if err != nil || ipNet.IP.To4() == nil { + log.Log.Error("onvif.targetsFromSubnets(): invalid CIDR '" + subnet + "'") + continue + } + if ones, bits := ipNet.Mask.Size(); bits != 32 || ones < 22 { + log.Log.Error("onvif.targetsFromSubnets(): range '" + subnet + "' is too large to scan (use /22 or smaller)") + continue + } + for _, host := range hostsInNetwork(ipNet) { + if _, exists := seen[host]; exists { + continue + } + seen[host] = struct{}{} + targets = append(targets, host) + } + } + return targets +} + +// localScanTargets enumerates every usable IPv4 host address on the local +// network interfaces. To keep scans bounded we only expand subnets with a mask +// of /22 or smaller (at most ~1022 hosts per interface). +func localScanTargets() []string { + seen := make(map[string]struct{}) + var targets []string + + interfaces, err := net.Interfaces() + if err != nil { + log.Log.Error("onvif.localScanTargets(): " + err.Error()) + return targets + } + + for _, iface := range interfaces { + if iface.Flags&net.FlagUp == 0 || iface.Flags&net.FlagLoopback != 0 { + continue + } + addrs, addrErr := iface.Addrs() + if addrErr != nil { + continue + } + for _, addr := range addrs { + ipNet, ok := addr.(*net.IPNet) + if !ok || ipNet.IP.To4() == nil { + continue + } + ones, bits := ipNet.Mask.Size() + if bits != 32 || ones < 22 { + // Skip huge or non-IPv4 ranges to avoid endless scans. + continue + } + for _, host := range hostsInNetwork(ipNet) { + if _, exists := seen[host]; exists { + continue + } + seen[host] = struct{}{} + targets = append(targets, host) + } + } + } + return targets +} + +// hostsInNetwork returns all assignable host addresses in the given network, +// excluding the network and broadcast addresses. +func hostsInNetwork(ipNet *net.IPNet) []string { + var hosts []string + network := ipNet.IP.Mask(ipNet.Mask).To4() + if network == nil { + return hosts + } + + for ip := cloneIP(network); ipNet.Contains(ip); incrementIP(ip) { + hosts = append(hosts, ip.String()) + } + // Drop network + broadcast addresses when present. + if len(hosts) > 2 { + hosts = hosts[1 : len(hosts)-1] + } + return hosts +} + +func cloneIP(ip net.IP) net.IP { + dup := make(net.IP, len(ip)) + copy(dup, ip) + return dup +} + +func incrementIP(ip net.IP) { + for i := len(ip) - 1; i >= 0; i-- { + ip[i]++ + if ip[i] != 0 { + break + } + } +} + +// hostFromXAddr extracts the host (IP) part from an ONVIF XAddr URL such as +// "http://192.168.1.69:8000/onvif/device_service". +func hostFromXAddr(xaddr string) string { + parsed, err := url.Parse(xaddr) + if err != nil { + return "" + } + host := parsed.Hostname() + if host == "" { + // Fall back to a naive split for values without a scheme. + host = strings.TrimPrefix(xaddr, "//") + if idx := strings.IndexAny(host, ":/"); idx >= 0 { + host = host[:idx] + } + } + return host +} + +// readARPTable parses /proc/net/arp (Linux) and returns a map of IP -> MAC. On +// non-Linux platforms or when the file is unavailable it returns an empty map. +func readARPTable() map[string]string { + table := make(map[string]string) + file, err := os.Open("/proc/net/arp") + if err != nil { + return table + } + defer file.Close() + + scanner := bufio.NewScanner(file) + // Skip the header line. + if scanner.Scan() { + _ = scanner.Text() + } + for scanner.Scan() { + fields := strings.Fields(scanner.Text()) + if len(fields) < 4 { + continue + } + ip := fields[0] + mac := fields[3] + if mac == "00:00:00:00:00:00" || mac == "" { + continue + } + table[ip] = strings.ToLower(mac) + } + return table +} + +// vendorFromMAC resolves a MAC address to a known camera vendor using its OUI. +func vendorFromMAC(mac string) string { + normalized := strings.ToUpper(strings.NewReplacer(":", "", "-", "", ".", "").Replace(mac)) + if len(normalized) < 6 { + return "" + } + // Try a longer prefix first (some vendors share the first 3 octets). + if len(normalized) >= 8 { + if vendor, ok := ouiVendors[normalized[:8]]; ok { + return vendor + } + } + if vendor, ok := ouiVendors[normalized[:6]]; ok { + return vendor + } + return "" +} + +// reverseDNS performs a best-effort, time-bounded reverse DNS lookup. +func reverseDNS(ip string, timeout time.Duration) string { + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + var resolver net.Resolver + names, err := resolver.LookupAddr(ctx, ip) + if err != nil || len(names) == 0 { + return "" + } + return strings.TrimSuffix(names[0], ".") +} + +// perHostTimeout derives a short per-connection dial timeout from the overall +// discovery timeout, clamped to a sensible range. +func perHostTimeout(timeout time.Duration) time.Duration { + dialTimeout := timeout / 4 + if dialTimeout < 300*time.Millisecond { + dialTimeout = 300 * time.Millisecond + } + if dialTimeout > 1500*time.Millisecond { + dialTimeout = 1500 * time.Millisecond + } + return dialTimeout +} + +func mergeSortedInts(existing, added []int) []int { + set := make(map[int]struct{}, len(existing)+len(added)) + for _, value := range existing { + set[value] = struct{}{} + } + for _, value := range added { + set[value] = struct{}{} + } + merged := make([]int, 0, len(set)) + for value := range set { + merged = append(merged, value) + } + sort.Ints(merged) + return merged +} + +func mergeUniqueStrings(existing, added []string) []string { + set := make(map[string]struct{}, len(existing)+len(added)) + merged := make([]string, 0, len(existing)+len(added)) + for _, value := range append(append([]string{}, existing...), added...) { + if _, ok := set[value]; ok { + continue + } + set[value] = struct{}{} + merged = append(merged, value) + } + return merged +} + +// ipLess compares two IPv4 address strings numerically. +func ipLess(a, b string) bool { + ipA := net.ParseIP(a).To4() + ipB := net.ParseIP(b).To4() + if ipA == nil || ipB == nil { + return a < b + } + for i := 0; i < 4; i++ { + if ipA[i] != ipB[i] { + return ipA[i] < ipB[i] + } + } + return false +} diff --git a/machinery/src/onvif/fingerprint.go b/machinery/src/onvif/fingerprint.go new file mode 100644 index 0000000..cd1e696 --- /dev/null +++ b/machinery/src/onvif/fingerprint.go @@ -0,0 +1,267 @@ +package onvif + +import ( + "bufio" + "net" + "strconv" + "strings" + "time" +) + +// deviceFingerprint holds the identifying information we can gather from a host +// without any credentials. It is populated by grabbing the RTSP and HTTP +// service banners and is then distilled into a manufacturer, model and a +// human-readable device type (e.g. "IP Camera", "DVR/NVR"). +type deviceFingerprint struct { + Manufacturer string + Model string + Type string + Server string + // realm is the WWW-Authenticate realm advertised by the HTTP service. Many + // cameras expose their model or vendor here (e.g. realm="Hikvision"). + realm string + // IsCamera is set when the collected evidence confidently identifies the + // device as a camera, NVR or DVR. + IsCamera bool +} + +// bannerVendors maps a lower-cased substring commonly found in RTSP/HTTP +// service banners or auth realms to a manufacturer. The list is ordered so the +// most specific matches win. This mirrors how tools such as Fing or ONVIF +// Device Manager fingerprint a device from its network banners. +var bannerVendors = []struct { + Match string + Vendor string + IsCamera bool +}{ + {"hikvision", "Hikvision", true}, + {"dahua", "Dahua", true}, + {"axis", "Axis", true}, + {"reolink", "Reolink", true}, + {"amcrest", "Amcrest", true}, + {"vivotek", "Vivotek", true}, + {"mobotix", "Mobotix", true}, + {"hanwha", "Hanwha", true}, + {"wisenet", "Hanwha", true}, + {"bosch", "Bosch", true}, + {"foscam", "Foscam", true}, + {"ubiquiti", "Ubiquiti", true}, + {"unifi", "Ubiquiti", true}, + {"uniview", "Uniview", true}, + {"tp-link", "TP-Link", true}, + {"tapo", "TP-Link", true}, + {"hipcam", "Hipcam", true}, + {"h264dvr", "Generic DVR", true}, + {"dvrdvs", "Hikvision", true}, + {"webs", "", false}, // generic embedded web server, no vendor + {"rtsp server", "", true}, + {"gstreamer", "", true}, + {"live555", "", true}, +} + +// genericRealms are auth realms that carry no useful model/vendor information. +var genericRealms = map[string]struct{}{ + "": {}, + "ip camera": {}, + "ipcamera": {}, + "camera": {}, + "login": {}, + "index": {}, + "streaming": {}, + "realm": {}, + "network video": {}, + "web": {}, + "protected": {}, + "authorized users only": {}, +} + +// fingerprintHost grabs the RTSP and HTTP banners for the given host (based on +// the ports found open during the scan) and classifies the device. It performs +// at most two lightweight, unauthenticated requests and is safe to run +// concurrently for every host. +func fingerprintHost(ip string, openPorts []int, timeout time.Duration) deviceFingerprint { + var fp deviceFingerprint + + // 1) RTSP OPTIONS on the first open RTSP port. The Server response header of + // most camera RTSP stacks reveals the device (e.g. "Dahua Rtsp Server", + // "Hipcam RealServer/V1.0", "H264DVR 1.0"). + for _, port := range openPorts { + if port == 554 || port == 8554 { + if banner := rtspServerBanner(ip, port, timeout); banner != "" { + fp.Server = banner + } + break + } + } + + // 2) HTTP banner + auth realm on the first open HTTP/ONVIF port. Cameras + // frequently expose their vendor/model in the Server header or the + // WWW-Authenticate realm. + for _, port := range openPorts { + if port == 80 || port == 8080 || port == 8000 { + server, realm := httpBanner(ip, port, timeout) + if fp.Server == "" { + fp.Server = server + } + fp.realm = realm + break + } + } + + classifyFingerprint(&fp, openPorts) + return fp +} + +// rtspServerBanner issues an unauthenticated RTSP OPTIONS request and returns +// the value of the Server response header (empty when the host does not answer +// or exposes no banner). +func rtspServerBanner(ip string, port int, timeout time.Duration) string { + address := net.JoinHostPort(ip, strconv.Itoa(port)) + conn, err := net.DialTimeout("tcp", address, timeout) + if err != nil { + return "" + } + defer conn.Close() + + _ = conn.SetDeadline(time.Now().Add(timeout)) + request := "OPTIONS rtsp://" + address + " RTSP/1.0\r\nCSeq: 1\r\nUser-Agent: KerberosDiscovery\r\n\r\n" + if _, err := conn.Write([]byte(request)); err != nil { + return "" + } + headers := readBannerHeaders(conn) + return headers["server"] +} + +// httpBanner issues an unauthenticated HTTP HEAD request and returns the Server +// header and the WWW-Authenticate realm (both best-effort, empty when absent). +func httpBanner(ip string, port int, timeout time.Duration) (server string, realm string) { + address := net.JoinHostPort(ip, strconv.Itoa(port)) + conn, err := net.DialTimeout("tcp", address, timeout) + if err != nil { + return "", "" + } + defer conn.Close() + + _ = conn.SetDeadline(time.Now().Add(timeout)) + request := "HEAD / HTTP/1.0\r\nHost: " + ip + "\r\nUser-Agent: KerberosDiscovery\r\nAccept: */*\r\n\r\n" + if _, err := conn.Write([]byte(request)); err != nil { + return "", "" + } + headers := readBannerHeaders(conn) + return headers["server"], parseRealm(headers["www-authenticate"]) +} + +// readBannerHeaders reads a status line followed by header lines from an +// RTSP/HTTP response and returns the headers keyed by their lower-cased name. +// Only the first occurrence of a header is kept. +func readBannerHeaders(conn net.Conn) map[string]string { + headers := make(map[string]string) + reader := bufio.NewReader(conn) + + // Discard the status line (e.g. "RTSP/1.0 200 OK" or "HTTP/1.1 401 ..."). + if _, err := reader.ReadString('\n'); err != nil { + return headers + } + for { + line, err := reader.ReadString('\n') + if err != nil { + break + } + line = strings.TrimRight(line, "\r\n") + if line == "" { + break + } + idx := strings.Index(line, ":") + if idx <= 0 { + continue + } + key := strings.ToLower(strings.TrimSpace(line[:idx])) + value := strings.TrimSpace(line[idx+1:]) + if _, exists := headers[key]; !exists { + headers[key] = value + } + } + return headers +} + +// parseRealm extracts the realm token from a WWW-Authenticate header value such +// as `Digest realm="Hikvision", nonce="..."`. +func parseRealm(header string) string { + lower := strings.ToLower(header) + marker := "realm=" + idx := strings.Index(lower, marker) + if idx < 0 { + return "" + } + value := header[idx+len(marker):] + value = strings.TrimSpace(value) + if strings.HasPrefix(value, "\"") { + value = value[1:] + if end := strings.Index(value, "\""); end >= 0 { + value = value[:end] + } + } else if end := strings.IndexAny(value, ", "); end >= 0 { + value = value[:end] + } + return strings.TrimSpace(value) +} + +// classifyFingerprint distils the collected banners and open ports into a +// manufacturer, model and device type. It also decides whether the evidence is +// strong enough to consider the host a camera/NVR. +func classifyFingerprint(fp *deviceFingerprint, openPorts []int) { + haystack := strings.ToLower(fp.Server + " " + fp.realm) + + // Manufacturer from the banner/realm. + for _, entry := range bannerVendors { + if !strings.Contains(haystack, entry.Match) { + continue + } + if entry.Vendor != "" && fp.Manufacturer == "" { + fp.Manufacturer = entry.Vendor + } + if entry.IsCamera { + fp.IsCamera = true + } + if fp.Manufacturer != "" { + break + } + } + + // Model from the auth realm when it looks specific (not a generic word). + if fp.Model == "" && fp.realm != "" { + if _, generic := genericRealms[strings.ToLower(fp.realm)]; !generic { + if !strings.EqualFold(fp.realm, fp.Manufacturer) { + fp.Model = fp.realm + } + } + } + + // Device type from ports and banners. + hasRTSP := containsInt(openPorts, 554) || containsInt(openPorts, 8554) + hasONVIF := containsInt(openPorts, 8000) || containsInt(openPorts, 8899) + hasDVRPort := containsInt(openPorts, 37777) || containsInt(openPorts, 34567) + + switch { + case strings.Contains(haystack, "nvr"): + fp.Type = "NVR" + fp.IsCamera = true + case strings.Contains(haystack, "dvr") || hasDVRPort: + fp.Type = "DVR/NVR" + fp.IsCamera = true + case hasRTSP || hasONVIF: + fp.Type = "IP Camera" + fp.IsCamera = true + case fp.IsCamera: + fp.Type = "IP Camera" + } +} + +func containsInt(values []int, target int) bool { + for _, value := range values { + if value == target { + return true + } + } + return false +} diff --git a/machinery/src/onvif/main.go b/machinery/src/onvif/main.go index 45623d7..bfe95f5 100644 --- a/machinery/src/onvif/main.go +++ b/machinery/src/onvif/main.go @@ -10,7 +10,6 @@ import ( "strings" "time" - onvifc "github.com/cedricve/go-onvif" "github.com/gin-gonic/gin" "github.com/kerberos-io/agent/machinery/src/log" "github.com/kerberos-io/agent/machinery/src/models" @@ -24,20 +23,66 @@ import ( xsdonvif "github.com/kerberos-io/onvif/xsd/onvif" ) -func Discover(timeout time.Duration) { - log.Log.Info("onvif.Discover(): Discovering devices") - log.Log.Info("Waiting for " + timeout.String()) - devices, err := onvifc.StartDiscovery(timeout) - if err != nil { - log.Log.Error("onvif.Discover(): " + err.Error()) - } else { - for _, device := range devices { - hostname, _ := device.GetHostname() - log.Log.Info("onvif.Discover(): " + hostname.Name + " (" + device.XAddr + ")") +// Discover performs an advanced Fing/WiFiman-style scan of the local network +// (ONVIF WS-Discovery + active port scan + MAC/vendor lookup) and prints a +// human readable summary of everything it finds. It is used by the +// `-action discover` CLI command. Optional subnets (CIDR, e.g. +// "192.168.1.0/24") override the auto-detected local subnets. +func Discover(timeout time.Duration, subnets ...string) { + log.Log.Info("onvif.Discover(): starting advanced network discovery") + log.Log.Info("onvif.Discover(): this may take up to " + timeout.String() + " for the ONVIF probe plus the port scan") + + devices := DiscoverDevices(timeout, subnets...) + if len(devices) == 0 { + log.Log.Info("onvif.Discover(): no devices discovered on the local network") + return + } + + cameraCount := 0 + for _, device := range devices { + if device.IsCamera { + cameraCount++ } - if len(devices) == 0 { - log.Log.Info("onvif.Discover(): No devices descovered\n") + } + log.Log.Info("onvif.Discover(): found " + strconv.Itoa(len(devices)) + " device(s), " + strconv.Itoa(cameraCount) + " likely camera(s)") + + for _, device := range devices { + label := "device" + if device.IsCamera { + label = "camera" } + summary := "onvif.Discover(): [" + label + "] " + device.IP + if device.Hostname != "" { + summary += " (" + device.Hostname + ")" + } + if device.MAC != "" { + summary += " mac=" + device.MAC + } + if device.Vendor != "" { + summary += " vendor=" + device.Vendor + } + if device.Type != "" { + summary += " type=" + device.Type + } + if device.Manufacturer != "" { + summary += " manufacturer=" + device.Manufacturer + } + if device.Model != "" { + summary += " model=" + device.Model + } + if device.Server != "" { + summary += " server=\"" + device.Server + "\"" + } + if device.ONVIF { + summary += " onvif=" + device.ONVIFXAddr + } + if len(device.Services) > 0 { + summary += " services=[" + strings.Join(device.Services, ", ") + "]" + } + if device.RTSPURL != "" { + summary += " rtsp=" + device.RTSPURL + } + log.Log.Info(summary) } } diff --git a/machinery/src/routers/http/methods.go b/machinery/src/routers/http/methods.go index b8d0aa1..303b970 100644 --- a/machinery/src/routers/http/methods.go +++ b/machinery/src/routers/http/methods.go @@ -1,6 +1,10 @@ package http import ( + "strconv" + "strings" + "time" + "github.com/gin-gonic/gin" "github.com/kerberos-io/agent/machinery/src/log" "github.com/kerberos-io/agent/machinery/src/models" @@ -17,6 +21,38 @@ import ( // @Success 200 {object} models.Authorization func Login() {} +// DiscoverCameras godoc +// @Router /api/camera/discover [get] +// @ID camera-discover +// @Tags onvif +// @Param timeout query int false "Discovery timeout in milliseconds (default 2000)" +// @Param subnet query string false "Optional subnet(s) to scan, e.g. '192.168.1.0/24' (comma-separated). Defaults to the local interfaces." +// @Summary Discover cameras and other devices on the local network. +// @Description Runs an advanced Fing/WiFiman-style scan (ONVIF WS-Discovery + TCP port scan + MAC/vendor lookup) and returns the devices found on the local network. +// @Success 200 {object} models.APIResponse +func DiscoverCameras(c *gin.Context) { + timeout := 2000 * time.Millisecond + if raw := c.Query("timeout"); raw != "" { + if milliseconds, err := strconv.Atoi(raw); err == nil && milliseconds > 0 { + timeout = time.Duration(milliseconds) * time.Millisecond + } + } + + var subnets []string + if raw := c.Query("subnet"); raw != "" { + for _, part := range strings.Split(raw, ",") { + if trimmed := strings.TrimSpace(part); trimmed != "" { + subnets = append(subnets, trimmed) + } + } + } + + devices := onvif.DiscoverDevices(timeout, subnets...) + c.JSON(200, models.APIResponse{ + Data: devices, + }) +} + // LoginToOnvif godoc // @Router /api/camera/onvif/login [post] // @ID camera-onvif-login diff --git a/machinery/src/routers/http/routes.go b/machinery/src/routers/http/routes.go index 15c9f16..f46c43f 100644 --- a/machinery/src/routers/http/routes.go +++ b/machinery/src/routers/http/routes.go @@ -96,6 +96,7 @@ func AddRoutes(r *gin.Engine, authMiddleware *jwt.GinJWTMiddleware, configDirect }) // Onvif specific methods. + api.GET("/camera/discover", DiscoverCameras) api.POST("/camera/onvif/verify", onvif.VerifyOnvifConnection) api.POST("/camera/onvif/login", LoginToOnvif) api.POST("/camera/onvif/capabilities", GetOnvifCapabilities) diff --git a/machinery/src/routers/mqtt/main.go b/machinery/src/routers/mqtt/main.go index 51fed31..694039f 100644 --- a/machinery/src/routers/mqtt/main.go +++ b/machinery/src/routers/mqtt/main.go @@ -14,7 +14,10 @@ import ( "sync" "time" + "context" + mqtt "github.com/eclipse/paho.mqtt.golang" + "github.com/kerberos-io/agent/machinery/src/capture" configService "github.com/kerberos-io/agent/machinery/src/config" "github.com/kerberos-io/agent/machinery/src/encryption" "github.com/kerberos-io/agent/machinery/src/log" @@ -338,6 +341,8 @@ func MQTTListenerHandler(mqttClient mqtt.Client, hubKey string, configDirectory go HandleNavigatePTZ(mqttClient, hubKey, payload, configuration, communication) case "request-config": go HandleRequestConfig(mqttClient, hubKey, payload, configuration, communication) + case "verify-stream": + go HandleVerifyStream(mqttClient, hubKey, payload, configuration, communication) case "update-config": go HandleUpdateConfig(mqttClient, hubKey, payload, configDirectory, configuration, communication) case "request-sd-stream": @@ -546,6 +551,98 @@ func HandleRequestConfig(mqttClient mqtt.Client, hubKey string, payload models.P } } +// HandleVerifyStream probes an RTSP stream (the one supplied in the request, or +// the currently configured main/sub stream) and reports back whether it can be +// connected to and decoded, along with the discovered codec/resolution/fps. +func HandleVerifyStream(mqttClient mqtt.Client, hubKey string, payload models.Payload, configuration *models.Configuration, communication *models.Communication) { + value := payload.Value + + // Convert map[string]interface{} to VerifyStreamPayload + jsonData, _ := json.Marshal(value) + var verifyPayload models.VerifyStreamPayload + json.Unmarshal(jsonData, &verifyPayload) + + if verifyPayload.Timestamp == 0 { + return + } + + stream := verifyPayload.Stream + if stream != "sub" { + stream = "main" + } + + // Resolve which RTSP url to verify: prefer the one supplied in the request + // (so users can verify unsaved edits), otherwise fall back to the configured + // stream url for the requested stream type. + rtspUrl := verifyPayload.RTSP + if rtspUrl == "" { + if stream == "sub" { + rtspUrl = configuration.Config.Capture.IPCamera.SubRTSP + } else { + rtspUrl = configuration.Config.Capture.IPCamera.RTSP + } + } + + success := false + errMsg := "" + width := 0 + height := 0 + codec := "" + fps := 0.0 + + if rtspUrl == "" { + errMsg = "No RTSP url configured for this stream." + } else { + // Probe the stream with a bounded timeout so a dead/unreachable camera + // can't hang the handler goroutine. + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + rtspClient := &capture.Golibrtsp{Url: rtspUrl} + errConnect := rtspClient.Connect(ctx, ctx) + if errConnect != nil { + errMsg = errConnect.Error() + } else { + videoStreams, errStreams := rtspClient.GetVideoStreams() + if errStreams != nil || len(videoStreams) == 0 { + errMsg = "Connected, but no decodable video stream was found." + } else { + success = true + vs := videoStreams[0] + width = vs.Width + height = vs.Height + codec = vs.Name + fps = vs.FPS + } + } + // Always release the connection. + rtspClient.Close(ctx) + } + + message := models.Message{ + Payload: models.Payload{ + Action: "verify-stream-result", + DeviceId: configuration.Config.Key, + Value: map[string]interface{}{ + "timestamp": verifyPayload.Timestamp, + "stream": stream, + "success": success, + "error": errMsg, + "width": width, + "height": height, + "codec": codec, + "fps": fps, + }, + }, + } + packagedPayload, err := models.PackageMQTTMessage(configuration, message) + if err == nil { + mqttClient.Publish("kerberos/hub/"+hubKey, 2, false, packagedPayload) + } else { + log.Log.Info("routers.mqtt.main.HandleVerifyStream(): something went wrong while sending result to hub: " + string(packagedPayload)) + } +} + func HandleUpdateConfig(mqttClient mqtt.Client, hubKey string, payload models.Payload, configDirectory string, configuration *models.Configuration, communication *models.Communication) { value := payload.Value