From d681fd13450f47b4f27efc0e8fd46383df0ee08a Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 28 Aug 2026 16:09:52 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=9E=E7=8E=B0=20MediaMTX=20?= =?UTF-8?q?=E5=88=86=E7=89=87=E5=88=86=E9=85=8D=E4=B8=8E=E6=95=85=E9=9A=9C?= =?UTF-8?q?=E8=8C=83=E5=9B=B4=E5=8F=AF=E8=A7=82=E5=AF=9F=20(#77)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Sense/config/mediamtx-shards.example.json | 9 + Sense/config/sense.demo.env.example | 2 + Sense/config/sense.env.example | 4 + Sense/scripts/build/build-windows.ps1 | 1 + Sense/scripts/build/test-package.ps1 | 2 +- Sense/scripts/runtime/sense-common.ps1 | 8 +- .../app/admin/router/sense_media_shard.go | 20 ++ Sense/server/app/sense/media/runtime.go | 23 ++ Sense/server/app/sense/media/service.go | 71 ++++- Sense/server/app/sense/media_shard/apis.go | 91 ++++++ Sense/server/app/sense/media_shard/audit.go | 19 ++ Sense/server/app/sense/media_shard/config.go | 80 ++++++ .../app/sense/media_shard/config_test.go | 46 ++++ Sense/server/app/sense/media_shard/dto.go | 72 +++++ Sense/server/app/sense/media_shard/models.go | 59 ++++ Sense/server/app/sense/media_shard/service.go | 259 ++++++++++++++++++ .../app/sense/media_shard/service_test.go | 130 +++++++++ .../version/2026082814000_media_shard.go | 57 ++++ .../version/2026082814000_media_shard_test.go | 86 ++++++ Sense/server/config/credential.env.example | 2 + Sense/server/tests/media/integration_test.go | 80 +++++- Sense/tests/package/run-tests.ps1 | 3 + Sense/ui/src/api/sense/media-shard.js | 13 + .../ui/src/views/sense/media-shard/index.vue | 65 +++++ .../sense/media-shard/mediaShardState.js | 6 + .../ui/tests/unit/sense/mediaShardApi.spec.js | 8 + .../tests/unit/sense/mediaShardState.spec.js | 7 + docs/02-architecture-and-code-map.md | 15 +- docs/03-business-rules-and-glossary.md | 15 +- docs/04-local-development-and-verification.md | 42 ++- docs/delivery/deployment-and-operations.md | 16 +- 31 files changed, 1293 insertions(+), 18 deletions(-) create mode 100644 Sense/config/mediamtx-shards.example.json create mode 100644 Sense/server/app/admin/router/sense_media_shard.go create mode 100644 Sense/server/app/sense/media_shard/apis.go create mode 100644 Sense/server/app/sense/media_shard/audit.go create mode 100644 Sense/server/app/sense/media_shard/config.go create mode 100644 Sense/server/app/sense/media_shard/config_test.go create mode 100644 Sense/server/app/sense/media_shard/dto.go create mode 100644 Sense/server/app/sense/media_shard/models.go create mode 100644 Sense/server/app/sense/media_shard/service.go create mode 100644 Sense/server/app/sense/media_shard/service_test.go create mode 100644 Sense/server/cmd/migrate/migration/version/2026082814000_media_shard.go create mode 100644 Sense/server/cmd/migrate/migration/version/2026082814000_media_shard_test.go create mode 100644 Sense/ui/src/api/sense/media-shard.js create mode 100644 Sense/ui/src/views/sense/media-shard/index.vue create mode 100644 Sense/ui/src/views/sense/media-shard/mediaShardState.js create mode 100644 Sense/ui/tests/unit/sense/mediaShardApi.spec.js create mode 100644 Sense/ui/tests/unit/sense/mediaShardState.spec.js diff --git a/Sense/config/mediamtx-shards.example.json b/Sense/config/mediamtx-shards.example.json new file mode 100644 index 0000000..7445c6c --- /dev/null +++ b/Sense/config/mediamtx-shards.example.json @@ -0,0 +1,9 @@ +[ + { + "id": "secondary", + "name": "备用媒体分片", + "mode": "external", + "controlAPI": "http://127.0.0.1:19997", + "capacity": 24 + } +] diff --git a/Sense/config/sense.demo.env.example b/Sense/config/sense.demo.env.example index 42884c1..f3a8c97 100644 --- a/Sense/config/sense.demo.env.example +++ b/Sense/config/sense.demo.env.example @@ -13,6 +13,8 @@ SENSE_MEDIAMTX_MODE=disabled SENSE_MEDIAMTX_BINARY= SENSE_MEDIAMTX_CONFIG= SENSE_MEDIAMTX_API=http://127.0.0.1:9997 +SENSE_MEDIAMTX_CAPACITY= +SENSE_MEDIAMTX_SHARDS_FILE= SENSE_WEB_ROOT=web SENSE_AUTO_MIGRATE=true SENSE_POSTGRES_BIN= diff --git a/Sense/config/sense.env.example b/Sense/config/sense.env.example index 1bfef50..4da51f9 100644 --- a/Sense/config/sense.env.example +++ b/Sense/config/sense.env.example @@ -13,6 +13,10 @@ SENSE_MEDIAMTX_MODE=managed SENSE_MEDIAMTX_BINARY=bin\mediamtx.exe SENSE_MEDIAMTX_CONFIG=config\mediamtx.yml SENSE_MEDIAMTX_API=http://127.0.0.1:9997 +# Leave capacity empty to use the current database quota. Optional extra +# shards are read from a repository-external JSON file and must use loopback APIs. +SENSE_MEDIAMTX_CAPACITY= +SENSE_MEDIAMTX_SHARDS_FILE= SENSE_WEB_ROOT=web SENSE_AUTO_MIGRATE=true SENSE_POSTGRES_BIN= diff --git a/Sense/scripts/build/build-windows.ps1 b/Sense/scripts/build/build-windows.ps1 index a08194f..b1881fa 100644 --- a/Sense/scripts/build/build-windows.ps1 +++ b/Sense/scripts/build/build-windows.ps1 @@ -81,6 +81,7 @@ try { Copy-Item -LiteralPath (Join-Path $senseRoot 'config\sense.demo.env.example') -Destination (Join-Path $staging 'config\sense.demo.env.example') Copy-Item -LiteralPath (Join-Path $senseRoot 'config\sense.demo.env.example') -Destination (Join-Path $staging 'config\sense.demo.env') Copy-Item -LiteralPath (Join-Path $senseRoot 'config\mediamtx.yml') -Destination (Join-Path $staging 'config\mediamtx.yml') + Copy-Item -LiteralPath (Join-Path $senseRoot 'config\mediamtx-shards.example.json') -Destination (Join-Path $staging 'config\mediamtx-shards.example.json') Copy-Item -LiteralPath (Join-Path $serverRoot 'config\db.sql') -Destination (Join-Path $staging 'config\db.sql') Copy-Item -LiteralPath (Join-Path $serverRoot 'config\pg.sql') -Destination (Join-Path $staging 'config\pg.sql') Copy-Item -LiteralPath (Join-Path $senseRoot 'README-WINDOWS.md') -Destination (Join-Path $staging 'README-WINDOWS.md') diff --git a/Sense/scripts/build/test-package.ps1 b/Sense/scripts/build/test-package.ps1 index 2035141..6399fce 100644 --- a/Sense/scripts/build/test-package.ps1 +++ b/Sense/scripts/build/test-package.ps1 @@ -9,7 +9,7 @@ $required = @( 'migrate-sense.bat', 'backup-sense.bat', 'restore-sense.bat', 'initialize-admin.bat', 'README-WINDOWS.md', 'config\sense.env', 'config\sense.env.example', 'config\sense.demo.env', - 'config\mediamtx.yml', 'config\db.sql', 'config\pg.sql', + 'config\mediamtx.yml', 'config\mediamtx-shards.example.json', 'config\db.sql', 'config\pg.sql', 'web\index.html', 'scripts\runtime\sense-common.ps1' ) foreach ($relative in $required) { diff --git a/Sense/scripts/runtime/sense-common.ps1 b/Sense/scripts/runtime/sense-common.ps1 index 03b32bc..dc7f8a0 100644 --- a/Sense/scripts/runtime/sense-common.ps1 +++ b/Sense/scripts/runtime/sense-common.ps1 @@ -7,7 +7,8 @@ $script:SenseAllowedEnvironment = @( 'SENSE_CREDENTIAL_KEY', 'SENSE_ONVIF_DISCOVERY_IP', 'SENSE_ONVIF_ALLOWED_CIDRS', 'SENSE_MEDIAMTX_MODE', 'SENSE_MEDIAMTX_BINARY', 'SENSE_MEDIAMTX_CONFIG', - 'SENSE_MEDIAMTX_API', 'SENSE_WEB_ROOT', 'SENSE_AUTO_MIGRATE', + 'SENSE_MEDIAMTX_API', 'SENSE_MEDIAMTX_CAPACITY', + 'SENSE_MEDIAMTX_SHARDS_FILE', 'SENSE_WEB_ROOT', 'SENSE_AUTO_MIGRATE', 'SENSE_POSTGRES_BIN' ) @@ -185,6 +186,7 @@ function Initialize-SenseRuntime { } $mediaBinary = Resolve-SenseConfiguredPath -PackageRoot $PackageRoot -Value (Get-SenseEnvironmentValue -Name 'SENSE_MEDIAMTX_BINARY') $mediaConfig = Resolve-SenseConfiguredPath -PackageRoot $PackageRoot -Value (Get-SenseEnvironmentValue -Name 'SENSE_MEDIAMTX_CONFIG') + $mediaShardsFile = Resolve-SenseConfiguredPath -PackageRoot $PackageRoot -Value (Get-SenseEnvironmentValue -Name 'SENSE_MEDIAMTX_SHARDS_FILE') if ($mediaMode -eq 'managed') { if (-not (Test-Path -LiteralPath $mediaBinary -PathType Leaf)) { throw 'Managed MediaMTX binary not found. Set SENSE_MEDIAMTX_BINARY to mediamtx.exe.' } if (-not (Test-Path -LiteralPath $mediaConfig -PathType Leaf)) { throw 'Managed MediaMTX configuration not found. Set SENSE_MEDIAMTX_CONFIG.' } @@ -192,6 +194,9 @@ function Initialize-SenseRuntime { if ($mediaMode -eq 'external' -and -not (Test-SenseTcpEndpoint -HostName $apiUri.Host -Port $apiUri.Port)) { throw "External MediaMTX Control API is unreachable at $($apiUri.Host):$($apiUri.Port)." } + if (-not [string]::IsNullOrWhiteSpace($mediaShardsFile) -and -not (Test-Path -LiteralPath $mediaShardsFile -PathType Leaf)) { + throw 'SENSE_MEDIAMTX_SHARDS_FILE does not exist.' + } $webRoot = Resolve-SenseConfiguredPath -PackageRoot $PackageRoot -Value (Get-SenseEnvironmentValue -Name 'SENSE_WEB_ROOT' -Default 'web') if (-not (Test-Path -LiteralPath (Join-Path $webRoot 'index.html') -PathType Leaf)) { throw 'Sense web assets are missing. Rebuild or replace the delivery package.' } @@ -205,6 +210,7 @@ function Initialize-SenseRuntime { [Environment]::SetEnvironmentVariable('SENSE_MEDIAMTX_CONFIG', '', 'Process') } [Environment]::SetEnvironmentVariable('SENSE_MEDIAMTX_API', $mediaAPI, 'Process') + [Environment]::SetEnvironmentVariable('SENSE_MEDIAMTX_SHARDS_FILE', $mediaShardsFile, 'Process') $runtimeDir = Join-Path $PackageRoot 'data\runtime' $logDir = Join-Path $PackageRoot 'logs' diff --git a/Sense/server/app/admin/router/sense_media_shard.go b/Sense/server/app/admin/router/sense_media_shard.go new file mode 100644 index 0000000..6fad368 --- /dev/null +++ b/Sense/server/app/admin/router/sense_media_shard.go @@ -0,0 +1,20 @@ +package router + +import ( + "github.com/gin-gonic/gin" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" + "git.ilapage.cn/ila/yovision/Sense/server/common/actions" + "git.ilapage.cn/ila/yovision/Sense/server/common/middleware" +) + +func init() { routerCheckRole = append(routerCheckRole, registerSenseMediaShardRouter) } + +func registerSenseMediaShardRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) { + api := &media_shard.API{} + r := v1.Group("/media-shards").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction()) + r.GET("", api.List) + r.GET("/:id", api.Get) + r.GET("/:id/migration-preflight", api.Preflight) +} diff --git a/Sense/server/app/sense/media/runtime.go b/Sense/server/app/sense/media/runtime.go index e34fc14..669ba17 100644 --- a/Sense/server/app/sense/media/runtime.go +++ b/Sense/server/app/sense/media/runtime.go @@ -9,6 +9,8 @@ import ( "time" "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" ) var runtimeState struct { @@ -41,6 +43,27 @@ func StartRuntime(parent context.Context, db *gorm.DB) error { return err } } + shards := media_shard.NewService(db, func(probeCtx context.Context, endpoint string) error { + probe, probeErr := NewHTTPController(endpoint) + if probeErr != nil { + return probeErr + } + return probe.Health(probeCtx) + }) + specs, err := media_shard.SpecsFromEnvironment(db, mode, config.APIBase) + if err != nil { + cancel() + return err + } + if err = shards.SyncSpecs(ctx, specs); err != nil { + cancel() + return err + } + if err = shards.RefreshAll(ctx); err != nil { + cancel() + return err + } + service.WithShardService(shards) runtimeState.Lock() if runtimeState.cancel != nil { runtimeState.cancel() diff --git a/Sense/server/app/sense/media/service.go b/Sense/server/app/sense/media/service.go index 409dc16..3a36625 100644 --- a/Sense/server/app/sense/media/service.go +++ b/Sense/server/app/sense/media/service.go @@ -13,6 +13,7 @@ import ( "gorm.io/gorm/clause" "git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential" + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" "git.ilapage.cn/ila/yovision/Sense/server/app/sense/reconcile" ) @@ -26,10 +27,16 @@ type Service struct { controller Controller process *Supervisor config RuntimeConfig + shards *media_shard.Service now func() time.Time mu sync.Mutex } +func (s *Service) WithShardService(shards *media_shard.Service) *Service { + s.shards = shards + return s +} + 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} } @@ -76,11 +83,19 @@ func (s *Service) ensureDevice(ctx context.Context, deviceID string, reactivate return err } } + if s.shards != nil { + if _, err := s.shards.EnsureAssignment(ctx, id); err != nil { + return fmt.Errorf("assign media route to shard: %w", err) + } + } } return nil } func (s *Service) Run(ctx context.Context) { + if s.shards != nil { + _ = s.shards.RefreshAll(ctx) + } _ = s.EnsureAllVerified(ctx) _ = s.ReconcileDue(ctx) interval := s.config.PollInterval @@ -94,6 +109,9 @@ func (s *Service) Run(ctx context.Context) { case <-ctx.Done(): return case <-ticker.C: + if s.shards != nil { + _ = s.shards.RefreshAll(ctx) + } _ = s.ReconcileDue(ctx) } } @@ -127,10 +145,11 @@ func (s *Service) Refresh(ctx context.Context, id string) (RouteResponse, error) if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil { return RouteResponse{}, ErrNotFound } - if err := s.ensureControl(ctx); err != nil { + controller, err := s.controllerForRoute(ctx, route) + if err != nil { return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err) } - status, err := s.controller.Status(ctx, route.Path) + status, err := controller.Status(ctx, route.Path) if err != nil { return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err) } @@ -159,12 +178,13 @@ func (s *Service) Reconcile(ctx context.Context, id string) (RouteResponse, erro return RouteResponse{}, err } if route.Desired == DesiredStopped { - if s.controller != nil { - _ = s.controller.Delete(ctx, route.Path) + if controller, err := s.controllerForRoute(ctx, route); err == nil { + _ = controller.Delete(ctx, route.Path) } return s.saveSuccess(ctx, route, "stopped", false, 0, "媒体路径已停止") } - if err := s.ensureControl(ctx); err != nil { + controller, err := s.controllerForRoute(ctx, route) + if err != nil { return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err) } var profile admissionProfile @@ -175,10 +195,10 @@ func (s *Service) Reconcile(ctx context.Context, id string) (RouteResponse, erro 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 { + if err = 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) + status, err := controller.Status(ctx, route.Path) if err != nil { return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err) } @@ -191,6 +211,43 @@ func (s *Service) Reconcile(ctx context.Context, id string) (RouteResponse, erro return s.saveSuccess(ctx, route, "waiting", false, status.Readers, "等待播放器连接并按需拉流") } +func (s *Service) controllerForRoute(ctx context.Context, route Route) (Controller, error) { + if s.shards == nil { + if err := s.ensureControl(ctx); err != nil { + return nil, err + } + return s.controller, nil + } + shard, err := s.shards.ResolveRoute(ctx, route.ID) + if errors.Is(err, media_shard.ErrNotFound) { + shard, err = s.shards.EnsureAssignment(ctx, route.ID) + } + if err != nil { + return nil, err + } + if shard.Status != media_shard.StatusRunning { + return nil, ErrRuntimeUnavailable + } + if shard.ID == "primary" { + if err = s.ensureControl(ctx); err != nil { + return nil, err + } + return s.controller, nil + } + endpoint, err := s.shards.ControlAPI(ctx, shard.ID) + if err != nil { + return nil, err + } + controller, err := NewHTTPController(endpoint) + if err != nil { + return nil, err + } + if err = controller.Health(ctx); err != nil { + return nil, err + } + return controller, nil +} + func (s *Service) ensureControl(ctx context.Context) error { if s.controller == nil { return ErrRuntimeUnavailable diff --git a/Sense/server/app/sense/media_shard/apis.go b/Sense/server/app/sense/media_shard/apis.go new file mode 100644 index 0000000..3876daf --- /dev/null +++ b/Sense/server/app/sense/media_shard/apis.go @@ -0,0 +1,91 @@ +package media_shard + +import ( + "errors" + "net/http" + "time" + + "github.com/gin-gonic/gin" + "github.com/go-admin-team/go-admin-core/sdk/api" + "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user" + coreService "github.com/go-admin-team/go-admin-core/sdk/service" + + "git.ilapage.cn/ila/yovision/Sense/server/common" +) + +const auditSuccess, auditFailure = "1", "2" + +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 NewService(base.Orm, nil), nil +} + +func (e *API) List(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + result, err := service.List(c.Request.Context()) + if err != nil { + e.audit(c, service, "List", auditFailure, "媒体分片列表查询失败") + e.writeError(err) + return + } + e.audit(c, service, "List", auditSuccess, "读取媒体分片列表") + e.OK(result, "查询成功") +} + +func (e *API) Get(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + result, err := service.Get(c.Request.Context(), c.Param("id")) + if err != nil { + e.audit(c, service, "Get", auditFailure, "媒体分片详情查询失败") + e.writeError(err) + return + } + e.audit(c, service, "Get", auditSuccess, "读取媒体分片详情 "+result.ID) + e.OK(result, "查询成功") +} + +func (e *API) Preflight(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + result, err := service.Preflight(c.Request.Context(), c.Param("id")) + if err != nil { + e.audit(c, service, "Preflight", auditFailure, "媒体分片迁移预检失败") + e.writeError(err) + return + } + e.audit(c, service, "Preflight", auditSuccess, "执行媒体分片只读迁移预检 "+result.SourceShardID) + e.OK(result, "预检完成") +} + +func (e *API) audit(c *gin.Context, service *Service, action, status, remark string) { + if err := WriteAudit(service.DB, Audit{Action: action, Method: c.Request.Method, Status: status, Username: user.GetUserName(c), UserID: user.GetUserId(c), ClientIP: common.GetClientIP(c), Route: c.FullPath(), Remark: remark, At: time.Now()}); err != nil { + api.GetRequestLogger(c).Errorf("media shard audit failed: %s", err.Error()) + } +} + +func (e *API) writeError(err error) { + switch { + case errors.Is(err, ErrInvalidID), errors.Is(err, ErrInvalidConfig): + e.Error(http.StatusBadRequest, err, err.Error()) + case errors.Is(err, ErrNotFound): + e.Error(http.StatusNotFound, err, err.Error()) + default: + e.Error(http.StatusInternalServerError, err, "媒体分片查询失败") + } +} diff --git a/Sense/server/app/sense/media_shard/audit.go b/Sense/server/app/sense/media_shard/audit.go new file mode 100644 index 0000000..9b00841 --- /dev/null +++ b/Sense/server/app/sense/media_shard/audit.go @@ -0,0 +1,19 @@ +package media_shard + +import ( + "time" + + adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models" + "gorm.io/gorm" +) + +type Audit struct { + Action, Method, Status, Username, ClientIP, Route, Remark string + UserID int + At time.Time +} + +func WriteAudit(db *gorm.DB, input Audit) error { + model := adminModels.SysOperaLog{Title: "媒体分片", BusinessType: "other", Method: "media_shard.API." + input.Action, RequestMethod: input.Method, OperatorType: "1", OperName: input.Username, OperUrl: input.Route, OperIp: input.ClientIP, Status: input.Status, OperTime: input.At.UTC(), Remark: input.Remark, CreatedAt: input.At.UTC(), UpdatedAt: input.At.UTC()} + return db.Create(&model).Error +} diff --git a/Sense/server/app/sense/media_shard/config.go b/Sense/server/app/sense/media_shard/config.go new file mode 100644 index 0000000..0e13508 --- /dev/null +++ b/Sense/server/app/sense/media_shard/config.go @@ -0,0 +1,80 @@ +package media_shard + +import ( + "encoding/json" + "errors" + "fmt" + "net" + "net/url" + "os" + "strconv" + "strings" + + "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota" +) + +var ErrInvalidConfig = errors.New("MediaMTX 分片配置不符合要求") + +func SpecsFromEnvironment(db *gorm.DB, mode, primaryAPI string) ([]Spec, error) { + capacity, err := quota.ReadLimit(db) + if raw := strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_CAPACITY")); raw != "" { + capacity, err = strconv.Atoi(raw) + } + if err != nil || capacity < 1 { + return nil, fmt.Errorf("%w: primary capacity", ErrInvalidConfig) + } + primary := Spec{ID: "primary", Name: "主媒体分片", Mode: mode, ControlAPI: strings.TrimSpace(primaryAPI), Capacity: capacity} + if err = validateSpec(primary, true); err != nil { + return nil, err + } + result := []Spec{primary} + path := strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_SHARDS_FILE")) + if path == "" { + return result, nil + } + data, err := os.ReadFile(path) + if err != nil { + return nil, fmt.Errorf("read MediaMTX shards file: %w", err) + } + var extra []Spec + if err = json.Unmarshal(data, &extra); err != nil { + return nil, fmt.Errorf("%w: shards JSON", ErrInvalidConfig) + } + seen := map[string]bool{"primary": true} + for _, item := range extra { + item.ID, item.Name, item.Mode, item.ControlAPI = strings.TrimSpace(item.ID), strings.TrimSpace(item.Name), strings.ToLower(strings.TrimSpace(item.Mode)), strings.TrimSpace(item.ControlAPI) + if item.Mode == "" { + item.Mode = "external" + } + if seen[item.ID] || item.Mode != "external" { + return nil, fmt.Errorf("%w: duplicate id or non-external extra shard", ErrInvalidConfig) + } + if err = validateSpec(item, false); err != nil { + return nil, err + } + seen[item.ID] = true + result = append(result, item) + } + return result, nil +} + +func validateSpec(spec Spec, primary bool) error { + if spec.ID == "" || len(spec.ID) > 64 || spec.Name == "" || len([]rune(spec.Name)) > 128 || spec.Capacity < 1 || spec.Capacity > 100000 { + return fmt.Errorf("%w: shard identity or capacity", ErrInvalidConfig) + } + if primary && spec.Mode != "managed" && spec.Mode != "external" { + return fmt.Errorf("%w: primary mode", ErrInvalidConfig) + } + parsed, err := url.Parse(spec.ControlAPI) + if err != nil || parsed.Scheme != "http" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" || parsed.Path != "" { + return fmt.Errorf("%w: control API", ErrInvalidConfig) + } + host := parsed.Hostname() + ip := net.ParseIP(host) + if !strings.EqualFold(host, "localhost") && (ip == nil || !ip.IsLoopback()) { + return fmt.Errorf("%w: control API must be loopback", ErrInvalidConfig) + } + return nil +} diff --git a/Sense/server/app/sense/media_shard/config_test.go b/Sense/server/app/sense/media_shard/config_test.go new file mode 100644 index 0000000..de66619 --- /dev/null +++ b/Sense/server/app/sense/media_shard/config_test.go @@ -0,0 +1,46 @@ +package media_shard + +import ( + "os" + "path/filepath" + "testing" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota" +) + +func TestSpecsFromEnvironmentUsesDatabaseQuotaAndExternalFile(t *testing.T) { + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate("a.Setting{}); err != nil { + t.Fatal(err) + } + if err = db.Create("a.Setting{ID: quota.SettingID, Limit: 23, Source: "database", Version: 1}).Error; err != nil { + t.Fatal(err) + } + path := filepath.Join(t.TempDir(), "shards.json") + if err = os.WriteFile(path, []byte(`[{"id":"secondary","name":"备用分片","controlAPI":"http://127.0.0.1:19997","capacity":11}]`), 0600); err != nil { + t.Fatal(err) + } + t.Setenv("SENSE_MEDIAMTX_CAPACITY", "") + t.Setenv("SENSE_MEDIAMTX_SHARDS_FILE", path) + specs, err := SpecsFromEnvironment(db, "managed", "http://127.0.0.1:9997") + if err != nil { + t.Fatal(err) + } + if len(specs) != 2 || specs[0].Capacity != 23 || specs[1].Capacity != 11 || specs[1].Mode != "external" { + t.Fatalf("unexpected specs: %+v", specs) + } +} + +func TestSpecsRejectsNonLoopbackControlAPI(t *testing.T) { + t.Setenv("SENSE_MEDIAMTX_CAPACITY", "5") + t.Setenv("SENSE_MEDIAMTX_SHARDS_FILE", "") + if _, err := SpecsFromEnvironment(nil, "external", "http://192.0.2.1:9997"); err == nil { + t.Fatal("expected unsafe control API to be rejected") + } +} diff --git a/Sense/server/app/sense/media_shard/dto.go b/Sense/server/app/sense/media_shard/dto.go new file mode 100644 index 0000000..cb793cf --- /dev/null +++ b/Sense/server/app/sense/media_shard/dto.go @@ -0,0 +1,72 @@ +package media_shard + +import "time" + +type Spec struct { + ID string `json:"id"` + Name string `json:"name"` + Mode string `json:"mode"` + ControlAPI string `json:"controlAPI"` + Capacity int `json:"capacity"` +} + +type Summary struct { + Total int `json:"total"` + Available int `json:"available"` + Failed int `json:"failed"` + Configured int `json:"configuredCapacity"` + Assigned int `json:"assignedPaths"` + Impacted int `json:"impactedPaths"` +} + +type Response struct { + ID string `json:"id"` + Name string `json:"name"` + Mode string `json:"mode"` + Status string `json:"status"` + Detail string `json:"detail"` + Capacity int `json:"capacity"` + AssignedPaths int `json:"assignedPaths"` + Remaining int `json:"remaining"` + LastProbeAt *time.Time `json:"lastProbeAt,omitempty"` + Stale bool `json:"stale"` +} + +type Impact struct { + RouteID string `json:"routeId"` + DeviceID string `json:"deviceId"` + DeviceName string `json:"deviceName"` + Location string `json:"location"` + ProfileToken string `json:"profileToken"` + Path string `json:"path"` + Desired string `json:"desired"` + Actual string `json:"actual"` +} + +type DetailResponse struct { + Response + Impact []Impact `json:"impact"` +} + +type PageResponse struct { + Summary Summary `json:"summary"` + List []Response `json:"list"` + Count int64 `json:"count"` +} + +type Check struct { + Name string `json:"name"` + Passed bool `json:"passed"` + Detail string `json:"detail"` +} + +type PreflightResponse struct { + SourceShardID string `json:"sourceShardId"` + TargetShardID string `json:"targetShardId,omitempty"` + TargetShardName string `json:"targetShardName,omitempty"` + ImpactedPaths int `json:"impactedPaths"` + Ready bool `json:"ready"` + ExecutionAuthorized bool `json:"executionAuthorized"` + RequiresIssue bool `json:"requiresSeparateHighRiskIssue"` + Checks []Check `json:"checks"` +} diff --git a/Sense/server/app/sense/media_shard/models.go b/Sense/server/app/sense/media_shard/models.go new file mode 100644 index 0000000..7268dc5 --- /dev/null +++ b/Sense/server/app/sense/media_shard/models.go @@ -0,0 +1,59 @@ +package media_shard + +import "time" + +const ( + StatusUnknown = "unknown" + StatusRunning = "running" + StatusFailed = "failed" + StatusDisabled = "disabled" +) + +// Shard stores only a loopback control endpoint. It is deliberately omitted +// from API response DTOs so credentials and control-plane topology cannot leak. +type Shard struct { + ID string `gorm:"size:64;primaryKey"` + Name string `gorm:"size:128;not null"` + Mode string `gorm:"size:16;not null"` + ControlAPI string `gorm:"size:512;not null"` + Capacity int `gorm:"not null"` + Status string `gorm:"size:16;not null;index"` + Detail string `gorm:"size:256;not null"` + LastProbeAt *time.Time `gorm:"index"` + ConfigVersion int64 `gorm:"not null;default:1"` + CreatedAt time.Time + UpdatedAt time.Time +} + +func (Shard) TableName() string { return "sense_media_shards" } + +// Assignment is the stable, Sense-owned mapping between an existing media +// route and a MediaMTX shard. Failures never rewrite this record automatically. +type Assignment struct { + RouteID string `gorm:"size:96;primaryKey"` + ShardID string `gorm:"size:64;not null;index"` + Algorithm string `gorm:"size:32;not null"` + CreatedAt time.Time `gorm:"not null"` + UpdatedAt time.Time `gorm:"not null"` +} + +func (Assignment) TableName() string { return "sense_media_shard_assignments" } + +type routeProjection struct { + ID string `gorm:"column:id"` + DeviceID string `gorm:"column:device_id"` + ProfileToken string `gorm:"column:profile_token"` + Path string `gorm:"column:path"` + Desired string `gorm:"column:desired"` + Actual string `gorm:"column:actual"` +} + +func (routeProjection) TableName() string { return "sense_media_routes" } + +type deviceProjection struct { + ID string `gorm:"column:id"` + Name string `gorm:"column:name"` + Location string `gorm:"column:location"` +} + +func (deviceProjection) TableName() string { return "sense_devices" } diff --git a/Sense/server/app/sense/media_shard/service.go b/Sense/server/app/sense/media_shard/service.go new file mode 100644 index 0000000..6728fc0 --- /dev/null +++ b/Sense/server/app/sense/media_shard/service.go @@ -0,0 +1,259 @@ +package media_shard + +import ( + "context" + "crypto/sha256" + "encoding/binary" + "errors" + "fmt" + "sort" + "strings" + "time" + + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +var ( + ErrNotFound = errors.New("媒体分片不存在") + ErrNoCapacity = errors.New("没有可用的媒体分片容量") + ErrInvalidID = errors.New("媒体分片标识不符合要求") +) + +const staleAfter = 30 * time.Second + +type Probe func(context.Context, string) error + +type Service struct { + DB *gorm.DB + Probe Probe + Now func() time.Time +} + +func NewService(db *gorm.DB, probe Probe) *Service { + return &Service{DB: db, Probe: probe, Now: time.Now} +} + +func (s *Service) SyncSpecs(ctx context.Context, specs []Spec) error { + return s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + for _, spec := range specs { + var current Shard + err := tx.First(¤t, "id = ?", spec.ID).Error + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + if errors.Is(err, gorm.ErrRecordNotFound) { + current = Shard{ID: spec.ID, Status: StatusUnknown, Detail: "等待健康探测", ConfigVersion: 1} + } else if current.Name != spec.Name || current.Mode != spec.Mode || current.ControlAPI != spec.ControlAPI || current.Capacity != spec.Capacity { + current.ConfigVersion++ + } + current.Name, current.Mode, current.ControlAPI, current.Capacity = spec.Name, spec.Mode, spec.ControlAPI, spec.Capacity + if err = tx.Save(¤t).Error; err != nil { + return err + } + } + ids := make([]string, 0, len(specs)) + for _, spec := range specs { + ids = append(ids, spec.ID) + } + return tx.Model(&Shard{}).Where("id NOT IN ?", ids).Updates(map[string]any{"status": StatusDisabled, "detail": "分片已从当前运行配置移除", "updated_at": s.now()}).Error + }) +} + +func (s *Service) RefreshAll(ctx context.Context) error { + var shards []Shard + if err := s.DB.WithContext(ctx).Where("status <> ?", StatusDisabled).Find(&shards).Error; err != nil { + return err + } + var result error + for _, shard := range shards { + status, detail := StatusRunning, "Control API 正常" + if s.Probe == nil || s.Probe(ctx, shard.ControlAPI) != nil { + status, detail = StatusFailed, "Control API 不可用" + } + now := s.now() + if err := s.DB.WithContext(ctx).Model(&Shard{}).Where("id = ?", shard.ID).Updates(map[string]any{"status": status, "detail": detail, "last_probe_at": &now, "updated_at": now}).Error; err != nil { + result = errors.Join(result, err) + } + } + return result +} + +func (s *Service) EnsureAssignment(ctx context.Context, routeID string) (Shard, error) { + routeID = strings.TrimSpace(routeID) + if routeID == "" { + return Shard{}, ErrInvalidID + } + var selected Shard + err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var existing Assignment + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&existing, "route_id = ?", routeID).Error; err == nil { + return tx.First(&selected, "id = ?", existing.ShardID).Error + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + var candidates []Shard + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("status IN ?", []string{StatusRunning, StatusUnknown}).Find(&candidates).Error; err != nil { + return err + } + counts, err := assignmentCounts(tx) + if err != nil { + return err + } + eligible := candidates[:0] + for _, item := range candidates { + if counts[item.ID] < int64(item.Capacity) { + eligible = append(eligible, item) + } + } + if len(eligible) == 0 { + return ErrNoCapacity + } + sort.SliceStable(eligible, func(i, j int) bool { return rendezvous(routeID, eligible[i].ID) > rendezvous(routeID, eligible[j].ID) }) + selected = eligible[0] + now := s.now() + assignment := Assignment{RouteID: routeID, ShardID: selected.ID, Algorithm: "rendezvous-v1", CreatedAt: now, UpdatedAt: now} + if err := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&assignment).Error; err != nil { + return err + } + if assignment.RouteID != "" { // reload also handles a concurrent winner on PostgreSQL. + if err := tx.First(&existing, "route_id = ?", routeID).Error; err != nil { + return err + } + return tx.First(&selected, "id = ?", existing.ShardID).Error + } + return nil + }) + return selected, err +} + +func (s *Service) ResolveRoute(ctx context.Context, routeID string) (Shard, error) { + var shard Shard + err := s.DB.WithContext(ctx).Table("sense_media_shards AS s").Select("s.*").Joins("JOIN sense_media_shard_assignments a ON a.shard_id = s.id").Where("a.route_id = ?", routeID).Scan(&shard).Error + if err != nil { + return Shard{}, err + } + if shard.ID == "" { + return Shard{}, ErrNotFound + } + return shard, nil +} + +func (s *Service) List(ctx context.Context) (PageResponse, error) { + var shards []Shard + if err := s.DB.WithContext(ctx).Order("name ASC, id ASC").Find(&shards).Error; err != nil { + return PageResponse{}, err + } + counts, err := assignmentCounts(s.DB.WithContext(ctx)) + if err != nil { + return PageResponse{}, err + } + result := PageResponse{List: make([]Response, 0, len(shards)), Count: int64(len(shards))} + for _, shard := range shards { + item := s.response(shard, int(counts[shard.ID])) + result.List = append(result.List, item) + result.Summary.Total++ + result.Summary.Configured += item.Capacity + result.Summary.Assigned += item.AssignedPaths + if item.Status == StatusRunning { + result.Summary.Available++ + } + if item.Status == StatusFailed { + result.Summary.Failed++ + result.Summary.Impacted += item.AssignedPaths + } + } + return result, nil +} + +func (s *Service) Get(ctx context.Context, id string) (DetailResponse, error) { + id = strings.TrimSpace(id) + if id == "" || len(id) > 64 { + return DetailResponse{}, ErrInvalidID + } + var shard Shard + if err := s.DB.WithContext(ctx).First(&shard, "id = ?", id).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return DetailResponse{}, ErrNotFound + } + return DetailResponse{}, err + } + impact, err := s.impact(ctx, id) + if err != nil { + return DetailResponse{}, err + } + return DetailResponse{Response: s.response(shard, len(impact)), Impact: impact}, nil +} + +func (s *Service) Preflight(ctx context.Context, id string) (PreflightResponse, error) { + detail, err := s.Get(ctx, id) + if err != nil { + return PreflightResponse{}, err + } + result := PreflightResponse{SourceShardID: id, ImpactedPaths: len(detail.Impact), ExecutionAuthorized: false, RequiresIssue: true} + result.Checks = append(result.Checks, Check{Name: "源分片状态", Passed: detail.Status == StatusRunning && !detail.Stale, Detail: detail.Detail}) + var candidates []Shard + if err = s.DB.WithContext(ctx).Where("id <> ? AND status = ?", id, StatusRunning).Find(&candidates).Error; err != nil { + return PreflightResponse{}, err + } + counts, err := assignmentCounts(s.DB.WithContext(ctx)) + if err != nil { + return PreflightResponse{}, err + } + for _, item := range candidates { + if item.Capacity-int(counts[item.ID]) >= len(detail.Impact) && (result.TargetShardID == "" || rendezvous(id, item.ID) > rendezvous(id, result.TargetShardID)) { + result.TargetShardID, result.TargetShardName = item.ID, item.Name + } + } + result.Checks = append(result.Checks, Check{Name: "目标容量", Passed: result.TargetShardID != "", Detail: map[bool]string{true: "存在可容纳全部受影响路径的健康目标分片", false: "没有可容纳全部受影响路径的健康目标分片"}[result.TargetShardID != ""]}) + result.Checks = append(result.Checks, Check{Name: "实施授权", Passed: false, Detail: "本工单只提供只读预检;实际迁移必须另建高风险工单并人工确认"}) + result.Ready = result.Checks[0].Passed && result.Checks[1].Passed + return result, nil +} + +func (s *Service) impact(ctx context.Context, id string) ([]Impact, error) { + var rows []Impact + err := s.DB.WithContext(ctx).Table("sense_media_shard_assignments AS a").Select("r.id AS route_id, r.device_id, COALESCE(d.name, '') AS device_name, COALESCE(d.location, '') AS location, r.profile_token, r.path, r.desired, r.actual").Joins("JOIN sense_media_routes r ON r.id = a.route_id").Joins("LEFT JOIN sense_devices d ON d.id = r.device_id").Where("a.shard_id = ?", id).Order("d.name ASC, r.profile_token ASC").Scan(&rows).Error + return rows, err +} + +func (s *Service) response(shard Shard, assigned int) Response { + remaining := shard.Capacity - assigned + if remaining < 0 { + remaining = 0 + } + stale := shard.LastProbeAt == nil || s.now().Sub(shard.LastProbeAt.UTC()) > staleAfter + return Response{ID: shard.ID, Name: shard.Name, Mode: shard.Mode, Status: shard.Status, Detail: shard.Detail, Capacity: shard.Capacity, AssignedPaths: assigned, Remaining: remaining, LastProbeAt: shard.LastProbeAt, Stale: stale} +} + +func assignmentCounts(db *gorm.DB) (map[string]int64, error) { + type row struct { + ShardID string + Count int64 + } + var rows []row + err := db.Model(&Assignment{}).Select("shard_id, COUNT(*) AS count").Group("shard_id").Scan(&rows).Error + out := map[string]int64{} + for _, r := range rows { + out[r.ShardID] = r.Count + } + return out, err +} +func rendezvous(routeID, shardID string) uint64 { + sum := sha256.Sum256([]byte(routeID + "\x00" + shardID)) + return binary.BigEndian.Uint64(sum[:8]) +} +func (s *Service) now() time.Time { + if s.Now != nil { + return s.Now().UTC() + } + return time.Now().UTC() +} + +func (s *Service) ControlAPI(ctx context.Context, id string) (string, error) { + var shard Shard + if err := s.DB.WithContext(ctx).First(&shard, "id = ?", id).Error; err != nil { + return "", fmt.Errorf("get media shard control endpoint: %w", err) + } + return shard.ControlAPI, nil +} diff --git a/Sense/server/app/sense/media_shard/service_test.go b/Sense/server/app/sense/media_shard/service_test.go new file mode 100644 index 0000000..8d4ec04 --- /dev/null +++ b/Sense/server/app/sense/media_shard/service_test.go @@ -0,0 +1,130 @@ +package media_shard + +import ( + "context" + "errors" + "testing" + "time" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +type testRoute struct{ ID, DeviceID, ProfileToken, Path, Desired, Actual string } + +func (testRoute) TableName() string { return "sense_media_routes" } + +type testDevice struct{ ID, Name, Location string } + +func (testDevice) TableName() string { return "sense_devices" } + +func testService(t *testing.T) (*Service, *gorm.DB) { + t.Helper() + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&Shard{}, &Assignment{}, &testRoute{}, &testDevice{}); err != nil { + t.Fatal(err) + } + now := time.Date(2026, 8, 28, 8, 0, 0, 0, time.UTC) + service := NewService(db, func(context.Context, string) error { return nil }) + service.Now = func() time.Time { return now } + return service, db +} + +func TestStableAssignmentUsesConfiguredCapacity(t *testing.T) { + service, _ := testService(t) + ctx := context.Background() + if err := service.SyncSpecs(ctx, []Spec{{ID: "alpha", Name: "A", Mode: "external", ControlAPI: "http://127.0.0.1:9997", Capacity: 3}, {ID: "beta", Name: "B", Mode: "external", ControlAPI: "http://127.0.0.1:19997", Capacity: 2}}); err != nil { + t.Fatal(err) + } + if err := service.RefreshAll(ctx); err != nil { + t.Fatal(err) + } + first, err := service.EnsureAssignment(ctx, "route-1") + if err != nil { + t.Fatal(err) + } + second, err := service.EnsureAssignment(ctx, "route-1") + if err != nil { + t.Fatal(err) + } + if first.ID != second.ID { + t.Fatalf("assignment changed from %s to %s", first.ID, second.ID) + } + for _, id := range []string{"route-2", "route-3", "route-4", "route-5"} { + if _, err = service.EnsureAssignment(ctx, id); err != nil { + t.Fatalf("assign %s: %v", id, err) + } + } + if _, err = service.EnsureAssignment(ctx, "route-6"); !errors.Is(err, ErrNoCapacity) { + t.Fatalf("expected configured capacity error, got %v", err) + } +} + +func TestFailureKeepsOwnershipAndLocatesImpact(t *testing.T) { + service, db := testService(t) + ctx := context.Background() + if err := service.SyncSpecs(ctx, []Spec{{ID: "failed", Name: "故障分片", Mode: "external", ControlAPI: "http://127.0.0.1:9997", Capacity: 7}}); err != nil { + t.Fatal(err) + } + if err := service.RefreshAll(ctx); err != nil { + t.Fatal(err) + } + if err := db.Create(&testDevice{ID: "device-1", Name: "东门摄像机", Location: "东门"}).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&testRoute{ID: "device-1:main", DeviceID: "device-1", ProfileToken: "main", Path: "sense_abc", Desired: "running", Actual: "ready"}).Error; err != nil { + t.Fatal(err) + } + if _, err := service.EnsureAssignment(ctx, "device-1:main"); err != nil { + t.Fatal(err) + } + service.Probe = func(context.Context, string) error { return errors.New("offline") } + if err := service.RefreshAll(ctx); err != nil { + t.Fatal(err) + } + detail, err := service.Get(ctx, "failed") + if err != nil { + t.Fatal(err) + } + if detail.Status != StatusFailed || len(detail.Impact) != 1 || detail.Impact[0].DeviceName != "东门摄像机" || detail.Impact[0].ProfileToken != "main" { + t.Fatalf("unexpected detail: %+v", detail) + } + resolved, err := service.ResolveRoute(ctx, "device-1:main") + if err != nil || resolved.ID != "failed" { + t.Fatalf("failure rewrote ownership: %+v %v", resolved, err) + } +} + +func TestMigrationPreflightIsReadOnly(t *testing.T) { + service, db := testService(t) + ctx := context.Background() + if err := service.SyncSpecs(ctx, []Spec{{ID: "source", Name: "源", Mode: "external", ControlAPI: "http://127.0.0.1:9997", Capacity: 1}, {ID: "target", Name: "目标", Mode: "external", ControlAPI: "http://127.0.0.1:19997", Capacity: 4}}); err != nil { + t.Fatal(err) + } + if err := service.RefreshAll(ctx); err != nil { + t.Fatal(err) + } + if err := db.Create(&testRoute{ID: "route", DeviceID: "d", ProfileToken: "p", Path: "path", Desired: "running", Actual: "ready"}).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&Assignment{RouteID: "route", ShardID: "source", Algorithm: "rendezvous-v1", CreatedAt: service.now(), UpdatedAt: service.now()}).Error; err != nil { + t.Fatal(err) + } + before := Assignment{} + db.First(&before, "route_id = ?", "route") + result, err := service.Preflight(ctx, "source") + if err != nil { + t.Fatal(err) + } + after := Assignment{} + db.First(&after, "route_id = ?", "route") + if !result.Ready || result.ExecutionAuthorized || !result.RequiresIssue || result.TargetShardID != "target" { + t.Fatalf("unexpected preflight: %+v", result) + } + if before.ShardID != after.ShardID || !before.UpdatedAt.Equal(after.UpdatedAt) { + t.Fatalf("preflight mutated assignment: before=%+v after=%+v", before, after) + } +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard.go b/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard.go new file mode 100644 index 0000000..bb14e94 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard.go @@ -0,0 +1,57 @@ +package version + +import ( + "fmt" + "runtime" + + "gorm.io/gorm" + "gorm.io/gorm/clause" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" + "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), migrateSenseMediaShard) +} + +func migrateSenseMediaShard(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := tx.AutoMigrate(&media_shard.Shard{}, &media_shard.Assignment{}); err != nil { + return err + } + root, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: senseLayoutMenuName, Title: "视频感知", Icon: "video-camera", Path: "/sense", MenuType: "M", Component: "Layout", Sort: 5, Visible: "0", IsFrame: "1"}) + if err != nil { + return err + } + page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaShard", Title: "媒体分片", Icon: "connection", Path: "media-shard", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:media-shard:list", ParentId: root.MenuId, Component: "/sense/media-shard/index", Sort: 10, Visible: "0", IsFrame: "1"}) + if err != nil { + return err + } + detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaShardDetail", Title: "查看分片详情", MenuType: "F", Action: "GET", Permission: "sense:media-shard:detail", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 1, Visible: "1", IsFrame: "1"}) + if err != nil { + return err + } + preflight, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaShardPreflight", Title: "迁移预检", MenuType: "F", Action: "GET", Permission: "sense:media-shard:preflight", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 2, Visible: "1", IsFrame: "1"}) + if err != nil { + return err + } + for _, role := range []string{"implementation_operator", "site_admin", "viewer"} { + if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, detail, preflight}); err != nil { + return err + } + for _, policy := range [][2]string{{"/api/v1/media-shards", "GET"}, {"/api/v1/media-shards/:id", "GET"}, {"/api/v1/media-shards/:id/migration-preflight", "GET"}} { + 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 + } + } + } + if err = rebuildSenseMenuPaths(tx, root.MenuId, "/0"); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard_test.go b/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard_test.go new file mode 100644 index 0000000..23b8800 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082814000_media_shard_test.go @@ -0,0 +1,86 @@ +package version + +import ( + "os" + "testing" + + "gorm.io/driver/postgres" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" + migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models" + common "git.ilapage.cn/ila/yovision/Sense/server/common/models" +) + +func TestMediaShardMigrationAddsReadOnlySurfaceWithoutFixtures(t *testing.T) { + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil { + t.Fatal(err) + } + 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) + } + } + const version = "2026082814000_media_shard.go" + if err = migrateSenseMediaShard(db, version); err != nil { + t.Fatal(err) + } + if !db.Migrator().HasTable(&media_shard.Shard{}) || !db.Migrator().HasTable(&media_shard.Assignment{}) { + t.Fatal("media shard tables missing") + } + var menus, reads, writes, fixtures, applied int64 + db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseMediaShard%").Count(&menus) + db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 = ?", "/api/v1/media-shards%", "GET").Count(&reads) + db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 <> ?", "/api/v1/media-shards%", "GET").Count(&writes) + db.Model(&media_shard.Shard{}).Count(&fixtures) + db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied) + if menus != 3 || reads != 9 || writes != 0 || fixtures != 0 || applied != 1 { + t.Fatalf("menus=%d reads=%d writes=%d fixtures=%d applied=%d", menus, reads, writes, fixtures, applied) + } +} + +func TestMediaShardMigrationOnPostgres(t *testing.T) { + dsn := os.Getenv("SENSE_MEDIA_SHARD_MIGRATION_TEST_DATABASE_URL") + if dsn == "" { + t.Skip("set SENSE_MEDIA_SHARD_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) + } + const schema = "sense_media_shard_77_test" + if err = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error; err != nil { + t.Fatal(err) + } + if err = db.Exec("CREATE SCHEMA " + schema).Error; err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error }) + sqlDB, err := db.DB() + if err != nil { + t.Fatal(err) + } + sqlDB.SetMaxOpenConns(1) + if err = db.Exec("SET search_path TO " + schema).Error; err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil { + t.Fatal(err) + } + 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 = migrateSenseMediaShard(db, "2026082814000_media_shard.go"); err != nil { + t.Fatal(err) + } + if !db.Migrator().HasTable(&media_shard.Shard{}) || !db.Migrator().HasTable(&media_shard.Assignment{}) { + t.Fatal("media shard tables missing") + } +} diff --git a/Sense/server/config/credential.env.example b/Sense/server/config/credential.env.example index 1489952..fc20b4e 100644 --- a/Sense/server/config/credential.env.example +++ b/Sense/server/config/credential.env.example @@ -9,3 +9,5 @@ SENSE_ONVIF_ALLOWED_CIDRS= SENSE_MEDIAMTX_BINARY= SENSE_MEDIAMTX_CONFIG= SENSE_MEDIAMTX_API=http://127.0.0.1:9997 +SENSE_MEDIAMTX_CAPACITY= +SENSE_MEDIAMTX_SHARDS_FILE= diff --git a/Sense/server/tests/media/integration_test.go b/Sense/server/tests/media/integration_test.go index de48959..f1bcaa0 100644 --- a/Sense/server/tests/media/integration_test.go +++ b/Sense/server/tests/media/integration_test.go @@ -10,6 +10,9 @@ import ( "time" "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media" + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard" + "gorm.io/driver/sqlite" + "gorm.io/gorm" ) func freeAddress(t *testing.T) string { @@ -25,6 +28,81 @@ func freeAddress(t *testing.T) string { return address } +func startTestMediaMTX(t *testing.T, binary string) string { + t.Helper() + apiAddress, rtspAddress := freeAddress(t), freeAddress(t) + configPath := filepath.Join(t.TempDir(), "mediamtx.yml") + config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtspTransports: [tcp]\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nmoq: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress) + if err := os.WriteFile(configPath, []byte(config), 0o600); 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) + }) + controller, err := media.NewHTTPController("http://" + apiAddress) + if err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(8 * time.Second) + for { + if err = controller.Health(context.Background()); err == nil { + break + } + if time.Now().After(deadline) { + t.Fatalf("MediaMTX did not become ready: %v", err) + } + time.Sleep(100 * time.Millisecond) + } + return "http://" + apiAddress +} + +func TestTwoRealMediaMTXShardsAreIndependentlyUsable(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") + } + first, second := startTestMediaMTX(t, binary), startTestMediaMTX(t, binary) + db, err := gorm.Open(sqlite.Open("file:real_media_shards?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&media_shard.Shard{}, &media_shard.Assignment{}); err != nil { + t.Fatal(err) + } + service := media_shard.NewService(db, func(ctx context.Context, endpoint string) error { + controller, buildErr := media.NewHTTPController(endpoint) + if buildErr != nil { + return buildErr + } + return controller.Health(ctx) + }) + specs := []media_shard.Spec{{ID: "one", Name: "测试分片一", Mode: "external", ControlAPI: first, Capacity: 2}, {ID: "two", Name: "测试分片二", Mode: "external", ControlAPI: second, Capacity: 2}} + if err = service.SyncSpecs(context.Background(), specs); err != nil { + t.Fatal(err) + } + if err = service.RefreshAll(context.Background()); err != nil { + t.Fatal(err) + } + for _, routeID := range []string{"route-a", "route-b", "route-c", "route-d"} { + if _, err = service.EnsureAssignment(context.Background(), routeID); err != nil { + t.Fatalf("assign %s: %v", routeID, err) + } + } + result, err := service.List(context.Background()) + if err != nil { + t.Fatal(err) + } + if result.Summary.Available != 2 || result.Summary.Assigned != 4 || result.Summary.Configured != 4 { + t.Fatalf("unexpected real shard summary: %+v", result.Summary) + } +} + func TestRealMediaMTXControlLifecycle(t *testing.T) { binary := os.Getenv("SENSE_MEDIAMTX_TEST_BINARY") if binary == "" { @@ -32,7 +110,7 @@ func TestRealMediaMTXControlLifecycle(t *testing.T) { } 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) + config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtspTransports: [tcp]\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nmoq: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress) if err := os.WriteFile(configPath, []byte(config), 0o600); err != nil { t.Fatal(err) } diff --git a/Sense/tests/package/run-tests.ps1 b/Sense/tests/package/run-tests.ps1 index 814d07f..0e85177 100644 --- a/Sense/tests/package/run-tests.ps1 +++ b/Sense/tests/package/run-tests.ps1 @@ -34,6 +34,7 @@ try { [IO.File]::WriteAllText((Join-Path $temporary 'web\js\runtime.fixture.js'), 'fixture', (New-Object Text.UTF8Encoding($false))) [IO.File]::WriteAllText((Join-Path $temporary 'bin\mediamtx.exe'), 'fixture', (New-Object Text.UTF8Encoding($false))) [IO.File]::WriteAllText((Join-Path $temporary 'config\mediamtx.yml'), 'api: true', (New-Object Text.UTF8Encoding($false))) + [IO.File]::WriteAllText((Join-Path $temporary 'config\mediamtx-shards.json'), '[]', (New-Object Text.UTF8Encoding($false))) $webAssetAudit = Join-Path $senseRoot 'scripts\build\assert-web-assets.ps1' & $webAssetAudit -WebRoot (Join-Path $temporary 'web') @@ -55,6 +56,7 @@ try { 'SENSE_MEDIAMTX_BINARY=bin\mediamtx.exe', 'SENSE_MEDIAMTX_CONFIG=config\mediamtx.yml', 'SENSE_MEDIAMTX_API=http://127.0.0.1:9997', + 'SENSE_MEDIAMTX_SHARDS_FILE=config\mediamtx-shards.json', 'SENSE_WEB_ROOT=web' ) -join "`n" [IO.File]::WriteAllText((Join-Path $temporary 'config\sense.env'), $envText, (New-Object Text.UTF8Encoding($false))) @@ -68,6 +70,7 @@ try { Assert-True $generatedSettings.Contains('timeout: 2592000') 'generated runtime settings must keep login valid for 30 days' Assert-True ((Get-SenseEnvironmentValue -Name 'SENSE_DATABASE_URL').Contains('p#&;=x')) 'special characters must survive env parsing' Assert-True ($generatedSettings.Contains('password=p#\u0026;=x')) 'special characters must be safely JSON-escaped in YAML' + Assert-Equal (Join-Path $temporary 'config\mediamtx-shards.json') (Get-SenseEnvironmentValue -Name 'SENSE_MEDIAMTX_SHARDS_FILE') 'relative shard config path must resolve inside package root' Assert-True (-not ((Get-SenseDatabaseInfo (Get-SenseEnvironmentValue -Name 'SENSE_DATABASE_URL')).Sanitized.Contains('password='))) 'PostgreSQL tool arguments must not contain password' [Environment]::SetEnvironmentVariable('SENSE_PORT', $null, 'Process') diff --git a/Sense/ui/src/api/sense/media-shard.js b/Sense/ui/src/api/sense/media-shard.js new file mode 100644 index 0000000..a7c28a2 --- /dev/null +++ b/Sense/ui/src/api/sense/media-shard.js @@ -0,0 +1,13 @@ +import request from '@/utils/request' + +export function listMediaShards() { + return request({ url: '/api/v1/media-shards', method: 'get' }) +} + +export function getMediaShard(id) { + return request({ url: `/api/v1/media-shards/${encodeURIComponent(id)}`, method: 'get' }) +} + +export function preflightMediaShard(id) { + return request({ url: `/api/v1/media-shards/${encodeURIComponent(id)}/migration-preflight`, method: 'get' }) +} diff --git a/Sense/ui/src/views/sense/media-shard/index.vue b/Sense/ui/src/views/sense/media-shard/index.vue new file mode 100644 index 0000000..2f6dce9 --- /dev/null +++ b/Sense/ui/src/views/sense/media-shard/index.vue @@ -0,0 +1,65 @@ + + + + + diff --git a/Sense/ui/src/views/sense/media-shard/mediaShardState.js b/Sense/ui/src/views/sense/media-shard/mediaShardState.js new file mode 100644 index 0000000..fa4772c --- /dev/null +++ b/Sense/ui/src/views/sense/media-shard/mediaShardState.js @@ -0,0 +1,6 @@ +const labels = { running: '运行正常', failed: '运行异常', unknown: '等待探测', disabled: '已停用', managed: 'Sense 托管', external: '外部运行' } +export function stateLabel(value) { return labels[value] || value || '未知' } +export function stateType(value) { return ({ running: 'success', failed: 'danger', unknown: 'warning', disabled: 'info' })[value] || 'info' } +export function usagePercent(item) { const capacity = Number(item?.capacity) || 0; return capacity > 0 ? Math.min(100, Math.round((Number(item.assignedPaths) || 0) * 100 / capacity)) : 0 } +export function freshnessText(item) { if (!item?.lastProbeAt) return '尚未探测'; return item.stale ? '状态已陈旧' : '刚刚更新' } +export function preflightConclusion(result) { if (!result) return ''; if (!result.ready) return '当前不满足迁移前置条件'; return '预检通过,但仍需单独建立高风险工单并人工确认' } diff --git a/Sense/ui/tests/unit/sense/mediaShardApi.spec.js b/Sense/ui/tests/unit/sense/mediaShardApi.spec.js new file mode 100644 index 0000000..0d64423 --- /dev/null +++ b/Sense/ui/tests/unit/sense/mediaShardApi.spec.js @@ -0,0 +1,8 @@ +jest.mock('@/utils/request', () => jest.fn(config => config)) +import request from '@/utils/request' +import { getMediaShard, listMediaShards, preflightMediaShard } from '@/api/sense/media-shard' + +describe('Sense media shard API', () => { + beforeEach(() => request.mockClear()) + test('exposes only read operations and safely encodes ids', () => { expect(listMediaShards()).toEqual({ url: '/api/v1/media-shards', method: 'get' }); expect(getMediaShard('a/b')).toEqual({ url: '/api/v1/media-shards/a%2Fb', method: 'get' }); expect(preflightMediaShard('a b')).toEqual({ url: '/api/v1/media-shards/a%20b/migration-preflight', method: 'get' }); expect(request).toHaveBeenCalledTimes(3) }) +}) diff --git a/Sense/ui/tests/unit/sense/mediaShardState.spec.js b/Sense/ui/tests/unit/sense/mediaShardState.spec.js new file mode 100644 index 0000000..af26723 --- /dev/null +++ b/Sense/ui/tests/unit/sense/mediaShardState.spec.js @@ -0,0 +1,7 @@ +import { freshnessText, preflightConclusion, stateLabel, stateType, usagePercent } from '@/views/sense/media-shard/mediaShardState' + +describe('Sense media shard presentation state', () => { + test('uses clear operator-facing labels and non-color-only status', () => { expect(stateLabel('failed')).toBe('运行异常'); expect(stateType('failed')).toBe('danger'); expect(stateLabel('external')).toBe('外部运行') }) + test('derives usage from configured capacity instead of a fixed tier', () => { expect(usagePercent({ assignedPaths: 7, capacity: 23 })).toBe(30); expect(usagePercent({ assignedPaths: 99, capacity: 0 })).toBe(0) }) + test('makes stale data and read-only preflight explicit', () => { expect(freshnessText({ lastProbeAt: '2026-08-28T00:00:00Z', stale: true })).toBe('状态已陈旧'); expect(preflightConclusion({ ready: true })).toContain('高风险工单'); expect(preflightConclusion({ ready: false })).toContain('不满足') }) +}) diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 3bb4a70..0f5b809 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: 5c2b5df080b160c1174be3e935fffba420885bff -synchronized_at: 2026-08-28T06:14:33Z +wiki_revision: 578ddbaae3d7e846037d958085e40609cd398bef +synchronized_at: 2026-08-28T08:02:32Z # 架构与代码地图 @@ -225,3 +225,14 @@ PostgreSQL 表 `sense_provisioning_batches` 保存幂等键、配额快照和汇 - go-admin-ui 页面:`Sense/ui/src/views/sense/edge-node/`;API 封装:`Sense/ui/src/api/sense/edge-node.js`。页面复用 BasicLayout、Element Plus 表格/分页/Dialog/Tag/Alert 和权限指令。 - 内部 `ApplyHeartbeat` 是未来 Sense 适配器的投影写入口,本工单不把它暴露为 HTTP API,也不定义跨项目契约。 + + +## Sense MediaMTX 分片代码路径 + +- 分片模型、容量分配、健康探测、故障影响和只读迁移预检位于 `Sense/server/app/sense/media_shard/`;既有媒体路由仍由 `Sense/server/app/sense/media/` 管理。 +- PostgreSQL 表 `sense_media_shards` 保存运行配置投影和健康状态,`sense_media_shard_assignments` 保存路径到分片的稳定归属;`sense_media_routes` 继续作为路径状态事实源。 +- 新路径在健康且有剩余配置容量的分片间使用 `rendezvous-v1` 确定性分配。已有归属在分片故障时保持不变,不自动迁移。 +- 主分片继续复用既有 MediaMTX Supervisor;额外分片通过仓库外 JSON 配置接入本机回环 Control API,并按 external 实例管理,不由 Sense 启停。 +- GoAdmin 路由位于 `Sense/server/app/admin/router/sense_media_shard.go`,只开放列表、详情和迁移预检三个 GET 接口。go-admin-ui 页面位于 `Sense/ui/src/views/sense/media-shard/`,复用 BasicLayout、Element Plus 表格、进度、Dialog、Tag、Alert 和权限指令。 +- 迁移 `2026082814000_media_shard.go` 建表、注册动态菜单并为 implementation_operator、site_admin、viewer 建立只读权限;生产迁移和启动不写入合成分片。 + diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 6ff9036..9ecc800 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: 8f81b925af671209341f1282e128c7f048bca8ec -synchronized_at: 2026-08-28T06:14:44Z +wiki_revision: 999eb1aee3558ff75cfa929acac77841af20b742 +synchronized_at: 2026-08-28T08:02:42Z # 业务规则与术语 @@ -199,3 +199,14 @@ synchronized_at: 2026-08-28T06:14:44Z - 状态转换写入不可变事件;列表和详情读取写入脱敏操作审计,不记录凭据或原始心跳载荷。 - 合成夹具必须显式调用,生产迁移和启动不得自动写入节点。 + + +## Sense MediaMTX 分片规则 + +- **媒体分片**是一个独立 MediaMTX Control API 及其配置容量的 Sense 本地投影,不是 Brain 分片,也不是跨产品共享节点身份。 +- 主分片容量优先读取显式 `SENSE_MEDIAMTX_CAPACITY`;未设置时读取 Sense PostgreSQL 的当前统一配额。额外分片容量逐项来自仓库外 JSON 配置。16 和 128 均不得作为分配算法硬上限。 +- 新媒体路径只在状态为运行正常或等待首次探测、且仍有配置容量的分片中确定性分配;同一路径重复处理必须保持原归属。 +- 分片异常、停用或状态陈旧都不得自动改写已有路径归属。详情必须能定位受影响的设备、位置、Profile 和媒体路径,同时不返回 Control API、RTSP URI 或凭据。 +- 跨分片迁移预检是只读操作,只检查源状态、全部受影响路径和候选目标容量;即使预检通过也不授予执行权限。实际迁移必须另建高风险工单并取得人工确认。 +- 额外分片只能配置为 external,Control API 必须使用无用户信息、无查询参数、无路径的本机回环 HTTP 地址。Sense 不停止外部实例。 + diff --git a/docs/04-local-development-and-verification.md b/docs/04-local-development-and-verification.md index 35fa9e4..9f0c4f7 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: cfec18809b3cb0d938003102c9caa5e6d6f4e818 -synchronized_at: 2026-08-28T06:25:12Z +wiki_revision: 05435d53528271a866c525655486689afc198762 +synchronized_at: 2026-08-28T08:02:52Z # 本地开发与验证 @@ -477,3 +477,41 @@ go test ./cmd/migrate/migration/version -run TestEdgeNodeMigrationOnPostgres -co 验证应覆盖在线、超过 90 秒的离线、心跳恢复但通道未收敛、最后已知状态与陈旧时长、负载和回填队列、不可变状态事件、只读权限、脱敏读取审计,以及 Brain/Bell 不运行时的独立性。合成节点只能由测试或显式开发调用;没有隔离 PostgreSQL 连接时必须记录迁移真库测试未执行。 + + +### Sense MediaMTX 分片验证 + +定向与全量验证: + +```powershell +cd Sense/server +go test -race ./app/sense/media ./app/sense/media_shard ./cmd/migrate/migration/version +go test ./... +go vet ./... + +cd ../ui +corepack pnpm@9.15.1 lint +corepack pnpm@9.15.1 test:unit +corepack pnpm@9.15.1 build:prod + +cd ../.. +powershell.exe -NoProfile -File .\Sense\tests\package\run-tests.ps1 +``` + +使用真实本机 MediaMTX 二进制验证两个独立测试实例;测试会分配临时回环端口并在结束时停止其自有进程: + +```powershell +cd Sense/server +$env:SENSE_MEDIAMTX_TEST_BINARY = '' +go test ./tests/media -run 'TestTwoRealMediaMTXShardsAreIndependentlyUsable|TestRealMediaMTXControlLifecycle' -count=1 -v +``` + +PostgreSQL 迁移测试仅使用专用隔离连接: + +```powershell +$env:SENSE_MEDIA_SHARD_MIGRATION_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>' +go test ./cmd/migrate/migration/version -run TestMediaShardMigrationOnPostgres -count=1 -v +``` + +验证应覆盖任意配置容量、稳定重复分配、容量耗尽、分片故障后归属不变、设备/Profile/路径影响范围、只读迁移预检、只读 RBAC、Control API 不出现在响应,以及 Brain/Bell 均不运行。没有专用 PostgreSQL 连接时必须记录真库迁移测试未执行。 + diff --git a/docs/delivery/deployment-and-operations.md b/docs/delivery/deployment-and-operations.md index f96afd0..79defa2 100644 --- a/docs/delivery/deployment-and-operations.md +++ b/docs/delivery/deployment-and-operations.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Deployment-and-Operations wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Deployment-and-Operations.- -wiki_revision: 34b85f1c3fc2c3e0e7fa0de3df056940fcb71333 -synchronized_at: 2026-08-28T06:17:09Z +wiki_revision: ce246054849b5b797dc4da0a55116238a4c00d7e +synchronized_at: 2026-08-28T08:04:52Z # YoVision 部署与运维 @@ -100,3 +100,15 @@ Sense\start_sense.bat 节点超过 90 秒没有心跳即显示离线。离线时页面继续显示最后已知的隧道、视频、负载和回填状态,并标明陈旧时长,这不表示通道仍实时可用。心跳恢复但通道未收敛时显示“恢复中”,应分别检查控制隧道、视频数据面和回填队列。当前版本只提供 Sense 内部投影入口,不包含外部心跳接入协议或跨节点回填执行;没有节点数据时不会自动生成演示节点。 + + +## Sense MediaMTX 分片配置与排错 + +主 MediaMTX 继续使用 `SENSE_MEDIAMTX_MODE`、`SENSE_MEDIAMTX_BINARY`、`SENSE_MEDIAMTX_CONFIG` 和 `SENSE_MEDIAMTX_API`。可选 `SENSE_MEDIAMTX_CAPACITY` 为正整数;留空时使用数据库中的当前统一配额。 + +额外分片通过 `SENSE_MEDIAMTX_SHARDS_FILE` 指向仓库外 JSON 文件。相对路径按 Windows 交付包根目录解析;示例结构见包内 `config\mediamtx-shards.example.json`。每项必须包含唯一 id、名称、本机回环 Control API 和正容量,mode 只能为 external。文件不得包含摄像头账号、密码、数据库连接或其他秘密。 + +额外 MediaMTX 进程由部署方独立启动并为各实例分配不冲突的 API、RTSP 及其他监听端口;Sense 只做健康探测和路径管理,不启停这些外部进程。升级后看不到“媒体分片”时,确认迁移 `2026082814000_media_shard.go` 成功并重新登录刷新动态菜单。 + +页面“运行异常”表示最近一次 Control API 探测失败;“状态已陈旧”表示运行循环已超过 30 秒没有更新探测结果。故障时先在详情定位设备/Profile/路径,不要手工改数据库归属。迁移预检不会执行迁移;任何实际跨分片迁移都必须另建高风险工单和回退方案。 + -- 2.34.1