Compare commits

...

14 Commits

Author SHA1 Message Date
Cédric Verstraeten
94b26cf096 Merge pull request #305 from kerberos-io/feature/add-motion-detection-pixel-changes
feature/add-motion-detection-pixel-changes
2026-07-16 15:33:02 +02:00
Cédric Verstraeten
57ef7ebaaf Enhance ONVIF fingerprinting and brand profiles
Expand discovery and classification to better identify device vendors and stream paths by adding new brand aliases/profiles (including D-Link, Trendnet, Lorex, Honeywell, Pelco, and TOA), improved realm matching, and hostname-based brand hints. Add HTTP body fingerprinting for OEM/rebadged devices, introduce audio-device detection with a new `is_audio` API field, and prevent camera RTSP guessing/fallback URLs for audio-only devices while updating discovery logging labels.
2026-07-16 14:39:55 +02:00
Cédric Verstraeten
ddf58fe633 Add Linksys ONVIF fingerprint and RTSP paths
Adds Linksys camera detection across ONVIF heuristics by introducing a dedicated brand profile, realm aliases, and banner fingerprinting. It also prioritizes Linksys-specific RTSP endpoints (including /ONVIF/channel1 and /ONVIF/channel2) and includes fallback stream paths used by Cisco/Linksys models.
2026-07-16 13:20:47 +02:00
Cédric Verstraeten
6f2d35cdf1 Add RTSP brand probing to ONVIF discovery
Extend device discovery to generate RTSP stream candidates using built-in brand profiles, RTSP DESCRIBE probing, auth-realm parsing, and port hints. Add `RTSPStreams`/`RTSPStream` to API responses, prefer verified stream URLs as primary `RTSPURL`, and let stronger RTSP-derived brand/model signals refine detected camera metadata. Also add focused unit tests for discriminating vs non-discriminating devices, realm-based brand/model detection, and generic fallback behavior.
2026-07-16 09:09:27 +02:00
Cédric Verstraeten
c836cef28d 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.
2026-07-15 23:17:23 +02:00
Cédric Verstraeten
c97bb70cb5 Update main.go 2026-07-15 12:08:10 +02:00
Cédric Verstraeten
96b145b046 Reset config.json to clean defaults
Removes sensitive credentials, private keys, and personal configuration values from config.json. Clears RTSP URLs, hub keys, encryption/signing keys, kstorage credentials, and resets various settings to neutral defaults.
2026-07-14 21:57:54 +02:00
Cédric Verstraeten
09a697e00b Make PixelChangeThreshold a pointer to distinguish unset vs 0
Changes PixelChangeThreshold from int to *int so nil (unset) defaults to 150, while 0 explicitly disables motion detection. Updates ProcessMotion to handle the new three-state logic and also emits pixelChangeThreshold in MQTT motion messages for live view visualization.
2026-07-14 21:52:24 +02:00
Cédric Verstraeten
1d0714f199 Merge pull request #304 from kerberos-io/feature/tweak-remote-recording
feature/tweak-remote-recording
2026-07-13 21:17:50 +02:00
Cédric Verstraeten
42e91867ec Refactor JSON handling and enhance motion detection
Replace unsafe fmt.Sprintf JSON formatting with proper struct marshalling in cloud.go. Add rawJSONOrEmptyArray() helper to safely handle json.RawMessage with fallback to empty arrays. Enhance motion detection overlay data in computervision by adding main stream dimensions (mainWidth/mainHeight) alongside motion frame dimensions, enabling accurate live-view scaling of detected regions.
2026-07-13 21:10:55 +02:00
Cédric Verstraeten
155c4a7e44 Support motion regions in continuous recording mode
Refactor motion detection logic to enable motion region visualization during continuous recording. Motion detection now runs in continuous mode when a motion region is configured, allowing live-view overlay display without triggering motion-based recording. In continuous mode without regions, motion detection is skipped as before. Updated conditional logic and added clarifying comments explaining the different code paths.
2026-07-13 17:59:34 +02:00
Cédric Verstraeten
4fe4977559 Update main.go 2026-07-13 17:32:18 +02:00
Cédric Verstraeten
67e66e863a Add continuousRecording flag to Hub status
Report whether a camera is in continuous recording mode (24/7) to the Hub. This allows the Hub live view to disable the manual record button, which is a no-op when the camera is already recording continuously.
2026-07-13 16:51:46 +02:00
Cédric Verstraeten
bd34e9d836 Auto-stop stale manual recordings
Add heartbeat-aware lifecycle management for manual/live-view recordings. The agent now tracks manual recording start time and viewer heartbeats, auto-stops recordings when heartbeats lapse or a max duration is reached, and clears state on stop/restart. MQTT recording payloads gain a `heartbeat` flag so keep-alives refresh active sessions without unintentionally restarting recordings after auto-stop.
2026-07-13 16:29:46 +02:00
20 changed files with 2523 additions and 148 deletions

View File

@@ -122,4 +122,4 @@
"signing": {},
"realtimeprocessing": "false",
"realtimeprocessing_topic": ""
}
}

View File

@@ -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":
{

View File

@@ -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 {

View File

@@ -50,6 +50,55 @@ func publishRecordingState(mqttClient mqtt.Client, hubKey string, configuration
}
}
const (
// manualRecordingHeartbeatTimeout is how long the agent keeps a manual
// (live-view / remote) recording alive after the LAST viewer heartbeat. The
// frontend re-sends the record command every ~15s while the user stays on the
// page; if several heartbeats are missed (the viewer closed the tab, went idle
// or lost connectivity) the recorder auto-stops the recording so the camera
// doesn't record forever when the "stop" message never arrives.
manualRecordingHeartbeatTimeout = 45 * time.Second
// manualRecordingMaxDuration caps a single manual recording so a forgotten
// record button can't record indefinitely even while the viewer keeps sending
// heartbeats. After this the recording auto-stops and the viewer must press
// record again to continue.
manualRecordingMaxDuration = 5 * time.Minute
)
// manualRecordingExpired reports whether an active manual (live-view) recording
// has outlived its viewer heartbeat window or the maximum duration cap. When it
// has, it clears the manual-recording state (so the motion recorder lets the
// current clip close normally and broadcasts recording:false) and returns true.
// It is a no-op returning false when no manual recording is active.
func manualRecordingExpired(communication *models.Communication, now int64) bool {
if communication.IsRecordingManual.IsNotSet() {
return false
}
manualStart := communication.RecordingManualStart.Load()
maxDurationReached := manualStart > 0 && now-manualStart > manualRecordingMaxDuration.Milliseconds()
// The heartbeat timeout only applies once the viewer has proven it supports
// heartbeats (an older frontend that starts a recording but never heartbeats
// still records up to the max-duration cap instead of being cut off early).
heartbeatExpired := false
if communication.RecordingManualHeartbeatSeen.IsSet() {
lastHeartbeat := communication.RecordingManualHeartbeat.Load()
heartbeatExpired = lastHeartbeat > 0 && now-lastHeartbeat > manualRecordingHeartbeatTimeout.Milliseconds()
}
if !heartbeatExpired && !maxDurationReached {
return false
}
if heartbeatExpired {
log.Log.Info("capture.main.HandleRecordStream(motiondetection): auto-stopping manual recording, no viewer heartbeat within timeout.")
} else {
log.Log.Info("capture.main.HandleRecordStream(motiondetection): auto-stopping manual recording, maximum duration reached.")
}
communication.IsRecordingManual.UnSet()
communication.RecordingManualHeartbeat.Store(0)
communication.RecordingManualStart.Store(0)
communication.RecordingManualHeartbeatSeen.UnSet()
return true
}
func CleanupRecordingDirectory(configDirectory string, configuration *models.Configuration) {
autoClean := configuration.Config.AutoClean
if autoClean != "true" {
@@ -227,6 +276,9 @@ func HandleRecordStream(queue *packets.Queue, configDirectory string, configurat
// Start each capture session with manual recording off, so a leftover
// request from before a restart/reconnect doesn't silently persist.
communication.IsRecordingManual.UnSet()
communication.RecordingManualHeartbeat.Store(0)
communication.RecordingManualStart.Store(0)
communication.RecordingManualHeartbeatSeen.UnSet()
if config.Capture.Recording == "false" {
log.Log.Info("capture.main.HandleRecordStream(): disabled, we will not record anything.")
@@ -674,7 +726,10 @@ func HandleRecordStream(queue *packets.Queue, configDirectory string, configurat
// motion timestamp every iteration so the post-recording timeout
// never fires. The clip still rolls over at maxRecordingPeriod and
// is restarted below, until the viewer stops the manual recording.
if communication.IsRecordingManual.IsSet() {
// It also auto-stops when the viewer's heartbeat lapses (closed page
// or idle) or the max remote-recording duration is reached, so a
// missed "stop" message can't keep the camera recording forever.
if communication.IsRecordingManual.IsSet() && !manualRecordingExpired(communication, now) {
motionTimestamp = now
}
@@ -751,8 +806,10 @@ func HandleRecordStream(queue *packets.Queue, configDirectory string, configurat
// If the viewer still has a manual recording running, this clip just
// rolled over at the max length — immediately kick off the next
// segment so recording stays continuous until they stop it.
if communication.IsRecordingManual.IsSet() {
// segment so recording stays continuous until they stop it. Skip the
// restart when the recording has expired (heartbeat lapsed or max
// duration reached), so it ends here instead of recording forever.
if communication.IsRecordingManual.IsSet() && !manualRecordingExpired(communication, time.Now().UnixMilli()) {
select {
case communication.HandleMotion <- models.MotionDataPartial{Timestamp: time.Now().Unix(), NumberOfChanges: 100000000}:
default:

View File

@@ -6,7 +6,6 @@ import (
"crypto/tls"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"os"
"strings"
@@ -219,6 +218,16 @@ func GetSystemInfo() (models.System, error) {
return system, nil
}
// rawJSONOrEmptyArray returns pre-marshalled JSON bytes as a json.RawMessage.
// When the input is empty it falls back to an empty JSON array so the
// surrounding payload always stays valid JSON.
func rawJSONOrEmptyArray(b []byte) json.RawMessage {
if len(b) == 0 {
return json.RawMessage("[]")
}
return json.RawMessage(b)
}
func HandleHeartBeat(configuration *models.Configuration, communication *models.Communication, uptimeStart time.Time) {
log.Log.Debug("cloud.HandleHeartBeat(): started")
@@ -472,6 +481,14 @@ loop:
hasBackChannel = "true"
}
// Whether this camera records continuously (24/7) rather than on
// motion. The Hub live view uses this to disable the manual record
// button, which is a no-op in continuous mode (already recording).
continuousRecording := "false"
if config.Capture.Continuous == "true" {
continuousRecording = "true"
}
hub_encryption := "false"
if config.HubEncryption == "true" {
hub_encryption = "true"
@@ -496,48 +513,98 @@ loop:
// We need a hub URI and hub public key before we will send a heartbeat
if hubURI != "" && key != "" {
var object = fmt.Sprintf(`{
"key" : "%s",
"version" : "%s",
"hub_encryption": "%s",
"e2e_encryption": "%s",
"release" : "%s",
"cpuid" : "%s",
"clouduser" : "%s",
"cloudpublickey" : "%s",
"cameraname" : "%s",
"enterprise" : %t,
"hostname" : "%s",
"architecture" : "%s",
"totalMemory" : "%d",
"usedMemory" : "%d",
"freeMemory" : "%d",
"processMemory" : "%d",
"mac_list" : %s,
"ip_list" : %s,
"board" : "",
"disk1size" : "%s",
"disk3size" : "%s",
"diskvdasize" : "%s",
"uptime" : "%s",
"boot_time" : "%s",
"siteID" : "%s",
"onvif" : "%s",
"onvif_zoom" : "%s",
"onvif_pantilt" : "%s",
"onvif_presets": "%s",
"onvif_presets_list": %s,
"onvif_events_list": %s,
"cameraConnected": "%s",
"hasBackChannel": "%s",
"livePreviewHttp": true,
"numberoffiles" : "33",
"timestamp" : 1564747908,
"cameratype" : "IPCamera",
"docker" : true,
"kios" : false,
"raspberrypi" : false
}`, config.Key, kerberosAgentVersion, hub_encryption, e2e_encryption, system.Version, system.CPUId, username, key, name, isEnterprise, system.Hostname, system.Architecture, system.TotalMemory, system.UsedMemory, system.FreeMemory, system.ProcessUsedMemory, macs, ips, "0", "0", "0", uptimeString, boottimeString, config.HubSite, onvifEnabled, onvifZoom, onvifPanTilt, onvifPresets, onvifPresetsList, onvifEventsList, cameraConnected, hasBackChannel)
heartbeat := struct {
Key string `json:"key"`
Version string `json:"version"`
HubEncryption string `json:"hub_encryption"`
E2EEncryption string `json:"e2e_encryption"`
Release string `json:"release"`
CPUId string `json:"cpuid"`
CloudUser string `json:"clouduser"`
CloudPublicKey string `json:"cloudpublickey"`
CameraName string `json:"cameraname"`
Enterprise bool `json:"enterprise"`
Hostname string `json:"hostname"`
Architecture string `json:"architecture"`
TotalMemory string `json:"totalMemory"`
UsedMemory string `json:"usedMemory"`
FreeMemory string `json:"freeMemory"`
ProcessMemory string `json:"processMemory"`
MacList json.RawMessage `json:"mac_list"`
IPList json.RawMessage `json:"ip_list"`
Board string `json:"board"`
Disk1Size string `json:"disk1size"`
Disk3Size string `json:"disk3size"`
DiskVdaSize string `json:"diskvdasize"`
Uptime string `json:"uptime"`
BootTime string `json:"boot_time"`
SiteID string `json:"siteID"`
Onvif string `json:"onvif"`
OnvifZoom string `json:"onvif_zoom"`
OnvifPanTilt string `json:"onvif_pantilt"`
OnvifPresets string `json:"onvif_presets"`
OnvifPresetsList json.RawMessage `json:"onvif_presets_list"`
OnvifEventsList json.RawMessage `json:"onvif_events_list"`
CameraConnected string `json:"cameraConnected"`
HasBackChannel string `json:"hasBackChannel"`
ContinuousRecording string `json:"continuousRecording"`
LivePreviewHTTP bool `json:"livePreviewHttp"`
NumberOfFiles string `json:"numberoffiles"`
Timestamp int64 `json:"timestamp"`
CameraType string `json:"cameratype"`
Docker bool `json:"docker"`
Kios bool `json:"kios"`
RaspberryPi bool `json:"raspberrypi"`
}{
Key: config.Key,
Version: kerberosAgentVersion,
HubEncryption: hub_encryption,
E2EEncryption: e2e_encryption,
Release: system.Version,
CPUId: system.CPUId,
CloudUser: username,
CloudPublicKey: key,
CameraName: name,
Enterprise: isEnterprise,
Hostname: system.Hostname,
Architecture: system.Architecture,
TotalMemory: strconv.FormatUint(system.TotalMemory, 10),
UsedMemory: strconv.FormatUint(system.UsedMemory, 10),
FreeMemory: strconv.FormatUint(system.FreeMemory, 10),
ProcessMemory: strconv.FormatUint(system.ProcessUsedMemory, 10),
MacList: rawJSONOrEmptyArray(macs),
IPList: rawJSONOrEmptyArray(ips),
Board: "",
Disk1Size: "0",
Disk3Size: "0",
DiskVdaSize: "0",
Uptime: uptimeString,
BootTime: boottimeString,
SiteID: config.HubSite,
Onvif: onvifEnabled,
OnvifZoom: onvifZoom,
OnvifPanTilt: onvifPanTilt,
OnvifPresets: onvifPresets,
OnvifPresetsList: rawJSONOrEmptyArray(onvifPresetsList),
OnvifEventsList: rawJSONOrEmptyArray(onvifEventsList),
CameraConnected: cameraConnected,
HasBackChannel: hasBackChannel,
ContinuousRecording: continuousRecording,
LivePreviewHTTP: true,
NumberOfFiles: "33",
Timestamp: 1564747908,
CameraType: "IPCamera",
Docker: true,
Kios: false,
RaspberryPi: false,
}
objectBytes, err := json.Marshal(heartbeat)
if err != nil {
log.Log.Error("cloud.HandleHeartBeat(): error while marshalling heartbeat: " + err.Error())
objectBytes = []byte("{}")
}
object := string(objectBytes)
// Get the private key to encrypt the data using symmetric encryption: AES.
privateKey := config.HubPrivateKey
@@ -551,11 +618,21 @@ loop:
// Base64 encode the encrypted data.
encryptedBase64 := base64.StdEncoding.EncodeToString(encrypted)
object = fmt.Sprintf(`{
"cloudpublicKey": "%s",
"encrypted" : %t,
"encryptedData" : "%s"
}`, config.HubKey, true, encryptedBase64)
encryptedPayload := struct {
CloudPublicKey string `json:"cloudpublicKey"`
Encrypted bool `json:"encrypted"`
EncryptedData string `json:"encryptedData"`
}{
CloudPublicKey: config.HubKey,
Encrypted: true,
EncryptedData: encryptedBase64,
}
encryptedBytes, err := json.Marshal(encryptedPayload)
if err != nil {
log.Log.Error("cloud.HandleHeartBeat(): error while marshalling encrypted heartbeat: " + err.Error())
encryptedBytes = []byte("{}")
}
object = string(encryptedBytes)
}
var jsonStr = []byte(object)
@@ -586,43 +663,86 @@ loop:
secretAccessKey := config.KStorage.SecretAccessKey
if vaultURI != "" && accessKey != "" && secretAccessKey != "" {
var object = fmt.Sprintf(`{
"key" : "%s",
"version" : "%s",
"release" : "%s",
"cpuid" : "%s",
"clouduser" : "%s",
"cloudpublickey" : "%s",
"cameraname" : "%s",
"enterprise" : %t,
"hostname" : "%s",
"architecture" : "%s",
"totalMemory" : "%d",
"usedMemory" : "%d",
"freeMemory" : "%d",
"processMemory" : "%d",
"mac_list" : %s,
"ip_list" : %s,
"board" : "",
"disk1size" : "%s",
"disk3size" : "%s",
"diskvdasize" : "%s",
"uptime" : "%s",
"boot_time" : "%s",
"siteID" : "%s",
"onvif" : "%s",
"onvif_zoom" : "%s",
"onvif_pantilt" : "%s",
"onvif_presets": "%s",
"onvif_presets_list": %s,
"cameraConnected": "%s",
"numberoffiles" : "33",
"timestamp" : 1564747908,
"cameratype" : "IPCamera",
"docker" : true,
"kios" : false,
"raspberrypi" : false
}`, config.Key, kerberosAgentVersion, system.Version, system.CPUId, username, key, name, isEnterprise, system.Hostname, system.Architecture, system.TotalMemory, system.UsedMemory, system.FreeMemory, system.ProcessUsedMemory, macs, ips, "0", "0", "0", uptimeString, boottimeString, config.HubSite, onvifEnabled, onvifZoom, onvifPanTilt, onvifPresets, onvifPresetsList, cameraConnected)
heartbeat := struct {
Key string `json:"key"`
Version string `json:"version"`
Release string `json:"release"`
CPUId string `json:"cpuid"`
CloudUser string `json:"clouduser"`
CloudPublicKey string `json:"cloudpublickey"`
CameraName string `json:"cameraname"`
Enterprise bool `json:"enterprise"`
Hostname string `json:"hostname"`
Architecture string `json:"architecture"`
TotalMemory string `json:"totalMemory"`
UsedMemory string `json:"usedMemory"`
FreeMemory string `json:"freeMemory"`
ProcessMemory string `json:"processMemory"`
MacList json.RawMessage `json:"mac_list"`
IPList json.RawMessage `json:"ip_list"`
Board string `json:"board"`
Disk1Size string `json:"disk1size"`
Disk3Size string `json:"disk3size"`
DiskVdaSize string `json:"diskvdasize"`
Uptime string `json:"uptime"`
BootTime string `json:"boot_time"`
SiteID string `json:"siteID"`
Onvif string `json:"onvif"`
OnvifZoom string `json:"onvif_zoom"`
OnvifPanTilt string `json:"onvif_pantilt"`
OnvifPresets string `json:"onvif_presets"`
OnvifPresetsList json.RawMessage `json:"onvif_presets_list"`
CameraConnected string `json:"cameraConnected"`
NumberOfFiles string `json:"numberoffiles"`
Timestamp int64 `json:"timestamp"`
CameraType string `json:"cameratype"`
Docker bool `json:"docker"`
Kios bool `json:"kios"`
RaspberryPi bool `json:"raspberrypi"`
}{
Key: config.Key,
Version: kerberosAgentVersion,
Release: system.Version,
CPUId: system.CPUId,
CloudUser: username,
CloudPublicKey: key,
CameraName: name,
Enterprise: isEnterprise,
Hostname: system.Hostname,
Architecture: system.Architecture,
TotalMemory: strconv.FormatUint(system.TotalMemory, 10),
UsedMemory: strconv.FormatUint(system.UsedMemory, 10),
FreeMemory: strconv.FormatUint(system.FreeMemory, 10),
ProcessMemory: strconv.FormatUint(system.ProcessUsedMemory, 10),
MacList: rawJSONOrEmptyArray(macs),
IPList: rawJSONOrEmptyArray(ips),
Board: "",
Disk1Size: "0",
Disk3Size: "0",
DiskVdaSize: "0",
Uptime: uptimeString,
BootTime: boottimeString,
SiteID: config.HubSite,
Onvif: onvifEnabled,
OnvifZoom: onvifZoom,
OnvifPanTilt: onvifPanTilt,
OnvifPresets: onvifPresets,
OnvifPresetsList: rawJSONOrEmptyArray(onvifPresetsList),
CameraConnected: cameraConnected,
NumberOfFiles: "33",
Timestamp: 1564747908,
CameraType: "IPCamera",
Docker: true,
Kios: false,
RaspberryPi: false,
}
objectBytes, err := json.Marshal(heartbeat)
if err != nil {
log.Log.Error("cloud.HandleHeartBeat(): error while marshalling vault heartbeat: " + err.Error())
objectBytes = []byte("{}")
}
object := string(objectBytes)
var jsonStr = []byte(object)
buffy := bytes.NewBuffer(jsonStr)

View File

@@ -76,6 +76,9 @@ func Bootstrap(ctx context.Context, configDirectory string, configuration *model
communication.HandleLiveHLS = make(chan string, 1)
communication.IsConfiguring = abool.New()
communication.IsRecordingManual = abool.New()
communication.RecordingManualHeartbeat = &atomic.Int64{}
communication.RecordingManualStart = &atomic.Int64{}
communication.RecordingManualHeartbeatSeen = abool.New()
cameraSettings := &models.Camera{}

View File

@@ -22,20 +22,44 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
var isPixelChangeThresholdReached = false
var changesToReturn = 0
var motionRectangle models.MotionRectangle
var motionRectangles []models.MotionRectangle
pixelThreshold := config.Capture.PixelChangeThreshold
// Might not be set in the config file, so set it to 150
if pixelThreshold == 0 {
pixelThreshold = 150
// Resolve the motion sensitivity (pixel-change threshold):
// nil (unset) -> default 150
// 0 -> motion detection DISABLED (temporary off switch from the UI)
// > 0 -> trigger when the number of changed pixels exceeds it
pixelThreshold := 150
motionDisabled := false
if config.Capture.PixelChangeThreshold != nil {
pixelThreshold = *config.Capture.PixelChangeThreshold
if pixelThreshold <= 0 {
motionDisabled = true
}
}
// In motion mode we always run detection. In CONTINUOUS mode recording is
// 24/7 so motion detection is normally skipped, BUT if a motion region is
// configured we still run it so the live view can visualise the motion boxes
// + region. In that case we only emit the motion EVENT — no motion-triggered
// recording (continuous already records, and the recorder's motion branch
// isn't draining HandleMotion in continuous mode).
continuousMode := config.Capture.Continuous == "true"
hasMotionRegion := config.Region != nil && len(config.Region.Polygon) > 0
if config.Capture.Continuous == "true" {
if motionDisabled {
log.Log.Info("computervision.main.ProcessMotion(): you've enabled continuous recording, so no motion detection required.")
log.Log.Info("computervision.main.ProcessMotion(): motion detection disabled (pixelChangeThreshold set to 0), skipping.")
} else if continuousMode && !hasMotionRegion {
log.Log.Info("computervision.main.ProcessMotion(): continuous recording enabled and no motion region configured, so no motion detection required.")
} else {
log.Log.Info("computervision.main.ProcessMotion(): motion detected is enabled, so starting the motion detection.")
if continuousMode {
log.Log.Info("computervision.main.ProcessMotion(): continuous recording enabled with a motion region, running motion detection for live-view visualisation only (no motion-triggered recording).")
} else {
log.Log.Info("computervision.main.ProcessMotion(): motion detected is enabled, so starting the motion detection.")
}
hubKey := config.HubKey
deviceKey := config.Key
@@ -100,20 +124,47 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
}
}
// Frame dimensions + the motion region polygon(s) in image space, shipped
// with each motion event so the live view can draw a motion-debug overlay
// (the boxes below + the detection region).
var imageCols, imageRows int
var regionPolygons [][]map[string]int
if config.Region != nil {
for _, polygon := range config.Region.Polygon {
var pts []map[string]int
for _, c := range polygon.Coordinates {
pts = append(pts, map[string]int{
"x": int(c.X * baseWidthRatio),
"y": int(c.Y * baseHeightRatio),
})
}
if len(pts) > 0 {
regionPolygons = append(regionPolygons, pts)
}
}
}
img := imageArray[0]
var coordinatesToCheck []int
var coordinatesPerRegion [][]int
totalCoordinates := 0
if img != nil {
bounds := img.Bounds()
rows := bounds.Dy()
cols := bounds.Dx()
imageCols = cols
imageRows = rows
// Make fixed size array of uinty8
// Build a SEPARATE coordinate list per region. Motion is evaluated
// independently per region: pixels are NOT shared between regions, so
// the threshold must be exceeded within a single region to trigger.
coordinatesPerRegion = make([][]int, len(polyObjects))
for y := 0; y < rows; y++ {
for x := 0; x < cols; x++ {
for _, poly := range polyObjects {
point := geo.NewPoint(float64(x), float64(y))
point := geo.NewPoint(float64(x), float64(y))
for idx, poly := range polyObjects {
if poly.Contains(point) {
coordinatesToCheck = append(coordinatesToCheck, y*cols+x)
coordinatesPerRegion[idx] = append(coordinatesPerRegion[idx], y*cols+x)
totalCoordinates++
}
}
}
@@ -121,7 +172,7 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
}
// If no region is set, we'll skip the motion detection
if len(coordinatesToCheck) > 0 {
if totalCoordinates > 0 {
// Start the motion detection
i := 0
@@ -146,12 +197,16 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
log.Log.Debug("computervision.main.ProcessMotion(): " + err.Error() + ".")
}
if config.Capture.Motion != "false" {
// Run detection when motion is enabled, OR when we're in continuous
// mode with a region: there config.Capture.Motion (the motion-RECORDING
// switch) is irrelevant, so the configured region alone is enough to
// emit motion events for the live-view overlay.
if config.Capture.Motion != "false" || continuousMode {
if detectMotion {
// Remember additional information about the result of findmotion
isPixelChangeThresholdReached, changesToReturn, motionRectangle = FindMotion(imageArray, coordinatesToCheck, pixelThreshold)
isPixelChangeThresholdReached, changesToReturn, motionRectangle, motionRectangles = FindMotion(imageArray, coordinatesPerRegion, pixelThreshold)
if isPixelChangeThresholdReached {
// If offline mode is disabled, send a message to the hub
@@ -164,7 +219,24 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
DeviceId: configuration.Config.Key,
Value: map[string]interface{}{
"timestamp": time.Now().Unix(),
},
// Live-view motion-debug overlay data. The boxes/region
// are in the MOTION frame's pixel space (width/height =
// the stream motion ran on, i.e. the sub stream when
// set). mainWidth/mainHeight are the MAIN stream's
// dimensions so the live view can extrapolate the
// boxes/region onto the high-res main view it shows —
// we know both, so no guessing from the <video> element.
"width": imageCols,
"height": imageRows,
"mainWidth": configuration.Config.Capture.IPCamera.Width,
"mainHeight": configuration.Config.Capture.IPCamera.Height,
"regions": motionRectangles,
"polygon": regionPolygons, // Motion sensitivity = the pixel-change threshold that must
// be exceeded before motion triggers. The live view renders
// a reference square of sqrt(threshold) px (in this MOTION
// frame's pixel space) so the user can visually gauge how
// large a moving object must be before it is detected.
"pixelChangeThreshold": pixelThreshold, },
},
}
payload, err := models.PackageMQTTMessage(configuration, message)
@@ -179,7 +251,12 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
}
}
if config.Capture.Recording != "false" {
// Trigger motion-based recording — but NOT in continuous mode:
// there the recorder runs the continuous branch and does not
// drain HandleMotion, so a (blocking) send would hang the motion
// loop. In continuous mode we only publish the motion event above
// for the live-view overlay.
if config.Capture.Recording != "false" && !continuousMode {
dataToPass := models.MotionDataPartial{
Timestamp: time.Now().Unix(),
NumberOfChanges: changesToReturn,
@@ -205,17 +282,73 @@ func ProcessMotion(motionCursor *packets.QueueCursor, configuration *models.Conf
log.Log.Debug("computervision.main.ProcessMotion(): stop the motion detection.")
}
func FindMotion(imageArray [3]*image.Gray, coordinatesToCheck []int, pixelChangeThreshold int) (thresholdReached bool, changesDetected int, motionRectangle models.MotionRectangle) {
func FindMotion(imageArray [3]*image.Gray, coordinatesPerRegion [][]int, pixelChangeThreshold int) (thresholdReached bool, changesDetected int, motionRectangle models.MotionRectangle, motionRectangles []models.MotionRectangle) {
image1 := imageArray[0]
image2 := imageArray[1]
image3 := imageArray[2]
threshold := 60
changes, motionRectangle := AbsDiffBitwiseAndThreshold(image1, image2, image3, threshold, coordinatesToCheck)
return changes > pixelChangeThreshold, changes, motionRectangle
// Evaluate each region INDEPENDENTLY — pixels are not shared between regions,
// so the threshold must be exceeded within a single region to trigger. The
// overall rectangle (recording metadata) and the per-cluster rectangles
// (live-view overlay) are aggregated across all regions.
var combinedRectangles []models.MotionRectangle
var overall models.MotionRectangle
haveOverall := false
totalChanges := 0
for _, coordinatesToCheck := range coordinatesPerRegion {
if len(coordinatesToCheck) == 0 {
continue
}
changes, rect, rects := AbsDiffBitwiseAndThreshold(image1, image2, image3, threshold, coordinatesToCheck)
totalChanges += changes
if changes > pixelChangeThreshold {
thresholdReached = true
}
combinedRectangles = append(combinedRectangles, rects...)
if changes > 0 {
if !haveOverall {
overall = rect
haveOverall = true
} else {
overall = unionMotionRectangle(overall, rect)
}
}
}
return thresholdReached, totalChanges, overall, combinedRectangles
}
func AbsDiffBitwiseAndThreshold(img1 *image.Gray, img2 *image.Gray, img3 *image.Gray, threshold int, coordinatesToCheck []int) (int, models.MotionRectangle) {
// unionMotionRectangle returns the smallest rectangle that contains both a and b.
func unionMotionRectangle(a, b models.MotionRectangle) models.MotionRectangle {
minX := a.X
if b.X < minX {
minX = b.X
}
minY := a.Y
if b.Y < minY {
minY = b.Y
}
maxX := a.X + a.Width
if b.X+b.Width > maxX {
maxX = b.X + b.Width
}
maxY := a.Y + a.Height
if b.Y+b.Height > maxY {
maxY = b.Y + b.Height
}
return models.MotionRectangle{
X: minX,
Y: minY,
Width: maxX - minX,
Height: maxY - minY,
}
}
func AbsDiffBitwiseAndThreshold(img1 *image.Gray, img2 *image.Gray, img3 *image.Gray, threshold int, coordinatesToCheck []int) (int, models.MotionRectangle, []models.MotionRectangle) {
changes := 0
cols := img1.Bounds().Dx()
rows := img1.Bounds().Dy()
var pixelList [][]int
for i := 0; i < len(coordinatesToCheck); i++ {
pixel := coordinatesToCheck[i]
@@ -224,7 +357,7 @@ func AbsDiffBitwiseAndThreshold(img1 *image.Gray, img2 *image.Gray, img3 *image.
if (diff > threshold || diff < -threshold) && (diff2 > threshold || diff2 < -threshold) {
changes++
// Store the pixel coordinates where the change is detected
pixelList = append(pixelList, []int{pixel % img1.Bounds().Dx(), pixel / img1.Bounds().Dx()})
pixelList = append(pixelList, []int{pixel % cols, pixel / cols})
}
}
@@ -258,5 +391,118 @@ func AbsDiffBitwiseAndThreshold(img1 *image.Gray, img2 *image.Gray, img3 *image.
}
log.Log.Debugf("Motion rectangle: %+v", motionRectangle)
}
return changes, motionRectangle
// Cluster the changed pixels into separate bounding boxes so the live view can
// visualise WHERE motion happened (a single overall rectangle is useless when
// two objects move in opposite corners). Cheap grid-based connected components.
motionRectangles := clusterMotionRectangles(pixelList, cols, rows)
return changes, motionRectangle, motionRectangles
}
// clusterMotionRectangles groups the changed-pixel coordinates into a handful of
// bounding boxes using connected-components on a coarse grid (8-connectivity).
// It is intentionally lightweight — it runs only when the motion threshold is
// reached and the boxes are meant for a debug overlay, not precise detection.
func clusterMotionRectangles(pixelList [][]int, cols, rows int) []models.MotionRectangle {
if len(pixelList) == 0 || cols <= 0 || rows <= 0 {
return nil
}
// ~40 cells across the longest side keeps the grid small (cheap to cluster)
// while still separating distinct motion blobs.
const gridDim = 40
cellW := cols / gridDim
if cellW < 1 {
cellW = 1
}
cellH := rows / gridDim
if cellH < 1 {
cellH = 1
}
gCols := (cols + cellW - 1) / cellW
gRows := (rows + cellH - 1) / cellH
grid := make([]bool, gCols*gRows)
for _, p := range pixelList {
cx := p[0] / cellW
cy := p[1] / cellH
if cx >= 0 && cx < gCols && cy >= 0 && cy < gRows {
grid[cy*gCols+cx] = true
}
}
visited := make([]bool, gCols*gRows)
var rectangles []models.MotionRectangle
const maxBoxes = 12
stack := make([][2]int, 0, 64)
for cy := 0; cy < gRows; cy++ {
for cx := 0; cx < gCols; cx++ {
idx := cy*gCols + cx
if !grid[idx] || visited[idx] {
continue
}
// Flood-fill this component (8-connectivity) and track its extent.
minX, minY, maxX, maxY := cx, cy, cx, cy
cellCount := 0
stack = stack[:0]
stack = append(stack, [2]int{cx, cy})
visited[idx] = true
for len(stack) > 0 {
cur := stack[len(stack)-1]
stack = stack[:len(stack)-1]
ccx, ccy := cur[0], cur[1]
cellCount++
if ccx < minX {
minX = ccx
}
if ccy < minY {
minY = ccy
}
if ccx > maxX {
maxX = ccx
}
if ccy > maxY {
maxY = ccy
}
for dy := -1; dy <= 1; dy++ {
for dx := -1; dx <= 1; dx++ {
nx, ny := ccx+dx, ccy+dy
if nx < 0 || ny < 0 || nx >= gCols || ny >= gRows {
continue
}
nIdx := ny*gCols + nx
if grid[nIdx] && !visited[nIdx] {
visited[nIdx] = true
stack = append(stack, [2]int{nx, ny})
}
}
}
}
// Skip single-cell specks (sensor noise) unless it's the only motion.
if cellCount < 2 && len(pixelList) > 4 {
continue
}
x := minX * cellW
y := minY * cellH
w := (maxX - minX + 1) * cellW
h := (maxY - minY + 1) * cellH
if x+w > cols {
w = cols - x
}
if y+h > rows {
h = rows - y
}
rectangles = append(rectangles, models.MotionRectangle{X: x, Y: y, Width: w, Height: h})
if len(rectangles) >= maxBoxes {
return rectangles
}
}
}
return rectangles
}

View File

@@ -401,7 +401,7 @@ func applyAgentEnvVars(configuration *models.Configuration, prefix string, apply
case "AGENT_CAPTURE_PIXEL_CHANGE":
count, err := strconv.Atoi(value)
if err == nil {
configuration.Config.Capture.PixelChangeThreshold = count
configuration.Config.Capture.PixelChangeThreshold = &count
}
break
case "AGENT_CAPTURE_FRAGMENTED":

View File

@@ -19,6 +19,44 @@ 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"`
RTSPStreams []RTSPStream `json:"rtsp_streams,omitempty" bson:"rtsp_streams"`
IsCamera bool `json:"is_camera" bson:"is_camera"`
// IsAudio marks audio-only devices (e.g. IP speakers / intercoms such as
// TOA) that expose RTSP to receive/stream audio rather than video.
IsAudio bool `json:"is_audio" bson:"is_audio"`
}
// RTSPStream is a candidate RTSP stream URL for a discovered camera, derived
// from a built-in brand -> RTSP path mapping. When Verified is true the path was
// confirmed to exist on the device via an unauthenticated RTSP DESCRIBE probe
// (a 200 OK or a 401/403 "auth required" both prove the path is valid).
type RTSPStream struct {
Brand string `json:"brand,omitempty" bson:"brand"`
Stream string `json:"stream,omitempty" bson:"stream"` // "main" or "sub"
Path string `json:"path" bson:"path"`
URL string `json:"url" bson:"url"`
Verified bool `json:"verified" bson:"verified"`
RequiresAuth bool `json:"requires_auth,omitempty" bson:"requires_auth"`
}
type OnvifPanTilt struct {
OnvifCredentials OnvifCredentials `json:"onvif_credentials,omitempty" bson:"onvif_credentials"`
Pan float64 `json:"pan,omitempty" bson:"pan"`

View File

@@ -52,12 +52,30 @@ type Communication struct {
// recorder keeps recording (it does not auto-close on the post-recording
// timeout) until the viewer stops it again. It is independent of motion
// detection so it also works when nothing is moving.
IsRecordingManual *abool.AtomicBool
Queue *packets.Queue
SubQueue *packets.Queue
Image string
CameraConnected bool
MainStreamConnected bool
SubStreamConnected bool
HasBackChannel bool
IsRecordingManual *abool.AtomicBool
// RecordingManualHeartbeat holds the unix-milliseconds timestamp of the last
// heartbeat received from the live view while a manual recording is active.
// The frontend re-sends the record command every few seconds while the user
// stays on the page; if the heartbeats stop (the viewer closed the tab, went
// idle or lost connectivity) the recorder auto-stops the manual recording so
// it can't record forever when the "stop" message never arrives.
RecordingManualHeartbeat *atomic.Int64
// RecordingManualStart holds the unix-milliseconds timestamp at which the
// current manual recording started. It bounds a manual recording to a maximum
// duration (see capture.manualRecordingMaxDuration) so a forgotten record
// button can't record indefinitely even while the viewer stays active.
RecordingManualStart *atomic.Int64
// RecordingManualHeartbeatSeen is set once the current manual recording has
// received at least one heartbeat, i.e. the viewer proved it supports
// heartbeating. Only then does the recorder enforce the heartbeat timeout; a
// viewer that starts a recording but never heartbeats (an older frontend)
// still records up to the max-duration cap instead of being cut off early.
RecordingManualHeartbeatSeen *abool.AtomicBool
Queue *packets.Queue
SubQueue *packets.Queue
Image string
CameraConnected bool
MainStreamConnected bool
SubStreamConnected bool
HasBackChannel bool
}

View File

@@ -74,7 +74,7 @@ type Capture struct {
ForwardWebRTC string `json:"forwardwebrtc"`
Fragmented string `json:"fragmented,omitempty" bson:"fragmented,omitempty"`
FragmentedDuration int64 `json:"fragmentedduration,omitempty" bson:"fragmentedduration,omitempty"`
PixelChangeThreshold int `json:"pixelChangeThreshold,omitempty"`
PixelChangeThreshold *int `json:"pixelChangeThreshold,omitempty"`
}
// IPCamera configuration, such as the RTSP url of the IPCamera and the FPS.

View File

@@ -154,6 +154,15 @@ type RecordPayload struct {
// recording (and keeps it running), false stops it. Older clients that only
// send a timestamp default to false; the live view always sets it explicitly.
Recording bool `json:"recording"`
// Heartbeat marks a keep-alive re-send (with Recording=true) from a viewer
// that supports heartbeating, as opposed to the initial start (the record
// button). While a user stays on the page the live view re-sends the record
// command every few seconds; the agent uses this flag to (a) refresh the
// recording's keep-alive without restarting an already auto-stopped clip from
// a stray heartbeat, and (b) only enable the heartbeat-timeout auto-stop once
// it has actually seen a heartbeat — so older viewers that never heartbeat
// still record up to the max-duration cap instead of being cut off early.
Heartbeat bool `json:"heartbeat"`
}
// We received a preset position request, we'll request it through onvif and send it back.
@@ -166,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.

View File

@@ -0,0 +1,523 @@
package onvif
import (
"bufio"
"net"
"strconv"
"strings"
"time"
"github.com/kerberos-io/agent/machinery/src/models"
)
// brandProfile describes a camera brand together with the RTSP URL path
// templates it exposes for its main (high quality) and sub (low quality)
// streams. The paths are the well-known, widely documented defaults for each
// vendor and are used both to identify the brand (by probing which path the
// device recognises) and to pre-fill a working RTSP URL for the user.
//
// The order of the list matters: more specific / more common brands come first
// so that when we actively probe a device the first matching profile wins.
type brandProfile struct {
Brand string
// aliases are lower-cased tokens that, when seen in a banner/realm/MAC
// vendor, map onto this brand.
Aliases []string
MainPath string
SubPath string
// extraMainPaths are alternative main-stream paths tried during active
// probing when the primary MainPath is not recognised.
extraMainPaths []string
}
// brandProfiles is the built-in brand -> RTSP path mapping. It mirrors the
// tables used by tools such as ONVIF Device Manager, iSpy/Agent DVR and
// Blue Iris.
var brandProfiles = []brandProfile{
{
Brand: "Hikvision",
Aliases: []string{"hikvision", "dvrdvs", "ds-", "hik"},
MainPath: "/Streaming/Channels/101",
SubPath: "/Streaming/Channels/102",
extraMainPaths: []string{"/h264/ch1/main/av_stream", "/ISAPI/Streaming/Channels/101"},
},
{
Brand: "Dahua",
Aliases: []string{"dahua", "dh-"},
MainPath: "/cam/realmonitor?channel=1&subtype=0",
SubPath: "/cam/realmonitor?channel=1&subtype=1",
extraMainPaths: []string{"/live"},
},
{
Brand: "Amcrest",
Aliases: []string{"amcrest"},
MainPath: "/cam/realmonitor?channel=1&subtype=0",
SubPath: "/cam/realmonitor?channel=1&subtype=1",
},
{
Brand: "Axis",
Aliases: []string{"axis"},
MainPath: "/axis-media/media.amp",
SubPath: "/axis-media/media.amp?videocodec=h264&resolution=640x480",
extraMainPaths: []string{"/mpeg4/media.amp"},
},
{
Brand: "Reolink",
Aliases: []string{"reolink", "rlc", "rln", "rlc-", "rln-", "trackmix", "duo"},
MainPath: "/h264Preview_01_main",
SubPath: "/h264Preview_01_sub",
extraMainPaths: []string{"/Preview_01_main"},
},
{
Brand: "Hanwha",
Aliases: []string{"hanwha", "wisenet", "samsung techwin"},
MainPath: "/profile2/media.smp",
SubPath: "/profile3/media.smp",
extraMainPaths: []string{"/profile1/media.smp", "/onvif/profile2/media.smp"},
},
{
Brand: "Bosch",
Aliases: []string{"bosch"},
MainPath: "/rtsp_tunnel",
SubPath: "/rtsp_tunnel?inst=2",
extraMainPaths: []string{"/rtsp_tunnel?inst=1", "/?inst=1"},
},
{
Brand: "Vivotek",
Aliases: []string{"vivotek"},
MainPath: "/live.sdp",
SubPath: "/live2.sdp",
extraMainPaths: []string{"/live1s1.sdp"},
},
{
Brand: "Foscam",
Aliases: []string{"foscam"},
MainPath: "/videoMain",
SubPath: "/videoSub",
},
{
Brand: "Uniview",
Aliases: []string{"uniview", "unv"},
MainPath: "/media/video1",
SubPath: "/media/video2",
extraMainPaths: []string{"/unicast/c1/s0/live", "/unicast/c1/s1/live"},
},
{
Brand: "TP-Link",
Aliases: []string{"tp-link", "tplink", "tapo"},
MainPath: "/stream1",
SubPath: "/stream2",
},
{
Brand: "Mobotix",
Aliases: []string{"mobotix"},
MainPath: "/cam0/mjpeg",
SubPath: "/cam1/mjpeg",
extraMainPaths: []string{"/live.sdp"},
},
{
Brand: "Ubiquiti",
Aliases: []string{"ubiquiti", "unifi"},
MainPath: "/s0",
SubPath: "/s1",
extraMainPaths: []string{"/live/ch00_0"},
},
{
Brand: "Panasonic",
Aliases: []string{"panasonic", "i-pro", "ipro"},
MainPath: "/MediaInput/h264",
SubPath: "/MediaInput/h264/stream_2",
},
{
Brand: "Sony",
Aliases: []string{"sony"},
MainPath: "/media/video1",
SubPath: "/media/video2",
},
{
// D-Link mydlink IP cameras. Older models stream MJPEG over HTTP; the
// RTSP-capable ones expose SDP-named streams, newer DCS models use
// "/live/profile.0".
Brand: "D-Link",
Aliases: []string{"d-link", "dlink", "dcs-", "dcs"},
MainPath: "/live1.sdp",
SubPath: "/live2.sdp",
extraMainPaths: []string{"/live.sdp", "/live/profile.0", "/play1.sdp"},
},
{
// TRENDnet. Newer PoE bullet/dome models (TV-IPxxxPI) use a
// Hikvision-style path; older ones expose SDP streams.
Brand: "Trendnet",
Aliases: []string{"trendnet", "tv-ip"},
MainPath: "/Streaming/Channels/101",
SubPath: "/Streaming/Channels/102",
extraMainPaths: []string{"/play1.sdp", "/play2.sdp", "/ch0_0.h264", "/live/av0"},
},
{
// Lorex is built largely on Dahua hardware, so it shares Dahua's
// realmonitor path scheme.
Brand: "Lorex",
Aliases: []string{"lorex"},
MainPath: "/cam/realmonitor?channel=1&subtype=0",
SubPath: "/cam/realmonitor?channel=1&subtype=1",
extraMainPaths: []string{"/ch01/0"},
},
{
// Honeywell ships both Dahua-OEM models (realmonitor) and in-house
// firmwares exposing "/h264" or "/media".
Brand: "Honeywell",
Aliases: []string{"honeywell"},
MainPath: "/cam/realmonitor?channel=1&subtype=0",
SubPath: "/cam/realmonitor?channel=1&subtype=1",
extraMainPaths: []string{"/h264", "/media", "/live.sdp"},
},
{
Brand: "Pelco",
Aliases: []string{"pelco"},
MainPath: "/stream1",
SubPath: "/stream2",
extraMainPaths: []string{"/1/stream1"},
},
{
// TOA network audio devices (IP horn speakers / intercoms, banner
// "TOA rtsp server") expose their stream through ONVIF rather than a
// documented fixed RTSP path. These ONVIF-style paths are a best-effort
// default; the authoritative URL should come from an ONVIF GetStreamUri
// query with credentials.
Brand: "TOA",
Aliases: []string{"toa"},
MainPath: "/ONVIF/channel1",
SubPath: "/ONVIF/channel2",
extraMainPaths: []string{"/media/video1", "/live"},
},
{
// Linksys/Cisco IP cameras (e.g. LCAD03FLN, LCAB03VLNOD, LCAM0336OD)
// run a mini_httpd server and expose ONVIF-style stream paths with a
// capitalised "ONVIF" segment (distinct from the generic "/onvif1").
Brand: "Linksys",
Aliases: []string{"linksys", "lcad", "lcab", "lcam", "lcae"},
MainPath: "/ONVIF/channel1",
SubPath: "/ONVIF/channel2",
extraMainPaths: []string{"/img/media.sav", "/live"},
},
}
// genericRTSPPaths are last-resort, vendor-neutral RTSP paths used when the
// brand is unknown. Many ONVIF/embedded cameras answer on one of these.
var genericRTSPPaths = []string{
"/ONVIF/channel1", "/ONVIF/channel2", "/onvif1", "/live", "/live/ch0", "/11", "/12",
"/stream0", "/stream1", "/h264", "/media/video1", "/ch0_0.h264",
}
// brandProfileFor returns the profile whose aliases best match the given brand
// hint (from a banner, realm or MAC vendor). It returns nil when nothing
// matches.
func brandProfileFor(hint string) *brandProfile {
hint = strings.ToLower(strings.TrimSpace(hint))
if hint == "" {
return nil
}
for i := range brandProfiles {
for _, alias := range brandProfiles[i].Aliases {
if strings.Contains(hint, alias) {
return &brandProfiles[i]
}
}
}
return nil
}
// realmBrands maps a lower-cased substring of an RTSP/HTTP WWW-Authenticate
// realm to a manufacturer. The auth realm is one of the most reliable brand
// signals because a camera advertises it even when it refuses every
// unauthenticated request (e.g. Hikvision realm "IP Camera(E3669)", Dahua realm
// "Login to <serial>"). Ordered so the most specific matches win.
var realmBrands = []struct {
Match string
Vendor string
}{
{"login to", "Dahua"},
{"surveillance server", "Dahua"},
{"real time streaming", "Dahua"},
{"dahua", "Dahua"},
{"ip camera(", "Hikvision"},
{"hikvision", "Hikvision"},
{"ds-", "Hikvision"},
{"axis", "Axis"},
{"reolink", "Reolink"},
{"amcrest", "Amcrest"},
{"wisenet", "Hanwha"},
{"hanwha", "Hanwha"},
{"uniview", "Uniview"},
{"tp-link", "TP-Link"},
{"tapo", "TP-Link"},
{"foscam", "Foscam"},
{"vivotek", "Vivotek"},
{"mobotix", "Mobotix"},
{"bosch", "Bosch"},
{"please log in with a valid username", "Bosch"},
{"d-link", "D-Link"},
{"dcs-", "D-Link"},
{"trendnet", "Trendnet"},
{"lorex", "Lorex"},
{"honeywell", "Honeywell"},
{"pelco", "Pelco"},
{"linksys", "Linksys"},
{"lcad", "Linksys"},
{"lcab", "Linksys"},
{"lcam", "Linksys"},
}
// brandFromRealm resolves a manufacturer from an auth realm string.
func brandFromRealm(realm string) string {
r := strings.ToLower(strings.TrimSpace(realm))
if r == "" {
return ""
}
for _, entry := range realmBrands {
if strings.Contains(r, entry.Match) {
return entry.Vendor
}
}
return ""
}
// modelFromRealm extracts a model/device code embedded in an auth realm, e.g.
// Hikvision's realm="IP Camera(E3669)" -> "E3669".
func modelFromRealm(realm string) string {
open := strings.Index(realm, "(")
closeIdx := strings.Index(realm, ")")
if open >= 0 && closeIdx > open+1 {
return strings.TrimSpace(realm[open+1 : closeIdx])
}
return ""
}
// firstNonEmpty returns the first non-blank value.
func firstNonEmpty(values ...string) string {
for _, value := range values {
if strings.TrimSpace(value) != "" {
return value
}
}
return ""
}
// guessRTSPStreams determines the most likely RTSP stream URLs for a camera. It
// combines the brand hint discovered from banners/MAC with an active,
// unauthenticated RTSP DESCRIBE probe and the auth realm advertised by the
// device.
//
// Detection strategy (most reliable first):
// 1. Send a control DESCRIBE for a random, non-existent path. Its 401 response
// usually carries a WWW-Authenticate realm that reveals the brand
// (Hikvision "IP Camera(...)", Dahua "Login to ..."). The realm is the
// strongest signal and works even when the device challenges auth for every
// request. The control also tells us whether the device distinguishes valid
// from invalid paths.
// 2. If the device discriminates paths, probe each brand's main path (realm
// brand first); the first the device recognises (200 or 401/403) confirms a
// working URL.
// 3. Otherwise fall back to the realm / hint / port brand's default paths and
// return them as unverified suggestions.
//
// It returns the detected brand, an optional model code parsed from the realm,
// and the ordered list of candidate streams (verified first).
func guessRTSPStreams(ip string, port int, brandHint string, openPorts []int, timeout time.Duration) (brand string, model string, streams []models.RTSPStream) {
base := "rtsp://" + net.JoinHostPort(ip, strconv.Itoa(port))
build := func(profileBrand, stream, path string, verified, requiresAuth bool) models.RTSPStream {
return models.RTSPStream{
Brand: profileBrand,
Stream: stream,
Path: path,
URL: base + path,
Verified: verified,
RequiresAuth: requiresAuth,
}
}
// 1) Control probe: distinguish behaviour + capture the auth realm.
bogusPath := "/kerberos-probe-" + strconv.FormatInt(time.Now().UnixNano(), 36)
controlStatus, controlRealm, _ := rtspDescribe(ip, port, bogusPath, timeout)
controlExists := controlStatus == 200 || controlStatus == 401 || controlStatus == 403
controlAuth := controlStatus == 401 || controlStatus == 403
discriminates := !controlExists
realmBrand := brandFromRealm(controlRealm)
model = modelFromRealm(controlRealm)
// The realm brand (when present) is authoritative and probed first.
primaryHint := firstNonEmpty(realmBrand, brandHint)
var verified []models.RTSPStream
var unverified []models.RTSPStream
detected := ""
// 2) Trustworthy active per-brand probing (device discriminates paths).
if discriminates {
for _, profile := range orderedProfiles(primaryHint) {
mainCandidates := append([]string{profile.MainPath}, profile.extraMainPaths...)
matchedMain := ""
matchedAuth := false
for _, path := range mainCandidates {
ok, requiresAuth := rtspPathExists(ip, port, path, timeout)
if ok {
matchedMain = path
matchedAuth = requiresAuth
break
}
}
if matchedMain == "" {
continue
}
detected = profile.Brand
verified = append(verified, build(profile.Brand, "main", matchedMain, true, matchedAuth))
if profile.SubPath != "" {
subOK, subAuth := rtspPathExists(ip, port, profile.SubPath, timeout)
verified = append(verified, build(profile.Brand, "sub", profile.SubPath, subOK, subAuth || matchedAuth))
}
break
}
}
// 3) Fall back to unverified suggestions from realm / hint / port signals.
if len(verified) == 0 {
profile := brandProfileFor(primaryHint)
if profile == nil {
profile = brandProfileForPorts(openPorts)
}
if profile != nil {
detected = profile.Brand
unverified = append(unverified, build(profile.Brand, "main", profile.MainPath, false, controlAuth))
if profile.SubPath != "" {
unverified = append(unverified, build(profile.Brand, "sub", profile.SubPath, false, controlAuth))
}
} else {
for _, path := range genericRTSPPaths {
unverified = append(unverified, build("Generic", "main", path, false, controlAuth))
}
}
}
// The realm brand always wins for the manufacturer name.
if realmBrand != "" {
detected = realmBrand
}
return detected, model, append(verified, unverified...)
}
// brandProfileForPorts derives a brand from vendor-specific control ports that
// were found open during the scan (used when banners give no hint).
func brandProfileForPorts(openPorts []int) *brandProfile {
if containsInt(openPorts, 37777) {
return brandProfileByName("Dahua")
}
return nil
}
// brandProfileByName returns the profile with the given brand name (nil when
// absent).
func brandProfileByName(name string) *brandProfile {
for i := range brandProfiles {
if brandProfiles[i].Brand == name {
return &brandProfiles[i]
}
}
return nil
}
// orderedProfiles returns the brand profiles with the profile matching the
// brand hint (if any) moved to the front so it is probed first.
func orderedProfiles(brandHint string) []brandProfile {
match := brandProfileFor(brandHint)
if match == nil {
return brandProfiles
}
ordered := make([]brandProfile, 0, len(brandProfiles))
ordered = append(ordered, *match)
for i := range brandProfiles {
if brandProfiles[i].Brand != match.Brand {
ordered = append(ordered, brandProfiles[i])
}
}
return ordered
}
// rtspDescribe sends an unauthenticated RTSP DESCRIBE for the given path and
// returns the response status code together with the WWW-Authenticate realm and
// Server header (when present). status is 0 when the device does not answer.
func rtspDescribe(ip string, port int, path string, timeout time.Duration) (status int, realm string, server string) {
address := net.JoinHostPort(ip, strconv.Itoa(port))
conn, err := net.DialTimeout("tcp", address, timeout)
if err != nil {
return 0, "", ""
}
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(timeout))
request := "DESCRIBE rtsp://" + address + path + " RTSP/1.0\r\n" +
"CSeq: 1\r\n" +
"User-Agent: KerberosDiscovery\r\n" +
"Accept: application/sdp\r\n\r\n"
if _, err := conn.Write([]byte(request)); err != nil {
return 0, "", ""
}
status, headers := readRTSPResponse(conn)
return status, parseRealm(headers["www-authenticate"]), headers["server"]
}
// rtspPathExists reports whether the device recognises the given RTSP path. A
// 200 OK means the path is publicly accessible; a 401/403 means the path is
// valid but requires credentials (still a positive match). Any other status
// (404, 400, 455, ...) means the path is not recognised.
func rtspPathExists(ip string, port int, path string, timeout time.Duration) (exists bool, requiresAuth bool) {
status, _, _ := rtspDescribe(ip, port, path, timeout)
switch status {
case 200:
return true, false
case 401, 403:
return true, true
default:
return false, false
}
}
// readRTSPResponse reads and parses the status code and headers of an RTSP
// response. Only the first occurrence of each header is kept.
func readRTSPResponse(conn net.Conn) (status int, headers map[string]string) {
headers = make(map[string]string)
reader := bufio.NewReader(conn)
line, err := reader.ReadString('\n')
if err != nil {
return 0, headers
}
fields := strings.Fields(line)
if len(fields) >= 2 && strings.HasPrefix(strings.ToUpper(fields[0]), "RTSP/") {
status, _ = strconv.Atoi(fields[1])
}
for {
hline, err := reader.ReadString('\n')
if err != nil {
break
}
hline = strings.TrimRight(hline, "\r\n")
if hline == "" {
break
}
idx := strings.Index(hline, ":")
if idx <= 0 {
continue
}
key := strings.ToLower(strings.TrimSpace(hline[:idx]))
value := strings.TrimSpace(hline[idx+1:])
if _, exists := headers[key]; !exists {
headers[key] = value
}
}
return status, headers
}

View File

@@ -0,0 +1,171 @@
package onvif
import (
"bufio"
"net"
"strconv"
"strings"
"testing"
"time"
)
// mockRTSPServer starts a TCP listener that answers RTSP DESCRIBE requests. For
// each incoming request it extracts the path and calls respond(path) to obtain
// the numeric status code and optional auth realm to return. It returns the
// listener host, port and a cleanup function.
func mockRTSPServer(t *testing.T, respond func(path string) (int, string)) (string, int, func()) {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("failed to start mock RTSP server: %v", err)
}
go func() {
for {
conn, err := listener.Accept()
if err != nil {
return
}
go func(c net.Conn) {
defer c.Close()
_ = c.SetDeadline(time.Now().Add(2 * time.Second))
reader := bufio.NewReader(c)
line, err := reader.ReadString('\n')
if err != nil {
return
}
path := ""
fields := strings.Fields(line)
if len(fields) >= 2 {
url := fields[1]
url = strings.TrimPrefix(url, "rtsp://")
if idx := strings.Index(url, "/"); idx >= 0 {
path = url[idx:]
}
}
status, realm := respond(path)
reason := map[int]string{200: "OK", 401: "Unauthorized", 404: "Not Found"}[status]
response := "RTSP/1.0 " + strconv.Itoa(status) + " " + reason + "\r\nCSeq: 1\r\n"
if realm != "" {
response += "WWW-Authenticate: Digest realm=\"" + realm + "\", nonce=\"abc\"\r\n"
}
response += "\r\n"
_, _ = c.Write([]byte(response))
}(conn)
}
}()
host, portStr, _ := net.SplitHostPort(listener.Addr().String())
port, _ := strconv.Atoi(portStr)
return host, port, func() { listener.Close() }
}
// TestGuessRTSPStreams_DiscriminatingHikvision verifies that a device which
// distinguishes valid from invalid paths (returning 401 only for the Hikvision
// path) is correctly identified as Hikvision with a confirmed main/sub stream.
func TestGuessRTSPStreams_DiscriminatingHikvision(t *testing.T) {
host, port, cleanup := mockRTSPServer(t, func(path string) (int, string) {
if strings.HasPrefix(path, "/Streaming/Channels/") {
return 401, "" // valid path, needs auth
}
return 404, "" // everything else is unknown -> device discriminates
})
defer cleanup()
brand, _, streams := guessRTSPStreams(host, port, "", nil, 2*time.Second)
if brand != "Hikvision" {
t.Fatalf("expected brand Hikvision, got %q", brand)
}
if len(streams) == 0 || !streams[0].Verified {
t.Fatalf("expected a verified main stream, got %+v", streams)
}
if !streams[0].RequiresAuth {
t.Errorf("expected main stream to require auth")
}
if streams[0].Path != "/Streaming/Channels/101" {
t.Errorf("expected main path /Streaming/Channels/101, got %q", streams[0].Path)
}
}
// TestGuessRTSPStreams_ChallengesEverything verifies that a device which returns
// 401 for *any* path (including a bogus one) does NOT get mis-detected via path
// probing, and instead falls back to the port hint (Dahua control port 37777)
// with unverified suggestions.
func TestGuessRTSPStreams_ChallengesEverything(t *testing.T) {
host, port, cleanup := mockRTSPServer(t, func(path string) (int, string) {
return 401, "" // challenges auth before checking the path, no realm
})
defer cleanup()
brand, _, streams := guessRTSPStreams(host, port, "", []int{37777}, 2*time.Second)
if brand != "Dahua" {
t.Fatalf("expected fallback brand Dahua from port hint, got %q", brand)
}
if len(streams) == 0 {
t.Fatalf("expected suggested streams, got none")
}
if streams[0].Verified {
t.Errorf("expected unverified suggestion for a non-discriminating device")
}
if streams[0].Path != "/cam/realmonitor?channel=1&subtype=0" {
t.Errorf("expected Dahua main path, got %q", streams[0].Path)
}
}
// TestGuessRTSPStreams_RealmDetectsHikvision verifies that a device which
// challenges auth for every path (so path probing cannot help) is still
// identified from its RTSP auth realm, and the model code is extracted.
func TestGuessRTSPStreams_RealmDetectsHikvision(t *testing.T) {
host, port, cleanup := mockRTSPServer(t, func(path string) (int, string) {
return 401, "IP Camera(E3669)" // Hikvision realm signature, 401 for all paths
})
defer cleanup()
brand, model, streams := guessRTSPStreams(host, port, "", nil, 2*time.Second)
if brand != "Hikvision" {
t.Fatalf("expected brand Hikvision from realm, got %q", brand)
}
if model != "E3669" {
t.Errorf("expected model E3669 from realm, got %q", model)
}
if len(streams) == 0 || streams[0].Path != "/Streaming/Channels/101" {
t.Fatalf("expected Hikvision default main path, got %+v", streams)
}
if !streams[0].RequiresAuth {
t.Errorf("expected the suggestion to be marked auth-required")
}
}
// TestGuessRTSPStreams_RealmDetectsDahua verifies Dahua detection from its
// "Login to ..." realm.
func TestGuessRTSPStreams_RealmDetectsDahua(t *testing.T) {
host, port, cleanup := mockRTSPServer(t, func(path string) (int, string) {
return 401, "Login to 5df61a6057b10cc99d471769516d3c11"
})
defer cleanup()
brand, _, streams := guessRTSPStreams(host, port, "", nil, 2*time.Second)
if brand != "Dahua" {
t.Fatalf("expected brand Dahua from realm, got %q", brand)
}
if len(streams) == 0 || streams[0].Path != "/cam/realmonitor?channel=1&subtype=0" {
t.Fatalf("expected Dahua default main path, got %+v", streams)
}
}
// TestGuessRTSPStreams_UnknownFallsBackToGeneric verifies that an unknown device
// (discriminating but matching no brand) yields generic suggestions.
func TestGuessRTSPStreams_UnknownFallsBackToGeneric(t *testing.T) {
host, port, cleanup := mockRTSPServer(t, func(path string) (int, string) {
return 404, "" // discriminates, but nothing matches
})
defer cleanup()
brand, _, streams := guessRTSPStreams(host, port, "", nil, 2*time.Second)
if brand != "" {
t.Fatalf("expected no detected brand, got %q", brand)
}
if len(streams) == 0 || streams[0].Brand != "Generic" {
t.Fatalf("expected generic suggestions, got %+v", streams)
}
}

View File

@@ -0,0 +1,557 @@
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)
// Resolve a hostname now (ONVIF WS-Discovery may already have set
// one; otherwise fall back to reverse DNS). Camera hostnames often
// encode the model (e.g. Reolink "RLC-823S2"), which is a useful
// brand hint when the RTSP/HTTP banners are anonymous.
mutex.Lock()
hostname := ""
if existing, ok := devicesByIP[ip]; ok {
hostname = existing.Hostname
}
mutex.Unlock()
if hostname == "" {
hostname = reverseDNS(ip, dialTimeout)
}
// Guess (and actively confirm) the RTSP stream URLs from a built-in
// brand -> RTSP path mapping when an RTSP port is open.
var rtspPort int
for _, port := range openPorts {
if port == 554 || port == 8554 {
rtspPort = port
break
}
}
// The banner manufacturer is most reliable; fall back to the
// hostname (model code) so devices that only reveal themselves via
// their name (e.g. Reolink RLC-*) still get the right stream paths.
brandHint := fingerprint.Manufacturer
if brandHint == "" {
brandHint = hostname
}
var rtspStreams []models.RTSPStream
detectedBrand := ""
detectedModel := ""
if rtspPort != 0 && !fingerprint.IsAudio {
detectedBrand, detectedModel, rtspStreams = guessRTSPStreams(ip, rtspPort, brandHint, openPorts, dialTimeout)
}
device := upsert(ip)
mutex.Lock()
device.OpenPorts = mergeSortedInts(device.OpenPorts, openPorts)
device.Services = mergeUniqueStrings(device.Services, services)
if hostname != "" && device.Hostname == "" {
device.Hostname = hostname
}
if isCamera || fingerprint.IsCamera {
device.IsCamera = true
}
if fingerprint.IsAudio {
device.IsAudio = true
device.IsCamera = false
}
if fingerprint.Manufacturer != "" {
device.Manufacturer = fingerprint.Manufacturer
}
// A brand derived from the RTSP auth realm, a confirmed path probe or
// a vendor-specific control port is more reliable than a banner
// string, so let it win.
if detectedBrand != "" && detectedBrand != "Generic" {
device.Manufacturer = detectedBrand
device.IsCamera = true
}
if fingerprint.Model != "" {
device.Model = fingerprint.Model
}
if device.Model == "" && detectedModel != "" {
device.Model = detectedModel
}
if fingerprint.Type != "" {
device.Type = fingerprint.Type
}
if fingerprint.Server != "" {
device.Server = fingerprint.Server
}
if len(rtspStreams) > 0 {
device.RTSPStreams = rtspStreams
// Prefer the first verified stream as the primary RTSP URL.
device.RTSPURL = rtspStreams[0].URL
for _, stream := range rtspStreams {
if stream.Verified {
device.RTSPURL = stream.URL
break
}
}
} else if rtspPort != 0 && !fingerprint.IsAudio {
device.RTSPURL = "rtsp://" + ip + ":" + strconv.Itoa(rtspPort) + "/"
}
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
}

View File

@@ -0,0 +1,380 @@
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
// body holds a lower-cased slice of the HTTP landing page, fetched only when
// the banners are anonymous. Rebadged/OEM cameras often reveal their vendor
// there (logo filenames, embedded scripts), e.g. ADI "Capture".
body string
// IsCamera is set when the collected evidence confidently identifies the
// device as a camera, NVR or DVR.
IsCamera bool
// IsAudio is set for audio-only devices (IP speakers / intercoms, e.g. TOA)
// that use RTSP for audio rather than video.
IsAudio 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},
{"linksys", "Linksys", true},
{"d-link", "D-Link", true},
{"dlink", "D-Link", true},
{"trendnet", "Trendnet", true},
{"lorex", "Lorex", true},
{"honeywell", "Honeywell", true},
{"pelco", "Pelco", true},
{"toa rtsp", "TOA", false},
{"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},
}
// bodyVendors maps a distinctive lower-cased substring found in a camera's HTML
// landing page (logo filename, embedded script, product string) to a
// manufacturer. Used only when the RTSP/HTTP banners are anonymous, so it can
// identify rebadged/OEM cameras (e.g. ADI "Capture") that hide their model
// behind a generic "httpd" server and an "RTSP" realm.
var bodyVendors = []struct {
Match string
Vendor string
IsCamera bool
}{
{"logo_white(capture)", "Capture", true},
{"logo_capture", "Capture", 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": {},
"please log in with a valid username.": {},
"please log in with a valid username": {},
}
// 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.
httpPort := 0
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
httpPort = port
break
}
}
// 3) When the banners are anonymous (generic server, no vendor realm), fetch
// a slice of the landing page. Rebadged/OEM cameras (e.g. ADI "Capture")
// only reveal their vendor in the HTML.
if httpPort != 0 && isGenericServer(fp.Server) {
fp.body = httpBody(ip, httpPort, timeout)
}
classifyFingerprint(&fp, openPorts)
return fp
}
// isGenericServer reports whether an HTTP Server header is a generic embedded
// web server that carries no vendor information (so the HTML body is worth a
// look).
func isGenericServer(server string) bool {
s := strings.ToLower(strings.TrimSpace(server))
if s == "" {
return true
}
for _, generic := range []string{"httpd", "webs", "boa", "lighttpd", "nginx", "gsoap", "mini_httpd", "thttpd", "apache"} {
if strings.Contains(s, generic) {
return true
}
}
return false
}
// httpBody issues an unauthenticated HTTP GET / and returns a lower-cased,
// size-bounded slice of the response (headers + body). Best-effort; empty on
// error.
func httpBody(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 := "GET / 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 ""
}
var builder strings.Builder
buf := make([]byte, 4096)
for builder.Len() < 65536 {
n, err := conn.Read(buf)
if n > 0 {
builder.Write(buf[:n])
}
if err != nil {
break
}
}
return strings.ToLower(builder.String())
}
// 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
}
}
}
// Vendor from the HTML landing page when the banners revealed nothing.
// Rebadged/OEM cameras (e.g. ADI "Capture") only identify themselves via
// logo filenames or embedded scripts.
if fp.Manufacturer == "" && fp.body != "" {
for _, entry := range bodyVendors {
if strings.Contains(fp.body, entry.Match) {
fp.Manufacturer = entry.Vendor
if entry.IsCamera {
fp.IsCamera = true
}
break
}
}
}
// 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)
// Audio devices (IP speakers / intercoms) also speak RTSP, but for audio
// rather than video, so classify them separately and never as a camera.
if fp.Manufacturer == "TOA" ||
strings.Contains(haystack, "speaker") ||
strings.Contains(haystack, "sip audio") ||
strings.Contains(haystack, "audio server") {
fp.IsAudio = true
fp.IsCamera = false
fp.Type = "IP Speaker/Audio"
return
}
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
}

View File

@@ -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,19 +23,81 @@ 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"
} else if device.IsAudio {
label = "speaker"
}
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)
// Detail the guessed RTSP stream URLs from the brand -> RTSP mapping.
for _, stream := range device.RTSPStreams {
status := "guess"
if stream.Verified {
status = "confirmed"
}
line := "onvif.Discover(): -> " + stream.Stream + " stream [" + status + "]"
if stream.RequiresAuth {
line += " (auth required)"
}
line += ": " + stream.URL
log.Log.Info(line)
}
}
}

View File

@@ -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

View File

@@ -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)

View File

@@ -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":
@@ -381,23 +386,46 @@ func HandleRecording(mqttClient mqtt.Client, hubKey string, payload models.Paylo
}
if recordPayload.Recording {
// Start a manual recording from the live view (record button). Keep it
// running until the viewer stops it again — the motion recorder honours
// communication.IsRecordingManual and won't auto-close on the
// post-recording timeout while it's set. We also inject a motion event
// so the recording starts immediately, even when nothing is moving.
log.Log.Info("routers.mqtt.main.HandleRecording(): manual recording started.")
communication.IsRecordingManual.Set()
select {
case communication.HandleMotion <- models.MotionDataPartial{Timestamp: timestamp, NumberOfChanges: 100000000}:
default:
log.Log.Warning("routers.mqtt.main.HandleRecording(): motion channel full, manual recording start not queued.")
now := time.Now().UnixMilli()
if recordPayload.Heartbeat {
// Keep-alive from a viewer that supports heartbeats. Only refresh while
// a manual recording is actually running; if it already auto-stopped
// (heartbeat timeout / max duration) we IGNORE it so a stray heartbeat
// can't restart a recording we just ended. Seeing a heartbeat also arms
// the recorder's heartbeat-timeout auto-stop.
if communication.IsRecordingManual.IsSet() {
communication.RecordingManualHeartbeat.Store(now)
communication.RecordingManualHeartbeatSeen.Set()
log.Log.Debug("routers.mqtt.main.HandleRecording(): manual recording heartbeat received.")
} else {
log.Log.Debug("routers.mqtt.main.HandleRecording(): ignoring heartbeat, no active manual recording.")
}
} else {
// Explicit start from the live view (record button). Start a manual
// recording and keep it running — the motion recorder honours
// communication.IsRecordingManual and won't auto-close on the
// post-recording timeout while it's set. We also inject a motion event
// so the recording starts immediately, even when nothing is moving.
communication.RecordingManualHeartbeat.Store(now)
if communication.IsRecordingManual.SetToIf(false, true) {
communication.RecordingManualStart.Store(now)
communication.RecordingManualHeartbeatSeen.UnSet()
log.Log.Info("routers.mqtt.main.HandleRecording(): manual recording started.")
select {
case communication.HandleMotion <- models.MotionDataPartial{Timestamp: timestamp, NumberOfChanges: 100000000}:
default:
log.Log.Warning("routers.mqtt.main.HandleRecording(): motion channel full, manual recording start not queued.")
}
}
}
} else {
// Stop the manual recording; the motion recorder closes the clip once the
// post-recording window elapses.
// post-recording window elapses. Clear the heartbeat/start markers too.
log.Log.Info("routers.mqtt.main.HandleRecording(): manual recording stopped.")
communication.IsRecordingManual.UnSet()
communication.RecordingManualHeartbeat.Store(0)
communication.RecordingManualStart.Store(0)
communication.RecordingManualHeartbeatSeen.UnSet()
}
}
@@ -523,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