[T-006] 完成 Sense 单实机与五路混合源集成验收 #25
+26
-1
@@ -1,6 +1,6 @@
|
||||
# Sense M1 骨架
|
||||
|
||||
本目录是 YoVision Sense 的无实机接入骨架。数据库保存期望态,ONVIF 和 MediaMTX 通过端口隔离;当前只能用确定性 fake、MediaMTX 假 HTTP 服务及可选合成 RTSP 源验证,不能据此宣称任何真实摄像头兼容性。
|
||||
本目录是 YoVision Sense 的 M1 接入骨架。数据库保存期望态,ONVIF 和 MediaMTX 通过端口隔离;默认关闭真实 ONVIF,显式设置 `SENSE_ONVIF_MODE=standard` 后才启用标准 SOAP/WS-Security 适配器。T-006 的真实样机结论仅覆盖已批准的精确海康基线,不能据此宣称多品牌兼容。
|
||||
|
||||
## 常用命令
|
||||
|
||||
@@ -26,6 +26,21 @@ Unix 将构建产物改为 `bin/sense-api`。服务默认监听 `127.0.0.1:8080`
|
||||
| `SENSE_MEDIAMTX_URL` | `http://127.0.0.1:9997` | MediaMTX 控制 API;不得包含 userinfo |
|
||||
| `SENSE_RECONCILE_INTERVAL` | `5s` | 对账周期 |
|
||||
| `SENSE_PROBE_INTERVAL` | `10s` | path 探活周期 |
|
||||
| `SENSE_ONVIF_MODE` | `disabled` | `disabled` 或 `standard`;默认不访问真实摄像头 |
|
||||
| `SENSE_ONVIF_RTSP_REWRITE_HOST` | 空 | NAT 或故障代理场景下重写 ONVIF 返回的 RTSP 主机 |
|
||||
| `SENSE_ONVIF_RTSP_REWRITE_PORT` | `0` | 非零时重写 ONVIF 返回的 RTSP 端口 |
|
||||
| `SENSE_ONVIF_RTSP_STRIP_QUERY` | `false` | 仅在已验证设备返回不可用查询串时显式移除;默认保留标准 URI 语义 |
|
||||
|
||||
设备台账只保存 `env://<key>` 凭据引用。真实适配器从进程环境读取以下变量,不把秘密写入 SQLite、日志或 MediaMTX 错误:
|
||||
|
||||
```text
|
||||
SENSE_CREDENTIAL_<KEY>_ONVIF_USERNAME
|
||||
SENSE_CREDENTIAL_<KEY>_ONVIF_PASSWORD
|
||||
SENSE_CREDENTIAL_<KEY>_RTSP_USERNAME
|
||||
SENSE_CREDENTIAL_<KEY>_RTSP_PASSWORD
|
||||
```
|
||||
|
||||
`cmd/sense-lab` 是回环实验室播种与脱敏收敛查询工具,不是已冻结的公共设备管理 API。`cmd/rtsp-fault-proxy` 只用于 T-006 控制真实上游网络路径故障。
|
||||
|
||||
MediaMTX `v1.19.3` 应作为独立二进制启动并只在可信网络开放 API。获取与 SHA-256 校验值见 `docs/03-tech-stack.md`。生成客户端使用固定版本工具和 vendored 官方 OpenAPI;`internal/mtx/generated/client.gen.go` 不可手改。
|
||||
|
||||
@@ -40,3 +55,13 @@ Expand-Archive "$env:TEMP\$asset" -DestinationPath "$env:TEMP\yovision-mediamtx-
|
||||
```
|
||||
|
||||
Linux amd64 使用同版 `mediamtx_v1.19.3_linux_amd64.tar.gz`,SHA-256 为 `a7ba21268fccda3ebc43fdad76b87fddb85ce77e725b5cb637bca724b5394fbe`。不要把下载的二进制或摄像头凭据提交到仓库。
|
||||
|
||||
## T-006 Windows 集成验证
|
||||
|
||||
脚本会启动两套独立 MediaMTX、4 个独立 FFmpeg publisher、真实摄像头网络故障代理和 Sense,在临时目录播种 5 条期望态,执行四类恢复后再连续观察 30 分钟。脚本只输出脱敏计数与时间,不保存视频:
|
||||
|
||||
```powershell
|
||||
./Sense/scripts/t006-integration.ps1 -CameraEnv D:\path\to\ip_camera.env
|
||||
```
|
||||
|
||||
`ip_camera.env` 必须保持在 Git 忽略范围内。调试时可把 `-ObservationMinutes` 降为 1;正式 T-006 证据必须使用默认 30 分钟,且最终 `maximum_unconverged`、`final_unconverged` 都为 0。
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net"
|
||||
"os"
|
||||
"os/signal"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
)
|
||||
|
||||
func main() {
|
||||
listenAddress := flag.String("listen", "127.0.0.1:10554", "local listen address")
|
||||
upstreamAddress := flag.String("upstream", "", "upstream host:port")
|
||||
flag.Parse()
|
||||
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
|
||||
if *upstreamAddress == "" {
|
||||
logger.Error("upstream is required")
|
||||
os.Exit(2)
|
||||
}
|
||||
if _, _, err := net.SplitHostPort(*upstreamAddress); err != nil {
|
||||
logger.Error("upstream must be host:port")
|
||||
os.Exit(2)
|
||||
}
|
||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||
defer stop()
|
||||
if err := serve(ctx, *listenAddress, *upstreamAddress); err != nil {
|
||||
logger.Error("RTSP fault proxy stopped", "error", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
func serve(ctx context.Context, listenAddress, upstreamAddress string) error {
|
||||
listener, err := net.Listen("tcp", listenAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("listen: %w", err)
|
||||
}
|
||||
defer listener.Close()
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
_ = listener.Close()
|
||||
}()
|
||||
var connections sync.WaitGroup
|
||||
defer connections.Wait()
|
||||
for {
|
||||
client, acceptErr := listener.Accept()
|
||||
if acceptErr != nil {
|
||||
if ctx.Err() != nil {
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("accept: %w", acceptErr)
|
||||
}
|
||||
connections.Add(1)
|
||||
go func() {
|
||||
defer connections.Done()
|
||||
proxy(client, upstreamAddress)
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
||||
func proxy(client net.Conn, upstreamAddress string) {
|
||||
defer client.Close()
|
||||
upstream, err := net.DialTimeout("tcp", upstreamAddress, 5*time.Second)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer upstream.Close()
|
||||
done := make(chan struct{}, 2)
|
||||
copyOneWay := func(destination, source net.Conn) {
|
||||
_, _ = io.Copy(destination, source)
|
||||
done <- struct{}{}
|
||||
}
|
||||
go copyOneWay(upstream, client)
|
||||
go copyOneWay(client, upstream)
|
||||
<-done
|
||||
}
|
||||
@@ -48,9 +48,17 @@ func run(logger *slog.Logger) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// T-003 deliberately has no real camera adapter. T-006 replaces this port
|
||||
// only after the device whitelist and five-camera evidence are available.
|
||||
reconciler := reconcile.New(repository, onvif.UnavailableAdapter{}, mediaClient)
|
||||
credentials := onvif.EnvCredentials{}
|
||||
var cameraAdapter onvif.Adapter = onvif.UnavailableAdapter{}
|
||||
if cfg.ONVIFMode == "standard" {
|
||||
cameraAdapter = onvif.NewHTTPAdapter(credentials, nil, onvif.HTTPOptions{
|
||||
RTSPRewriteHost: cfg.RTSPRewriteHost,
|
||||
RTSPRewritePort: cfg.RTSPRewritePort,
|
||||
StripRTSPQuery: cfg.RTSPStripQuery,
|
||||
})
|
||||
}
|
||||
discovery := onvif.NewRouter(cameraAdapter, credentials)
|
||||
reconciler := reconcile.New(repository, discovery, mediaClient)
|
||||
checker := probe.New(repository, mediaClient)
|
||||
report := func(err error) {
|
||||
// Domain and MediaMTX errors intentionally omit stream URIs and credentials.
|
||||
|
||||
@@ -0,0 +1,143 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
|
||||
"yovision/sense/internal/device"
|
||||
"yovision/sense/internal/store"
|
||||
)
|
||||
|
||||
type manifest struct {
|
||||
Site manifestSite `json:"site"`
|
||||
Devices []manifestDevice `json:"devices"`
|
||||
}
|
||||
|
||||
type manifestSite struct {
|
||||
TenantID string `json:"tenant_id"`
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
MaxVideoChannels int `json:"max_video_channels"`
|
||||
}
|
||||
|
||||
type manifestDevice struct {
|
||||
ID string `json:"id"`
|
||||
TenantID string `json:"tenant_id"`
|
||||
SiteID string `json:"site_id"`
|
||||
SerialNumber string `json:"serial_number"`
|
||||
Name string `json:"name"`
|
||||
Capabilities []string `json:"capabilities"`
|
||||
EndpointRef string `json:"endpoint_ref"`
|
||||
CredentialRef string `json:"credential_ref"`
|
||||
PathName string `json:"path_name"`
|
||||
}
|
||||
|
||||
func main() {
|
||||
if err := run(os.Args[1:], os.Stdout); err != nil {
|
||||
_, _ = fmt.Fprintln(os.Stderr, err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
func run(args []string, output io.Writer) error {
|
||||
if len(args) == 0 {
|
||||
return errors.New("usage: sense-lab <seed|status>")
|
||||
}
|
||||
switch args[0] {
|
||||
case "seed":
|
||||
return seed(args[1:], output)
|
||||
case "status":
|
||||
return status(args[1:], output)
|
||||
default:
|
||||
return fmt.Errorf("unknown command %q", args[0])
|
||||
}
|
||||
}
|
||||
|
||||
func seed(args []string, output io.Writer) error {
|
||||
flags := flag.NewFlagSet("seed", flag.ContinueOnError)
|
||||
flags.SetOutput(io.Discard)
|
||||
dsn := flags.String("db", "", "SQLite DSN")
|
||||
manifestPath := flags.String("manifest", "", "manifest JSON path")
|
||||
if err := flags.Parse(args); err != nil {
|
||||
return err
|
||||
}
|
||||
if *dsn == "" || *manifestPath == "" {
|
||||
return errors.New("seed requires -db and -manifest")
|
||||
}
|
||||
file, err := os.Open(*manifestPath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("open manifest: %w", err)
|
||||
}
|
||||
defer file.Close()
|
||||
var value manifest
|
||||
decoder := json.NewDecoder(io.LimitReader(file, 1<<20))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(&value); err != nil {
|
||||
return fmt.Errorf("decode manifest: %w", err)
|
||||
}
|
||||
repository, err := store.OpenSQLite(context.Background(), *dsn)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer repository.Close()
|
||||
ctx := context.Background()
|
||||
if err := repository.EnsureSite(ctx, device.Site{
|
||||
TenantID: value.Site.TenantID, ID: value.Site.ID, Name: value.Site.Name,
|
||||
MaxVideoChannels: value.Site.MaxVideoChannels,
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, input := range value.Devices {
|
||||
capabilities := make([]device.Capability, 0, len(input.Capabilities))
|
||||
for _, capability := range input.Capabilities {
|
||||
capabilities = append(capabilities, device.Capability(capability))
|
||||
}
|
||||
if err := repository.CreateDevice(ctx, device.Device{
|
||||
ID: input.ID, TenantID: input.TenantID, SiteID: input.SiteID,
|
||||
SerialNumber: input.SerialNumber, Name: input.Name, Modality: device.ModalityVideo,
|
||||
Capabilities: capabilities, DesiredState: device.DesiredEnabled, ActualState: device.ActualPending,
|
||||
EndpointRef: input.EndpointRef, CredentialRef: input.CredentialRef, PathName: input.PathName,
|
||||
}); err != nil {
|
||||
return fmt.Errorf("create device %s: %w", input.ID, err)
|
||||
}
|
||||
}
|
||||
return json.NewEncoder(output).Encode(map[string]int{"seeded": len(value.Devices)})
|
||||
}
|
||||
|
||||
func status(args []string, output io.Writer) error {
|
||||
flags := flag.NewFlagSet("status", flag.ContinueOnError)
|
||||
flags.SetOutput(io.Discard)
|
||||
dsn := flags.String("db", "", "SQLite DSN")
|
||||
expect := flags.Int("expect", -1, "expected device count")
|
||||
requireConverged := flags.Bool("require-converged", false, "fail when unconverged is non-zero")
|
||||
if err := flags.Parse(args); err != nil {
|
||||
return err
|
||||
}
|
||||
if *dsn == "" {
|
||||
return errors.New("status requires -db")
|
||||
}
|
||||
repository, err := store.OpenSQLite(context.Background(), *dsn)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer repository.Close()
|
||||
snapshot, err := repository.ConvergenceSnapshot(context.Background())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := json.NewEncoder(output).Encode(snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
if *expect >= 0 && snapshot.Total != *expect {
|
||||
return fmt.Errorf("expected %d devices, got %d", *expect, snapshot.Total)
|
||||
}
|
||||
if *requireConverged && snapshot.Unconverged != 0 {
|
||||
return fmt.Errorf("unconverged devices: %d", snapshot.Unconverged)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
# T-006-only publisher fixture. This is a separate MediaMTX instance; Sense
|
||||
# never edits these source paths. Each path has one independent FFmpeg process.
|
||||
logLevel: warn
|
||||
rtspAddress: 127.0.0.1:8555
|
||||
rtspTransports: [tcp]
|
||||
api: false
|
||||
metrics: false
|
||||
rtmp: false
|
||||
hls: false
|
||||
webrtc: false
|
||||
srt: false
|
||||
moq: false
|
||||
paths:
|
||||
synthetic-1:
|
||||
source: publisher
|
||||
synthetic-2:
|
||||
source: publisher
|
||||
synthetic-3:
|
||||
source: publisher
|
||||
synthetic-4:
|
||||
source: publisher
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"net/url"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -16,6 +17,7 @@ const (
|
||||
defaultMediaMTXURL = "http://127.0.0.1:9997"
|
||||
defaultReconcilePeriod = 5 * time.Second
|
||||
defaultProbePeriod = 10 * time.Second
|
||||
defaultONVIFMode = "disabled"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
@@ -25,6 +27,10 @@ type Config struct {
|
||||
MediaMTXURL string
|
||||
ReconcileInterval time.Duration
|
||||
ProbeInterval time.Duration
|
||||
ONVIFMode string
|
||||
RTSPRewriteHost string
|
||||
RTSPRewritePort int
|
||||
RTSPStripQuery bool
|
||||
}
|
||||
|
||||
func Load() (Config, error) {
|
||||
@@ -40,6 +46,14 @@ func Load() (Config, error) {
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
rewritePort, err := intEnv("SENSE_ONVIF_RTSP_REWRITE_PORT", 0)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
stripQuery, err := boolEnv("SENSE_ONVIF_RTSP_STRIP_QUERY", false)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
|
||||
cfg := Config{
|
||||
HTTPAddress: stringEnv("SENSE_HTTP_ADDR", defaultHTTPAddress),
|
||||
@@ -48,6 +62,10 @@ func Load() (Config, error) {
|
||||
MediaMTXURL: stringEnv("SENSE_MEDIAMTX_URL", defaultMediaMTXURL),
|
||||
ReconcileInterval: reconcilePeriod,
|
||||
ProbeInterval: probePeriod,
|
||||
ONVIFMode: stringEnv("SENSE_ONVIF_MODE", defaultONVIFMode),
|
||||
RTSPRewriteHost: stringEnv("SENSE_ONVIF_RTSP_REWRITE_HOST", ""),
|
||||
RTSPRewritePort: rewritePort,
|
||||
RTSPStripQuery: stripQuery,
|
||||
}
|
||||
if err := cfg.Validate(); err != nil {
|
||||
return Config{}, err
|
||||
@@ -78,6 +96,18 @@ func (c Config) Validate() error {
|
||||
if c.ReconcileInterval <= 0 || c.ProbeInterval <= 0 {
|
||||
return fmt.Errorf("loop intervals must be positive")
|
||||
}
|
||||
if c.ONVIFMode != "" && c.ONVIFMode != "disabled" && c.ONVIFMode != "standard" {
|
||||
return fmt.Errorf("SENSE_ONVIF_MODE must be disabled or standard")
|
||||
}
|
||||
if c.RTSPRewritePort < 0 || c.RTSPRewritePort > 65535 {
|
||||
return fmt.Errorf("SENSE_ONVIF_RTSP_REWRITE_PORT must be between 0 and 65535")
|
||||
}
|
||||
if c.RTSPRewriteHost != "" {
|
||||
if strings.TrimSpace(c.RTSPRewriteHost) != c.RTSPRewriteHost ||
|
||||
strings.ContainsAny(c.RTSPRewriteHost, "/@") {
|
||||
return fmt.Errorf("invalid SENSE_ONVIF_RTSP_REWRITE_HOST")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -111,3 +141,15 @@ func durationEnv(name string, fallback time.Duration) (time.Duration, error) {
|
||||
}
|
||||
return parsed, nil
|
||||
}
|
||||
|
||||
func intEnv(name string, fallback int) (int, error) {
|
||||
value, ok := os.LookupEnv(name)
|
||||
if !ok {
|
||||
return fallback, nil
|
||||
}
|
||||
parsed, err := strconv.Atoi(value)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("invalid %s: %w", name, err)
|
||||
}
|
||||
return parsed, nil
|
||||
}
|
||||
|
||||
@@ -33,3 +33,40 @@ func TestValidateRejectsCredentialsInMediaMTXURL(t *testing.T) {
|
||||
t.Fatal("expected credentials in MediaMTX URL to be rejected")
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateONVIFModeAndRewritePort(t *testing.T) {
|
||||
t.Parallel()
|
||||
cfg := Config{
|
||||
HTTPAddress: "127.0.0.1:8080",
|
||||
DatabaseDSN: "file:test.db",
|
||||
MediaMTXURL: "http://127.0.0.1:9997",
|
||||
ReconcileInterval: 1,
|
||||
ProbeInterval: 1,
|
||||
ONVIFMode: "standard",
|
||||
RTSPRewriteHost: "127.0.0.1",
|
||||
RTSPRewritePort: 10554,
|
||||
}
|
||||
if err := cfg.Validate(); err != nil {
|
||||
t.Fatalf("valid ONVIF configuration failed: %v", err)
|
||||
}
|
||||
cfg.ONVIFMode = "vendor"
|
||||
if err := cfg.Validate(); err == nil {
|
||||
t.Fatal("unknown ONVIF mode must be rejected")
|
||||
}
|
||||
cfg.ONVIFMode = "standard"
|
||||
cfg.RTSPRewritePort = 65536
|
||||
if err := cfg.Validate(); err == nil {
|
||||
t.Fatal("invalid RTSP rewrite port must be rejected")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadRTSPStripQueryOptIn(t *testing.T) {
|
||||
t.Setenv("SENSE_ONVIF_RTSP_STRIP_QUERY", "true")
|
||||
cfg, err := Load()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !cfg.RTSPStripQuery {
|
||||
t.Fatal("explicit RTSP query stripping was not loaded")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
package onvif
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
type Credentials struct {
|
||||
ONVIFUsername string
|
||||
ONVIFPassword string
|
||||
RTSPUsername string
|
||||
RTSPPassword string
|
||||
}
|
||||
|
||||
type CredentialProvider interface {
|
||||
Resolve(reference string) (Credentials, error)
|
||||
}
|
||||
|
||||
type EnvCredentials struct {
|
||||
LookupEnv func(string) (string, bool)
|
||||
}
|
||||
|
||||
var credentialKey = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
|
||||
|
||||
func (p EnvCredentials) Resolve(reference string) (Credentials, error) {
|
||||
parsed, err := url.Parse(reference)
|
||||
if err != nil || parsed.Scheme != "env" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" {
|
||||
return Credentials{}, fmt.Errorf("credential reference must use env://<key>")
|
||||
}
|
||||
key := strings.Trim(strings.TrimSpace(parsed.Host+parsed.Path), "/")
|
||||
if !credentialKey.MatchString(key) {
|
||||
return Credentials{}, fmt.Errorf("credential reference contains an invalid key")
|
||||
}
|
||||
lookup := p.LookupEnv
|
||||
if lookup == nil {
|
||||
lookup = os.LookupEnv
|
||||
}
|
||||
prefix := "SENSE_CREDENTIAL_" + strings.ToUpper(strings.ReplaceAll(key, "-", "_"))
|
||||
read := func(suffix string) string {
|
||||
value, _ := lookup(prefix + suffix)
|
||||
return value
|
||||
}
|
||||
result := Credentials{
|
||||
ONVIFUsername: read("_ONVIF_USERNAME"),
|
||||
ONVIFPassword: read("_ONVIF_PASSWORD"),
|
||||
RTSPUsername: read("_RTSP_USERNAME"),
|
||||
RTSPPassword: read("_RTSP_PASSWORD"),
|
||||
}
|
||||
if result.ONVIFUsername == "" || result.ONVIFPassword == "" {
|
||||
return Credentials{}, fmt.Errorf("ONVIF credentials are not configured for reference")
|
||||
}
|
||||
if result.RTSPUsername == "" {
|
||||
result.RTSPUsername = result.ONVIFUsername
|
||||
}
|
||||
if result.RTSPPassword == "" {
|
||||
result.RTSPPassword = result.ONVIFPassword
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
@@ -0,0 +1,342 @@
|
||||
package onvif
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"crypto/sha1"
|
||||
"encoding/base64"
|
||||
"encoding/xml"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
deviceNamespace = "http://www.onvif.org/ver10/device/wsdl"
|
||||
mediaNamespace = "http://www.onvif.org/ver10/media/wsdl"
|
||||
)
|
||||
|
||||
type HTTPOptions struct {
|
||||
RTSPRewriteHost string
|
||||
RTSPRewritePort int
|
||||
StripRTSPQuery bool
|
||||
}
|
||||
|
||||
type HTTPAdapter struct {
|
||||
credentials CredentialProvider
|
||||
client *http.Client
|
||||
options HTTPOptions
|
||||
now func() time.Time
|
||||
random io.Reader
|
||||
}
|
||||
|
||||
func NewHTTPAdapter(credentials CredentialProvider, client *http.Client, options HTTPOptions) *HTTPAdapter {
|
||||
if client == nil {
|
||||
client = &http.Client{Timeout: 10 * time.Second}
|
||||
}
|
||||
return &HTTPAdapter{
|
||||
credentials: credentials,
|
||||
client: client,
|
||||
options: options,
|
||||
now: time.Now,
|
||||
random: rand.Reader,
|
||||
}
|
||||
}
|
||||
|
||||
func (a *HTTPAdapter) Probe(ctx context.Context, target Target) (ProbeResult, error) {
|
||||
endpoint, credentials, err := a.target(target)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
infoBody, err := a.call(ctx, endpoint, deviceNamespace+"/GetDeviceInformation",
|
||||
`<tds:GetDeviceInformation xmlns:tds="`+deviceNamespace+`"/>`, credentials)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
var info deviceInformationEnvelope
|
||||
if err := xml.Unmarshal(infoBody, &info); err != nil {
|
||||
return ProbeResult{}, invalidResponse("decode device information")
|
||||
}
|
||||
|
||||
servicesBody, err := a.call(ctx, endpoint, deviceNamespace+"/GetServices",
|
||||
`<tds:GetServices xmlns:tds="`+deviceNamespace+`"><tds:IncludeCapability>false</tds:IncludeCapability></tds:GetServices>`, credentials)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
var services servicesEnvelope
|
||||
if err := xml.Unmarshal(servicesBody, &services); err != nil {
|
||||
return ProbeResult{}, invalidResponse("decode services")
|
||||
}
|
||||
mediaEndpoint, err := externalMediaEndpoint(endpoint, services.Body.Response.Services)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
|
||||
profilesBody, err := a.call(ctx, mediaEndpoint, mediaNamespace+"/GetProfiles",
|
||||
`<trt:GetProfiles xmlns:trt="`+mediaNamespace+`"/>`, credentials)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
var profilesResponse profilesEnvelope
|
||||
if err := xml.Unmarshal(profilesBody, &profilesResponse); err != nil {
|
||||
return ProbeResult{}, invalidResponse("decode profiles")
|
||||
}
|
||||
profiles := make([]Profile, 0, len(profilesResponse.Body.Response.Profiles))
|
||||
selectedToken := ""
|
||||
for _, value := range profilesResponse.Body.Response.Profiles {
|
||||
video := value.VideoEncoder != nil
|
||||
profiles = append(profiles, Profile{Token: value.Token, Name: value.Name, VideoEncoder: video})
|
||||
if selectedToken == "" && video && value.Token != "" {
|
||||
selectedToken = value.Token
|
||||
}
|
||||
}
|
||||
if selectedToken == "" {
|
||||
return ProbeResult{}, invalidResponse("no video profile")
|
||||
}
|
||||
|
||||
streamRequest := `<trt:GetStreamUri xmlns:trt="` + mediaNamespace + `" xmlns:tt="http://www.onvif.org/ver10/schema">` +
|
||||
`<trt:StreamSetup><tt:Stream>RTP-Unicast</tt:Stream><tt:Transport><tt:Protocol>RTSP</tt:Protocol></tt:Transport></trt:StreamSetup>` +
|
||||
`<trt:ProfileToken>` + escapeXML(selectedToken) + `</trt:ProfileToken></trt:GetStreamUri>`
|
||||
streamBody, err := a.call(ctx, mediaEndpoint, mediaNamespace+"/GetStreamUri", streamRequest, credentials)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
var streamResponse streamURIEnvelope
|
||||
if err := xml.Unmarshal(streamBody, &streamResponse); err != nil {
|
||||
return ProbeResult{}, invalidResponse("decode stream URI")
|
||||
}
|
||||
streamURI, err := a.rewriteStreamURI(endpoint, streamResponse.Body.Response.MediaURI.URI, credentials)
|
||||
if err != nil {
|
||||
return ProbeResult{}, err
|
||||
}
|
||||
return ProbeResult{
|
||||
Manufacturer: info.Body.Response.Manufacturer,
|
||||
Model: info.Body.Response.Model,
|
||||
FirmwareVersion: info.Body.Response.FirmwareVersion,
|
||||
SerialNumber: info.Body.Response.SerialNumber,
|
||||
Profiles: profiles,
|
||||
StreamURI: streamURI,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (a *HTTPAdapter) SetSystemDateAndTime(ctx context.Context, target Target, value time.Time) error {
|
||||
endpoint, credentials, err := a.target(target)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
utc := value.UTC()
|
||||
body := `<tds:SetSystemDateAndTime xmlns:tds="` + deviceNamespace + `" xmlns:tt="http://www.onvif.org/ver10/schema">` +
|
||||
`<tds:DateTimeType>Manual</tds:DateTimeType><tds:DaylightSavings>false</tds:DaylightSavings>` +
|
||||
`<tds:UTCDateTime><tt:Time><tt:Hour>` + strconv.Itoa(utc.Hour()) + `</tt:Hour><tt:Minute>` + strconv.Itoa(utc.Minute()) +
|
||||
`</tt:Minute><tt:Second>` + strconv.Itoa(utc.Second()) + `</tt:Second></tt:Time><tt:Date><tt:Year>` + strconv.Itoa(utc.Year()) +
|
||||
`</tt:Year><tt:Month>` + strconv.Itoa(int(utc.Month())) + `</tt:Month><tt:Day>` + strconv.Itoa(utc.Day()) +
|
||||
`</tt:Day></tt:Date></tds:UTCDateTime></tds:SetSystemDateAndTime>`
|
||||
_, err = a.call(ctx, endpoint, deviceNamespace+"/SetSystemDateAndTime", body, credentials)
|
||||
return err
|
||||
}
|
||||
|
||||
func (a *HTTPAdapter) target(target Target) (*url.URL, Credentials, error) {
|
||||
endpoint, err := url.Parse(target.EndpointRef)
|
||||
if err != nil || endpoint.Host == "" || endpoint.User != nil || (endpoint.Scheme != "http" && endpoint.Scheme != "https") {
|
||||
return nil, Credentials{}, invalidResponse("invalid ONVIF endpoint")
|
||||
}
|
||||
credentials, err := a.credentials.Resolve(target.CredentialRef)
|
||||
if err != nil {
|
||||
return nil, Credentials{}, &Error{Code: ErrorAuthentication, Err: err}
|
||||
}
|
||||
return endpoint, credentials, nil
|
||||
}
|
||||
|
||||
func (a *HTTPAdapter) call(ctx context.Context, endpoint *url.URL, action, body string, credentials Credentials) ([]byte, error) {
|
||||
nonce := make([]byte, 20)
|
||||
if _, err := io.ReadFull(a.random, nonce); err != nil {
|
||||
return nil, &Error{Code: ErrorUnavailable, Err: fmt.Errorf("create authentication nonce")}
|
||||
}
|
||||
created := a.now().UTC().Format("2006-01-02T15:04:05Z")
|
||||
digestInput := append(append(append([]byte{}, nonce...), []byte(created)...), []byte(credentials.ONVIFPassword)...)
|
||||
digest := sha1.Sum(digestInput)
|
||||
envelope := `<?xml version="1.0" encoding="UTF-8"?>` +
|
||||
`<s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope" xmlns:wsse="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-wssecurity-secext-1.0.xsd" xmlns:wsu="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-wssecurity-utility-1.0.xsd">` +
|
||||
`<s:Header><wsse:Security s:mustUnderstand="1"><wsse:UsernameToken><wsse:Username>` + escapeXML(credentials.ONVIFUsername) +
|
||||
`</wsse:Username><wsse:Password Type="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-username-token-profile-1.0#PasswordDigest">` +
|
||||
base64.StdEncoding.EncodeToString(digest[:]) + `</wsse:Password><wsse:Nonce EncodingType="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-soap-message-security-1.0#Base64Binary">` +
|
||||
base64.StdEncoding.EncodeToString(nonce) + `</wsse:Nonce><wsu:Created>` + created +
|
||||
`</wsu:Created></wsse:UsernameToken></wsse:Security></s:Header><s:Body>` + body + `</s:Body></s:Envelope>`
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint.String(), strings.NewReader(envelope))
|
||||
if err != nil {
|
||||
return nil, &Error{Code: ErrorUnavailable, Err: fmt.Errorf("create ONVIF request")}
|
||||
}
|
||||
request.Header.Set("Content-Type", `application/soap+xml; charset=utf-8; action="`+action+`"`)
|
||||
response, err := a.client.Do(request)
|
||||
if err != nil {
|
||||
code := ErrorUnavailable
|
||||
if errors.Is(err, context.DeadlineExceeded) || errors.Is(ctx.Err(), context.DeadlineExceeded) {
|
||||
code = ErrorTimeout
|
||||
}
|
||||
return nil, &Error{Code: code, Err: fmt.Errorf("ONVIF transport failed")}
|
||||
}
|
||||
defer response.Body.Close()
|
||||
responseBody, err := io.ReadAll(io.LimitReader(response.Body, 2<<20))
|
||||
if err != nil {
|
||||
return nil, &Error{Code: ErrorUnavailable, Err: fmt.Errorf("read ONVIF response")}
|
||||
}
|
||||
if response.StatusCode == http.StatusUnauthorized || response.StatusCode == http.StatusForbidden {
|
||||
return nil, &Error{Code: ErrorAuthentication, Err: fmt.Errorf("ONVIF authorization failed")}
|
||||
}
|
||||
if fault := soapFault(responseBody); fault != "" {
|
||||
code := ErrorInvalidReply
|
||||
lower := strings.ToLower(fault)
|
||||
if strings.Contains(lower, "authoriz") || strings.Contains(lower, "notauthorized") {
|
||||
code = ErrorAuthentication
|
||||
}
|
||||
return nil, &Error{Code: code, Err: fmt.Errorf("ONVIF SOAP fault")}
|
||||
}
|
||||
if response.StatusCode != http.StatusOK {
|
||||
return nil, &Error{Code: ErrorUnavailable, Err: fmt.Errorf("ONVIF returned HTTP status %d", response.StatusCode)}
|
||||
}
|
||||
return responseBody, nil
|
||||
}
|
||||
|
||||
func externalMediaEndpoint(deviceEndpoint *url.URL, services []service) (*url.URL, error) {
|
||||
for _, value := range services {
|
||||
if value.Namespace != mediaNamespace || value.XAddr == "" {
|
||||
continue
|
||||
}
|
||||
mediaEndpoint, err := url.Parse(value.XAddr)
|
||||
if err != nil || mediaEndpoint.Host == "" {
|
||||
return nil, invalidResponse("invalid media service address")
|
||||
}
|
||||
mediaEndpoint.Scheme = deviceEndpoint.Scheme
|
||||
mediaEndpoint.Host = deviceEndpoint.Host
|
||||
mediaEndpoint.User = nil
|
||||
return mediaEndpoint, nil
|
||||
}
|
||||
return nil, invalidResponse("media service is unavailable")
|
||||
}
|
||||
|
||||
func (a *HTTPAdapter) rewriteStreamURI(deviceEndpoint *url.URL, raw string, credentials Credentials) (string, error) {
|
||||
stream, err := url.Parse(raw)
|
||||
if err != nil || stream.Host == "" || (stream.Scheme != "rtsp" && stream.Scheme != "rtsps") {
|
||||
return "", invalidResponse("invalid stream URI")
|
||||
}
|
||||
host := a.options.RTSPRewriteHost
|
||||
if host == "" {
|
||||
host = deviceEndpoint.Hostname()
|
||||
}
|
||||
port := a.options.RTSPRewritePort
|
||||
if port == 0 {
|
||||
if parsedPort := stream.Port(); parsedPort != "" {
|
||||
value, parseErr := strconv.Atoi(parsedPort)
|
||||
if parseErr != nil {
|
||||
return "", invalidResponse("invalid stream port")
|
||||
}
|
||||
port = value
|
||||
}
|
||||
}
|
||||
if port > 0 {
|
||||
stream.Host = net.JoinHostPort(host, strconv.Itoa(port))
|
||||
} else {
|
||||
stream.Host = host
|
||||
}
|
||||
stream.User = url.UserPassword(credentials.RTSPUsername, credentials.RTSPPassword)
|
||||
if a.options.StripRTSPQuery {
|
||||
stream.RawQuery = ""
|
||||
stream.ForceQuery = false
|
||||
}
|
||||
return stream.String(), nil
|
||||
}
|
||||
|
||||
func invalidResponse(message string) error {
|
||||
return &Error{Code: ErrorInvalidReply, Err: fmt.Errorf("%s", message)}
|
||||
}
|
||||
|
||||
func escapeXML(value string) string {
|
||||
var buffer bytes.Buffer
|
||||
_ = xml.EscapeText(&buffer, []byte(value))
|
||||
return buffer.String()
|
||||
}
|
||||
|
||||
func soapFault(body []byte) string {
|
||||
decoder := xml.NewDecoder(bytes.NewReader(body))
|
||||
inFault := false
|
||||
for {
|
||||
token, err := decoder.Token()
|
||||
if errors.Is(err, io.EOF) {
|
||||
return ""
|
||||
}
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
switch value := token.(type) {
|
||||
case xml.StartElement:
|
||||
if value.Name.Local == "Fault" {
|
||||
inFault = true
|
||||
}
|
||||
if inFault && (value.Name.Local == "Text" || value.Name.Local == "faultstring") {
|
||||
var message string
|
||||
if decoder.DecodeElement(&message, &value) == nil {
|
||||
return message
|
||||
}
|
||||
}
|
||||
case xml.EndElement:
|
||||
if value.Name.Local == "Fault" {
|
||||
return "SOAP fault"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type deviceInformationEnvelope struct {
|
||||
Body struct {
|
||||
Response struct {
|
||||
Manufacturer string `xml:"Manufacturer"`
|
||||
Model string `xml:"Model"`
|
||||
FirmwareVersion string `xml:"FirmwareVersion"`
|
||||
SerialNumber string `xml:"SerialNumber"`
|
||||
} `xml:"GetDeviceInformationResponse"`
|
||||
} `xml:"Body"`
|
||||
}
|
||||
|
||||
type service struct {
|
||||
Namespace string `xml:"Namespace"`
|
||||
XAddr string `xml:"XAddr"`
|
||||
}
|
||||
|
||||
type servicesEnvelope struct {
|
||||
Body struct {
|
||||
Response struct {
|
||||
Services []service `xml:"Service"`
|
||||
} `xml:"GetServicesResponse"`
|
||||
} `xml:"Body"`
|
||||
}
|
||||
|
||||
type profileResponse struct {
|
||||
Token string `xml:"token,attr"`
|
||||
Name string `xml:"Name"`
|
||||
VideoEncoder *struct{} `xml:"VideoEncoderConfiguration"`
|
||||
}
|
||||
|
||||
type profilesEnvelope struct {
|
||||
Body struct {
|
||||
Response struct {
|
||||
Profiles []profileResponse `xml:"Profiles"`
|
||||
} `xml:"GetProfilesResponse"`
|
||||
} `xml:"Body"`
|
||||
}
|
||||
|
||||
type streamURIEnvelope struct {
|
||||
Body struct {
|
||||
Response struct {
|
||||
MediaURI struct {
|
||||
URI string `xml:"Uri"`
|
||||
} `xml:"MediaUri"`
|
||||
} `xml:"GetStreamUriResponse"`
|
||||
} `xml:"Body"`
|
||||
}
|
||||
@@ -0,0 +1,136 @@
|
||||
package onvif
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type staticCredentials struct {
|
||||
value Credentials
|
||||
err error
|
||||
}
|
||||
|
||||
func (s staticCredentials) Resolve(string) (Credentials, error) {
|
||||
return s.value, s.err
|
||||
}
|
||||
|
||||
func TestHTTPAdapterDiscoversMediaAndRewritesNATStream(t *testing.T) {
|
||||
t.Parallel()
|
||||
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
|
||||
if !strings.Contains(request.Header.Get("Content-Type"), "action=") {
|
||||
t.Fatal("SOAP action is required")
|
||||
}
|
||||
body := ""
|
||||
switch {
|
||||
case strings.Contains(request.Header.Get("Content-Type"), "GetDeviceInformation"):
|
||||
body = `<tds:GetDeviceInformationResponse xmlns:tds="http://www.onvif.org/ver10/device/wsdl"><tds:Manufacturer>HIKVISION</tds:Manufacturer><tds:Model>camera</tds:Model><tds:FirmwareVersion>v1</tds:FirmwareVersion><tds:SerialNumber>serial</tds:SerialNumber></tds:GetDeviceInformationResponse>`
|
||||
case strings.Contains(request.Header.Get("Content-Type"), "GetServices"):
|
||||
body = `<tds:GetServicesResponse xmlns:tds="http://www.onvif.org/ver10/device/wsdl"><tds:Service><tds:Namespace>http://www.onvif.org/ver10/media/wsdl</tds:Namespace><tds:XAddr>http://192.0.2.10/onvif/Media</tds:XAddr></tds:Service></tds:GetServicesResponse>`
|
||||
case strings.Contains(request.Header.Get("Content-Type"), "GetProfiles"):
|
||||
body = `<trt:GetProfilesResponse xmlns:trt="http://www.onvif.org/ver10/media/wsdl"><trt:Profiles token="main"><trt:Name>Main</trt:Name><trt:VideoEncoderConfiguration/></trt:Profiles></trt:GetProfilesResponse>`
|
||||
case strings.Contains(request.Header.Get("Content-Type"), "GetStreamUri"):
|
||||
body = `<trt:GetStreamUriResponse xmlns:trt="http://www.onvif.org/ver10/media/wsdl"><trt:MediaUri><trt:Uri>rtsp://192.0.2.10:554/Streaming/Channels/101?transportmode=unicast&profile=Profile_1</trt:Uri></trt:MediaUri></trt:GetStreamUriResponse>`
|
||||
default:
|
||||
http.Error(writer, "unexpected action", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
writer.Header().Set("Content-Type", "application/soap+xml")
|
||||
_, _ = fmt.Fprintf(writer, `<s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope"><s:Body>%s</s:Body></s:Envelope>`, body)
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
credentials := Credentials{
|
||||
ONVIFUsername: "onvif-user", ONVIFPassword: "onvif-password",
|
||||
RTSPUsername: "rtsp-user", RTSPPassword: "rtsp-password",
|
||||
}
|
||||
adapter := NewHTTPAdapter(staticCredentials{value: credentials}, server.Client(), HTTPOptions{
|
||||
RTSPRewriteHost: "127.0.0.1", RTSPRewritePort: 10554, StripRTSPQuery: true,
|
||||
})
|
||||
result, err := adapter.Probe(context.Background(), Target{
|
||||
EndpointRef: server.URL + "/onvif/device_service", CredentialRef: "env://camera",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if result.Manufacturer != "HIKVISION" || result.Model != "camera" || len(result.Profiles) != 1 {
|
||||
t.Fatalf("unexpected probe result: %+v", result)
|
||||
}
|
||||
stream, err := url.Parse(result.StreamURI)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if stream.Host != "127.0.0.1:10554" || stream.Path != "/Streaming/Channels/101" || stream.RawQuery != "" {
|
||||
t.Fatalf("unexpected rewritten stream address: host=%s path=%s", stream.Host, stream.Path)
|
||||
}
|
||||
password, _ := stream.User.Password()
|
||||
if stream.User.Username() != credentials.RTSPUsername || password != credentials.RTSPPassword {
|
||||
t.Fatal("RTSP credentials were not injected")
|
||||
}
|
||||
}
|
||||
|
||||
func TestHTTPAdapterPreservesRTSPQueryByDefault(t *testing.T) {
|
||||
t.Parallel()
|
||||
adapter := NewHTTPAdapter(staticCredentials{}, nil, HTTPOptions{})
|
||||
streamURI, err := adapter.rewriteStreamURI(
|
||||
&url.URL{Scheme: "http", Host: "camera.example:8008"},
|
||||
"rtsp://192.0.2.10:554/live?profile=main",
|
||||
Credentials{RTSPUsername: "user", RTSPPassword: "secret"},
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stream, err := url.Parse(streamURI)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if stream.RawQuery != "profile=main" {
|
||||
t.Fatalf("RTSP query was unexpectedly changed: %q", stream.RawQuery)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHTTPAdapterMapsAuthorizationFaultWithoutLeakingSecret(t *testing.T) {
|
||||
t.Parallel()
|
||||
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, _ *http.Request) {
|
||||
writer.WriteHeader(http.StatusBadRequest)
|
||||
_, _ = writer.Write([]byte(`<s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope"><s:Body><s:Fault><s:Reason><s:Text>The action requires authorization</s:Text></s:Reason></s:Fault></s:Body></s:Envelope>`))
|
||||
}))
|
||||
defer server.Close()
|
||||
secret := "not-for-errors"
|
||||
adapter := NewHTTPAdapter(staticCredentials{value: Credentials{
|
||||
ONVIFUsername: "user", ONVIFPassword: secret, RTSPUsername: "user", RTSPPassword: secret,
|
||||
}}, server.Client(), HTTPOptions{})
|
||||
_, err := adapter.Probe(context.Background(), Target{EndpointRef: server.URL, CredentialRef: "env://camera"})
|
||||
if CodeOf(err) != ErrorAuthentication || strings.Contains(err.Error(), secret) {
|
||||
t.Fatalf("expected redacted authentication error, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnvCredentialsAndDirectRTSPRouting(t *testing.T) {
|
||||
t.Parallel()
|
||||
values := map[string]string{
|
||||
"SENSE_CREDENTIAL_CAMERA_ONVIF_USERNAME": "onvif",
|
||||
"SENSE_CREDENTIAL_CAMERA_ONVIF_PASSWORD": "onvif-secret",
|
||||
"SENSE_CREDENTIAL_CAMERA_RTSP_USERNAME": "rtsp",
|
||||
"SENSE_CREDENTIAL_CAMERA_RTSP_PASSWORD": "rtsp-secret",
|
||||
}
|
||||
provider := EnvCredentials{LookupEnv: func(name string) (string, bool) {
|
||||
value, ok := values[name]
|
||||
return value, ok
|
||||
}}
|
||||
resolved, err := provider.Resolve("env://camera")
|
||||
if err != nil || resolved.RTSPUsername != "rtsp" {
|
||||
t.Fatalf("resolve credentials: %+v err=%v", resolved, err)
|
||||
}
|
||||
router := NewRouter(UnavailableAdapter{}, provider)
|
||||
result, err := router.Probe(context.Background(), Target{
|
||||
EndpointRef: "rtsp://127.0.0.1:8555/synthetic-1",
|
||||
})
|
||||
if err != nil || result.StreamURI != "rtsp://127.0.0.1:8555/synthetic-1" {
|
||||
t.Fatalf("route direct RTSP: %+v err=%v", result, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package onvif
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Router struct {
|
||||
camera Adapter
|
||||
credentials CredentialProvider
|
||||
}
|
||||
|
||||
func NewRouter(camera Adapter, credentials CredentialProvider) *Router {
|
||||
return &Router{camera: camera, credentials: credentials}
|
||||
}
|
||||
|
||||
func (r *Router) Probe(ctx context.Context, target Target) (ProbeResult, error) {
|
||||
endpoint, err := url.Parse(target.EndpointRef)
|
||||
if err != nil || endpoint.Host == "" || endpoint.User != nil {
|
||||
return ProbeResult{}, &Error{Code: ErrorInvalidReply, Err: fmt.Errorf("endpoint reference is invalid")}
|
||||
}
|
||||
switch endpoint.Scheme {
|
||||
case "http", "https":
|
||||
return r.camera.Probe(ctx, target)
|
||||
case "rtsp", "rtsps":
|
||||
if target.CredentialRef != "" {
|
||||
credentials, resolveErr := r.credentials.Resolve(target.CredentialRef)
|
||||
if resolveErr != nil {
|
||||
return ProbeResult{}, &Error{Code: ErrorAuthentication, Err: resolveErr}
|
||||
}
|
||||
endpoint.User = url.UserPassword(credentials.RTSPUsername, credentials.RTSPPassword)
|
||||
}
|
||||
return ProbeResult{
|
||||
Profiles: []Profile{{Token: "direct", Name: "direct", VideoEncoder: true}},
|
||||
StreamURI: endpoint.String(),
|
||||
}, nil
|
||||
default:
|
||||
return ProbeResult{}, &Error{Code: ErrorInvalidReply, Err: fmt.Errorf("unsupported endpoint scheme")}
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Router) SetSystemDateAndTime(ctx context.Context, target Target, value time.Time) error {
|
||||
endpoint, err := url.Parse(target.EndpointRef)
|
||||
if err != nil || (endpoint.Scheme != "http" && endpoint.Scheme != "https") {
|
||||
return &Error{Code: ErrorInvalidReply, Err: fmt.Errorf("clock sync requires an ONVIF endpoint")}
|
||||
}
|
||||
return r.camera.SetSystemDateAndTime(ctx, target, value)
|
||||
}
|
||||
@@ -15,6 +15,7 @@ const defaultBatchSize = 128
|
||||
type Repository interface {
|
||||
ListEnabledVideoDevices(ctx context.Context, limit int) ([]device.Device, error)
|
||||
UpdateActualState(ctx context.Context, id string, state device.ActualState, now time.Time) error
|
||||
RequestReconcile(ctx context.Context, id string, now time.Time) error
|
||||
}
|
||||
|
||||
type RuntimePaths interface {
|
||||
@@ -50,7 +51,13 @@ func (c *Checker) RunOnce(ctx context.Context) error {
|
||||
if probeErr == nil && ready {
|
||||
state = device.ActualOnline
|
||||
}
|
||||
if updateErr := c.repository.UpdateActualState(ctx, value.ID, state, c.now().UTC()); updateErr != nil {
|
||||
now := c.now().UTC()
|
||||
if probeErr != nil {
|
||||
if requestErr := c.repository.RequestReconcile(ctx, value.ID, now); requestErr != nil {
|
||||
runErrors = append(runErrors, fmt.Errorf("request device %s reconciliation: %w", value.ID, requestErr))
|
||||
}
|
||||
}
|
||||
if updateErr := c.repository.UpdateActualState(ctx, value.ID, state, now); updateErr != nil {
|
||||
runErrors = append(runErrors, fmt.Errorf("update device %s health: %w", value.ID, updateErr))
|
||||
}
|
||||
if probeErr != nil {
|
||||
|
||||
@@ -10,8 +10,17 @@ import (
|
||||
)
|
||||
|
||||
type fakeRepository struct {
|
||||
devices []device.Device
|
||||
states map[string]device.ActualState
|
||||
devices []device.Device
|
||||
states map[string]device.ActualState
|
||||
requested map[string]int
|
||||
}
|
||||
|
||||
func (f *fakeRepository) RequestReconcile(_ context.Context, id string, _ time.Time) error {
|
||||
if f.requested == nil {
|
||||
f.requested = make(map[string]int)
|
||||
}
|
||||
f.requested[id]++
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeRepository) ListEnabledVideoDevices(context.Context, int) ([]device.Device, error) {
|
||||
@@ -50,4 +59,7 @@ func TestCheckerMapsReadyAndUnavailablePaths(t *testing.T) {
|
||||
if repository.states["online"] != device.ActualOnline || repository.states["offline"] != device.ActualOffline {
|
||||
t.Fatalf("unexpected actual states: %+v", repository.states)
|
||||
}
|
||||
if repository.requested["offline"] != 1 || repository.requested["online"] != 0 {
|
||||
t.Fatalf("unexpected reconcile requests: %+v", repository.requested)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -25,6 +25,25 @@ type ReconcileCandidate struct {
|
||||
NextAttempt *time.Time
|
||||
}
|
||||
|
||||
type DeviceConvergence struct {
|
||||
ID string `json:"id"`
|
||||
PathName string `json:"path_name"`
|
||||
DesiredState device.DesiredState `json:"desired_state"`
|
||||
ActualState device.ActualState `json:"actual_state"`
|
||||
Generation int64 `json:"generation"`
|
||||
ObservedGeneration int64 `json:"observed_generation"`
|
||||
FailureCount int `json:"failure_count"`
|
||||
NextAttemptAt *time.Time `json:"next_attempt_at,omitempty"`
|
||||
LastErrorCode string `json:"last_error_code,omitempty"`
|
||||
Converged bool `json:"converged"`
|
||||
}
|
||||
|
||||
type ConvergenceSnapshot struct {
|
||||
Total int `json:"total"`
|
||||
Unconverged int `json:"unconverged"`
|
||||
Devices []DeviceConvergence `json:"devices"`
|
||||
}
|
||||
|
||||
type SQLite struct {
|
||||
db *sql.DB
|
||||
}
|
||||
@@ -521,6 +540,76 @@ func (s *SQLite) UpdateActualState(ctx context.Context, id string, state device.
|
||||
return nil
|
||||
}
|
||||
|
||||
// RequestReconcile invalidates the observed generation without changing the
|
||||
// desired state or retry backoff. Runtime probes use it when MediaMTX loses a
|
||||
// configured path, including after a MediaMTX process restart.
|
||||
func (s *SQLite) RequestReconcile(ctx context.Context, id string, now time.Time) error {
|
||||
result, err := s.db.ExecContext(ctx, `
|
||||
UPDATE sense_reconcile_state
|
||||
SET observed_generation = 0, updated_at = ?
|
||||
WHERE device_id = ?`, formatTime(now), id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("request device reconciliation: %w", err)
|
||||
}
|
||||
if affected, _ := result.RowsAffected(); affected != 1 {
|
||||
return ErrNotFound
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ConvergenceSnapshot returns only identifiers, state and counters. Endpoint
|
||||
// and credential references are deliberately excluded from diagnostics.
|
||||
func (s *SQLite) ConvergenceSnapshot(ctx context.Context) (ConvergenceSnapshot, error) {
|
||||
rows, err := s.db.QueryContext(ctx, `
|
||||
SELECT d.id, d.path_name, d.desired_state, d.actual_state, d.generation,
|
||||
r.observed_generation, r.failure_count, r.next_attempt_at, r.last_error_code
|
||||
FROM sense_devices d
|
||||
JOIN sense_reconcile_state r ON r.device_id = d.id
|
||||
WHERE d.desired_state = 'enabled'
|
||||
AND EXISTS (
|
||||
SELECT 1 FROM sense_device_capabilities c
|
||||
WHERE c.device_id = d.id AND c.capability = 'video_capture'
|
||||
)
|
||||
ORDER BY d.id`)
|
||||
if err != nil {
|
||||
return ConvergenceSnapshot{}, fmt.Errorf("query convergence snapshot: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
snapshot := ConvergenceSnapshot{Devices: make([]DeviceConvergence, 0)}
|
||||
for rows.Next() {
|
||||
var value DeviceConvergence
|
||||
var nextAttempt, lastError sql.NullString
|
||||
if err := rows.Scan(
|
||||
&value.ID, &value.PathName, &value.DesiredState, &value.ActualState,
|
||||
&value.Generation, &value.ObservedGeneration, &value.FailureCount,
|
||||
&nextAttempt, &lastError,
|
||||
); err != nil {
|
||||
return ConvergenceSnapshot{}, fmt.Errorf("scan convergence snapshot: %w", err)
|
||||
}
|
||||
if nextAttempt.Valid {
|
||||
parsed, parseErr := parseTime(nextAttempt.String)
|
||||
if parseErr != nil {
|
||||
return ConvergenceSnapshot{}, parseErr
|
||||
}
|
||||
value.NextAttemptAt = &parsed
|
||||
}
|
||||
if lastError.Valid {
|
||||
value.LastErrorCode = lastError.String
|
||||
}
|
||||
value.Converged = value.ObservedGeneration == value.Generation &&
|
||||
value.FailureCount == 0 && value.ActualState == device.ActualOnline
|
||||
if !value.Converged {
|
||||
snapshot.Unconverged++
|
||||
}
|
||||
snapshot.Devices = append(snapshot.Devices, value)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return ConvergenceSnapshot{}, fmt.Errorf("iterate convergence snapshot: %w", err)
|
||||
}
|
||||
snapshot.Total = len(snapshot.Devices)
|
||||
return snapshot, nil
|
||||
}
|
||||
|
||||
const deviceColumns = `d.id, d.tenant_id, d.site_id, d.serial_number, d.name, d.modality,
|
||||
d.desired_state, d.actual_state, d.endpoint_ref, d.credential_ref,
|
||||
d.path_name, d.generation, d.created_at, d.updated_at`
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"yovision/sense/internal/device"
|
||||
)
|
||||
@@ -144,6 +145,42 @@ func TestLowerQuotaDoesNotDisableExistingStreams(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConvergenceSnapshotAndRuntimeReconcileRequest(t *testing.T) {
|
||||
t.Parallel()
|
||||
store := openTestStore(t)
|
||||
ctx := context.Background()
|
||||
if err := store.EnsureSite(ctx, device.Site{TenantID: "tenant", ID: "site", Name: "Site"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := store.CreateDevice(ctx, videoDevice(1, "tenant", "site")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now := time.Date(2026, 8, 7, 0, 0, 0, 0, time.UTC)
|
||||
if err := store.MarkReconciled(ctx, "camera-001", 1, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := store.UpdateActualState(ctx, "camera-001", device.ActualOnline, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
snapshot, err := store.ConvergenceSnapshot(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if snapshot.Total != 1 || snapshot.Unconverged != 0 {
|
||||
t.Fatalf("expected converged snapshot, got %+v", snapshot)
|
||||
}
|
||||
if err := store.RequestReconcile(ctx, "camera-001", now.Add(time.Second)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
snapshot, err = store.ConvergenceSnapshot(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if snapshot.Unconverged != 1 || snapshot.Devices[0].ObservedGeneration != 0 {
|
||||
t.Fatalf("runtime loss must invalidate convergence: %+v", snapshot)
|
||||
}
|
||||
}
|
||||
|
||||
func openTestStore(t *testing.T) *SQLite {
|
||||
t.Helper()
|
||||
dsn := "file:" + filepath.ToSlash(filepath.Join(t.TempDir(), "sense.db"))
|
||||
|
||||
@@ -0,0 +1,441 @@
|
||||
[CmdletBinding()]
|
||||
param(
|
||||
[Parameter(Mandatory = $true)]
|
||||
[string]$CameraEnv,
|
||||
[string]$RuntimeRoot = (Join-Path ([IO.Path]::GetTempPath()) 'yovision-t006'),
|
||||
[ValidateRange(1, 1440)]
|
||||
[int]$ObservationMinutes = 30,
|
||||
[switch]$KeepSession
|
||||
)
|
||||
|
||||
$ErrorActionPreference = 'Stop'
|
||||
$ProgressPreference = 'SilentlyContinue'
|
||||
|
||||
function Read-EnvFile([string]$Path) {
|
||||
$result = @{}
|
||||
Get-Content -LiteralPath $Path | ForEach-Object {
|
||||
$line = $_.Trim()
|
||||
if (-not $line -or $line.StartsWith('#') -or -not $line.Contains('=')) {
|
||||
return
|
||||
}
|
||||
$parts = $line -split '=', 2
|
||||
$value = $parts[1].Trim()
|
||||
if ($value.Length -ge 2 -and (($value.StartsWith('"') -and $value.EndsWith('"')) -or ($value.StartsWith("'") -and $value.EndsWith("'")))) {
|
||||
$value = $value.Substring(1, $value.Length - 2)
|
||||
}
|
||||
$result[$parts[0].Trim().ToLowerInvariant()] = $value
|
||||
}
|
||||
return $result
|
||||
}
|
||||
|
||||
function Require-Keys([hashtable]$Config, [string[]]$Keys) {
|
||||
foreach ($key in $Keys) {
|
||||
if (-not $Config.ContainsKey($key) -or [string]::IsNullOrWhiteSpace($Config[$key])) {
|
||||
throw "camera environment is missing required key: $key"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function Assert-PortFree([int]$Port) {
|
||||
$client = [Net.Sockets.TcpClient]::new()
|
||||
try {
|
||||
$task = $client.ConnectAsync('127.0.0.1', $Port)
|
||||
if ($task.Wait(250) -and $client.Connected) {
|
||||
throw "required local port is already in use: $Port"
|
||||
}
|
||||
}
|
||||
catch [AggregateException] {
|
||||
return
|
||||
}
|
||||
catch [Net.Sockets.SocketException] {
|
||||
return
|
||||
}
|
||||
finally {
|
||||
$client.Dispose()
|
||||
}
|
||||
}
|
||||
|
||||
function Wait-Port([int]$Port, [int]$TimeoutSeconds = 30) {
|
||||
$watch = [Diagnostics.Stopwatch]::StartNew()
|
||||
while ($watch.Elapsed.TotalSeconds -lt $TimeoutSeconds) {
|
||||
$client = [Net.Sockets.TcpClient]::new()
|
||||
try {
|
||||
$task = $client.ConnectAsync('127.0.0.1', $Port)
|
||||
if ($task.Wait(500) -and $client.Connected) {
|
||||
return
|
||||
}
|
||||
}
|
||||
catch {
|
||||
}
|
||||
finally {
|
||||
$client.Dispose()
|
||||
}
|
||||
Start-Sleep -Milliseconds 250
|
||||
}
|
||||
throw "local port did not become ready: $Port"
|
||||
}
|
||||
|
||||
function Start-ManagedProcess(
|
||||
[string]$Name,
|
||||
[string]$FilePath,
|
||||
[string[]]$Arguments,
|
||||
[hashtable]$Environment = @{}
|
||||
) {
|
||||
$start = [Diagnostics.ProcessStartInfo]::new()
|
||||
$start.FileName = $FilePath
|
||||
$start.WorkingDirectory = $session
|
||||
$start.UseShellExecute = $false
|
||||
$start.CreateNoWindow = $true
|
||||
$start.RedirectStandardOutput = $true
|
||||
$start.RedirectStandardError = $true
|
||||
foreach ($argument in $Arguments) {
|
||||
$start.ArgumentList.Add($argument)
|
||||
}
|
||||
foreach ($entry in $Environment.GetEnumerator()) {
|
||||
$start.Environment[$entry.Key] = [string]$entry.Value
|
||||
}
|
||||
$process = [Diagnostics.Process]::new()
|
||||
$process.StartInfo = $start
|
||||
if (-not $process.Start()) {
|
||||
throw "failed to start process: $Name"
|
||||
}
|
||||
return [pscustomobject]@{
|
||||
Name = $Name
|
||||
Process = $process
|
||||
Stdout = $process.StandardOutput.ReadToEndAsync()
|
||||
Stderr = $process.StandardError.ReadToEndAsync()
|
||||
StartedAt = [DateTimeOffset]::UtcNow
|
||||
}
|
||||
}
|
||||
|
||||
function Stop-ManagedProcess($Managed) {
|
||||
if ($null -eq $Managed -or $null -eq $Managed.Process) {
|
||||
return
|
||||
}
|
||||
if (-not $Managed.Process.HasExited) {
|
||||
# Every fixture is started as a leaf process. Killing only that exact
|
||||
# process avoids Windows process-tree edge cases during fault tests.
|
||||
$Managed.Process.Kill()
|
||||
$Managed.Process.WaitForExit(10000) | Out-Null
|
||||
}
|
||||
}
|
||||
|
||||
function Assert-Alive($Managed) {
|
||||
if ($null -eq $Managed -or $Managed.Process.HasExited) {
|
||||
$code = if ($null -eq $Managed) { 'not-started' } else { $Managed.Process.ExitCode }
|
||||
throw "required process exited: $($Managed.Name), code=$code"
|
||||
}
|
||||
}
|
||||
|
||||
function Invoke-LabStatus {
|
||||
$lastDiagnostic = ''
|
||||
foreach ($attempt in 1..5) {
|
||||
$raw = @(& $labBinary status -db $databaseDSN 2>&1)
|
||||
if ($LASTEXITCODE -eq 0) {
|
||||
try {
|
||||
return (($raw -join "`n") | ConvertFrom-Json)
|
||||
}
|
||||
catch {
|
||||
$lastDiagnostic = 'invalid JSON response'
|
||||
}
|
||||
}
|
||||
else {
|
||||
$lastDiagnostic = (($raw -join "`n") -split "`r?`n" | Select-Object -Last 2) -join ' | '
|
||||
}
|
||||
Start-Sleep -Milliseconds 250
|
||||
}
|
||||
throw "sense-lab status failed after retries: $lastDiagnostic"
|
||||
}
|
||||
|
||||
function Wait-Converged([int]$TimeoutSeconds = 180) {
|
||||
$watch = [Diagnostics.Stopwatch]::StartNew()
|
||||
while ($watch.Elapsed.TotalSeconds -lt $TimeoutSeconds) {
|
||||
try {
|
||||
$snapshot = Invoke-LabStatus
|
||||
if ([int]$snapshot.total -eq 5 -and [int]$snapshot.unconverged -eq 0) {
|
||||
return [Math]::Round($watch.Elapsed.TotalSeconds, 1)
|
||||
}
|
||||
}
|
||||
catch {
|
||||
}
|
||||
Start-Sleep -Seconds 1
|
||||
}
|
||||
throw 'five devices did not converge before timeout'
|
||||
}
|
||||
|
||||
function Wait-DeviceState([string]$ID, [string]$State, [int]$TimeoutSeconds = 90) {
|
||||
$watch = [Diagnostics.Stopwatch]::StartNew()
|
||||
while ($watch.Elapsed.TotalSeconds -lt $TimeoutSeconds) {
|
||||
try {
|
||||
$snapshot = Invoke-LabStatus
|
||||
$match = @($snapshot.devices | Where-Object { $_.id -eq $ID })
|
||||
if ($match.Count -eq 1 -and $match[0].actual_state -eq $State) {
|
||||
return [Math]::Round($watch.Elapsed.TotalSeconds, 1)
|
||||
}
|
||||
}
|
||||
catch {
|
||||
}
|
||||
Start-Sleep -Seconds 1
|
||||
}
|
||||
throw "device did not reach expected state: $ID/$State"
|
||||
}
|
||||
|
||||
function Wait-SenseHealth($Managed, [int]$TimeoutSeconds = 30) {
|
||||
$watch = [Diagnostics.Stopwatch]::StartNew()
|
||||
while ($watch.Elapsed.TotalSeconds -lt $TimeoutSeconds) {
|
||||
if ($Managed.Process.HasExited) {
|
||||
$output = @(
|
||||
$Managed.Stdout.GetAwaiter().GetResult()
|
||||
$Managed.Stderr.GetAwaiter().GetResult()
|
||||
) -join "`n"
|
||||
$summary = (($output -split "`r?`n") | Where-Object { $_ } | Select-Object -Last 3) -join ' | '
|
||||
throw "Sense exited before health check, code=$($Managed.Process.ExitCode), output=$summary"
|
||||
}
|
||||
try {
|
||||
$response = Invoke-RestMethod -Method Get -Uri 'http://127.0.0.1:18080/healthz' -TimeoutSec 2 -NoProxy
|
||||
if ($response.status -eq 'ok') {
|
||||
return
|
||||
}
|
||||
}
|
||||
catch {
|
||||
}
|
||||
Start-Sleep -Milliseconds 500
|
||||
}
|
||||
$listeners = @(Get-NetTCPConnection -State Listen -OwningProcess $Managed.Process.Id -ErrorAction SilentlyContinue |
|
||||
ForEach-Object { "$($_.LocalAddress):$($_.LocalPort)" })
|
||||
throw "Sense health endpoint did not become ready; process listeners=$($listeners -join ',')"
|
||||
}
|
||||
|
||||
function Start-Publisher([int]$Index) {
|
||||
return Start-ManagedProcess "publisher-$Index" $ffmpeg @(
|
||||
'-hide_banner', '-loglevel', 'warning', '-re',
|
||||
'-f', 'lavfi', '-i', "testsrc2=size=640x360:rate=10",
|
||||
'-c:v', 'libx264', '-preset', 'ultrafast', '-tune', 'zerolatency',
|
||||
'-pix_fmt', 'yuv420p', '-g', '10', '-an',
|
||||
'-f', 'rtsp', '-rtsp_transport', 'tcp',
|
||||
"rtsp://127.0.0.1:8555/synthetic-$Index"
|
||||
)
|
||||
}
|
||||
|
||||
function Start-Proxy {
|
||||
return Start-ManagedProcess 'real-camera-network-proxy' $proxyBinary @(
|
||||
'-listen', '127.0.0.1:10554', '-upstream', "$($camera.host):$($camera.rtspport)"
|
||||
)
|
||||
}
|
||||
|
||||
function Start-Sense {
|
||||
$environment = @{
|
||||
SENSE_HTTP_ADDR = '127.0.0.1:18080'
|
||||
SENSE_DB_DSN = $databaseDSN
|
||||
SENSE_MEDIAMTX_URL = 'http://127.0.0.1:9997'
|
||||
SENSE_RECONCILE_INTERVAL = '1s'
|
||||
SENSE_PROBE_INTERVAL = '1s'
|
||||
SENSE_ONVIF_MODE = 'standard'
|
||||
SENSE_ONVIF_RTSP_REWRITE_HOST = '127.0.0.1'
|
||||
SENSE_ONVIF_RTSP_REWRITE_PORT = '10554'
|
||||
SENSE_ONVIF_RTSP_STRIP_QUERY = 'true'
|
||||
SENSE_CREDENTIAL_CAMERA_ONVIF_USERNAME = $camera.onvifuser
|
||||
SENSE_CREDENTIAL_CAMERA_ONVIF_PASSWORD = $camera.onvifpwd
|
||||
SENSE_CREDENTIAL_CAMERA_RTSP_USERNAME = $camera.username
|
||||
SENSE_CREDENTIAL_CAMERA_RTSP_PASSWORD = $camera.password
|
||||
}
|
||||
return Start-ManagedProcess -Name 'sense-api' -FilePath $senseBinary -Arguments @() -Environment $environment
|
||||
}
|
||||
|
||||
function Start-ProductionMediaMTX {
|
||||
return Start-ManagedProcess 'mediamtx-production' $mediaMTX @($productionConfig)
|
||||
}
|
||||
|
||||
$repoRoot = (Resolve-Path (Join-Path $PSScriptRoot '..\..')).Path
|
||||
$cameraPath = (Resolve-Path -LiteralPath $CameraEnv).Path
|
||||
$camera = Read-EnvFile $cameraPath
|
||||
Require-Keys $camera @('host', 'username', 'password', 'rtspport', 'onvif', 'onvifuser', 'onvifpwd')
|
||||
|
||||
foreach ($port in 8554, 8555, 9997, 10554, 18080) {
|
||||
Assert-PortFree $port
|
||||
}
|
||||
|
||||
New-Item -ItemType Directory -Path $RuntimeRoot -Force | Out-Null
|
||||
$mediaDirectory = Join-Path $RuntimeRoot 'mediamtx-v1.19.3'
|
||||
$mediaMTX = Join-Path $mediaDirectory 'mediamtx.exe'
|
||||
if (-not (Test-Path -LiteralPath $mediaMTX)) {
|
||||
$zip = Join-Path $RuntimeRoot 'mediamtx_v1.19.3_windows_amd64.zip'
|
||||
Invoke-WebRequest 'https://github.com/bluenviron/mediamtx/releases/download/v1.19.3/mediamtx_v1.19.3_windows_amd64.zip' -OutFile $zip
|
||||
$actualHash = (Get-FileHash -LiteralPath $zip -Algorithm SHA256).Hash.ToLowerInvariant()
|
||||
if ($actualHash -ne '5d82148d1032a6a190d9909a2997d9989457aaadf49af87dd02cd4512d31bebe') {
|
||||
throw 'MediaMTX checksum mismatch'
|
||||
}
|
||||
New-Item -ItemType Directory -Path $mediaDirectory -Force | Out-Null
|
||||
Expand-Archive -LiteralPath $zip -DestinationPath $mediaDirectory -Force
|
||||
}
|
||||
|
||||
$ffmpeg = (Get-Command ffmpeg -ErrorAction Stop).Source
|
||||
$session = Join-Path $RuntimeRoot ('session-' + [Guid]::NewGuid().ToString('N'))
|
||||
New-Item -ItemType Directory -Path $session | Out-Null
|
||||
$senseBinary = Join-Path $session 'sense-api.exe'
|
||||
$labBinary = Join-Path $session 'sense-lab.exe'
|
||||
$proxyBinary = Join-Path $session 'rtsp-fault-proxy.exe'
|
||||
$productionConfig = Join-Path $session 'mediamtx-production.yml'
|
||||
$syntheticConfig = Join-Path $session 'mediamtx-synthetic.yml'
|
||||
[IO.File]::Copy((Join-Path $repoRoot 'Sense\deploy\mediamtx.yml'), $productionConfig)
|
||||
[IO.File]::Copy((Join-Path $repoRoot 'Sense\deploy\mediamtx-synthetic.yml'), $syntheticConfig)
|
||||
|
||||
& go -C (Join-Path $repoRoot 'Sense') build -o $senseBinary ./cmd/sense-api
|
||||
if ($LASTEXITCODE -ne 0) { throw 'build sense-api failed' }
|
||||
& go -C (Join-Path $repoRoot 'Sense') build -o $labBinary ./cmd/sense-lab
|
||||
if ($LASTEXITCODE -ne 0) { throw 'build sense-lab failed' }
|
||||
& go -C (Join-Path $repoRoot 'Sense') build -o $proxyBinary ./cmd/rtsp-fault-proxy
|
||||
if ($LASTEXITCODE -ne 0) { throw 'build RTSP fault proxy failed' }
|
||||
|
||||
$databasePath = Join-Path $session 'sense.db'
|
||||
$databaseDSN = 'file:' + $databasePath.Replace('\', '/')
|
||||
$manifestPath = Join-Path $session 'manifest.json'
|
||||
$onvifEndpoint = if ($camera.onvif -match '^https?://') {
|
||||
$camera.onvif
|
||||
} else {
|
||||
"http://$($camera.host):$($camera.onvif)/onvif/device_service"
|
||||
}
|
||||
$devices = @(
|
||||
[ordered]@{ id = 'camera-real'; tenant_id = 'lab'; site_id = 'site'; serial_number = 'e9ed6a555ae0'; name = 'Approved real camera'; capabilities = @('video_capture', 'audio_capture'); endpoint_ref = $onvifEndpoint; credential_ref = 'env://camera'; path_name = 'sense/lab/site/camera-real' }
|
||||
)
|
||||
foreach ($index in 1..4) {
|
||||
$devices += [ordered]@{ id = "synthetic-$index"; tenant_id = 'lab'; site_id = 'site'; serial_number = "synthetic-$index"; name = "Synthetic source $index"; capabilities = @('video_capture'); endpoint_ref = "rtsp://127.0.0.1:8555/synthetic-$index"; credential_ref = ''; path_name = "sense/lab/site/synthetic-$index" }
|
||||
}
|
||||
[ordered]@{
|
||||
site = [ordered]@{ tenant_id = 'lab'; id = 'site'; name = 'T-006 Lab'; max_video_channels = 16 }
|
||||
devices = $devices
|
||||
} | ConvertTo-Json -Depth 8 | Set-Content -LiteralPath $manifestPath -Encoding UTF8
|
||||
|
||||
$managed = [Collections.Generic.List[object]]::new()
|
||||
$publishers = @{}
|
||||
$events = [Collections.Generic.List[object]]::new()
|
||||
$observationSamples = 0
|
||||
$maxUnconverged = 0
|
||||
$success = $false
|
||||
$stage = 'starting fixtures'
|
||||
$failureMessage = $null
|
||||
try {
|
||||
$sourceMedia = Start-ManagedProcess 'mediamtx-synthetic' $mediaMTX @($syntheticConfig)
|
||||
$managed.Add($sourceMedia)
|
||||
Wait-Port 8555
|
||||
$productionMedia = Start-ProductionMediaMTX
|
||||
$managed.Add($productionMedia)
|
||||
Wait-Port 9997
|
||||
|
||||
foreach ($index in 1..4) {
|
||||
$publishers[$index] = Start-Publisher $index
|
||||
$managed.Add($publishers[$index])
|
||||
}
|
||||
$stage = 'checking synthetic publishers'
|
||||
Start-Sleep -Seconds 3
|
||||
foreach ($publisher in $publishers.Values) { Assert-Alive $publisher }
|
||||
|
||||
$proxy = Start-Proxy
|
||||
$managed.Add($proxy)
|
||||
Wait-Port 10554
|
||||
|
||||
& $labBinary seed -db $databaseDSN -manifest $manifestPath | Out-Null
|
||||
if ($LASTEXITCODE -ne 0) { throw 'seed device ledger failed' }
|
||||
$stage = 'initial convergence'
|
||||
$sense = Start-Sense
|
||||
$managed.Add($sense)
|
||||
Wait-SenseHealth $sense
|
||||
$initialSeconds = Wait-Converged
|
||||
$events.Add([ordered]@{ event = 'initial_convergence'; seconds = $initialSeconds; unconverged = 0 })
|
||||
|
||||
$stage = 'real camera network recovery'
|
||||
Stop-ManagedProcess $proxy
|
||||
$offlineSeconds = Wait-DeviceState 'camera-real' 'offline'
|
||||
$proxy = Start-Proxy
|
||||
$managed.Add($proxy)
|
||||
Wait-Port 10554
|
||||
$recoverySeconds = Wait-Converged
|
||||
$events.Add([ordered]@{ event = 'real_camera_network'; offline_detect_seconds = $offlineSeconds; recovery_seconds = $recoverySeconds; unconverged = 0 })
|
||||
|
||||
$stage = 'synthetic publisher recovery'
|
||||
Stop-ManagedProcess $publishers[2]
|
||||
$offlineSeconds = Wait-DeviceState 'synthetic-2' 'offline'
|
||||
$publishers[2] = Start-Publisher 2
|
||||
$managed.Add($publishers[2])
|
||||
$recoverySeconds = Wait-Converged
|
||||
$events.Add([ordered]@{ event = 'synthetic_publisher'; offline_detect_seconds = $offlineSeconds; recovery_seconds = $recoverySeconds; unconverged = 0 })
|
||||
|
||||
$stage = 'Sense restart recovery'
|
||||
Stop-ManagedProcess $sense
|
||||
$restartWatch = [Diagnostics.Stopwatch]::StartNew()
|
||||
$sense = Start-Sense
|
||||
$managed.Add($sense)
|
||||
Wait-SenseHealth $sense
|
||||
$recoverySeconds = Wait-Converged
|
||||
$events.Add([ordered]@{ event = 'sense_restart'; process_ready_seconds = [Math]::Round($restartWatch.Elapsed.TotalSeconds, 1); recovery_seconds = $recoverySeconds; unconverged = 0 })
|
||||
|
||||
$stage = 'MediaMTX restart recovery'
|
||||
Stop-ManagedProcess $productionMedia
|
||||
Start-Sleep -Seconds 3
|
||||
$productionMedia = Start-ProductionMediaMTX
|
||||
$managed.Add($productionMedia)
|
||||
Wait-Port 9997
|
||||
$recoverySeconds = Wait-Converged 240
|
||||
$events.Add([ordered]@{ event = 'mediamtx_restart'; recovery_seconds = $recoverySeconds; unconverged = 0 })
|
||||
|
||||
$stage = 'checking MediaMTX path count'
|
||||
$configured = Invoke-RestMethod -Method Get -Uri 'http://127.0.0.1:9997/v3/config/paths/list' -TimeoutSec 5 -NoProxy
|
||||
if ([int]$configured.itemCount -ne 5) {
|
||||
throw "expected 5 MediaMTX paths, got $($configured.itemCount)"
|
||||
}
|
||||
|
||||
$stage = 'stability observation'
|
||||
$observation = [Diagnostics.Stopwatch]::StartNew()
|
||||
$targetSeconds = $ObservationMinutes * 60
|
||||
while ($observation.Elapsed.TotalSeconds -lt $targetSeconds) {
|
||||
Assert-Alive $sourceMedia
|
||||
Assert-Alive $productionMedia
|
||||
Assert-Alive $sense
|
||||
Assert-Alive $proxy
|
||||
foreach ($publisher in $publishers.Values) { Assert-Alive $publisher }
|
||||
$snapshot = Invoke-LabStatus
|
||||
$observationSamples++
|
||||
$maxUnconverged = [Math]::Max($maxUnconverged, [int]$snapshot.unconverged)
|
||||
if ([int]$snapshot.total -ne 5 -or [int]$snapshot.unconverged -ne 0) {
|
||||
throw "observation detected unconverged devices: $($snapshot.unconverged)"
|
||||
}
|
||||
Start-Sleep -Seconds 10
|
||||
}
|
||||
$observationSeconds = [Math]::Round($observation.Elapsed.TotalSeconds, 1)
|
||||
$final = Invoke-LabStatus
|
||||
$success = $true
|
||||
[ordered]@{
|
||||
success = $true
|
||||
mediamtx_version = 'v1.19.3'
|
||||
mediamtx_sha256 = '5d82148d1032a6a190d9909a2997d9989457aaadf49af87dd02cd4512d31bebe'
|
||||
source_count = 5
|
||||
synthetic_publishers = 4
|
||||
configured_paths = [int]$configured.itemCount
|
||||
recovery = $events
|
||||
observation_seconds = $observationSeconds
|
||||
observation_samples = $observationSamples
|
||||
maximum_unconverged = $maxUnconverged
|
||||
final_unconverged = [int]$final.unconverged
|
||||
} | ConvertTo-Json -Depth 8
|
||||
}
|
||||
catch {
|
||||
$failureMessage = "T-006 stage '$stage' failed: $($_.Exception.Message)"
|
||||
}
|
||||
finally {
|
||||
foreach ($item in @($managed)) {
|
||||
Stop-ManagedProcess $item
|
||||
}
|
||||
if (-not $KeepSession) {
|
||||
Get-ChildItem -LiteralPath $session -Recurse -Force -ErrorAction SilentlyContinue | ForEach-Object { $_.Attributes = 'Normal' }
|
||||
if (Test-Path -LiteralPath $session) {
|
||||
(Get-Item -LiteralPath $session -Force).Attributes = 'Directory'
|
||||
Remove-Item -LiteralPath $session -Recurse -ErrorAction SilentlyContinue
|
||||
}
|
||||
}
|
||||
if (-not $success) {
|
||||
Write-Warning 'T-006 integration did not complete; no success evidence was emitted.'
|
||||
}
|
||||
}
|
||||
if ($failureMessage) {
|
||||
throw $failureMessage
|
||||
}
|
||||
@@ -1,11 +1,11 @@
|
||||
# 当前实现状态
|
||||
|
||||
> 快照日期:2026-08-04。只记录仓库现实与 blocker;任务实时状态到 Gitea Issue 查看。
|
||||
> 快照日期:2026-08-07。只记录仓库现实与 blocker;任务实时状态到 Gitea Issue 查看。
|
||||
|
||||
## 当前阶段
|
||||
|
||||
- 阶段:M0 已收敛为首期指定摄像头型号准入,实机验证尚未执行;M1 Sense 无实机软件骨架已完成,但不得提前宣称 M0/M1 出口完成。
|
||||
- 生产代码:Sense M1 无实机骨架已建立,包含可构建进程、SQLite 台账、ONVIF port/fake、MediaMTX 生成客户端、最小对账与探活;真实 ONVIF adapter 和单实机混合源验收仍未开始。
|
||||
- 阶段:M0 指定摄像头型号准入已完成;M1 的“一实机 + 四合成源”实验室软件闭环已通过,但五条独立真实上游和生产 SLA 尚未验收。
|
||||
- 生产代码:Sense 已包含可构建进程、SQLite 台账、标准 ONVIF SOAP/WS-Security adapter、凭据引用、MediaMTX 生成客户端、对账、探活和 MediaMTX 重启重建;公共设备管理 API 与生产部署能力仍未实现。
|
||||
- 默认容量:16 路;单站点本阶段上限 128 路,必须横向分片。
|
||||
|
||||
## 仓库现实
|
||||
@@ -13,7 +13,8 @@
|
||||
- `Sense/` 已有 Go module 与 `cmd/sense-api`;`Brain/`、`Bell/` 仍只有目录占位。
|
||||
- Sense 设备模型使用 `modality + capabilities`,SQLite 执行 v1 migration;视频配额默认 16、允许 1~128,17/128/129、新增/启用和“降低配额不关闭已有流”均有测试。
|
||||
- MediaMTX 固定为独立二进制 `v1.19.3`,官方 OpenAPI 已按 SHA-256 vendoring,并由固定 `oapi-codegen v2.8.0` 生成客户端;手写薄封装有 create/read/delete、幂等 ensure 与探活假 HTTP 测试。
|
||||
- T-003 对账进度与指数退避持久化,覆盖取消和 SQLite 重启恢复;当前不枚举/删除孤儿,也不包含真实摄像头 adapter。
|
||||
- T-003 对账进度与指数退避持久化,覆盖取消和 SQLite 重启恢复;T-006 增加真实 ONVIF adapter、RTSP router、实验室播种/状态工具、故障代理和五路自动验收。当前仍不枚举/删除孤儿。
|
||||
- T-006 正式使用 1 台准入实机和 4 个独立合成 publisher 连续观察 `1806.6 s` / 180 次采样,四类恢复均通过,最大与最终 `unconverged` 均为 0;详细证据见 `docs/research/sense-5-stream-integration.md`。
|
||||
- `docs/raw/01`~`08` 已记录需求、分析、方案、客户场景、事件比对和三系统职责。
|
||||
- `docs/raw/contracts/event-v0.1.schema.json` 已冻结,并有多份示例与语义说明。
|
||||
- harness coding 文档、上下文清单、Gitea Issue/PR 模板和治理脚本已接入。
|
||||
@@ -49,7 +50,7 @@ Sense 默认监听 `127.0.0.1:8080`,提供 `/healthz` 与 `/readyz` 运维探
|
||||
## 当前 blocker / 待确认
|
||||
|
||||
- T-001 已完成一台 HIKVISION `DS-2CD3321FD-IW1-T`、硬件 `0x0`、固件 `V5.5.61 build 180929` 的 ONVIF/RTSP 指定基线准入;主/子码流、错误鉴权、三个 ONVIF 核心操作和时间漂移均有脱敏证据。项目负责人豁免了无法现场执行的三次真实断网恢复,该项没有原始时间线,单样机结论也不代表批次或多品牌兼容。
|
||||
- 后续本地开发统一使用现有一台 Hikvision IP Camera;T-006 以该实机 + 至少 4 条独立合成 RTSP 源完成五路软件闭环,不要求当前购买更多摄像头。合成 publisher 的具体工具和版本在 T-006 领取后冻结。
|
||||
- 后续本地开发统一使用现有一台 Hikvision IP Camera;T-006 已用该实机 + 4 条独立 FFmpeg 8.1.2 合成 RTSP 源完成五路软件闭环,不要求当前购买更多摄像头。
|
||||
- T-007 保留至少 5 条独立真实上游的客户现场门禁,可使用客户授权、借用或租赁设备;它不阻塞本地开发/M1 实验室出口,但继续阻塞生产试点和真实多路 SLA。
|
||||
- S2 真实生产试点的未成年人影像、公共安全视频法规适用性和最终留存政策仍需客户/法务确认,阻塞 M3 上线但不阻塞 M1 实验室骨架。
|
||||
- 人脸方向已延后至 M5 的 S4 成人园区候选试点;必要性/PIP 影响评估、单独同意与替代方式、合法底库来源和删除流程未完成,阻塞人脸能力上线。
|
||||
@@ -59,7 +60,7 @@ Sense 默认监听 `127.0.0.1:8080`,提供 `/healthz` 与 `/readyz` 运维探
|
||||
|
||||
## 下一步
|
||||
|
||||
下一步领取 T-006,使用 T-001 已准入的一台实机和至少 4 条独立合成 RTSP 上游完成五路软件闭环,无需采购额外摄像头。客户/借用/租赁条件具备后再执行 T-007 真实多路现场门禁。实时领取状态仍以 Gitea 为准。
|
||||
下一步按 Gitea 依赖关系选择后续 M1 任务;客户授权、借用或租赁条件具备后执行 T-007 五条独立真实上游现场门禁。T-006 的合成结果不解除 T-007,也不形成容量或生产 SLA 承诺。实时领取状态仍以 Gitea 为准。
|
||||
|
||||
## 已知风险
|
||||
|
||||
|
||||
@@ -0,0 +1,102 @@
|
||||
# Sense 单实机与五路混合源集成记录
|
||||
|
||||
> 执行日期:2026-08-07
|
||||
>
|
||||
> 结论:PASS(M1 实验室集成 smoke)
|
||||
|
||||
## 1. 验收范围与结论
|
||||
|
||||
本轮使用 T-001 已批准的 1 台真实 HIKVISION `DS-2CD3321FD-IW1-T`(硬件 `0x0`、固件 `V5.5.61 build 180929`)和 4 个可独立启停的 FFmpeg 合成 RTSP publisher,完成 5 路期望态、MediaMTX path、运行态探活和自动恢复闭环。
|
||||
|
||||
正式观察持续 `1806.6` 秒,共读取 `180` 个脱敏收敛快照;`maximum_unconverged = 0`、`final_unconverged = 0`。真实摄像头网络路径、单个合成 publisher、Sense 进程和 MediaMTX 进程四类恢复均通过,无需人工改 MediaMTX 配置或重建 SQLite。
|
||||
|
||||
该结论只证明“一台已准入实机 + 四条无人物合成源”的 M1 软件闭环,不证明五台真实设备兼容、故障隔离、默认 16 路或最大 128 路容量,也不构成生产 SLA。五条独立真实上游仍由 T-007 验收。
|
||||
|
||||
## 2. 实测架构
|
||||
|
||||
```text
|
||||
真实摄像头 -- RTSP/TCP -- 故障代理 :10554 --+
|
||||
|
|
||||
4 × FFmpeg -- publisher --> 源 MediaMTX :8555 +--> 生产 MediaMTX :8554
|
||||
^
|
||||
SQLite 期望态 --> Sense reconcile --> ONVIF/router ----+ API :9997
|
||||
^ |
|
||||
+---- Sense probe +---- 运行态反馈/重新对账
|
||||
```
|
||||
|
||||
- SQLite 是期望态真相源;设备台账只保存 `env://camera` 凭据引用。
|
||||
- HTTP/HTTPS endpoint 进入标准 ONVIF SOAP 1.2 + WS-Security PasswordDigest adapter;合成 `rtsp://` endpoint 由 router 直接映射,不伪装 ONVIF。
|
||||
- ONVIF adapter 获取设备信息、服务、profile 和 stream URI,在内存中注入 RTSP 凭据,并支持显式 NAT 主机/端口重写。错误不回显 endpoint、URI 或凭据。
|
||||
- 生产 MediaMTX path 只能由 Sense 对账通过 API 创建;`sense-lab` 只负责实验室播种和读取脱敏状态,不是冻结的公共设备管理 API。
|
||||
- probe 发现 MediaMTX path 丢失时使已观察代失效,触发重新对账,因此 MediaMTX 空配置重启后能从 SQLite 自动恢复五条 path。
|
||||
- 真实网络故障由本地 TCP 代理进程启停注入;不修改摄像头配置,也不继承 T-001 的物理断网豁免。
|
||||
|
||||
## 3. 固定版本与产物
|
||||
|
||||
| 组件 | 实测版本 / 摘要 |
|
||||
|---|---|
|
||||
| Sense | `dev`,Go `1.23.0 windows/amd64`;同源重建 `sense-api.exe` SHA-256 `5d49b91ff6c76802da3096e9155f6496187059a6a97114fbeb4e9b5a39128983` |
|
||||
| MediaMTX | `v1.19.3`,Windows amd64 发布包 SHA-256 `5d82148d1032a6a190d9909a2997d9989457aaadf49af87dd02cd4512d31bebe` |
|
||||
| FFmpeg | `8.1.2-full_build-www.gyan.dev`,4 个独立 `libx264` publisher |
|
||||
| SQLite driver | `modernc.org/sqlite v1.38.2`,随 Sense module 锁定 |
|
||||
|
||||
二进制、数据库、临时清单、日志和媒体都位于系统临时目录并在运行后删除;仓库不保存摄像头主机、内网地址、账号、密码、完整 RTSP/ONVIF URI、人物画面或可复用 token。
|
||||
|
||||
## 4. 可重复步骤
|
||||
|
||||
1. 将 T-001 格式的 `ip_camera.env` 放在仓库外或 Git 忽略路径;字段为 `host`、`username`、`password`、`rtspport`、`onvif`、`onvifuser`、`onvifpwd`。
|
||||
2. 确保本机 `8554`、`8555`、`9997`、`10554`、`18080` 未被占用,并安装 Go 1.23 与带 `libx264` 的 FFmpeg。
|
||||
3. 执行正式脚本;默认观察 30 分钟:
|
||||
|
||||
```powershell
|
||||
./Sense/scripts/t006-integration.ps1 -CameraEnv D:\path\to\ip_camera.env
|
||||
```
|
||||
|
||||
4. 脚本校验 MediaMTX 发布包 SHA-256,构建三个 Go 命令,启动两套 MediaMTX 和四个独立 publisher,播种一条真实设备和四条合成设备,再依次执行四类故障。
|
||||
5. 成功输出必须同时满足:`source_count = 5`、`synthetic_publishers = 4`、`configured_paths = 5`、每类恢复 `unconverged = 0`、观察至少 1800 秒、`maximum_unconverged = 0`、`final_unconverged = 0`。
|
||||
|
||||
调试时可加 `-ObservationMinutes 1`,但短窗口不能替代正式证据。本轮先完成 `60.6` 秒、6 次采样的 smoke,随后才执行正式窗口。
|
||||
|
||||
## 5. 正式结果
|
||||
|
||||
| 项目 | 结果 | 实测值 |
|
||||
|---|---|---:|
|
||||
| 初始五路自动收敛 | PASS | `8.2 s` |
|
||||
| 真实摄像头网络路径离线探测 | PASS | `1.0 s` |
|
||||
| 真实摄像头网络路径恢复 | PASS | `5.1 s` |
|
||||
| 单个合成 publisher 离线探测 | PASS | `1.0 s` |
|
||||
| 单个合成 publisher 恢复 | PASS | `5.1 s` |
|
||||
| Sense 重启到健康 | PASS | `0.5 s` |
|
||||
| Sense 重启后重新收敛 | PASS | `0.0 s`(SQLite 已保持一致) |
|
||||
| MediaMTX 空配置重启后恢复 | PASS | `2.1 s` |
|
||||
| 自动配置 path 数 | PASS | `5` |
|
||||
| 连续观察 | PASS | `1806.6 s` / `180` 次采样 |
|
||||
| 观察期最大 / 最终未收敛数 | PASS | `0 / 0` |
|
||||
|
||||
一分钟 smoke 的恢复值为:初始 `5.1 s`、真实网络 `5.1 s`、合成 publisher `5.1 s`、Sense `0.0 s`、MediaMTX `3.1 s`,观察 `60.6 s` / 6 次采样,最大和最终未收敛数均为 0。
|
||||
|
||||
## 6. 实机兼容发现与裁决
|
||||
|
||||
该旧固件通过 ONVIF 返回的主码流路径正确,但 URI 携带 `transportmode/profile` 查询串;使用正确 RTSP 账号请求该完整 URI仍返回 `401 Unauthorized`。相同凭据、主机、端口和通道路径仅移除查询串后,可立即解码 H.264 1920×1080。
|
||||
|
||||
因此实现增加 `SENSE_ONVIF_RTSP_STRIP_QUERY`,默认 `false`,仅对已经实测需要该兼容行为的部署显式设为 `true`。默认仍保留 ONVIF URI 查询语义;代码不按厂商名或型号硬编码,也不把该发现外推到其他海康型号、固件或品牌。
|
||||
|
||||
## 7. 失败记录与修正
|
||||
|
||||
| 现象 | 原因 | 修正与防回归 |
|
||||
|---|---|---|
|
||||
| 第二套 MediaMTX API 未启动 | 合成实例同时占用默认 UDP、RTMP/HLS/WebRTC/SRT/MoQ 端口 | 合成 fixture 限制为 RTSP over TCP,并关闭无关协议 |
|
||||
| Sense 已监听但 PowerShell 健康检查超时 | 测试机配置全局 HTTP 代理,本地请求未显式绕过 | 本地 `Invoke-RestMethod` 使用 `-NoProxy` |
|
||||
| Sense 未读取实验环境变量 | PowerShell 空数组位置参数被省略,哈希表错绑为命令参数 | `Start-ManagedProcess` 改用显式命名参数 |
|
||||
| 真实 path 长期 offline | 准入旧固件返回的 RTSP 查询串触发 401 | 增加默认关闭、按部署启用的 strip-query 兼容开关和单测 |
|
||||
| 观察期一次状态读取非零 | 独立状态进程与单写者 SQLite 短暂竞争 | 状态读取增加 5 次、250 ms 有界重试;持续失败仍终止验收 |
|
||||
|
||||
失败样本未从矩阵删除;每次修正后均从头重跑四类故障。正式结果来自最终完整运行,不拼接前序成功片段。
|
||||
|
||||
## 8. 限制与后续
|
||||
|
||||
- 公网/NAT 样机只用于受控开发验证;生产环境必须使用 VPN、专网或来源白名单,不应长期裸露 ONVIF/RTSP。
|
||||
- 单台精确固件基线不能证明批次差异、多品牌兼容或五台真实设备的并发故障隔离。
|
||||
- 合成流是 640×360、10 fps、H.264 无音频测试图案,不代表真实码率、音频、夜视、弱网、存储或 AI 负载。
|
||||
- 本任务没有验证 16/64/128 路容量;默认 16、单站点最大 128 的产品配额语义保持不变,容量与 GPU 承诺需独立压测。
|
||||
- `/healthz` 只表示 Sense 进程存活;设备收敛必须继续查看脱敏快照/指标,不能把进程健康等同于业务健康。
|
||||
+17
-4
@@ -3,12 +3,12 @@ id: T-006
|
||||
title: 完成 Sense 单实机与五路混合源集成验收
|
||||
phase: 1
|
||||
deps: [T-001, T-003]
|
||||
status: TODO
|
||||
status: DONE
|
||||
created: 2026-08-04
|
||||
issue: 16
|
||||
context_ref: null
|
||||
claim_branch: null
|
||||
work_branch: null
|
||||
context_ref: 94ae2f00488bde184a9db4e5437232dfceb5f1bb
|
||||
claim_branch: claims/T-006
|
||||
work_branch: agent/codex/T-006
|
||||
write_paths:
|
||||
- docs/tasks/T-006.md
|
||||
- Sense/
|
||||
@@ -64,6 +64,19 @@ T-003 只用 fake ONVIF 和 MediaMTX 假服务建立无实机骨架,不能证
|
||||
|
||||
## 执行记录
|
||||
|
||||
### 2026-08-07 完成五路混合源验收
|
||||
|
||||
- 实现标准 ONVIF SOAP 1.2 / WS-Security adapter、`env://` 凭据解析、HTTP/RTSP router、NAT RTSP 重写,以及默认关闭的 RTSP query 兼容开关;凭据只进入进程环境和内存 URI,不写 SQLite 或证据。
|
||||
- 增加 `sense-lab` 播种/脱敏状态工具、真实 RTSP 故障代理、4 路独立 FFmpeg publisher fixture 和一键集成脚本;MediaMTX runtime path 丢失会使观察代失效并自动触发重新对账。
|
||||
- 正式使用 1 台 T-001 准入实机和 4 个独立合成 publisher:5 路初始收敛 `8.2 s`;真实网络恢复 `5.1 s`、单 publisher 恢复 `5.1 s`、Sense 重启恢复 `0.0 s`、MediaMTX 重启恢复 `2.1 s`,各阶段 `unconverged = 0`。
|
||||
- 连续观察 `1806.6 s`、采样 180 次,`maximum_unconverged = 0`、`final_unconverged = 0`;自动配置 path 数为 5。结果、版本、失败记录和限制见 `docs/research/sense-5-stream-integration.md`。
|
||||
- 该结论仅关闭 M1 单实机混合源 smoke;五条独立真实上游、真实多设备故障隔离与生产 SLA 仍由 T-007 验收,16/64/128 路容量不在本任务结论内。
|
||||
|
||||
### 2026-08-07 领取任务
|
||||
|
||||
- dispatcher `ila` 将任务分配给 `codex`;`context_ref` 为 `94ae2f00488bde184a9db4e5437232dfceb5f1bb`,claim 为 `claims/T-006`,工作分支为 `agent/codex/T-006`。
|
||||
- 接受既有写路径:`docs/tasks/T-006.md`、`Sense/`、`docs/research/sense-5-stream-integration.md`、`docs/current-state.md`。
|
||||
|
||||
### 2026-08-04 任务定义
|
||||
|
||||
- 项目负责人批准将 T-003 拆成无实机软件骨架,并由本任务保留真实摄像头和 5 路 M1 集成门禁。
|
||||
|
||||
Reference in New Issue
Block a user