From 08692e33e9d01cdecfe6b69241222ed92a934af6 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 7 Aug 2026 15:16:10 +0800 Subject: [PATCH 1/2] chore(task): claim T-006 --- docs/tasks/T-006.md | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/docs/tasks/T-006.md b/docs/tasks/T-006.md index cabd8d5..b26c2c6 100644 --- a/docs/tasks/T-006.md +++ b/docs/tasks/T-006.md @@ -3,12 +3,12 @@ id: T-006 title: 完成 Sense 单实机与五路混合源集成验收 phase: 1 deps: [T-001, T-003] -status: TODO +status: DOING 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,11 @@ T-003 只用 fake ONVIF 和 MediaMTX 假服务建立无实机骨架,不能证 ## 执行记录 +### 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 集成门禁。 From 95678380457cfb96dbc1802a82c2623de0e57446 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 7 Aug 2026 16:33:35 +0800 Subject: [PATCH 2/2] feat(sense): complete T-006 five-stream integration --- Sense/README.md | 27 +- Sense/cmd/rtsp-fault-proxy/main.go | 81 ++++ Sense/cmd/sense-api/main.go | 14 +- Sense/cmd/sense-lab/main.go | 143 +++++++ Sense/deploy/mediamtx-synthetic.yml | 21 + Sense/internal/config/config.go | 42 ++ Sense/internal/config/config_test.go | 37 ++ Sense/internal/onvif/credentials.go | 62 +++ Sense/internal/onvif/http.go | 342 +++++++++++++++ Sense/internal/onvif/http_test.go | 136 ++++++ Sense/internal/onvif/router.go | 50 +++ Sense/internal/probe/probe.go | 9 +- Sense/internal/probe/probe_test.go | 16 +- Sense/internal/store/sqlite.go | 89 ++++ Sense/internal/store/sqlite_test.go | 37 ++ Sense/scripts/t006-integration.ps1 | 441 ++++++++++++++++++++ docs/current-state.md | 13 +- docs/research/sense-5-stream-integration.md | 102 +++++ docs/tasks/T-006.md | 10 +- 19 files changed, 1658 insertions(+), 14 deletions(-) create mode 100644 Sense/cmd/rtsp-fault-proxy/main.go create mode 100644 Sense/cmd/sense-lab/main.go create mode 100644 Sense/deploy/mediamtx-synthetic.yml create mode 100644 Sense/internal/onvif/credentials.go create mode 100644 Sense/internal/onvif/http.go create mode 100644 Sense/internal/onvif/http_test.go create mode 100644 Sense/internal/onvif/router.go create mode 100644 Sense/scripts/t006-integration.ps1 create mode 100644 docs/research/sense-5-stream-integration.md diff --git a/Sense/README.md b/Sense/README.md index 0da49be..15c1f69 100644 --- a/Sense/README.md +++ b/Sense/README.md @@ -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://` 凭据引用。真实适配器从进程环境读取以下变量,不把秘密写入 SQLite、日志或 MediaMTX 错误: + +```text +SENSE_CREDENTIAL__ONVIF_USERNAME +SENSE_CREDENTIAL__ONVIF_PASSWORD +SENSE_CREDENTIAL__RTSP_USERNAME +SENSE_CREDENTIAL__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。 diff --git a/Sense/cmd/rtsp-fault-proxy/main.go b/Sense/cmd/rtsp-fault-proxy/main.go new file mode 100644 index 0000000..3366cfa --- /dev/null +++ b/Sense/cmd/rtsp-fault-proxy/main.go @@ -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 +} diff --git a/Sense/cmd/sense-api/main.go b/Sense/cmd/sense-api/main.go index 2fa69ea..3d30d0f 100644 --- a/Sense/cmd/sense-api/main.go +++ b/Sense/cmd/sense-api/main.go @@ -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. diff --git a/Sense/cmd/sense-lab/main.go b/Sense/cmd/sense-lab/main.go new file mode 100644 index 0000000..1bfac49 --- /dev/null +++ b/Sense/cmd/sense-lab/main.go @@ -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 ") + } + 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 +} diff --git a/Sense/deploy/mediamtx-synthetic.yml b/Sense/deploy/mediamtx-synthetic.yml new file mode 100644 index 0000000..97c4183 --- /dev/null +++ b/Sense/deploy/mediamtx-synthetic.yml @@ -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 diff --git a/Sense/internal/config/config.go b/Sense/internal/config/config.go index 1512177..4f7558f 100644 --- a/Sense/internal/config/config.go +++ b/Sense/internal/config/config.go @@ -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 +} diff --git a/Sense/internal/config/config_test.go b/Sense/internal/config/config_test.go index 480b850..58a6d1c 100644 --- a/Sense/internal/config/config_test.go +++ b/Sense/internal/config/config_test.go @@ -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") + } +} diff --git a/Sense/internal/onvif/credentials.go b/Sense/internal/onvif/credentials.go new file mode 100644 index 0000000..04be8e1 --- /dev/null +++ b/Sense/internal/onvif/credentials.go @@ -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 := 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 +} diff --git a/Sense/internal/onvif/http.go b/Sense/internal/onvif/http.go new file mode 100644 index 0000000..5b583ba --- /dev/null +++ b/Sense/internal/onvif/http.go @@ -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", + ``, 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", + `false`, 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", + ``, 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 := `` + + `RTP-UnicastRTSP` + + `` + escapeXML(selectedToken) + `` + 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 := `` + + `Manualfalse` + + `` + strconv.Itoa(utc.Hour()) + `` + strconv.Itoa(utc.Minute()) + + `` + strconv.Itoa(utc.Second()) + `` + strconv.Itoa(utc.Year()) + + `` + strconv.Itoa(int(utc.Month())) + `` + strconv.Itoa(utc.Day()) + + `` + _, 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 := `` + + `` + + `` + escapeXML(credentials.ONVIFUsername) + + `` + + base64.StdEncoding.EncodeToString(digest[:]) + `` + + base64.StdEncoding.EncodeToString(nonce) + `` + created + + `` + body + `` + 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"` +} diff --git a/Sense/internal/onvif/http_test.go b/Sense/internal/onvif/http_test.go new file mode 100644 index 0000000..f13e719 --- /dev/null +++ b/Sense/internal/onvif/http_test.go @@ -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 = `HIKVISIONcamerav1serial` + case strings.Contains(request.Header.Get("Content-Type"), "GetServices"): + body = `http://www.onvif.org/ver10/media/wsdlhttp://192.0.2.10/onvif/Media` + case strings.Contains(request.Header.Get("Content-Type"), "GetProfiles"): + body = `Main` + case strings.Contains(request.Header.Get("Content-Type"), "GetStreamUri"): + body = `rtsp://192.0.2.10:554/Streaming/Channels/101?transportmode=unicast&profile=Profile_1` + default: + http.Error(writer, "unexpected action", http.StatusBadRequest) + return + } + writer.Header().Set("Content-Type", "application/soap+xml") + _, _ = fmt.Fprintf(writer, `%s`, 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(`The action requires authorization`)) + })) + 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) + } +} diff --git a/Sense/internal/onvif/router.go b/Sense/internal/onvif/router.go new file mode 100644 index 0000000..841401e --- /dev/null +++ b/Sense/internal/onvif/router.go @@ -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) +} diff --git a/Sense/internal/probe/probe.go b/Sense/internal/probe/probe.go index b90d4c5..b511244 100644 --- a/Sense/internal/probe/probe.go +++ b/Sense/internal/probe/probe.go @@ -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 { diff --git a/Sense/internal/probe/probe_test.go b/Sense/internal/probe/probe_test.go index 6be8b07..9fca5df 100644 --- a/Sense/internal/probe/probe_test.go +++ b/Sense/internal/probe/probe_test.go @@ -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) + } } diff --git a/Sense/internal/store/sqlite.go b/Sense/internal/store/sqlite.go index 4cf93f6..cc36ea6 100644 --- a/Sense/internal/store/sqlite.go +++ b/Sense/internal/store/sqlite.go @@ -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` diff --git a/Sense/internal/store/sqlite_test.go b/Sense/internal/store/sqlite_test.go index ca2a273..a8df3a9 100644 --- a/Sense/internal/store/sqlite_test.go +++ b/Sense/internal/store/sqlite_test.go @@ -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")) diff --git a/Sense/scripts/t006-integration.ps1 b/Sense/scripts/t006-integration.ps1 new file mode 100644 index 0000000..aa3ae23 --- /dev/null +++ b/Sense/scripts/t006-integration.ps1 @@ -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 +} diff --git a/docs/current-state.md b/docs/current-state.md index 1a0af9c..ced0a63 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -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 为准。 ## 已知风险 diff --git a/docs/research/sense-5-stream-integration.md b/docs/research/sense-5-stream-integration.md new file mode 100644 index 0000000..4bef329 --- /dev/null +++ b/docs/research/sense-5-stream-integration.md @@ -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 进程存活;设备收敛必须继续查看脱敏快照/指标,不能把进程健康等同于业务健康。 diff --git a/docs/tasks/T-006.md b/docs/tasks/T-006.md index b26c2c6..e805eb1 100644 --- a/docs/tasks/T-006.md +++ b/docs/tasks/T-006.md @@ -3,7 +3,7 @@ id: T-006 title: 完成 Sense 单实机与五路混合源集成验收 phase: 1 deps: [T-001, T-003] -status: DOING +status: DONE created: 2026-08-04 issue: 16 context_ref: 94ae2f00488bde184a9db4e5437232dfceb5f1bb @@ -64,6 +64,14 @@ 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`。