From 19f9bfa1a0e660ee3786b2e76d03f796bdaf056b Mon Sep 17 00:00:00 2001
From: QiuSW <105186638@qq.com>
Date: Fri, 14 Aug 2026 17:41:58 +0800
Subject: [PATCH 1/4] =?UTF-8?q?feat:=20=E9=87=8D=E5=BB=BA=20MediaMTX=20?=
=?UTF-8?q?=E7=94=9F=E5=91=BD=E5=91=A8=E6=9C=9F=E4=B8=8E=E7=8A=B6=E6=80=81?=
=?UTF-8?q?=E5=AF=B9=E8=B4=A6=20(#67)?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
Sense/README.md | 2 +
Sense/server/app/admin/router/sense_media.go | 21 ++
Sense/server/app/sense/admission/apis.go | 4 +
Sense/server/app/sense/media/apis.go | 98 ++++++
Sense/server/app/sense/media/config.go | 91 ++++++
Sense/server/app/sense/media/controller.go | 168 ++++++++++
.../server/app/sense/media/controller_test.go | 67 ++++
Sense/server/app/sense/media/dto.go | 24 ++
Sense/server/app/sense/media/models.go | 37 +++
Sense/server/app/sense/media/runtime.go | 70 ++++
Sense/server/app/sense/media/service.go | 305 ++++++++++++++++++
Sense/server/app/sense/media/service_test.go | 128 ++++++++
Sense/server/app/sense/media/supervisor.go | 149 +++++++++
Sense/server/app/sense/reconcile/backoff.go | 17 +
.../app/sense/reconcile/backoff_test.go | 12 +
Sense/server/cmd/api/server.go | 18 ++
.../migration/version/2026081419000_media.go | 60 ++++
.../version/2026081419000_media_test.go | 54 ++++
Sense/server/config/credential.env.example | 3 +
.../config/mediamtx/mediamtx.yml.example | 7 +
Sense/server/tests/media/integration_test.go | 73 +++++
Sense/server/tests/media/postgres_test.go | 70 ++++
Sense/ui/src/api/sense/media.js | 7 +
Sense/ui/src/views/sense/media/index.vue | 55 ++++
Sense/ui/src/views/sense/media/mediaStatus.js | 10 +
Sense/ui/tests/unit/sense/mediaStatus.spec.js | 10 +
26 files changed, 1560 insertions(+)
create mode 100644 Sense/server/app/admin/router/sense_media.go
create mode 100644 Sense/server/app/sense/media/apis.go
create mode 100644 Sense/server/app/sense/media/config.go
create mode 100644 Sense/server/app/sense/media/controller.go
create mode 100644 Sense/server/app/sense/media/controller_test.go
create mode 100644 Sense/server/app/sense/media/dto.go
create mode 100644 Sense/server/app/sense/media/models.go
create mode 100644 Sense/server/app/sense/media/runtime.go
create mode 100644 Sense/server/app/sense/media/service.go
create mode 100644 Sense/server/app/sense/media/service_test.go
create mode 100644 Sense/server/app/sense/media/supervisor.go
create mode 100644 Sense/server/app/sense/reconcile/backoff.go
create mode 100644 Sense/server/app/sense/reconcile/backoff_test.go
create mode 100644 Sense/server/cmd/migrate/migration/version/2026081419000_media.go
create mode 100644 Sense/server/cmd/migrate/migration/version/2026081419000_media_test.go
create mode 100644 Sense/server/config/mediamtx/mediamtx.yml.example
create mode 100644 Sense/server/tests/media/integration_test.go
create mode 100644 Sense/server/tests/media/postgres_test.go
create mode 100644 Sense/ui/src/api/sense/media.js
create mode 100644 Sense/ui/src/views/sense/media/index.vue
create mode 100644 Sense/ui/src/views/sense/media/mediaStatus.js
create mode 100644 Sense/ui/tests/unit/sense/mediaStatus.spec.js
diff --git a/Sense/README.md b/Sense/README.md
index 445443a..190ab37 100644
--- a/Sense/README.md
+++ b/Sense/README.md
@@ -44,4 +44,6 @@ go run . server -c C:\secure-path\sense-settings.yml
使用 ONVIF 发现或手工接入前,还必须设置 `SENSE_ONVIF_DISCOVERY_IP` 和 `SENSE_ONVIF_ALLOWED_CIDRS`。前者只能是获准用于 WS-Discovery 的本机网卡地址;后者是获准访问的摄像头网段(多个 CIDR 用逗号分隔)。未配置时系统会给出可行动提示且不会扫描任意网卡;手工地址、Media XAddr 和 Stream URI 同样受该网段限制,并拒绝重定向或 URL 内凭据。
+MediaMTX 保持独立二进制。配置 `SENSE_MEDIAMTX_BINARY`、`SENSE_MEDIAMTX_CONFIG` 和只允许回环地址的 `SENSE_MEDIAMTX_API`。Sense 只生成无摄像头凭据的基础配置;路径和凭据在运行时通过回环 Control API 下发。模板见 `server/config/mediamtx/mediamtx.yml.example`。
+
仓库不提供默认账号、默认密码或可用密钥。首位管理员通过受仓库外 `SENSE_BOOTSTRAP_TOKEN` 保护的一次性初始化接口创建,详细步骤以项目 Wiki 的本地开发与验证页为准。
diff --git a/Sense/server/app/admin/router/sense_media.go b/Sense/server/app/admin/router/sense_media.go
new file mode 100644
index 0000000..0664a38
--- /dev/null
+++ b/Sense/server/app/admin/router/sense_media.go
@@ -0,0 +1,21 @@
+package router
+
+import (
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
+ "git.ilapage.cn/ila/yovision/Sense/server/common/actions"
+ "git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
+ "github.com/gin-gonic/gin"
+ jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
+)
+
+func init() { routerCheckRole = append(routerCheckRole, registerSenseMediaRouter) }
+
+func registerSenseMediaRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
+ api := &media.API{}
+ r := v1.Group("/media").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
+ r.GET("/routes", api.List)
+ r.GET("/process", api.Process)
+ r.POST("/reconcile", api.ReconcileAll)
+ r.POST("/routes/:id/reconcile", api.Reconcile)
+ r.POST("/routes/:id/stop", api.Stop)
+}
diff --git a/Sense/server/app/sense/admission/apis.go b/Sense/server/app/sense/admission/apis.go
index 52cd72c..ff01407 100644
--- a/Sense/server/app/sense/admission/apis.go
+++ b/Sense/server/app/sense/admission/apis.go
@@ -12,6 +12,7 @@ import (
"github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
)
@@ -53,6 +54,9 @@ func (e *API) Probe(c *gin.Context) {
e.writeError(err)
return
}
+ // Media route intent is best-effort and credential-free. A MediaMTX
+ // failure must never roll back the verified device/Profile transaction.
+ _ = media.EnsureDeviceRoutes(c.Request.Context(), service.Orm, request.DeviceID)
e.OK(result, "探测完成")
}
func (e *API) Get(c *gin.Context) {
diff --git a/Sense/server/app/sense/media/apis.go b/Sense/server/app/sense/media/apis.go
new file mode 100644
index 0000000..f9a8469
--- /dev/null
+++ b/Sense/server/app/sense/media/apis.go
@@ -0,0 +1,98 @@
+package media
+
+import (
+ "errors"
+ "net/http"
+
+ "github.com/gin-gonic/gin"
+ "github.com/go-admin-team/go-admin-core/sdk/api"
+ coreService "github.com/go-admin-team/go-admin-core/sdk/service"
+)
+
+type API struct{ api.Api }
+
+func (e *API) service(c *gin.Context) (*Service, error) {
+ base := coreService.Service{}
+ if err := e.MakeContext(c).MakeOrm().MakeService(&base).Errors; err != nil {
+ return nil, err
+ }
+ return serviceFor(base.Orm), nil
+}
+
+func (e *API) List(c *gin.Context) {
+ service, err := e.service(c)
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ items, err := service.List(c.Request.Context())
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ e.OK(gin.H{"list": items, "total": len(items)}, "查询成功")
+}
+
+func (e *API) Process(c *gin.Context) {
+ service, err := e.service(c)
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ e.OK(service.ProcessState(), "查询成功")
+}
+
+func (e *API) ReconcileAll(c *gin.Context) {
+ service, err := e.service(c)
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ if err = service.EnsureAllVerified(c.Request.Context()); err == nil {
+ err = service.ReconcileDue(c.Request.Context())
+ }
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ e.OK(nil, "对账完成")
+}
+
+func (e *API) Reconcile(c *gin.Context) {
+ service, err := e.service(c)
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ item, err := service.Reconcile(c.Request.Context(), c.Param("id"))
+ if err != nil && item.ID == "" {
+ e.writeError(err)
+ return
+ }
+ e.OK(item, "对账完成")
+}
+
+func (e *API) Stop(c *gin.Context) {
+ service, err := e.service(c)
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ item, err := service.StopRoute(c.Request.Context(), c.Param("id"))
+ if err != nil {
+ e.writeError(err)
+ return
+ }
+ e.OK(item, "媒体路径已停止")
+}
+
+func (e *API) writeError(err error) {
+ switch {
+ case errors.Is(err, ErrNotFound):
+ e.Error(http.StatusNotFound, err, err.Error())
+ case errors.Is(err, ErrRuntimeUnavailable), errors.Is(err, ErrUnsafeControlAPI):
+ e.Error(http.StatusServiceUnavailable, err, "视频服务运行配置不可用")
+ default:
+ e.Error(http.StatusInternalServerError, err, "视频服务操作失败")
+ }
+}
diff --git a/Sense/server/app/sense/media/config.go b/Sense/server/app/sense/media/config.go
new file mode 100644
index 0000000..80671a4
--- /dev/null
+++ b/Sense/server/app/sense/media/config.go
@@ -0,0 +1,91 @@
+package media
+
+import (
+ "errors"
+ "fmt"
+ "io"
+ "net"
+ "net/url"
+ "os"
+ "path/filepath"
+ "strings"
+ "time"
+)
+
+var ErrUnsafeControlAPI = errors.New("MediaMTX Control API 必须使用本机回环地址")
+
+type RuntimeConfig struct {
+ Binary string
+ ConfigPath string
+ APIBase string
+ PollInterval time.Duration
+ StartTimeout time.Duration
+}
+
+func ConfigFromEnvironment() (RuntimeConfig, error) {
+ c := RuntimeConfig{
+ Binary: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_BINARY")), ConfigPath: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_CONFIG")),
+ APIBase: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_API")), PollInterval: 5 * time.Second, StartTimeout: 8 * time.Second,
+ }
+ if c.APIBase == "" {
+ c.APIBase = "http://127.0.0.1:9997"
+ }
+ if err := validateControlAPI(c.APIBase); err != nil {
+ return RuntimeConfig{}, err
+ }
+ if c.Binary != "" && c.ConfigPath == "" {
+ return RuntimeConfig{}, errors.New("SENSE_MEDIAMTX_CONFIG is required when SENSE_MEDIAMTX_BINARY is configured")
+ }
+ return c, nil
+}
+
+func validateControlAPI(value string) error {
+ u, err := url.Parse(value)
+ if err != nil || u.Scheme != "http" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || u.Path != "" {
+ return ErrUnsafeControlAPI
+ }
+ host := u.Hostname()
+ ip := net.ParseIP(host)
+ if host != "localhost" && (ip == nil || !ip.IsLoopback()) {
+ return ErrUnsafeControlAPI
+ }
+ return nil
+}
+
+func RenderBaseConfig(w io.Writer, apiBase string) error {
+ if err := validateControlAPI(apiBase); err != nil {
+ return err
+ }
+ u, _ := url.Parse(apiBase)
+ _, err := fmt.Fprintf(w, "logLevel: info\napi: true\napiAddress: %s\nmetrics: false\npaths: {}\n", u.Host)
+ return err
+}
+
+func EnsureBaseConfig(path, apiBase string) error {
+ if path == "" {
+ return errors.New("MediaMTX config path is empty")
+ }
+ if _, err := os.Stat(path); err == nil {
+ return nil
+ } else if !errors.Is(err, os.ErrNotExist) {
+ return err
+ }
+ if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
+ return err
+ }
+ temporary := path + ".tmp"
+ f, err := os.OpenFile(temporary, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
+ if err != nil {
+ return err
+ }
+ writeErr := RenderBaseConfig(f, apiBase)
+ closeErr := f.Close()
+ if writeErr != nil || closeErr != nil {
+ _ = os.Remove(temporary)
+ return errors.Join(writeErr, closeErr)
+ }
+ if err = os.Rename(temporary, path); err != nil {
+ _ = os.Remove(temporary)
+ }
+ return err
+}
diff --git a/Sense/server/app/sense/media/controller.go b/Sense/server/app/sense/media/controller.go
new file mode 100644
index 0000000..15ae3a5
--- /dev/null
+++ b/Sense/server/app/sense/media/controller.go
@@ -0,0 +1,168 @@
+package media
+
+import (
+ "bytes"
+ "context"
+ "encoding/json"
+ "errors"
+ "fmt"
+ "io"
+ "net/http"
+ "net/url"
+ "regexp"
+ "strings"
+ "time"
+)
+
+var validPath = regexp.MustCompile(`^[A-Za-z0-9_-]{1,96}$`)
+
+type Source struct {
+ Path, URI, Username, Password string
+}
+
+type PathStatus struct {
+ Exists, Ready bool
+ Readers int
+}
+
+type Controller interface {
+ Health(context.Context) error
+ Apply(context.Context, Source) error
+ Delete(context.Context, string) error
+ Status(context.Context, string) (PathStatus, error)
+}
+
+type HTTPController struct {
+ base string
+ client *http.Client
+}
+
+func NewHTTPController(base string) (*HTTPController, error) {
+ if err := validateControlAPI(base); err != nil {
+ return nil, err
+ }
+ transport := http.DefaultTransport.(*http.Transport).Clone()
+ transport.Proxy = nil
+ return &HTTPController{base: strings.TrimRight(base, "/"), client: &http.Client{
+ Timeout: 5 * time.Second, Transport: transport,
+ CheckRedirect: func(*http.Request, []*http.Request) error {
+ return errors.New("MediaMTX Control API redirect rejected")
+ },
+ }}, nil
+}
+
+func (c *HTTPController) Health(ctx context.Context) error {
+ return c.request(ctx, http.MethodGet, "/v3/config/global/get", nil, nil)
+}
+
+func (c *HTTPController) Apply(ctx context.Context, source Source) error {
+ if !validPath.MatchString(source.Path) {
+ return errors.New("invalid MediaMTX path")
+ }
+ u, err := url.Parse(source.URI)
+ if err != nil || u.User != nil || u.Host == "" || (u.Scheme != "rtsp" && u.Scheme != "rtsps") {
+ return errors.New("invalid credential-free RTSP source")
+ }
+ payload := map[string]any{"source": u.String(), "sourceOnDemand": true, "rtspTransport": "tcp"}
+ if source.Username != "" {
+ // MediaMTX v1.19.3 has no sourceUser/sourcePass fields. Credentials
+ // are assembled only for this loopback request and are never stored,
+ // logged, or returned by a Sense endpoint.
+ u.User = url.UserPassword(source.Username, source.Password)
+ payload["source"] = u.String()
+ }
+ data, err := json.Marshal(payload)
+ if err != nil {
+ return err
+ }
+ configured, err := c.configured(ctx, source.Path)
+ if err != nil {
+ return err
+ }
+ method, action := http.MethodPost, "add"
+ if configured {
+ method, action = http.MethodPatch, "patch"
+ }
+ return c.request(ctx, method, "/v3/config/paths/"+action+"/"+url.PathEscape(source.Path), bytes.NewReader(data), nil)
+}
+
+func (c *HTTPController) configured(ctx context.Context, path string) (bool, error) {
+ status := 0
+ err := c.request(ctx, http.MethodGet, "/v3/config/paths/get/"+url.PathEscape(path), nil, &status)
+ if status == http.StatusNotFound {
+ return false, nil
+ }
+ return err == nil, err
+}
+
+func (c *HTTPController) Delete(ctx context.Context, path string) error {
+ if !validPath.MatchString(path) {
+ return errors.New("invalid MediaMTX path")
+ }
+ status := 0
+ err := c.request(ctx, http.MethodDelete, "/v3/config/paths/delete/"+url.PathEscape(path), nil, &status)
+ if status == http.StatusNotFound {
+ return nil
+ }
+ return err
+}
+
+func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, error) {
+ statusCode := 0
+ var raw struct {
+ Ready bool `json:"ready"`
+ Readers []any `json:"readers"`
+ }
+ err := c.requestJSON(ctx, http.MethodGet, "/v3/paths/get/"+url.PathEscape(path), &statusCode, &raw)
+ if statusCode == http.StatusNotFound {
+ return PathStatus{}, nil
+ }
+ if err != nil {
+ return PathStatus{}, err
+ }
+ return PathStatus{Exists: true, Ready: raw.Ready, Readers: len(raw.Readers)}, nil
+}
+
+func (c *HTTPController) request(ctx context.Context, method, path string, body io.Reader, statusOut *int) error {
+ return c.requestJSON(ctx, method, path, body, statusOut, nil)
+}
+
+func (c *HTTPController) requestJSON(ctx context.Context, method, path string, args ...any) error {
+ var body io.Reader
+ var statusOut *int
+ var target any
+ for _, arg := range args {
+ switch value := arg.(type) {
+ case io.Reader:
+ body = value
+ case *int:
+ statusOut = value
+ default:
+ target = value
+ }
+ }
+ req, err := http.NewRequestWithContext(ctx, method, c.base+path, body)
+ if err != nil {
+ return err
+ }
+ if body != nil {
+ req.Header.Set("Content-Type", "application/json")
+ }
+ res, err := c.client.Do(req)
+ if err != nil {
+ return fmt.Errorf("MediaMTX Control API unavailable: %w", err)
+ }
+ defer res.Body.Close()
+ if statusOut != nil {
+ *statusOut = res.StatusCode
+ }
+ if res.StatusCode < 200 || res.StatusCode >= 300 {
+ _, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
+ return fmt.Errorf("MediaMTX Control API returned %d", res.StatusCode)
+ }
+ if target == nil {
+ _, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
+ return nil
+ }
+ return json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(target)
+}
diff --git a/Sense/server/app/sense/media/controller_test.go b/Sense/server/app/sense/media/controller_test.go
new file mode 100644
index 0000000..cff4cee
--- /dev/null
+++ b/Sense/server/app/sense/media/controller_test.go
@@ -0,0 +1,67 @@
+package media
+
+import (
+ "context"
+ "encoding/json"
+ "net/http"
+ "net/http/httptest"
+ "strings"
+ "testing"
+)
+
+func TestControllerUsesCredentialsOnlyAtLoopbackBoundary(t *testing.T) {
+ var payload map[string]any
+ server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+ switch {
+ case r.URL.Path == "/v3/config/paths/get/sense_test":
+ http.NotFound(w, r)
+ case r.URL.Path == "/v3/config/paths/add/sense_test":
+ if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
+ t.Error(err)
+ }
+ w.WriteHeader(http.StatusOK)
+ default:
+ http.NotFound(w, r)
+ }
+ }))
+ defer server.Close()
+ controller, err := NewHTTPController(server.URL)
+ if err != nil {
+ t.Fatal(err)
+ }
+ source := Source{Path: "sense_test", URI: "rtsp://192.0.2.1/live", Username: "synthetic-user", Password: "synthetic-pass"}
+ if err = controller.Apply(context.Background(), source); err != nil {
+ t.Fatal(err)
+ }
+ if payload["source"] != "rtsp://synthetic-user:synthetic-pass@192.0.2.1/live" {
+ t.Fatalf("unexpected MediaMTX payload: %#v", payload)
+ }
+ if source.URI != "rtsp://192.0.2.1/live" || strings.Contains(source.URI, "synthetic") {
+ t.Fatalf("caller source was mutated: %#v", source)
+ }
+}
+
+func TestControlAPIAndGeneratedConfigAreLoopbackOnly(t *testing.T) {
+ if _, err := NewHTTPController("http://192.0.2.1:9997"); err != ErrUnsafeControlAPI {
+ t.Fatalf("error=%v", err)
+ }
+ var output strings.Builder
+ if err := RenderBaseConfig(&output, "http://127.0.0.1:9997"); err != nil {
+ t.Fatal(err)
+ }
+ if strings.Contains(strings.ToLower(output.String()), "password") || !strings.Contains(output.String(), "paths: {}") {
+ t.Fatalf("unsafe config: %s", output.String())
+ }
+}
+
+func TestExternalProcessUsesOrphanSafetyGate(t *testing.T) {
+ supervisor := NewSupervisor("", "")
+ supervisor.MarkExternal()
+ if err := supervisor.Stop(context.Background()); err != nil {
+ t.Fatal(err)
+ }
+ state := supervisor.State()
+ if !state.External || state.Owned || state.Phase != "running" {
+ t.Fatalf("unexpected state: %#v", state)
+ }
+}
diff --git a/Sense/server/app/sense/media/dto.go b/Sense/server/app/sense/media/dto.go
new file mode 100644
index 0000000..6ed788c
--- /dev/null
+++ b/Sense/server/app/sense/media/dto.go
@@ -0,0 +1,24 @@
+package media
+
+import "time"
+
+type RouteResponse struct {
+ ID string `json:"id"`
+ DeviceID string `json:"deviceId"`
+ ProfileToken string `json:"profileToken"`
+ Path string `json:"path"`
+ Desired string `json:"desired"`
+ Actual string `json:"actual"`
+ Readers int `json:"readers"`
+ SourceReady bool `json:"sourceReady"`
+ FailureCount int `json:"failureCount"`
+ NextRetryAt *time.Time `json:"nextRetryAt,omitempty"`
+ LastErrorCode string `json:"lastErrorCode,omitempty"`
+ Detail string `json:"detail"`
+ Version int64 `json:"version"`
+ UpdatedAt time.Time `json:"updatedAt"`
+}
+
+func routeResponse(r Route) RouteResponse {
+ return RouteResponse{ID: r.ID, DeviceID: r.DeviceID, ProfileToken: r.ProfileToken, Path: r.Path, Desired: r.Desired, Actual: r.Actual, Readers: r.Readers, SourceReady: r.SourceReady, FailureCount: r.FailureCount, NextRetryAt: r.NextRetryAt, LastErrorCode: r.LastErrorCode, Detail: r.Detail, Version: r.Version, UpdatedAt: r.UpdatedAt}
+}
diff --git a/Sense/server/app/sense/media/models.go b/Sense/server/app/sense/media/models.go
new file mode 100644
index 0000000..41b955f
--- /dev/null
+++ b/Sense/server/app/sense/media/models.go
@@ -0,0 +1,37 @@
+package media
+
+import "time"
+
+const (
+ DesiredRunning = "running"
+ DesiredStopped = "stopped"
+)
+
+type Route struct {
+ ID string `gorm:"size:96;primaryKey"`
+ DeviceID string `gorm:"size:36;not null;uniqueIndex:media_device_profile"`
+ ProfileToken string `gorm:"size:255;not null;uniqueIndex:media_device_profile"`
+ Path string `gorm:"size:96;not null;uniqueIndex"`
+ Desired string `gorm:"size:16;not null;index"`
+ Actual string `gorm:"size:32;not null;index"`
+ Readers int `gorm:"not null"`
+ SourceReady bool `gorm:"not null"`
+ FailureCount int `gorm:"not null"`
+ NextRetryAt *time.Time `gorm:"index"`
+ LastErrorCode string `gorm:"size:64;not null"`
+ Detail string `gorm:"size:512;not null"`
+ Version int64 `gorm:"not null"`
+ UpdatedAt time.Time
+}
+
+func (Route) TableName() string { return "sense_media_routes" }
+
+// admissionProfile intentionally maps only non-secret fields from #66.
+type admissionProfile struct {
+ DeviceID string `gorm:"column:device_id"`
+ Token string `gorm:"column:token"`
+ StreamURI string `gorm:"column:stream_uri"`
+ VerificationStatus string `gorm:"column:verification_status"`
+}
+
+func (admissionProfile) TableName() string { return "sense_admission_profiles" }
diff --git a/Sense/server/app/sense/media/runtime.go b/Sense/server/app/sense/media/runtime.go
new file mode 100644
index 0000000..defbf45
--- /dev/null
+++ b/Sense/server/app/sense/media/runtime.go
@@ -0,0 +1,70 @@
+package media
+
+import (
+ "context"
+ "errors"
+ "sync"
+ "time"
+
+ "gorm.io/gorm"
+)
+
+var runtimeState struct {
+ sync.RWMutex
+ service *Service
+ cancel context.CancelFunc
+}
+
+func StartRuntime(parent context.Context, db *gorm.DB) error {
+ if db == nil {
+ return errors.New("Sense database is unavailable for MediaMTX runtime")
+ }
+ config, err := ConfigFromEnvironment()
+ if err != nil {
+ return err
+ }
+ controller, err := NewHTTPController(config.APIBase)
+ if err != nil {
+ return err
+ }
+ service := NewService(db, controller, NewSupervisor(config.Binary, config.ConfigPath), config)
+ ctx, cancel := context.WithCancel(parent)
+ runtimeState.Lock()
+ if runtimeState.cancel != nil {
+ runtimeState.cancel()
+ }
+ runtimeState.service, runtimeState.cancel = service, cancel
+ runtimeState.Unlock()
+ go service.Run(ctx)
+ return nil
+}
+
+func ShutdownRuntime(ctx context.Context) error {
+ runtimeState.Lock()
+ service, cancel := runtimeState.service, runtimeState.cancel
+ runtimeState.service, runtimeState.cancel = nil, nil
+ runtimeState.Unlock()
+ if cancel != nil {
+ cancel()
+ }
+ if service != nil && service.process != nil {
+ return service.process.Stop(ctx)
+ }
+ return nil
+}
+
+func serviceFor(db *gorm.DB) *Service {
+ runtimeState.RLock()
+ service := runtimeState.service
+ runtimeState.RUnlock()
+ if service != nil {
+ return service
+ }
+ return NewService(db, nil, nil, RuntimeConfig{PollInterval: 5 * time.Second})
+}
+
+// EnsureDeviceRoutes is an internal post-admission port. It persists only
+// credential-free route intent and deliberately does not fail device probing.
+func EnsureDeviceRoutes(ctx context.Context, db *gorm.DB, deviceID string) error {
+ return serviceFor(db).EnsureDevice(ctx, deviceID)
+}
diff --git a/Sense/server/app/sense/media/service.go b/Sense/server/app/sense/media/service.go
new file mode 100644
index 0000000..409dc16
--- /dev/null
+++ b/Sense/server/app/sense/media/service.go
@@ -0,0 +1,305 @@
+package media
+
+import (
+ "context"
+ "crypto/sha256"
+ "errors"
+ "fmt"
+ "strings"
+ "sync"
+ "time"
+
+ "gorm.io/gorm"
+ "gorm.io/gorm/clause"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/reconcile"
+)
+
+var (
+ ErrNotFound = errors.New("媒体路由不存在")
+ ErrRuntimeUnavailable = errors.New("视频服务运行配置不可用")
+)
+
+type Service struct {
+ db *gorm.DB
+ controller Controller
+ process *Supervisor
+ config RuntimeConfig
+ now func() time.Time
+ mu sync.Mutex
+}
+
+func NewService(db *gorm.DB, controller Controller, process *Supervisor, config RuntimeConfig) *Service {
+ return &Service{db: db, controller: controller, process: process, config: config, now: time.Now}
+}
+
+func (s *Service) EnsureAllVerified(ctx context.Context) error {
+ var ids []string
+ if err := s.db.WithContext(ctx).Model(&admissionProfile{}).Where("verification_status = ?", "ready").Distinct().Pluck("device_id", &ids).Error; err != nil {
+ return err
+ }
+ for _, id := range ids {
+ if err := s.ensureDevice(ctx, id, false); err != nil {
+ return err
+ }
+ }
+ return nil
+}
+
+func (s *Service) EnsureDevice(ctx context.Context, deviceID string) error {
+ return s.ensureDevice(ctx, deviceID, true)
+}
+
+func (s *Service) ensureDevice(ctx context.Context, deviceID string, reactivate bool) error {
+ if strings.TrimSpace(deviceID) == "" {
+ return errors.New("device id is required")
+ }
+ var profiles []admissionProfile
+ if err := s.db.WithContext(ctx).Where("device_id = ? AND verification_status = ?", deviceID, "ready").Find(&profiles).Error; err != nil {
+ return err
+ }
+ now := s.now().UTC()
+ for _, profile := range profiles {
+ id := profile.DeviceID + ":" + profile.Token
+ digest := sha256.Sum256([]byte(id))
+ route := Route{ID: id, DeviceID: profile.DeviceID, ProfileToken: profile.Token, Path: fmt.Sprintf("sense_%x", digest[:12]), Desired: DesiredRunning, Actual: "pending", Detail: "等待视频服务对账", Version: 1, UpdatedAt: now}
+ if err := s.db.WithContext(ctx).Clauses(clause.OnConflict{Columns: []clause.Column{{Name: "id"}}, DoNothing: true}).Create(&route).Error; err != nil {
+ return err
+ }
+ var existing Route
+ if err := s.db.WithContext(ctx).First(&existing, "id = ?", id).Error; err != nil {
+ return err
+ }
+ if reactivate && existing.Desired != DesiredRunning {
+ if err := s.db.WithContext(ctx).Model(&existing).Updates(map[string]any{"desired": DesiredRunning, "actual": "pending", "detail": "等待视频服务对账", "next_retry_at": nil, "version": existing.Version + 1, "updated_at": now}).Error; err != nil {
+ return err
+ }
+ }
+ }
+ return nil
+}
+
+func (s *Service) Run(ctx context.Context) {
+ _ = s.EnsureAllVerified(ctx)
+ _ = s.ReconcileDue(ctx)
+ interval := s.config.PollInterval
+ if interval <= 0 {
+ interval = 5 * time.Second
+ }
+ ticker := time.NewTicker(interval)
+ defer ticker.Stop()
+ for {
+ select {
+ case <-ctx.Done():
+ return
+ case <-ticker.C:
+ _ = s.ReconcileDue(ctx)
+ }
+ }
+}
+
+func (s *Service) ReconcileDue(ctx context.Context) error {
+ now := s.now().UTC()
+ var routes []Route
+ if err := s.db.WithContext(ctx).Where("desired = ? AND (next_retry_at IS NULL OR next_retry_at <= ?)", DesiredRunning, now).Order("updated_at ASC").Limit(128).Find(&routes).Error; err != nil {
+ return err
+ }
+ var result error
+ for _, route := range routes {
+ var err error
+ if route.FailureCount == 0 && (route.Actual == "waiting" || route.Actual == "ready") {
+ _, err = s.Refresh(ctx, route.ID)
+ } else {
+ _, err = s.Reconcile(ctx, route.ID)
+ }
+ if err != nil {
+ result = errors.Join(result, err)
+ }
+ }
+ return result
+}
+
+func (s *Service) Refresh(ctx context.Context, id string) (RouteResponse, error) {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ var route Route
+ if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
+ return RouteResponse{}, ErrNotFound
+ }
+ if err := s.ensureControl(ctx); err != nil {
+ return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err)
+ }
+ status, err := s.controller.Status(ctx, route.Path)
+ if err != nil {
+ return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err)
+ }
+ if !status.Exists {
+ // A cold MediaMTX start begins with paths: {}; re-apply the route
+ // after releasing the service lock.
+ s.mu.Unlock()
+ response, reconcileErr := s.Reconcile(ctx, id)
+ s.mu.Lock()
+ return response, reconcileErr
+ }
+ if status.Ready {
+ return s.saveSuccess(ctx, route, "ready", true, status.Readers, "上游拉流正常")
+ }
+ return s.saveSuccess(ctx, route, "waiting", false, status.Readers, "等待播放器连接并按需拉流")
+}
+
+func (s *Service) Reconcile(ctx context.Context, id string) (RouteResponse, error) {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ var route Route
+ if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
+ if errors.Is(err, gorm.ErrRecordNotFound) {
+ return RouteResponse{}, ErrNotFound
+ }
+ return RouteResponse{}, err
+ }
+ if route.Desired == DesiredStopped {
+ if s.controller != nil {
+ _ = s.controller.Delete(ctx, route.Path)
+ }
+ return s.saveSuccess(ctx, route, "stopped", false, 0, "媒体路径已停止")
+ }
+ if err := s.ensureControl(ctx); err != nil {
+ return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err)
+ }
+ var profile admissionProfile
+ if err := s.db.WithContext(ctx).Where("device_id = ? AND token = ? AND verification_status = ?", route.DeviceID, route.ProfileToken, "ready").First(&profile).Error; err != nil {
+ return s.saveFailure(ctx, route, "profile_unavailable", "已验证 Profile 不可用", err)
+ }
+ value, err := credential.Read(s.db.WithContext(ctx), route.DeviceID, credential.PurposeRTSP)
+ if err != nil {
+ return s.saveFailure(ctx, route, "credential_unavailable", "RTSP 凭据不可用", err)
+ }
+ if err = s.controller.Apply(ctx, Source{Path: route.Path, URI: profile.StreamURI, Username: value.Username, Password: value.Password}); err != nil {
+ return s.saveFailure(ctx, route, "apply_failed", "媒体路径配置失败", err)
+ }
+ status, err := s.controller.Status(ctx, route.Path)
+ if err != nil {
+ return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err)
+ }
+ if !status.Exists {
+ return s.saveFailure(ctx, route, "path_missing", "媒体路径不存在", errors.New("MediaMTX path missing after apply"))
+ }
+ if status.Ready {
+ return s.saveSuccess(ctx, route, "ready", true, status.Readers, "上游拉流正常")
+ }
+ return s.saveSuccess(ctx, route, "waiting", false, status.Readers, "等待播放器连接并按需拉流")
+}
+
+func (s *Service) ensureControl(ctx context.Context) error {
+ if s.controller == nil {
+ return ErrRuntimeUnavailable
+ }
+ if err := s.controller.Health(ctx); err == nil {
+ if s.process != nil && !s.process.State().Owned {
+ s.process.MarkExternal()
+ }
+ return nil
+ }
+ if s.process == nil {
+ return ErrRuntimeUnavailable
+ }
+ if err := EnsureBaseConfig(s.config.ConfigPath, s.config.APIBase); err != nil {
+ s.process.MarkFailed("MediaMTX 基础配置不可用")
+ return err
+ }
+ if err := s.process.Start(ctx); err != nil {
+ return err
+ }
+ deadline := s.now().Add(s.config.StartTimeout)
+ for s.now().Before(deadline) {
+ if err := s.controller.Health(ctx); err == nil {
+ s.process.MarkReady()
+ return nil
+ }
+ select {
+ case <-ctx.Done():
+ return ctx.Err()
+ case <-time.After(100 * time.Millisecond):
+ }
+ }
+ s.process.MarkFailed("MediaMTX Control API 就绪超时")
+ return errors.New("MediaMTX readiness timeout")
+}
+
+func (s *Service) saveSuccess(ctx context.Context, route Route, actual string, sourceReady bool, readers int, detail string) (RouteResponse, error) {
+ if route.Actual == actual && route.SourceReady == sourceReady && route.Readers == readers && route.FailureCount == 0 && route.NextRetryAt == nil && route.LastErrorCode == "" && route.Detail == detail {
+ return routeResponse(route), nil
+ }
+ now := s.now().UTC()
+ updates := map[string]any{"actual": actual, "source_ready": sourceReady, "readers": readers, "failure_count": 0, "next_retry_at": nil, "last_error_code": "", "detail": detail, "version": route.Version + 1, "updated_at": now}
+ if err := s.db.WithContext(ctx).Model(&route).Updates(updates).Error; err != nil {
+ return RouteResponse{}, err
+ }
+ for key, value := range updates {
+ switch key {
+ case "actual":
+ route.Actual = value.(string)
+ case "source_ready":
+ route.SourceReady = value.(bool)
+ case "readers":
+ route.Readers = value.(int)
+ case "failure_count":
+ route.FailureCount = value.(int)
+ case "last_error_code":
+ route.LastErrorCode = value.(string)
+ case "detail":
+ route.Detail = value.(string)
+ case "version":
+ route.Version = value.(int64)
+ case "updated_at":
+ route.UpdatedAt = value.(time.Time)
+ }
+ }
+ route.NextRetryAt = nil
+ return routeResponse(route), nil
+}
+
+func (s *Service) saveFailure(ctx context.Context, route Route, code, detail string, cause error) (RouteResponse, error) {
+ now := s.now().UTC()
+ failures := route.FailureCount + 1
+ next := now.Add(reconcile.Backoff(failures))
+ updates := map[string]any{"actual": code, "source_ready": false, "readers": 0, "failure_count": failures, "next_retry_at": &next, "last_error_code": code, "detail": detail, "version": route.Version + 1, "updated_at": now}
+ if err := s.db.WithContext(ctx).Model(&route).Updates(updates).Error; err != nil {
+ return RouteResponse{}, errors.Join(cause, err)
+ }
+ route.Actual, route.SourceReady, route.Readers, route.FailureCount, route.NextRetryAt, route.LastErrorCode, route.Detail = code, false, 0, failures, &next, code, detail
+ route.Version, route.UpdatedAt = route.Version+1, now
+ return routeResponse(route), cause
+}
+
+func (s *Service) StopRoute(ctx context.Context, id string) (RouteResponse, error) {
+ var route Route
+ if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
+ return RouteResponse{}, ErrNotFound
+ }
+ route.Desired = DesiredStopped
+ if err := s.db.WithContext(ctx).Model(&route).Updates(map[string]any{"desired": DesiredStopped, "next_retry_at": nil, "version": route.Version + 1, "updated_at": s.now().UTC()}).Error; err != nil {
+ return RouteResponse{}, err
+ }
+ return s.Reconcile(ctx, id)
+}
+
+func (s *Service) List(ctx context.Context) ([]RouteResponse, error) {
+ var routes []Route
+ if err := s.db.WithContext(ctx).Order("updated_at DESC").Find(&routes).Error; err != nil {
+ return nil, err
+ }
+ items := make([]RouteResponse, 0, len(routes))
+ for _, route := range routes {
+ items = append(items, routeResponse(route))
+ }
+ return items, nil
+}
+
+func (s *Service) ProcessState() ProcessState {
+ if s.process == nil {
+ return ProcessState{Phase: "configuration_failed", Detail: "视频服务运行配置不可用", LastChanged: s.now().UTC()}
+ }
+ return s.process.State()
+}
diff --git a/Sense/server/app/sense/media/service_test.go b/Sense/server/app/sense/media/service_test.go
new file mode 100644
index 0000000..4345f96
--- /dev/null
+++ b/Sense/server/app/sense/media/service_test.go
@@ -0,0 +1,128 @@
+package media
+
+import (
+ "context"
+ "crypto/rand"
+ "encoding/base64"
+ "errors"
+ "os"
+ "testing"
+ "time"
+
+ "gorm.io/driver/sqlite"
+ "gorm.io/gorm"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
+)
+
+type fakeController struct {
+ healthErr, applyErr, statusErr error
+ applies int
+ status PathStatus
+}
+
+func (f *fakeController) Health(context.Context) error { return f.healthErr }
+func (f *fakeController) Apply(context.Context, Source) error { f.applies++; return f.applyErr }
+func (f *fakeController) Delete(context.Context, string) error { return nil }
+func (f *fakeController) Status(context.Context, string) (PathStatus, error) {
+ return f.status, f.statusErr
+}
+
+func mediaTestService(t *testing.T, controller Controller) (*Service, *gorm.DB) {
+ t.Helper()
+ db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err = db.AutoMigrate(&Route{}, &admissionProfile{}, &credential.DeviceCredential{}); err != nil {
+ t.Fatal(err)
+ }
+ key := make([]byte, 32)
+ if _, err = rand.Read(key); err != nil {
+ t.Fatal(err)
+ }
+ t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
+ vault, _ := credential.NewVault(key)
+ ciphertext, err := vault.Encrypt("device-1", credential.PurposeRTSP, "synthetic-user", "synthetic-password")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Create(&credential.DeviceCredential{DeviceID: "device-1", Purpose: credential.PurposeRTSP, Ciphertext: ciphertext, KeyVersion: credential.Version()}).Error; err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Create(&admissionProfile{DeviceID: "device-1", Token: "main", StreamURI: "rtsp://192.0.2.1/live", VerificationStatus: "ready"}).Error; err != nil {
+ t.Fatal(err)
+ }
+ service := NewService(db, controller, nil, RuntimeConfig{})
+ service.now = func() time.Time { return time.Date(2026, 8, 14, 0, 0, 0, 0, time.UTC) }
+ return service, db
+}
+
+func TestEnsureAndReconcileAreIdempotent(t *testing.T) {
+ controller := &fakeController{status: PathStatus{Exists: true, Ready: true, Readers: 2}}
+ service, db := mediaTestService(t, controller)
+ if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
+ t.Fatal(err)
+ }
+ if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
+ t.Fatal(err)
+ }
+ var count int64
+ if err := db.Model(&Route{}).Count(&count).Error; err != nil || count != 1 {
+ t.Fatalf("count=%d err=%v", count, err)
+ }
+ item, err := service.Reconcile(context.Background(), "device-1:main")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if item.Actual != "ready" || item.Readers != 2 || controller.applies != 1 {
+ t.Fatalf("item=%#v applies=%d", item, controller.applies)
+ }
+ if err = service.ReconcileDue(context.Background()); err != nil {
+ t.Fatal(err)
+ }
+ items, err := service.List(context.Background())
+ if err != nil || controller.applies != 1 || items[0].Version != item.Version {
+ t.Fatalf("steady route was rewritten: items=%#v applies=%d err=%v", items, controller.applies, err)
+ }
+}
+
+func TestColdStartDoesNotReactivateStoppedRoute(t *testing.T) {
+ controller := &fakeController{status: PathStatus{Exists: true}}
+ service, _ := mediaTestService(t, controller)
+ if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
+ t.Fatal(err)
+ }
+ if _, err := service.StopRoute(context.Background(), "device-1:main"); err != nil {
+ t.Fatal(err)
+ }
+ if err := service.EnsureAllVerified(context.Background()); err != nil {
+ t.Fatal(err)
+ }
+ items, err := service.List(context.Background())
+ if err != nil || len(items) != 1 || items[0].Desired != DesiredStopped {
+ t.Fatalf("stopped route was reactivated: %#v err=%v", items, err)
+ }
+}
+
+func TestFailurePersistsBackoffWithoutChangingProfile(t *testing.T) {
+ controller := &fakeController{applyErr: errors.New("synthetic apply failure")}
+ service, db := mediaTestService(t, controller)
+ if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
+ t.Fatal(err)
+ }
+ item, err := service.Reconcile(context.Background(), "device-1:main")
+ if err == nil || item.Actual != "apply_failed" || item.FailureCount != 1 || item.NextRetryAt == nil {
+ t.Fatalf("item=%#v err=%v", item, err)
+ }
+ var profile admissionProfile
+ if err = db.First(&profile, "device_id = ? AND token = ?", "device-1", "main").Error; err != nil {
+ t.Fatal(err)
+ }
+ if profile.VerificationStatus != "ready" {
+ t.Fatalf("profile changed: %#v", profile)
+ }
+ if value := os.Getenv(credential.EnvironmentKey); value == "" {
+ t.Fatal("test key unexpectedly missing")
+ }
+}
diff --git a/Sense/server/app/sense/media/supervisor.go b/Sense/server/app/sense/media/supervisor.go
new file mode 100644
index 0000000..4985b6c
--- /dev/null
+++ b/Sense/server/app/sense/media/supervisor.go
@@ -0,0 +1,149 @@
+package media
+
+import (
+ "context"
+ "errors"
+ "os"
+ "os/exec"
+ "path/filepath"
+ "sync"
+ "time"
+)
+
+type ProcessState struct {
+ Phase string `json:"phase"`
+ PID int `json:"pid,omitempty"`
+ Owned bool `json:"owned"`
+ External bool `json:"external"`
+ Restarts int `json:"restarts"`
+ Detail string `json:"detail"`
+ LastChanged time.Time `json:"lastChanged"`
+}
+
+type Process interface {
+ Start(context.Context) error
+ Stop(context.Context) error
+ State() ProcessState
+}
+
+type Supervisor struct {
+ binary, config string
+ mu sync.Mutex
+ command *exec.Cmd
+ state ProcessState
+}
+
+func NewSupervisor(binary, config string) *Supervisor {
+ phase, detail := "stopped", "MediaMTX 尚未启动"
+ if binary == "" {
+ phase, detail = "not_configured", "未配置 MediaMTX 二进制;可连接外部已启动实例"
+ }
+ return &Supervisor{binary: binary, config: config, state: ProcessState{Phase: phase, Detail: detail, LastChanged: time.Now().UTC()}}
+}
+
+func (s *Supervisor) MarkExternal() {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.state.Owned {
+ return
+ }
+ s.state.Phase, s.state.External, s.state.Detail, s.state.LastChanged = "running", true, "检测到外部 MediaMTX;孤儿安全闸禁止 Sense 停止该进程", time.Now().UTC()
+}
+
+func (s *Supervisor) Start(_ context.Context) error {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.state.Owned && s.command != nil {
+ return nil
+ }
+ if s.binary == "" {
+ return errors.New("SENSE_MEDIAMTX_BINARY is not configured")
+ }
+ args := []string{}
+ if s.config != "" {
+ args = append(args, s.config)
+ }
+ cmd := exec.Command(s.binary, args...)
+ if s.config != "" {
+ cmd.Dir = filepath.Dir(s.config)
+ }
+ if err := cmd.Start(); err != nil {
+ s.state.Phase, s.state.Detail, s.state.LastChanged = "failed", "MediaMTX 进程启动失败", time.Now().UTC()
+ return err
+ }
+ if s.state.Phase == "failed" {
+ s.state.Restarts++
+ }
+ s.command = cmd
+ s.state.Phase, s.state.PID, s.state.Owned, s.state.External = "starting", cmd.Process.Pid, true, false
+ s.state.Detail, s.state.LastChanged = "等待 MediaMTX Control API 就绪", time.Now().UTC()
+ go s.wait(cmd)
+ return nil
+}
+
+func (s *Supervisor) MarkReady() {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.state.Owned {
+ s.state.Phase, s.state.Detail, s.state.LastChanged = "running", "MediaMTX Control API 已就绪", time.Now().UTC()
+ }
+}
+
+func (s *Supervisor) MarkFailed(detail string) {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ s.state.Phase, s.state.Detail, s.state.LastChanged = "failed", detail, time.Now().UTC()
+}
+
+func (s *Supervisor) wait(cmd *exec.Cmd) {
+ err := cmd.Wait()
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.command != cmd {
+ return
+ }
+ s.command = nil
+ s.state.Phase, s.state.PID, s.state.Owned, s.state.External = "failed", 0, false, false
+ s.state.Detail = "MediaMTX 进程已退出"
+ if err == nil {
+ s.state.Phase, s.state.Detail = "stopped", "MediaMTX 进程已停止"
+ }
+ s.state.LastChanged = time.Now().UTC()
+}
+
+func (s *Supervisor) Stop(ctx context.Context) error {
+ s.mu.Lock()
+ cmd := s.command
+ owned := s.state.Owned
+ s.mu.Unlock()
+ if !owned || cmd == nil {
+ return nil
+ }
+ if err := cmd.Process.Signal(os.Interrupt); err != nil {
+ if killErr := cmd.Process.Kill(); killErr != nil {
+ return errors.Join(err, killErr)
+ }
+ }
+ ticker := time.NewTicker(50 * time.Millisecond)
+ defer ticker.Stop()
+ for {
+ select {
+ case <-ctx.Done():
+ _ = cmd.Process.Kill()
+ return ctx.Err()
+ case <-ticker.C:
+ s.mu.Lock()
+ finished := s.command != cmd
+ s.mu.Unlock()
+ if finished {
+ return nil
+ }
+ }
+ }
+}
+
+func (s *Supervisor) State() ProcessState {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ return s.state
+}
diff --git a/Sense/server/app/sense/reconcile/backoff.go b/Sense/server/app/sense/reconcile/backoff.go
new file mode 100644
index 0000000..e79d174
--- /dev/null
+++ b/Sense/server/app/sense/reconcile/backoff.go
@@ -0,0 +1,17 @@
+package reconcile
+
+import "time"
+
+func Backoff(failures int) time.Duration {
+ if failures <= 1 {
+ return time.Second
+ }
+ d := time.Second
+ for i := 1; i < failures; i++ {
+ if d >= 30*time.Second {
+ return time.Minute
+ }
+ d *= 2
+ }
+ return d
+}
diff --git a/Sense/server/app/sense/reconcile/backoff_test.go b/Sense/server/app/sense/reconcile/backoff_test.go
new file mode 100644
index 0000000..42fae39
--- /dev/null
+++ b/Sense/server/app/sense/reconcile/backoff_test.go
@@ -0,0 +1,12 @@
+package reconcile
+
+import (
+ "testing"
+ "time"
+)
+
+func TestBackoffIsBounded(t *testing.T) {
+ if Backoff(1) != time.Second || Backoff(4) != 8*time.Second || Backoff(20) != time.Minute {
+ t.Fatal("unexpected retry backoff")
+ }
+}
diff --git a/Sense/server/cmd/api/server.go b/Sense/server/cmd/api/server.go
index 896ffd5..1f0af69 100644
--- a/Sense/server/cmd/api/server.go
+++ b/Sense/server/cmd/api/server.go
@@ -20,6 +20,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/admin/router"
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/common/database"
"git.ilapage.cn/ila/yovision/Sense/server/common/global"
common "git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
@@ -85,6 +86,19 @@ func run() error {
for _, f := range AppRouters {
f()
}
+ runtimeCtx, runtimeCancel := context.WithCancel(context.Background())
+ defer runtimeCancel()
+ var runtimeDBFound bool
+ for _, db := range sdk.Runtime.GetDb() {
+ runtimeDBFound = true
+ if err := media.StartRuntime(runtimeCtx, db); err != nil {
+ log.Errorf("MediaMTX runtime unavailable: %v", err)
+ }
+ break
+ }
+ if !runtimeDBFound {
+ log.Error("MediaMTX runtime unavailable: Sense database is not initialized")
+ }
srv := &http.Server{
Addr: fmt.Sprintf("%s:%d", config.ApplicationConfig.Host, config.ApplicationConfig.Port),
@@ -139,6 +153,10 @@ func run() error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
log.Info("Shutdown Server ... ")
+ runtimeCancel()
+ if err := media.ShutdownRuntime(ctx); err != nil && !errors.Is(err, context.Canceled) {
+ log.Errorf("Shutdown MediaMTX runtime: %v", err)
+ }
if err := srv.Shutdown(ctx); err != nil {
log.Fatal("Server Shutdown:", err)
diff --git a/Sense/server/cmd/migrate/migration/version/2026081419000_media.go b/Sense/server/cmd/migrate/migration/version/2026081419000_media.go
new file mode 100644
index 0000000..98c503e
--- /dev/null
+++ b/Sense/server/cmd/migrate/migration/version/2026081419000_media.go
@@ -0,0 +1,60 @@
+package version
+
+import (
+ "runtime"
+
+ "gorm.io/gorm"
+ "gorm.io/gorm/clause"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
+ "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration"
+ migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
+ common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
+)
+
+func init() {
+ _, fileName, _, _ := runtime.Caller(0)
+ migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSenseMedia)
+}
+
+func migrateSenseMedia(db *gorm.DB, version string) error {
+ return db.Transaction(func(tx *gorm.DB) error {
+ if err := tx.AutoMigrate(&media.Route{}); err != nil {
+ return err
+ }
+ page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMedia", Title: "视频服务", Icon: "video-play", Path: "/sense/media", MenuType: "C", Permission: "sense:media:list", Component: "/sense/media/index", Sort: 7, Visible: "0", IsFrame: "1"})
+ if err != nil {
+ return err
+ }
+ reconcile, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaReconcile", Title: "立即对账", MenuType: "F", Action: "POST", Permission: "sense:media:reconcile", ParentId: page.MenuId, Sort: 1, Visible: "1", IsFrame: "1"})
+ if err != nil {
+ return err
+ }
+ stop, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaStop", Title: "停止路径", MenuType: "F", Action: "POST", Permission: "sense:media:stop", ParentId: page.MenuId, Sort: 2, Visible: "1", IsFrame: "1"})
+ if err != nil {
+ return err
+ }
+ for _, role := range []string{"implementation_operator", "site_admin"} {
+ if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, reconcile, stop}); err != nil {
+ return err
+ }
+ }
+ if err = attachDeviceRole(tx, "viewer", []migrationModels.SysMenu{page}); err != nil {
+ return err
+ }
+ read := [][2]string{{"/api/v1/media/routes", "GET"}, {"/api/v1/media/process", "GET"}}
+ write := [][2]string{{"/api/v1/media/reconcile", "POST"}, {"/api/v1/media/routes/:id/reconcile", "POST"}, {"/api/v1/media/routes/:id/stop", "POST"}}
+ for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
+ policies := append([][2]string{}, read...)
+ if role != "viewer" {
+ policies = append(policies, write...)
+ }
+ for _, policy := range policies {
+ if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}).Error; err != nil {
+ return err
+ }
+ }
+ }
+ return tx.Create(&common.Migration{Version: version}).Error
+ })
+}
diff --git a/Sense/server/cmd/migrate/migration/version/2026081419000_media_test.go b/Sense/server/cmd/migrate/migration/version/2026081419000_media_test.go
new file mode 100644
index 0000000..4d06942
--- /dev/null
+++ b/Sense/server/cmd/migrate/migration/version/2026081419000_media_test.go
@@ -0,0 +1,54 @@
+package version
+
+import (
+ "os"
+ "testing"
+
+ "gorm.io/driver/postgres"
+ "gorm.io/gorm"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
+ migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
+ common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
+)
+
+func TestMediaMigrationOnPostgres(t *testing.T) {
+ dsn := os.Getenv("SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL")
+ if dsn == "" {
+ t.Skip("set SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL to run the PostgreSQL migration test")
+ }
+ db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil {
+ t.Fatal(err)
+ }
+ t.Cleanup(func() {
+ db.Exec("DROP TABLE IF EXISTS sense_media_routes, sys_role_menu, sys_menu, sys_role, casbin_rule, sys_migration CASCADE")
+ })
+ for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
+ if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil {
+ t.Fatal(err)
+ }
+ }
+ if err = migrateSenseMedia(db, "2026081419000_media.go"); err != nil {
+ t.Fatal(err)
+ }
+ var routes, menus, policies, applied int64
+ if err = db.Model(&media.Route{}).Count(&routes).Error; err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseMedia%").Count(&menus).Error; err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Model(&deviceCasbinRule{}).Where("v1 LIKE ?", "/api/v1/media%").Count(&policies).Error; err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Model(&common.Migration{}).Where("version = ?", "2026081419000_media.go").Count(&applied).Error; err != nil {
+ t.Fatal(err)
+ }
+ if routes != 0 || menus != 3 || policies != 12 || applied != 1 {
+ t.Fatalf("routes=%d menus=%d policies=%d applied=%d", routes, menus, policies, applied)
+ }
+}
diff --git a/Sense/server/config/credential.env.example b/Sense/server/config/credential.env.example
index 99883ea..1489952 100644
--- a/Sense/server/config/credential.env.example
+++ b/Sense/server/config/credential.env.example
@@ -6,3 +6,6 @@ SENSE_CREDENTIAL_KEY=
# approved local interface/IP ranges; comma-separate multiple CIDRs.
SENSE_ONVIF_DISCOVERY_IP=
SENSE_ONVIF_ALLOWED_CIDRS=
+SENSE_MEDIAMTX_BINARY=
+SENSE_MEDIAMTX_CONFIG=
+SENSE_MEDIAMTX_API=http://127.0.0.1:9997
diff --git a/Sense/server/config/mediamtx/mediamtx.yml.example b/Sense/server/config/mediamtx/mediamtx.yml.example
new file mode 100644
index 0000000..2ea5ab1
--- /dev/null
+++ b/Sense/server/config/mediamtx/mediamtx.yml.example
@@ -0,0 +1,7 @@
+# Sense generates only this credential-free base configuration.
+# The Control API must stay on loopback; camera paths are applied at runtime.
+logLevel: info
+api: true
+apiAddress: 127.0.0.1:9997
+metrics: false
+paths: {}
diff --git a/Sense/server/tests/media/integration_test.go b/Sense/server/tests/media/integration_test.go
new file mode 100644
index 0000000..de48959
--- /dev/null
+++ b/Sense/server/tests/media/integration_test.go
@@ -0,0 +1,73 @@
+package media_test
+
+import (
+ "context"
+ "fmt"
+ "net"
+ "os"
+ "path/filepath"
+ "testing"
+ "time"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
+)
+
+func freeAddress(t *testing.T) string {
+ t.Helper()
+ listener, err := net.Listen("tcp", "127.0.0.1:0")
+ if err != nil {
+ t.Fatal(err)
+ }
+ address := listener.Addr().String()
+ if err = listener.Close(); err != nil {
+ t.Fatal(err)
+ }
+ return address
+}
+
+func TestRealMediaMTXControlLifecycle(t *testing.T) {
+ binary := os.Getenv("SENSE_MEDIAMTX_TEST_BINARY")
+ if binary == "" {
+ t.Skip("set SENSE_MEDIAMTX_TEST_BINARY to run the real MediaMTX integration")
+ }
+ apiAddress, rtspAddress := freeAddress(t), freeAddress(t)
+ configPath := filepath.Join(t.TempDir(), "mediamtx.yml")
+ config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress)
+ if err := os.WriteFile(configPath, []byte(config), 0o600); err != nil {
+ t.Fatal(err)
+ }
+ controller, err := media.NewHTTPController("http://" + apiAddress)
+ if err != nil {
+ t.Fatal(err)
+ }
+ supervisor := media.NewSupervisor(binary, configPath)
+ if err = supervisor.Start(context.Background()); err != nil {
+ t.Fatal(err)
+ }
+ t.Cleanup(func() {
+ ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
+ defer cancel()
+ _ = supervisor.Stop(ctx)
+ })
+ deadline := time.Now().Add(8 * time.Second)
+ for {
+ err = controller.Health(context.Background())
+ if err == nil {
+ break
+ }
+ if time.Now().After(deadline) {
+ t.Fatalf("MediaMTX did not become ready: %v", err)
+ }
+ time.Sleep(100 * time.Millisecond)
+ }
+ source := media.Source{Path: "sense_integration", URI: "rtsp://127.0.0.1:65530/test", Username: "synthetic-user", Password: "synthetic-password"}
+ if err = controller.Apply(context.Background(), source); err != nil {
+ t.Fatal(err)
+ }
+ if err = controller.Apply(context.Background(), source); err != nil {
+ t.Fatalf("replace must be idempotent: %v", err)
+ }
+ if err = controller.Delete(context.Background(), source.Path); err != nil {
+ t.Fatal(err)
+ }
+}
diff --git a/Sense/server/tests/media/postgres_test.go b/Sense/server/tests/media/postgres_test.go
new file mode 100644
index 0000000..6cab65b
--- /dev/null
+++ b/Sense/server/tests/media/postgres_test.go
@@ -0,0 +1,70 @@
+package media_test
+
+import (
+ "context"
+ "crypto/rand"
+ "encoding/base64"
+ "os"
+ "testing"
+
+ "gorm.io/driver/postgres"
+ "gorm.io/gorm"
+
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
+ "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
+)
+
+type readyController struct{}
+
+func (readyController) Health(context.Context) error { return nil }
+func (readyController) Apply(context.Context, media.Source) error { return nil }
+func (readyController) Delete(context.Context, string) error { return nil }
+func (readyController) Status(context.Context, string) (media.PathStatus, error) {
+ return media.PathStatus{Exists: true, Ready: true}, nil
+}
+
+func TestPostgresColdStartRestoresDesiredRoute(t *testing.T) {
+ dsn := os.Getenv("SENSE_MEDIA_TEST_DATABASE_URL")
+ if dsn == "" {
+ t.Skip("set SENSE_MEDIA_TEST_DATABASE_URL to run PostgreSQL media recovery")
+ }
+ db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err = db.AutoMigrate(&media.Route{}, &admission.Profile{}, &credential.DeviceCredential{}); err != nil {
+ t.Fatal(err)
+ }
+ t.Cleanup(func() {
+ db.Exec("DROP TABLE IF EXISTS sense_media_routes, sense_admission_profiles, sense_device_credentials")
+ })
+ key := make([]byte, 32)
+ if _, err = rand.Read(key); err != nil {
+ t.Fatal(err)
+ }
+ t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
+ vault, _ := credential.NewVault(key)
+ ciphertext, err := vault.Encrypt("device-pg", credential.PurposeRTSP, "synthetic-user", "synthetic-password")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Create(&credential.DeviceCredential{DeviceID: "device-pg", Purpose: credential.PurposeRTSP, Ciphertext: ciphertext, KeyVersion: credential.Version()}).Error; err != nil {
+ t.Fatal(err)
+ }
+ if err = db.Create(&admission.Profile{DeviceID: "device-pg", Token: "main", Name: "Main", StreamURI: "rtsp://192.0.2.1/live", VerificationStatus: "ready", VerificationDetail: "synthetic"}).Error; err != nil {
+ t.Fatal(err)
+ }
+ first := media.NewService(db, readyController{}, nil, media.RuntimeConfig{})
+ if err = first.EnsureDevice(context.Background(), "device-pg"); err != nil {
+ t.Fatal(err)
+ }
+ second := media.NewService(db.Session(&gorm.Session{NewDB: true}), readyController{}, nil, media.RuntimeConfig{})
+ if err = second.ReconcileDue(context.Background()); err != nil {
+ t.Fatal(err)
+ }
+ items, err := second.List(context.Background())
+ if err != nil || len(items) != 1 || items[0].Actual != "ready" {
+ t.Fatalf("items=%#v err=%v", items, err)
+ }
+}
diff --git a/Sense/ui/src/api/sense/media.js b/Sense/ui/src/api/sense/media.js
new file mode 100644
index 0000000..92bc6b9
--- /dev/null
+++ b/Sense/ui/src/api/sense/media.js
@@ -0,0 +1,7 @@
+import request from '@/utils/request'
+
+export function listMediaRoutes() { return request({ url: '/api/v1/media/routes', method: 'get' }) }
+export function getMediaProcess() { return request({ url: '/api/v1/media/process', method: 'get' }) }
+export function reconcileAllMedia() { return request({ url: '/api/v1/media/reconcile', method: 'post' }) }
+export function reconcileMediaRoute(id) { return request({ url: `/api/v1/media/routes/${encodeURIComponent(id)}/reconcile`, method: 'post' }) }
+export function stopMediaRoute(id) { return request({ url: `/api/v1/media/routes/${encodeURIComponent(id)}/stop`, method: 'post' }) }
diff --git a/Sense/ui/src/views/sense/media/index.vue b/Sense/ui/src/views/sense/media/index.vue
new file mode 100644
index 0000000..f8d8433
--- /dev/null
+++ b/Sense/ui/src/views/sense/media/index.vue
@@ -0,0 +1,55 @@
+
+
+
+
+
+
+
+ {{ processLabel(process.phase) }}
+ {{ process.owned ? 'Sense 启动' : process.external ? '外部启动(受保护)' : '未运行' }}
+ {{ process.pid || '—' }}
+ {{ process.detail || '尚未取得状态' }}
+
+
+
+
+ {{ scope.row.desired === 'running' ? '运行' : '停止' }}
+ {{ mediaStatusLabel(scope.row.actual) }}
+
+ {{ scope.row.nextRetryAt ? parseTime(scope.row.nextRetryAt) : '—' }}
+
+
+
+ 对账
+ 停止
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/Sense/ui/src/views/sense/media/mediaStatus.js b/Sense/ui/src/views/sense/media/mediaStatus.js
new file mode 100644
index 0000000..eee4889
--- /dev/null
+++ b/Sense/ui/src/views/sense/media/mediaStatus.js
@@ -0,0 +1,10 @@
+export function mediaStatusLabel(value) {
+ return ({ ready: '拉流正常', waiting: '等待拉流', pending: '等待对账', stopped: '已停止', process_unavailable: '进程不可用', apply_failed: '配置失败', status_unavailable: '状态未知', path_missing: '路径缺失', profile_unavailable: 'Profile 不可用', credential_unavailable: '凭据不可用' })[value] || value || '未知'
+}
+
+export function mediaStatusType(value) {
+ if (value === 'ready' || value === 'running') return 'success'
+ if (value === 'waiting' || value === 'pending' || value === 'starting') return 'warning'
+ if (value === 'stopped' || value === 'not_configured') return 'info'
+ return 'danger'
+}
diff --git a/Sense/ui/tests/unit/sense/mediaStatus.spec.js b/Sense/ui/tests/unit/sense/mediaStatus.spec.js
new file mode 100644
index 0000000..1b6f508
--- /dev/null
+++ b/Sense/ui/tests/unit/sense/mediaStatus.spec.js
@@ -0,0 +1,10 @@
+import { mediaStatusLabel, mediaStatusType } from '@/views/sense/media/mediaStatus'
+
+describe('Sense media status presentation', () => {
+ it('uses actionable labels for expected lifecycle states', () => {
+ expect(mediaStatusLabel('waiting')).toBe('等待拉流')
+ expect(mediaStatusLabel('process_unavailable')).toBe('进程不可用')
+ expect(mediaStatusType('ready')).toBe('success')
+ expect(mediaStatusType('apply_failed')).toBe('danger')
+ })
+})
--
2.34.1
From 4149cc4426c65c478b2a1da3cc7a883a52d44e5e Mon Sep 17 00:00:00 2001
From: QiuSW <105186638@qq.com>
Date: Fri, 14 Aug 2026 17:51:16 +0800
Subject: [PATCH 2/4] =?UTF-8?q?docs:=20=E8=AE=B0=E5=BD=95=20Sense=20?=
=?UTF-8?q?=E8=A7=86=E9=A2=91=E6=9C=8D=E5=8A=A1=E7=94=9F=E5=91=BD=E5=91=A8?=
=?UTF-8?q?=E6=9C=9F=20(#67)?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
docs/02-architecture-and-code-map.md | 8 ++++--
docs/03-business-rules-and-glossary.md | 15 +++++++++--
docs/04-local-development-and-verification.md | 27 +++++++++++++++++--
docs/06-troubleshooting.md | 18 +++++++++++--
docs/delivery/README.md | 14 ++++++++--
5 files changed, 72 insertions(+), 10 deletions(-)
diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md
index bdeda25..28b6205 100644
--- a/docs/02-architecture-and-code-map.md
+++ b/docs/02-architecture-and-code-map.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Architecture-and-Code-Map.-
-wiki_revision: 915b5ca9a9a097470cb67a80f73045d16d3d57be
-synchronized_at: 2026-08-14T08:53:04Z
+wiki_revision: c5482806a40428a436c202760eefdb6403406584
+synchronized_at: 2026-08-14T09:46:37Z
# 架构与代码地图
@@ -96,6 +96,10 @@ Sense JWT realm 固定为 `Sense`;浏览器令牌 Cookie 为 `Sense-Admin-Toke
工单 #66 在 Sense/server/app/sense/onvif/、rtsp/ 与 admission/ 建立视频接入边界:WS-Discovery 只能绑定 SENSE_ONVIF_DISCOVERY_IP 指定的本机网卡,所有 ONVIF、Media XAddr 与 RTSP Stream URI 都必须落在 SENSE_ONVIF_ALLOWED_CIDRS 明确授权的网段。HTTP 客户端禁止代理和重定向,并在每次连接时重新解析、校验和固定目标 IP,防止 DNS 重绑定;URL 用户信息及敏感查询参数被拒绝。
ONVIF 支持 Basic 与 MD5/SHA-256 Digest challenge,Profile 与无凭据 Stream URI 持久化到 PostgreSQL。接入失败会记录可行动状态但保留最后一次已验证 Profile;成功接入清除凭据更新触发的重试标记。前端继续复用 GoAdmin 动态菜单、权限链、BasicLayout 和 Element Plus 表单、Dialog、Table、Tag。
+工单 #67 在 `Sense/server/app/sense/media/` 与 `reconcile/` 建立 MediaMTX 管理面:`cmd/api/server.go` 随 Sense 生命周期启动后台对账并只停止本实例拥有的子进程;检测到外部实例时设置孤儿安全闸,不发送停止信号。MediaMTX Control API 只允许 loopback HTTP,禁用代理与重定向。
+
+媒体路由只保存设备/Profile 引用、无秘密路径名、期望态、实际态、reader、退避和下次重试;摄像头凭据从内部端口按需解密,仅在 loopback Control API 请求内临时组装,不写入路由表、基础配置、日志或 Sense 响应。MediaMTX 故障和退避不改变 #66 的设备/Profile 验证状态。
+
diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md
index 12c1110..ec1e5a6 100644
--- a/docs/03-business-rules-and-glossary.md
+++ b/docs/03-business-rules-and-glossary.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Business-Rules-and-Glossary
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Business-Rules-and-Glossary.-
-wiki_revision: d1471acadb4759b80e9bc20e5cec18bd21ba65cd
-synchronized_at: 2026-08-14T08:53:13Z
+wiki_revision: ea2c540a9ac3f8ccf40123f8389eb392fb5e172b
+synchronized_at: 2026-08-14T09:46:41Z
# 业务规则与术语
@@ -76,6 +76,17 @@ synchronized_at: 2026-08-14T08:53:13Z
- 登录成功/失败、登出、密码变更和鉴权拒绝必须留下身份审计;密码、令牌、Cookie、验证码、数据库连接和摄像头凭据不得进入审计正文。
- `site_admin` 可在账户维护流程中读取角色、部门、岗位和字典等必要支撑数据,但不能修改角色、菜单或系统配置;`implementation_operator` 与 `viewer` 不具备账户管理权限。未注册的配置和接口管理路由对所有角色返回 404。
+
+## Sense 视频服务规则
+
+- MediaMTX 始终是独立二进制;Sense 管理配置、进程生命周期、路径期望态和状态对账,不把媒体内核放入 GoAdmin handler 或 GORM model。
+- Control API 只能绑定回环地址。Sense 可启动配置的 MediaMTX,也可连接已由外部启动的实例;外部实例标记为非本实例所有,孤儿安全闸禁止 Sense 停止它。
+- 已验证 Profile 幂等形成媒体路径;数据库不保存带凭据 Stream URI。摄像头凭据只在 loopback Control API 请求边界临时使用,不进入基础配置、日志或 Sense API。
+- 路径状态区分 pending、waiting、ready、process_unavailable、apply_failed、status_unavailable、path_missing、stopped,并保存失败次数和有上限的下次重试时间。
+- 冷启动恢复 desired=running 路径;用户明确停止的路径保持 stopped,不因启动扫描自动重新启用。稳定路径只刷新状态,不重复下发配置或无意义增加版本。
+- MediaMTX 失败不得删除或降级设备台账与最后一次已验证 Profile。
+
+
## Sense 旧 MVP 规则状态
diff --git a/docs/04-local-development-and-verification.md b/docs/04-local-development-and-verification.md
index 5bbedd5..bd72188 100644
--- a/docs/04-local-development-and-verification.md
+++ b/docs/04-local-development-and-verification.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Local-Development-and-Verification
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Local-Development-and-Verification.-
-wiki_revision: c8afa08e1e189a9c2779da1cbd3ed75124fd7eae
-synchronized_at: 2026-08-14T08:53:20Z
+wiki_revision: e084a4ae3f5038c4ec7dae6028a2e3d7527680a8
+synchronized_at: 2026-08-14T09:46:47Z
# 本地开发与验证
@@ -167,6 +167,29 @@ $env:SENSE_CREDENTIAL_KEY = [Convert]::ToBase64String($keyBytes)
视频接入还需在仓库外配置 SENSE_ONVIF_DISCOVERY_IP(获准的本机网卡 IP)和 SENSE_ONVIF_ALLOWED_CIDRS(逗号分隔的获准摄像头网段)。不要使用 0.0.0.0/0 代替授权清单。
协议回归位于 app/sense/onvif、app/sense/rtsp、app/sense/admission;隔离 PostgreSQL 重启恢复测试通过 SENSE_ADMISSION_TEST_DATABASE_URL 显式启用。验证至少覆盖 Digest/Basic、无配置发现提示、URL 凭据和敏感查询拒绝、目标网段、重定向、Media/Stream 主机归一化、主子码流、失败重探保留已验证 Profile,以及 viewer 只读权限。
+MediaMTX 保持仓库外独立二进制。运行前在进程环境设置:
+
+```powershell
+$env:SENSE_MEDIAMTX_BINARY = ''
+$env:SENSE_MEDIAMTX_CONFIG = '<仓库外 mediamtx.yml>'
+$env:SENSE_MEDIAMTX_API = 'http://127.0.0.1:9997'
+```
+
+配置文件不存在时 Sense 只生成 loopback API 和空 `paths: {}` 的无凭据基础配置;模板位于 `Sense/server/config/mediamtx/mediamtx.yml.example`。Control API 不允许非回环地址。真实集成验证使用:
+
+```powershell
+$env:SENSE_MEDIAMTX_TEST_BINARY = ''
+go test ./tests/media -run TestRealMediaMTXControlLifecycle -v
+
+$env:SENSE_MEDIA_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
+go test ./tests/media -run TestPostgresColdStartRestoresDesiredRoute -v
+
+$env:SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
+go test ./cmd/migrate/migration/version -run TestMediaMigrationOnPostgres -v
+```
+
+测试必须使用隔离端口和数据库;结束后停止测试进程。不得输出连接串或摄像头凭据。
+
diff --git a/docs/06-troubleshooting.md b/docs/06-troubleshooting.md
index 6df6047..bce0f28 100644
--- a/docs/06-troubleshooting.md
+++ b/docs/06-troubleshooting.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Troubleshooting
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Troubleshooting
-wiki_revision: ae80dc2d34341ff3cddfff93f51373adb3bd0664
-synchronized_at: 2026-08-14T08:53:30Z
+wiki_revision: c615cd5235557a3888260beb5a682c6ed4e5b5dc
+synchronized_at: 2026-08-14T09:46:56Z
# 故障排查
@@ -73,4 +73,18 @@ synchronized_at: 2026-08-14T08:53:30Z
| 设备时间异常 | 在摄像机管理页或受控 NTP 环境校时后重新探测;Sense 不自动修改设备时间。 |
| 部分码流失败 | 查看逐 Profile 状态、设备 RTSP 权限和端口;最后一次已验证 Profile 会保留。 |
| 重定向已拒绝 | ONVIF 服务返回了 3xx;修正为摄像机最终服务地址,不允许 Sense 跟随到未知目标。 |
+
+## Sense 视频服务排错
+
+| 现象 | 原因与处理 |
+|---|---|
+| 进程状态“未配置” | 未设置 SENSE_MEDIAMTX_BINARY;如由外部服务管理,先确认 loopback Control API 已就绪,否则配置二进制和仓库外配置路径。 |
+| 进程启动失败 | 核对二进制存在、配置目录可写、MediaMTX 配置可解析,以及 RTSP/API 端口未被其他进程占用。 |
+| 等待拉流 | 路径已建立但 sourceOnDemand 尚无 reader;打开实时监看后再观察,不等同于接入失败。 |
+| 配置失败或路径缺失 | 在“视频服务”点击对账;检查 Control API 仍为 loopback、Profile 仍已验证、RTSP 凭据可用。 |
+| 显示外部启动(受保护) | Sense 检测到不是本实例启动的 MediaMTX;孤儿安全闸生效,Sense 关闭时不会停止它。 |
+| 持续自动重试 | 查看失败码、失败次数和下次重试时间;修正二进制、端口、凭据或上游后等待退避到期,或由有权限用户立即对账。 |
+| Sense 重启后路径未恢复 | 确认数据库 route 的 desired 为 running、迁移已执行、Control API 可达;明确停止的路径不会自动恢复。 |
+
+
diff --git a/docs/delivery/README.md b/docs/delivery/README.md
index e415ea3..443efcf 100644
--- a/docs/delivery/README.md
+++ b/docs/delivery/README.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Delivery-Documentation-Guide
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Delivery-Documentation-Guide.-
-wiki_revision: 07c8711ff4d75da5e4acbd61d029d27f01009747
-synchronized_at: 2026-08-14T08:54:16Z
+wiki_revision: 45c86f2a0e4d3252e8042df5ee725e633dd497c1
+synchronized_at: 2026-08-14T09:47:31Z
# 交付文档指南
@@ -119,4 +119,14 @@ Sense 面向网管、实施人员和非技术现场人员,菜单按日常任
部署人员必须先确认获准摄像头网段,再把本机对应网卡 IP 配置为 SENSE_ONVIF_DISCOVERY_IP,把获准网段配置为逗号分隔的 SENSE_ONVIF_ALLOWED_CIDRS。不得为了省事填写全网段。现场人员在“视频接入”选择已登记且已配置凭据的视频设备,可使用发现结果或手工填写不含账号密码的 ONVIF 地址。
验证结果区分可用、部分码流失败、认证失败、目标未获准、重定向拒绝、响应超时、设备时间异常和无法连接,并显示主/子码流及逐 Profile 状态。失败重试不会删除上次已验证 Profile;修改凭据后应重新验证。真实摄像机兼容性、网络 ACL 和设备校时仍需在客户授权环境完成。
+
+### Sense 视频服务交付说明
+
+“视频服务”面向实施、运维和站点管理员展示 MediaMTX 进程归属、媒体路径、拉流状态、观看数、失败原因和下次重试。只读用户只能查看;有操作权限的人员可立即对账或停止单一路径。
+
+交付时 MediaMTX 二进制和真实配置放在仓库外受控目录,Control API 只能绑定回环地址。Sense 关闭时只停止自己启动的 MediaMTX;外部启动进程会显示“外部启动(受保护)”。“等待拉流”表示按需路径尚无观看者,不等同于故障。停止路径不会删除设备或 Profile,后续重新接入可恢复期望态。
+
+任何客户文档、截图和日志都不得包含 Control API 请求体、摄像头凭据或带凭据 URI。现场至少验证启动失败、端口冲突、外部进程保护、重复对账、冷启动恢复和 reader 状态。
+
+
--
2.34.1
From 263a68c98dd5ce8cc65c7a27a378470a5de3d3aa Mon Sep 17 00:00:00 2001
From: QiuSW <105186638@qq.com>
Date: Fri, 14 Aug 2026 18:00:44 +0800
Subject: [PATCH 3/4] =?UTF-8?q?docs:=20=E5=BD=92=E6=A1=A3=E5=B7=A5?=
=?UTF-8?q?=E5=8D=95=20#67=20=E5=BE=85=E9=AA=8C=E6=94=B6=E8=AF=81=E6=8D=AE?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
...视频服务生命周期与状态对账.md | 66 +++++++++++++++++++
wiki-docs.json | 4 ++
2 files changed, 70 insertions(+)
create mode 100644 docs/task/67-Sense视频服务生命周期与状态对账.md
diff --git a/docs/task/67-Sense视频服务生命周期与状态对账.md b/docs/task/67-Sense视频服务生命周期与状态对账.md
new file mode 100644
index 0000000..cda1697
--- /dev/null
+++ b/docs/task/67-Sense视频服务生命周期与状态对账.md
@@ -0,0 +1,66 @@
+
+generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
+wiki_page: Task-67-Sense视频服务生命周期与状态对账
+wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Task-67-Sense%E8%A7%86%E9%A2%91%E6%9C%8D%E5%8A%A1%E7%94%9F%E5%91%BD%E5%91%A8%E6%9C%9F%E4%B8%8E%E7%8A%B6%E6%80%81%E5%AF%B9%E8%B4%A6.-
+wiki_revision: 8ce039b45a20e680e378932331d9a8da1cae1654
+synchronized_at: 2026-08-14T09:58:07Z
+
+
+# 67 Sense视频服务生命周期与状态对账
+
+- 类型:需求
+- 所属 Epic:#7
+- 所属 MVP / 版本:#8
+- 状态:待验收
+- 日期:2026-08-14
+- 工单:https://git.ilapage.cn/ila/yovision/issues/67
+- Pull Request:https://git.ilapage.cn/ila/yovision/pulls/85
+- 主项目:Sense
+
+## 背景与目标
+
+在 #66 已验证设备/Profile 上重建 MediaMTX 安全配置、独立进程生命周期、按需路径和持久化状态对账。Brain、Bell 不启动时可独立运行和验收。
+
+## 最终方案
+
+- MediaMTX 保持独立二进制;Sense 在 GoAdmin server 生命周期内启动后台对账,只停止自己启动的子进程。
+- Control API 只允许 loopback HTTP,禁用代理和重定向;检测到外部实例时标记“外部启动(受保护)”,孤儿安全闸禁止停止。
+- 路由只持久化 Device/Profile 引用、哈希路径、desired/actual、reader、失败次数、退避和下次重试,不保存带凭据 URI。
+- 摄像头凭据按需从 #65 内部端口读取,仅在 MediaMTX v1.19.3 loopback 请求边界临时组装,不进入基础配置、日志或 Sense API。
+- 接入成功后幂等建立路由;冷启动恢复 desired=running,明确 stopped 路径保持停止;稳定路径只刷新状态,不重复下发或增加版本。
+- 复用 GoAdmin JWT/Casbin/操作审计、迁移、动态菜单和 go-admin-ui BasicLayout、Element Plus Descriptions/Table/Tag/Button/MessageBox。
+
+## 验收结果
+
+| 标准 | 结果 |
+|---|---|
+| 已验证 Profile 幂等建立/更新路径 | 单元测试及接入内部端口通过 |
+| 进程未启动、启动中、失败和恢复可定位 | 状态机、真实二进制和页面反馈通过 |
+| 冷启动恢复 running 路径 | PostgreSQL 17 重建 Service 测试通过 |
+| reader、上游、退避、下次重试、孤儿闸可观察 | DTO、页面、退避和外部进程测试通过 |
+| 配置、日志和 Sense API 不泄漏凭据 | loopback 边界、无秘密模型/响应及扫描通过 |
+| MediaMTX 失败不破坏 Profile | 失败持久化测试确认 Profile 保持 ready |
+
+## 测试
+
+- `go test ./...`、`go vet ./...`、`go build ./...`:通过。
+- `go test -race ./app/sense/media ./app/sense/reconcile`:通过。
+- 前端 lint:0 error,32 条上游/目录命名 warning。
+- 前端单测:10 suites、35 tests 通过。
+- 前端生产构建:通过,6 条上游继承 warning。
+- 真实 MediaMTX v1.19.3:Control API 就绪、path add/patch/delete 通过。
+- PostgreSQL 17:desired route 保存后重建 Service 并恢复为 ready 通过。
+- PostgreSQL 17 定向迁移:3 个菜单、12 条角色策略、1 条迁移记录通过。
+- Wiki 镜像检查:通过。
+- 已处理测试安全问题:MediaMTX 自动 TLS 文件的工作目录已固定到外部配置目录;临时证书/私钥未进入最终提交或远端。
+- 未验证部分:未连接客户真实摄像机和现场网络;真实上游持续拉流、reader 变化、端口 ACL 和目标浏览器留待授权现场验收。
+- 已知相邻问题:空白 PostgreSQL 执行完整上游迁移链时,在到达 #67 前被旧 `sys_config` 初始化字段长度问题中止;#67 定向迁移已通过,空库安装链应由 #70/#71 单独复核,不在本工单混改上游初始化。
+
+## 回退
+
+回退 PR #85 可移除运行时、路由和页面入口;已验证 Device/Profile 不受影响。媒体路由表保留可审计状态,停止/回退不删除摄像机数据。
+
+## 相关提交
+
+- 19f9bfa feat: 重建 MediaMTX 生命周期与状态对账 (#67)
+- 4149cc4 docs: 记录 Sense 视频服务生命周期 (#67)
diff --git a/wiki-docs.json b/wiki-docs.json
index b98eb58..ee06431 100644
--- a/wiki-docs.json
+++ b/wiki-docs.json
@@ -115,6 +115,10 @@
{
"page": "Task-66-Sense视频接入与Profile",
"path": "docs/task/66-Sense视频接入与Profile.md"
+ },
+ {
+ "page": "Task-67-Sense视频服务生命周期与状态对账",
+ "path": "docs/task/67-Sense视频服务生命周期与状态对账.md"
}
]
}
--
2.34.1
From 35b9c7b8a4482407da841a85c5a878d38bebe987 Mon Sep 17 00:00:00 2001
From: QiuSW <105186638@qq.com>
Date: Fri, 14 Aug 2026 18:16:15 +0800
Subject: [PATCH 4/4] =?UTF-8?q?docs:=20=E8=AE=B0=E5=BD=95=E5=B7=A5?=
=?UTF-8?q?=E5=8D=95=20#67=20=E9=AA=8C=E6=94=B6=E9=80=9A=E8=BF=87?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
...7-Sense视频服务生命周期与状态对账.md | 12 +++++++++---
1 file changed, 9 insertions(+), 3 deletions(-)
diff --git a/docs/task/67-Sense视频服务生命周期与状态对账.md b/docs/task/67-Sense视频服务生命周期与状态对账.md
index cda1697..5f59a14 100644
--- a/docs/task/67-Sense视频服务生命周期与状态对账.md
+++ b/docs/task/67-Sense视频服务生命周期与状态对账.md
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Task-67-Sense视频服务生命周期与状态对账
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Task-67-Sense%E8%A7%86%E9%A2%91%E6%9C%8D%E5%8A%A1%E7%94%9F%E5%91%BD%E5%91%A8%E6%9C%9F%E4%B8%8E%E7%8A%B6%E6%80%81%E5%AF%B9%E8%B4%A6.-
-wiki_revision: 8ce039b45a20e680e378932331d9a8da1cae1654
-synchronized_at: 2026-08-14T09:58:07Z
+wiki_revision: 4642924705d7a3406874ef532579fb2d6a25a88b
+synchronized_at: 2026-08-14T10:13:23Z
# 67 Sense视频服务生命周期与状态对账
@@ -11,7 +11,7 @@ synchronized_at: 2026-08-14T09:58:07Z
- 类型:需求
- 所属 Epic:#7
- 所属 MVP / 版本:#8
-- 状态:待验收
+- 状态:已完成
- 日期:2026-08-14
- 工单:https://git.ilapage.cn/ila/yovision/issues/67
- Pull Request:https://git.ilapage.cn/ila/yovision/pulls/85
@@ -60,7 +60,13 @@ synchronized_at: 2026-08-14T09:58:07Z
回退 PR #85 可移除运行时、路由和页面入口;已验证 Device/Profile 不受影响。媒体路由表保留可审计状态,停止/回退不删除摄像机数据。
+## 验收
+
+- 用户于 2026-08-14 明确验收通过。
+
## 相关提交
- 19f9bfa feat: 重建 MediaMTX 生命周期与状态对账 (#67)
- 4149cc4 docs: 记录 Sense 视频服务生命周期 (#67)
+
+- 263a68c docs: 归档工单 #67 待验收证据
--
2.34.1