Compare commits

...
Author SHA1 Message Date
QiuSW d681fd1345 feat: 实现 MediaMTX 分片分配与故障范围可观察 (#77) 2026-08-28 16:09:52 +08:00
ila bdd78523b5 Merge PR #120: Sense 边缘节点资产与离线状态管理
Implements #76; remains pending user acceptance.
2026-08-28 14:29:55 +08:00
QiuSW 2db49fc955 feat: 实现边缘节点资产与离线状态管理 (#76) 2026-08-28 14:25:32 +08:00
ila a9fd8e8a47 Merge PR #119: Sense 容量与配额安全闸
合入 dev,工单 #75 保持待验收。
2026-08-28 11:55:51 +08:00
QiuSW 12f5419f21 feat: 实现容量与配额安全闸 (#75) 2026-08-28 11:55:11 +08:00
ila debf662138 Merge PR #118: Sense 运维中心 (#74)
将 #74 实现合入 dev;工单保持待验收。
2026-08-28 11:05:11 +08:00
QiuSW 184f101a86 feat: 实现 Sense 运维中心 (#74) 2026-08-28 11:02:34 +08:00
ila 4157d90e9f Merge PR #117: Sense 本地事件候选查看
关联 #73;合入 dev,工单保持待验收。
2026-08-28 10:26:58 +08:00
83 changed files with 4658 additions and 53 deletions
@@ -0,0 +1,9 @@
[
{
"id": "secondary",
"name": "备用媒体分片",
"mode": "external",
"controlAPI": "http://127.0.0.1:19997",
"capacity": 24
}
]
+2
View File
@@ -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=
+4
View File
@@ -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=
+1
View File
@@ -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')
+1 -1
View File
@@ -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) {
+7 -1
View File
@@ -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'
@@ -22,6 +22,7 @@ func registerSenseDeviceRouter(v1 *gin.RouterGroup, authMiddleware *jwt.GinJWTMi
r.POST("", api.Insert)
r.PUT("/:id", api.Update)
r.PUT("/:id/disable", api.Disable)
r.PUT("/:id/enable", api.Enable)
r.PUT("/:id/credentials", api.UpdateCredentials)
}
}
@@ -0,0 +1,19 @@
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/edge_node"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseEdgeNodeRouter) }
func registerSenseEdgeNodeRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &edge_node.API{}
r := v1.Group("/edge-nodes").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("", api.List)
r.GET("/:id", api.Get)
}
@@ -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)
}
@@ -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/operations"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseOperationsRouter) }
func registerSenseOperationsRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &operations.API{}
r := v1.Group("/operations").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("", api.List)
r.GET("/:id", api.Get)
r.POST("/:id/retry", api.Retry)
}
@@ -0,0 +1,31 @@
package router
import (
"net/http"
"testing"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
)
func TestSenseOperationsRoutes(t *testing.T) {
gin.SetMode(gin.TestMode)
engine := gin.New()
registerSenseOperationsRouter(engine.Group("/api/v1"), &jwt.GinJWTMiddleware{})
wanted := map[string]bool{
http.MethodGet + " /api/v1/operations": false,
http.MethodGet + " /api/v1/operations/:id": false,
http.MethodPost + " /api/v1/operations/:id/retry": false,
}
for _, route := range engine.Routes() {
key := route.Method + " " + route.Path
if _, ok := wanted[key]; ok {
wanted[key] = true
}
}
for route, found := range wanted {
if !found {
t.Fatalf("route not registered: %s", route)
}
}
}
@@ -0,0 +1,19 @@
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/quota"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseQuotaRouter) }
func registerSenseQuotaRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &quota.API{}
r := v1.Group("/quota").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("", api.Get)
r.PUT("", api.Update)
}
@@ -14,6 +14,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
deviceService "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service/dto"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
)
type Device struct{ api.Api }
@@ -106,6 +107,25 @@ func (e Device) Disable(c *gin.Context) {
e.OK(response, "设备已停用")
}
func (e Device) Enable(c *gin.Context) {
service := deviceService.Device{}
if err := e.MakeContext(c).MakeOrm().MakeService(&service.Service).Errors; err != nil {
e.Error(http.StatusInternalServerError, err, "服务初始化失败")
return
}
req := dto.EnableReq{ID: c.Param("id"), UpdateBy: user.GetUserId(c)}
if err := bindStrictJSON(c, &req); err != nil {
e.Error(http.StatusBadRequest, err, "请求内容格式不正确")
return
}
var response dto.DeviceResponse
if err := service.Enable(&req, &response); err != nil {
e.writeServiceError(err)
return
}
e.OK(response, "设备已启用,请重新完成视频接入验证")
}
func (e Device) UpdateCredentials(c *gin.Context) {
service := deviceService.Device{}
if err := e.MakeContext(c).MakeOrm().MakeService(&service.Service).Errors; err != nil {
@@ -133,6 +153,10 @@ func (e Device) writeServiceError(err error) {
e.Error(http.StatusNotFound, err, err.Error())
case errors.Is(err, deviceService.ErrVersionConflict):
e.Error(http.StatusConflict, err, err.Error())
case errors.Is(err, quota.ErrExceeded):
e.Error(http.StatusConflict, err, "当前配额已满,无法新增或启用设备")
case errors.Is(err, quota.ErrUnavailable):
e.Error(http.StatusServiceUnavailable, err, "配额配置不可读取,已拒绝新增或启用设备")
case errors.Is(err, credential.ErrKeyUnavailable):
e.Error(http.StatusServiceUnavailable, err, "摄像头凭据安全配置不可用")
default:
@@ -15,6 +15,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service/dto"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
)
var (
@@ -108,7 +109,12 @@ func (e *Device) Insert(req *dto.CreateReq, response *dto.DeviceResponse) error
}
model.CreateBy = req.CreateBy
model.UpdateBy = req.CreateBy
if err = e.Orm.Create(&model).Error; err != nil {
if err = quota.WithAvailableSlot(e.Orm, func(tx *gorm.DB) error {
return tx.Create(&model).Error
}); err != nil {
if errors.Is(err, quota.ErrUnavailable) || errors.Is(err, quota.ErrExceeded) {
return err
}
return fmt.Errorf("create device: %w", err)
}
return e.Get(model.ID, response)
@@ -155,6 +161,35 @@ func (e *Device) Disable(req *dto.DisableReq, response *dto.DeviceResponse) erro
return e.Get(req.ID, response)
}
func (e *Device) Enable(req *dto.EnableReq, response *dto.DeviceResponse) error {
if req.Version < 1 {
return ErrInvalidDevice
}
err := quota.WithAvailableSlot(e.Orm, func(tx *gorm.DB) error {
result := tx.Model(&models.Device{}).
Where("id = ? AND version = ? AND status = ?", req.ID, req.Version, models.StatusDisabled).
Updates(map[string]any{
"status": models.StatusPending,
"version": req.Version + 1, "update_by": req.UpdateBy, "updated_at": time.Now().UTC(),
})
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return e.notFoundOrConflictWith(tx, req.ID)
}
return nil
})
if err != nil {
if errors.Is(err, quota.ErrUnavailable) || errors.Is(err, quota.ErrExceeded) ||
errors.Is(err, ErrDeviceNotFound) || errors.Is(err, ErrVersionConflict) {
return err
}
return fmt.Errorf("enable device: %w", err)
}
return e.Get(req.ID, response)
}
func (e *Device) UpdateCredentials(req *dto.CredentialUpdateReq, response *dto.DeviceResponse) error {
if req.Version < 1 || strings.TrimSpace(req.ONVIFUsername) == "" || req.ONVIFPassword == "" {
return ErrInvalidDevice
@@ -13,6 +13,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service/dto"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
)
func testDeviceService(t *testing.T) (*Device, *gorm.DB) {
@@ -21,7 +22,10 @@ func testDeviceService(t *testing.T) (*Device, *gorm.DB) {
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&models.Device{}, &credential.DeviceCredential{}); err != nil {
if err = db.AutoMigrate(&models.Device{}, &credential.DeviceCredential{}, &quota.Setting{}, &quota.Change{}); err != nil {
t.Fatal(err)
}
if err = db.Create(&quota.Setting{ID: quota.SettingID, Limit: 32, Source: "test", Version: 1}).Error; err != nil {
t.Fatal(err)
}
key := make([]byte, 32)
@@ -36,6 +40,38 @@ func testDeviceService(t *testing.T) (*Device, *gorm.DB) {
return service, db
}
func TestCreateAndEnableUseQuotaSafetyGate(t *testing.T) {
service, db := testDeviceService(t)
if err := db.Model(&quota.Setting{}).Where("id = ?", quota.SettingID).Update("limit", 1).Error; err != nil {
t.Fatal(err)
}
var first dto.DeviceResponse
if err := service.Insert(&dto.CreateReq{Name: "第一路", Modality: "video", Capabilities: []string{"video"}}, &first); err != nil {
t.Fatal(err)
}
var rejected dto.DeviceResponse
if err := service.Insert(&dto.CreateReq{Name: "第二路", Modality: "video", Capabilities: []string{"video"}}, &rejected); err != quota.ErrExceeded {
t.Fatalf("second create error=%v", err)
}
var disabled dto.DeviceResponse
if err := service.Disable(&dto.DisableReq{ID: first.ID, Version: first.Version}, &disabled); err != nil {
t.Fatal(err)
}
var second dto.DeviceResponse
if err := service.Insert(&dto.CreateReq{Name: "第二路", Modality: "video", Capabilities: []string{"video"}}, &second); err != nil {
t.Fatal(err)
}
if err := service.Enable(&dto.EnableReq{ID: disabled.ID, Version: disabled.Version}, &rejected); err != quota.ErrExceeded {
t.Fatalf("enable over quota error=%v", err)
}
if err := db.Delete(&quota.Setting{}, quota.SettingID).Error; err != nil {
t.Fatal(err)
}
if err := service.Insert(&dto.CreateReq{Name: "第三路", Modality: "video", Capabilities: []string{"video"}}, &rejected); err != quota.ErrUnavailable {
t.Fatalf("create with unreadable quota error=%v", err)
}
}
func TestDeviceLifecycleUsesAllowlistedFieldsAndOptimisticVersion(t *testing.T) {
service, _ := testDeviceService(t)
var created dto.DeviceResponse
@@ -36,6 +36,12 @@ type DisableReq struct {
UpdateBy int `json:"-"`
}
type EnableReq struct {
ID string `json:"-"`
Version int64 `json:"version"`
UpdateBy int `json:"-"`
}
type CredentialUpdateReq struct {
ID string `json:"-"`
ONVIFUsername string `json:"onvifUsername"`
+81
View File
@@ -0,0 +1,81 @@
package edge_node
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
}
func (e *API) List(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := PageRequest{}
if err = e.MakeContext(c).Bind(&request).Errors; err != nil {
e.audit(c, service, "List", auditFailure, "边缘节点查询条件格式不正确")
e.Error(http.StatusBadRequest, err, "查询条件格式不正确")
return
}
response, err := service.List(request)
if err != nil {
e.audit(c, service, "List", auditFailure, "边缘节点列表查询失败")
e.writeError(err)
return
}
e.audit(c, service, "List", auditSuccess, "读取边缘节点列表")
e.OK(response, "查询成功")
}
func (e *API) Get(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
item, err := service.Get(c.Param("id"))
if err != nil {
e.audit(c, service, "Get", auditFailure, "边缘节点详情查询失败")
e.writeError(err)
return
}
e.audit(c, service, "Get", auditSuccess, "读取边缘节点详情 "+item.ID)
e.OK(item, "查询成功")
}
func (e *API) audit(c *gin.Context, service *Service, action, status, remark string) {
if err := WriteAudit(service.Orm, 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("edge node audit failed: %s", err.Error())
}
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrInvalidQuery):
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, "边缘节点查询失败")
}
}
+21
View File
@@ -0,0 +1,21 @@
package edge_node
import (
"time"
"gorm.io/gorm"
adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
)
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: "edge_node.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()}
model.CreateBy, model.UpdateBy = input.UserID, input.UserID
return db.Create(&model).Error
}
+65
View File
@@ -0,0 +1,65 @@
package edge_node
import (
"time"
commonDto "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
)
const HeartbeatTimeout = 90 * time.Second
type PageRequest struct {
commonDto.Pagination `search:"-"`
}
type Summary struct {
Total int `json:"total"`
Online int `json:"online"`
Offline int `json:"offline"`
Recovering int `json:"recovering"`
LoadUsed int `json:"loadUsed"`
LoadCapacity int `json:"loadCapacity"`
BackfillQueueDepth int `json:"backfillQueueDepth"`
HeartbeatTimeout int64 `json:"heartbeatTimeoutSeconds"`
}
type Response struct {
ID string `json:"id"`
Name string `json:"name"`
Location string `json:"location"`
Status string `json:"status"`
RuntimeVersion string `json:"runtimeVersion"`
ModelVersion string `json:"modelVersion"`
StartedAt time.Time `json:"startedAt"`
UptimeSeconds int64 `json:"uptimeSeconds"`
LastHeartbeatAt time.Time `json:"lastHeartbeatAt"`
LastCollectedAt time.Time `json:"lastCollectedAt"`
FreshnessSeconds int64 `json:"freshnessSeconds"`
Stale bool `json:"stale"`
LoadUsed int `json:"loadUsed"`
LoadCapacity int `json:"loadCapacity"`
ControlTunnelStatus string `json:"controlTunnelStatus"`
ControlTunnelDetail string `json:"controlTunnelDetail"`
ControlLatencyMs int `json:"controlLatencyMs"`
VideoPlaneStatus string `json:"videoPlaneStatus"`
VideoPlaneDetail string `json:"videoPlaneDetail"`
ActiveStreamCount int `json:"activeStreamCount"`
BackfillStatus string `json:"backfillStatus"`
BackfillQueueDepth int `json:"backfillQueueDepth"`
RecoveryPhase string `json:"recoveryPhase"`
ProjectionVersion int64 `json:"projectionVersion"`
Events []Event `json:"events,omitempty"`
}
type PageResponse struct {
Summary Summary `json:"summary"`
List []Response `json:"list"`
Count int64 `json:"count"`
}
// HeartbeatSample is an internal Sense adapter input. It is intentionally not
// exposed as an HTTP contract by this task.
type HeartbeatSample struct {
Node
CollectedAt time.Time
}
@@ -0,0 +1,20 @@
package edge_node
import (
"time"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
// SeedSyntheticFixture is explicit test/development data. Production startup
// and migrations never call it.
func SeedSyntheticFixture(db *gorm.DB, now time.Time) error {
now = now.UTC().Truncate(time.Second)
items := []Node{
{ID: "SEN-EDGE-01", Name: "主楼边缘节点", Location: "主楼机房", RuntimeVersion: "sense-edge 0.1.0", ModelVersion: "detector-v3", StartedAt: now.Add(-7 * 24 * time.Hour), LastHeartbeatAt: now.Add(-18 * time.Second), LastCollectedAt: now.Add(-18 * time.Second), LoadUsed: 6, LoadCapacity: 16, ControlTunnelStatus: ChannelReady, ControlTunnelDetail: "控制隧道正常", ControlLatencyMs: 24, VideoPlaneStatus: ChannelReady, VideoPlaneDetail: "6 路视频正常", ActiveStreamCount: 6, BackfillStatus: BackfillIdle, ProjectionVersion: 1},
{ID: "SEN-EDGE-02", Name: "仓库边缘节点", Location: "仓库弱电间", RuntimeVersion: "sense-edge 0.1.0", ModelVersion: "detector-v3", StartedAt: now.Add(-3 * 24 * time.Hour), LastHeartbeatAt: now.Add(-4 * time.Minute), LastCollectedAt: now.Add(-4 * time.Minute), LoadUsed: 4, LoadCapacity: 16, ControlTunnelStatus: ChannelReady, ControlTunnelDetail: "最后已知:控制隧道正常", ControlLatencyMs: 31, VideoPlaneStatus: ChannelReady, VideoPlaneDetail: "最后已知:4 路视频正常", ActiveStreamCount: 4, BackfillStatus: BackfillPending, BackfillQueueDepth: 12, ProjectionVersion: 1},
{ID: "SEN-EDGE-03", Name: "南门边缘节点", Location: "南门岗亭", RuntimeVersion: "sense-edge 0.1.0", ModelVersion: "detector-v3", StartedAt: now.Add(-10 * time.Hour), LastHeartbeatAt: now.Add(-12 * time.Second), LastCollectedAt: now.Add(-12 * time.Second), LoadUsed: 2, LoadCapacity: 8, ControlTunnelStatus: ChannelReady, ControlTunnelDetail: "控制隧道已恢复", ControlLatencyMs: 46, VideoPlaneStatus: ChannelUnavailable, VideoPlaneDetail: "视频数据面恢复中", ActiveStreamCount: 0, BackfillStatus: BackfillPending, BackfillQueueDepth: 5, RecoveryPhase: "video_reconnecting", ProjectionVersion: 2},
}
return db.Clauses(clause.OnConflict{DoNothing: true}).Create(&items).Error
}
@@ -0,0 +1,61 @@
package edge_node
import (
"time"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
const (
StatusOnline = "online"
StatusOffline = "offline"
StatusRecovering = "recovering"
ChannelReady = "ready"
ChannelUnavailable = "unavailable"
BackfillIdle = "idle"
BackfillPending = "pending"
)
// Node is the Sense-owned last-known projection of one edge runtime. It is
// deliberately not a shared identity or a cross-product protocol model.
type Node struct {
ID string `gorm:"size:64;primaryKey" json:"id"`
Name string `gorm:"size:128;not null" json:"name"`
Location string `gorm:"size:256;not null;default:''" json:"location"`
RuntimeVersion string `gorm:"size:64;not null;default:''" json:"runtimeVersion"`
ModelVersion string `gorm:"size:64;not null;default:''" json:"modelVersion"`
StartedAt time.Time `gorm:"not null" json:"startedAt"`
LastHeartbeatAt time.Time `gorm:"not null;index" json:"lastHeartbeatAt"`
LastCollectedAt time.Time `gorm:"not null" json:"lastCollectedAt"`
LoadUsed int `gorm:"not null;default:0" json:"loadUsed"`
LoadCapacity int `gorm:"not null;default:0" json:"loadCapacity"`
ControlTunnelStatus string `gorm:"size:32;not null;default:'unavailable'" json:"controlTunnelStatus"`
ControlTunnelDetail string `gorm:"size:256;not null;default:''" json:"controlTunnelDetail"`
ControlLatencyMs int `gorm:"not null;default:0" json:"controlLatencyMs"`
VideoPlaneStatus string `gorm:"size:32;not null;default:'unavailable'" json:"videoPlaneStatus"`
VideoPlaneDetail string `gorm:"size:256;not null;default:''" json:"videoPlaneDetail"`
ActiveStreamCount int `gorm:"not null;default:0" json:"activeStreamCount"`
BackfillStatus string `gorm:"size:32;not null;default:'idle'" json:"backfillStatus"`
BackfillQueueDepth int `gorm:"not null;default:0" json:"backfillQueueDepth"`
RecoveryPhase string `gorm:"size:32;not null;default:''" json:"recoveryPhase"`
ProjectionVersion int64 `gorm:"not null;default:1" json:"projectionVersion"`
common.ControlBy
common.ModelTime
}
func (Node) TableName() string { return "sense_edge_nodes" }
// Event records projection state transitions without storing credentials or
// raw heartbeat payloads.
type Event struct {
ID uint64 `gorm:"primaryKey;autoIncrement" json:"id"`
NodeID string `gorm:"size:64;not null;index" json:"nodeId"`
EventType string `gorm:"size:32;not null;index" json:"eventType"`
FromStatus string `gorm:"size:32;not null;default:''" json:"fromStatus"`
ToStatus string `gorm:"size:32;not null" json:"toStatus"`
Detail string `gorm:"size:256;not null;default:''" json:"detail"`
OccurredAt time.Time `gorm:"not null;index" json:"occurredAt"`
}
func (Event) TableName() string { return "sense_edge_node_events" }
+167
View File
@@ -0,0 +1,167 @@
package edge_node
import (
"context"
"errors"
"fmt"
"strings"
"time"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
var (
ErrInvalidQuery = errors.New("边缘节点查询条件不符合要求")
ErrNotFound = errors.New("边缘节点不存在")
ErrInvalidSample = errors.New("边缘节点心跳样本不符合要求")
)
type Service struct {
coreService.Service
Now func() time.Time
}
func NewService(db *gorm.DB) *Service {
return &Service{Service: coreService.Service{Orm: db}, Now: time.Now}
}
func (s *Service) List(request PageRequest) (PageResponse, error) {
if request.GetPageSize() > 100 {
return PageResponse{}, ErrInvalidQuery
}
var count int64
if err := s.Orm.Model(&Node{}).Count(&count).Error; err != nil {
return PageResponse{}, fmt.Errorf("count edge nodes: %w", err)
}
var models []Node
if err := s.Orm.Order("name ASC, id ASC").Limit(request.GetPageSize()).Offset((request.GetPageIndex() - 1) * request.GetPageSize()).Find(&models).Error; err != nil {
return PageResponse{}, fmt.Errorf("list edge nodes: %w", err)
}
result := PageResponse{List: make([]Response, 0, len(models)), Count: count}
result.Summary.HeartbeatTimeout = int64(HeartbeatTimeout / time.Second)
var all []Node
if err := s.Orm.Find(&all).Error; err != nil {
return PageResponse{}, fmt.Errorf("summarize edge nodes: %w", err)
}
result.Summary.Total = len(all)
for _, item := range all {
status := s.status(item)
switch status {
case StatusOffline:
result.Summary.Offline++
case StatusRecovering:
result.Summary.Recovering++
default:
result.Summary.Online++
}
result.Summary.LoadUsed += item.LoadUsed
result.Summary.LoadCapacity += item.LoadCapacity
result.Summary.BackfillQueueDepth += item.BackfillQueueDepth
}
for _, item := range models {
result.List = append(result.List, s.response(item, nil))
}
return result, nil
}
func (s *Service) Get(id string) (Response, error) {
id = strings.TrimSpace(id)
if id == "" || len(id) > 64 {
return Response{}, ErrInvalidQuery
}
var model Node
if err := s.Orm.First(&model, "id = ?", id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return Response{}, ErrNotFound
}
return Response{}, fmt.Errorf("get edge node: %w", err)
}
var events []Event
if err := s.Orm.Where("node_id = ?", id).Order("occurred_at DESC, id DESC").Limit(20).Find(&events).Error; err != nil {
return Response{}, fmt.Errorf("list edge node events: %w", err)
}
return s.response(model, events), nil
}
// ApplyHeartbeat updates the Sense-owned projection atomically and records
// only state transitions. Callers are internal adapters, not public APIs.
func (s *Service) ApplyHeartbeat(ctx context.Context, sample HeartbeatSample) error {
now := sample.CollectedAt.UTC()
if now.IsZero() {
now = s.now()
}
sample.ID = strings.TrimSpace(sample.ID)
if sample.ID == "" || len(sample.ID) > 64 || strings.TrimSpace(sample.Name) == "" || sample.LastHeartbeatAt.IsZero() || sample.StartedAt.IsZero() || sample.LoadUsed < 0 || sample.LoadCapacity < 0 || sample.LoadUsed > sample.LoadCapacity || sample.BackfillQueueDepth < 0 {
return ErrInvalidSample
}
sample.LastHeartbeatAt = sample.LastHeartbeatAt.UTC()
sample.LastCollectedAt = now
return s.Orm.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var previous Node
err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&previous, "id = ?", sample.ID).Error
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
return fmt.Errorf("lock edge node projection: %w", err)
}
from := ""
if err == nil {
from = statusAt(previous, now)
sample.CreatedAt = previous.CreatedAt
sample.ProjectionVersion = previous.ProjectionVersion + 1
} else {
sample.ProjectionVersion = 1
}
if err = tx.Save(&sample.Node).Error; err != nil {
return fmt.Errorf("save edge node projection: %w", err)
}
to := statusAt(sample.Node, now)
if from == to && from != "" {
return nil
}
if from == StatusOffline {
timeoutAt := previous.LastHeartbeatAt.UTC().Add(HeartbeatTimeout)
if err = tx.Create(&Event{NodeID: sample.ID, EventType: "heartbeat_timeout", FromStatus: StatusOnline, ToStatus: StatusOffline, Detail: "节点心跳超过 90 秒未更新", OccurredAt: timeoutAt}).Error; err != nil {
return err
}
}
eventType, detail := "heartbeat_received", "收到节点首次心跳"
if from == StatusOffline {
eventType, detail = "node_recovered", "节点心跳恢复,通道状态继续独立收敛"
} else if to == StatusRecovering {
eventType, detail = "recovery_started", "节点进入恢复阶段"
}
return tx.Create(&Event{NodeID: sample.ID, EventType: eventType, FromStatus: from, ToStatus: to, Detail: detail, OccurredAt: now}).Error
})
}
func (s *Service) response(model Node, events []Event) Response {
now := s.now()
freshness := int64(now.Sub(model.LastHeartbeatAt.UTC()).Seconds())
if freshness < 0 {
freshness = 0
}
uptime := int64(model.LastHeartbeatAt.UTC().Sub(model.StartedAt.UTC()).Seconds())
if uptime < 0 {
uptime = 0
}
return Response{ID: model.ID, Name: model.Name, Location: model.Location, Status: statusAt(model, now), RuntimeVersion: model.RuntimeVersion, ModelVersion: model.ModelVersion, StartedAt: model.StartedAt, UptimeSeconds: uptime, LastHeartbeatAt: model.LastHeartbeatAt, LastCollectedAt: model.LastCollectedAt, FreshnessSeconds: freshness, Stale: freshness > int64(HeartbeatTimeout/time.Second), LoadUsed: model.LoadUsed, LoadCapacity: model.LoadCapacity, ControlTunnelStatus: model.ControlTunnelStatus, ControlTunnelDetail: model.ControlTunnelDetail, ControlLatencyMs: model.ControlLatencyMs, VideoPlaneStatus: model.VideoPlaneStatus, VideoPlaneDetail: model.VideoPlaneDetail, ActiveStreamCount: model.ActiveStreamCount, BackfillStatus: model.BackfillStatus, BackfillQueueDepth: model.BackfillQueueDepth, RecoveryPhase: model.RecoveryPhase, ProjectionVersion: model.ProjectionVersion, Events: events}
}
func (s *Service) status(model Node) string { return statusAt(model, s.now()) }
func (s *Service) now() time.Time {
if s.Now != nil {
return s.Now().UTC()
}
return time.Now().UTC()
}
func statusAt(model Node, now time.Time) string {
if now.UTC().Sub(model.LastHeartbeatAt.UTC()) > HeartbeatTimeout {
return StatusOffline
}
if strings.TrimSpace(model.RecoveryPhase) != "" && model.RecoveryPhase != "converged" {
return StatusRecovering
}
return StatusOnline
}
@@ -0,0 +1,103 @@
package edge_node
import (
"context"
"errors"
"testing"
"time"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
func edgeNodeTestService(t *testing.T) (*Service, *gorm.DB, time.Time) {
t.Helper()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&Node{}, &Event{}); err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 28, 4, 0, 0, 0, time.UTC)
service := NewService(db)
service.Now = func() time.Time { return now }
return service, db, now
}
func TestListClassifiesOnlineOfflineAndRecoveringWithoutErasingLastKnownState(t *testing.T) {
service, db, now := edgeNodeTestService(t)
if err := SeedSyntheticFixture(db, now); err != nil {
t.Fatal(err)
}
page, err := service.List(PageRequest{})
if err != nil {
t.Fatal(err)
}
if page.Count != 3 || page.Summary.Online != 1 || page.Summary.Offline != 1 || page.Summary.Recovering != 1 || page.Summary.BackfillQueueDepth != 17 {
t.Fatalf("unexpected summary: %+v", page.Summary)
}
var offline Response
for _, item := range page.List {
if item.Status == StatusOffline {
offline = item
}
}
if !offline.Stale || offline.ControlTunnelStatus != ChannelReady || offline.VideoPlaneStatus != ChannelReady || offline.BackfillQueueDepth != 12 {
t.Fatalf("offline projection lost last-known state: %+v", offline)
}
}
func TestApplyHeartbeatRecordsRecoveryAndPreservesIndependentChannelConvergence(t *testing.T) {
service, db, now := edgeNodeTestService(t)
old := Node{ID: "EDGE-1", Name: "旧节点", StartedAt: now.Add(-time.Hour), LastHeartbeatAt: now.Add(-5 * time.Minute), LastCollectedAt: now.Add(-5 * time.Minute), LoadCapacity: 16, ControlTunnelStatus: ChannelReady, VideoPlaneStatus: ChannelReady, BackfillStatus: BackfillPending, BackfillQueueDepth: 8, ProjectionVersion: 3}
if err := db.Create(&old).Error; err != nil {
t.Fatal(err)
}
sample := HeartbeatSample{Node: Node{ID: "EDGE-1", Name: "旧节点", StartedAt: old.StartedAt, LastHeartbeatAt: now.Add(-5 * time.Second), LoadUsed: 4, LoadCapacity: 16, ControlTunnelStatus: ChannelReady, VideoPlaneStatus: ChannelUnavailable, BackfillStatus: BackfillPending, BackfillQueueDepth: 3, RecoveryPhase: "video_reconnecting"}, CollectedAt: now}
if err := service.ApplyHeartbeat(context.Background(), sample); err != nil {
t.Fatal(err)
}
item, err := service.Get("EDGE-1")
if err != nil {
t.Fatal(err)
}
if item.Status != StatusRecovering || item.ProjectionVersion != 4 || item.VideoPlaneStatus != ChannelUnavailable || item.BackfillQueueDepth != 3 {
t.Fatalf("unexpected recovery projection: %+v", item)
}
if len(item.Events) != 2 || item.Events[0].EventType != "node_recovered" || item.Events[0].FromStatus != StatusOffline || item.Events[0].ToStatus != StatusRecovering || item.Events[1].EventType != "heartbeat_timeout" || item.Events[1].ToStatus != StatusOffline {
t.Fatalf("unexpected events: %+v", item.Events)
}
}
func TestApplyHeartbeatRejectsInvalidSampleAndDoesNotPersist(t *testing.T) {
service, db, now := edgeNodeTestService(t)
err := service.ApplyHeartbeat(context.Background(), HeartbeatSample{Node: Node{ID: "EDGE-2", Name: "节点", StartedAt: now.Add(-time.Hour), LastHeartbeatAt: now, LoadUsed: 17, LoadCapacity: 16}, CollectedAt: now})
if !errors.Is(err, ErrInvalidSample) {
t.Fatalf("err=%v", err)
}
var count int64
db.Model(&Node{}).Count(&count)
if count != 0 {
t.Fatalf("count=%d", count)
}
}
func TestGetRejectsMissingOrMalformedID(t *testing.T) {
service, _, _ := edgeNodeTestService(t)
if _, err := service.Get(""); !errors.Is(err, ErrInvalidQuery) {
t.Fatalf("err=%v", err)
}
if _, err := service.Get("missing"); !errors.Is(err, ErrNotFound) {
t.Fatalf("err=%v", err)
}
}
func TestListRejectsOversizedPage(t *testing.T) {
service, _, _ := edgeNodeTestService(t)
request := PageRequest{}
request.PageSize = 101
if _, err := service.List(request); !errors.Is(err, ErrInvalidQuery) {
t.Fatalf("err=%v", err)
}
}
+23
View File
@@ -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()
+64 -7
View File
@@ -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
@@ -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, "媒体分片查询失败")
}
}
@@ -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
}
@@ -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
}
@@ -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(&quota.Setting{}); err != nil {
t.Fatal(err)
}
if err = db.Create(&quota.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")
}
}
+72
View File
@@ -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"`
}
@@ -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" }
@@ -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(&current, "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(&current).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
}
@@ -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)
}
}
+107
View File
@@ -0,0 +1,107 @@
package operations
import (
"errors"
"net/http"
"time"
"github.com/gin-gonic/gin"
"github.com/gin-gonic/gin/binding"
"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
}
func (e *API) List(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := PageRequest{}
if err = e.MakeContext(c).Bind(&request).Errors; err != nil {
e.audit(c, service, "List", auditFailure, "运维中心查询条件格式不正确")
e.Error(http.StatusBadRequest, err, "查询条件格式不正确")
return
}
response, err := service.List(request)
if err != nil {
e.audit(c, service, "List", auditFailure, "运维中心列表查询失败")
e.writeError(err)
return
}
e.audit(c, service, "List", auditSuccess, "读取运维中心问题列表")
e.OK(response, "查询成功")
}
func (e *API) Get(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
item, err := service.Get(c.Param("id"))
if err != nil {
e.audit(c, service, "Get", auditFailure, "运维问题详情查询失败")
e.writeError(err)
return
}
e.audit(c, service, "Get", auditSuccess, "读取运维问题详情 "+item.ID)
e.OK(item, "查询成功")
}
func (e *API) Retry(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := RetryRequest{}
if err = e.MakeContext(c).Bind(&request, binding.JSON).Errors; err != nil {
e.audit(c, service, "Retry", auditFailure, "受控重试请求格式不正确")
e.Error(http.StatusBadRequest, err, "请求格式不正确")
return
}
item, err := service.Retry(c.Request.Context(), c.Param("id"), request.ExpectedVersion, user.GetUserId(c))
if err != nil {
e.audit(c, service, "Retry", auditFailure, "受控重试被拒绝")
e.writeError(err)
return
}
e.audit(c, service, "Retry", auditSuccess, "受控重试已排队 "+item.ID)
e.OK(item, "重试任务已排队")
}
func (e *API) audit(c *gin.Context, service *Service, action, status, remark string) {
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()})
if err != nil {
api.GetRequestLogger(c).Errorf("operations audit failed: %s", err.Error())
}
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrInvalidFilter):
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, ErrVersionConflict), errors.Is(err, ErrRetryInProgress), errors.Is(err, ErrRetryNotAllowed):
e.Error(http.StatusConflict, err, err.Error())
case errors.Is(err, ErrProblemNotFound):
e.Error(http.StatusNotFound, err, err.Error())
default:
e.Error(http.StatusInternalServerError, err, "运维中心操作失败")
}
}
@@ -0,0 +1,26 @@
package operations
import (
"time"
"gorm.io/gorm"
adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
)
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: "operations.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(),
}
model.CreateBy, model.UpdateBy = input.UserID, input.UserID
return db.Create(&model).Error
}
+68
View File
@@ -0,0 +1,68 @@
package operations
import (
"time"
commonDto "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
)
const (
ObjectDevice = "device"
ObjectMedia = "media"
ObjectLocalInference = "local_inference"
ProblemAuthentication = "authentication_failed"
ProblemBackoff = "backoff_wait"
ProblemClockDrift = "clock_drift"
ProblemOrphan = "orphan_safety_gate"
ProblemUnavailable = "capability_unavailable"
ProblemUnready = "not_converged"
)
type PageRequest struct {
commonDto.Pagination `search:"-"`
ObjectType string `form:"objectType"`
ProblemType string `form:"problemType"`
Severity string `form:"severity"`
Keyword string `form:"keyword"`
}
type RetryRequest struct {
ExpectedVersion int64 `json:"expectedVersion"`
}
type Summary struct {
ManagedCount int `json:"managedCount"`
ConvergedCount int `json:"convergedCount"`
ActionableProblems int `json:"actionableProblems"`
LocalInference string `json:"localInference"`
}
type Problem struct {
ID string `json:"id"`
ObjectType string `json:"objectType"`
ObjectID string `json:"objectId"`
ObjectName string `json:"objectName"`
Location string `json:"location,omitempty"`
ProblemType string `json:"problemType"`
Severity string `json:"severity"`
Expected string `json:"expected"`
Actual string `json:"actual"`
Difference string `json:"difference"`
NextAction string `json:"nextAction"`
NextRetryAt *time.Time `json:"nextRetryAt,omitempty"`
LastSuccessAt *time.Time `json:"lastSuccessAt,omitempty"`
AttemptCount int `json:"attemptCount"`
Retryable bool `json:"retryable"`
RetryInProgress bool `json:"retryInProgress"`
Version int64 `json:"version"`
SafetyGate string `json:"safetyGate"`
OperationalOnly bool `json:"operationalOnly"`
UpdatedAt time.Time `json:"updatedAt"`
}
type PageResponse struct {
Summary Summary `json:"summary"`
List []Problem `json:"list"`
Count int64 `json:"count"`
}
@@ -0,0 +1,359 @@
package operations
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"time"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
)
var (
ErrInvalidFilter = errors.New("运维中心查询条件不符合要求")
ErrProblemNotFound = errors.New("运维问题不存在或已经收敛")
ErrVersionConflict = errors.New("状态已经变化,请刷新后重试")
ErrRetryInProgress = errors.New("该对象已有重试任务")
ErrRetryNotAllowed = errors.New("该问题不允许重试")
)
type Service struct {
DB *gorm.DB
Now func() time.Time
LocalInferenceConfigured bool
DeviceRetry func(context.Context, admission.ProbeRequest) error
}
type admissionResult struct {
DeviceID string
Address string
Status string
Detail string
CheckedAt time.Time
}
func (admissionResult) TableName() string { return "sense_admission_results" }
func NewService(db *gorm.DB) *Service {
service := &Service{DB: db, Now: time.Now, LocalInferenceConfigured: false}
service.DeviceRetry = func(ctx context.Context, request admission.ProbeRequest) error {
base := coreService.Service{}
base.Orm = db
runtimeService, err := admission.NewRuntime(base)
if err != nil {
return err
}
_, err = runtimeService.Probe(ctx, request)
if err == nil {
// Keep the existing admission boundary: route intent is best-effort and
// must not turn a successful device probe into a MediaMTX failure.
_ = media.EnsureDeviceRoutes(ctx, db, request.DeviceID)
}
return err
}
return service
}
func (s *Service) List(request PageRequest) (PageResponse, error) {
if err := validateFilter(request); err != nil {
return PageResponse{}, err
}
problems, summary, err := s.project()
if err != nil {
return PageResponse{}, err
}
filtered := make([]Problem, 0, len(problems))
keyword := strings.ToLower(strings.TrimSpace(request.Keyword))
for _, item := range problems {
if request.ObjectType != "" && item.ObjectType != request.ObjectType {
continue
}
if request.ProblemType != "" && item.ProblemType != request.ProblemType {
continue
}
if request.Severity != "" && item.Severity != request.Severity {
continue
}
searchable := strings.ToLower(strings.Join([]string{item.ObjectID, item.ObjectName, item.Location, item.Difference}, " "))
if keyword != "" && !strings.Contains(searchable, keyword) {
continue
}
filtered = append(filtered, item)
}
pageIndex, pageSize := request.GetPageIndex(), request.GetPageSize()
if pageSize > 100 {
pageSize = 100
}
start := (pageIndex - 1) * pageSize
if start > len(filtered) {
start = len(filtered)
}
end := start + pageSize
if end > len(filtered) {
end = len(filtered)
}
return PageResponse{Summary: summary, List: filtered[start:end], Count: int64(len(filtered))}, nil
}
func (s *Service) Get(id string) (Problem, error) {
problems, _, err := s.project()
if err != nil {
return Problem{}, err
}
for _, item := range problems {
if item.ID == id {
return item, nil
}
}
return Problem{}, ErrProblemNotFound
}
func (s *Service) Retry(ctx context.Context, id string, expectedVersion int64, userID int) (Problem, error) {
if expectedVersion < 1 {
return Problem{}, ErrVersionConflict
}
problem, err := s.Get(id)
if err != nil {
return Problem{}, err
}
if !problem.Retryable {
if problem.RetryInProgress {
return Problem{}, ErrRetryInProgress
}
return Problem{}, ErrRetryNotAllowed
}
now := s.now().UTC()
switch problem.ObjectType {
case ObjectDevice:
var admissionState admissionResult
if err := s.DB.First(&admissionState, "device_id = ?", problem.ObjectID).Error; err != nil || strings.TrimSpace(admissionState.Address) == "" {
return Problem{}, ErrRetryNotAllowed
}
result := s.DB.Model(&deviceModels.Device{}).
Where("id = ? AND version = ? AND retry_requested_at IS NULL", problem.ObjectID, expectedVersion).
Updates(map[string]any{"retry_requested_at": now, "version": expectedVersion + 1, "updated_at": now})
if result.Error != nil {
return Problem{}, fmt.Errorf("queue device retry: %w", result.Error)
}
if result.RowsAffected == 0 {
return Problem{}, s.retryConflict(ObjectDevice, problem.ObjectID, expectedVersion)
}
if s.DeviceRetry != nil {
err = s.DeviceRetry(ctx, admission.ProbeRequest{DeviceID: problem.ObjectID, Address: admissionState.Address, Version: expectedVersion + 1, UpdateBy: userID})
if err != nil {
clearErr := s.DB.Model(&deviceModels.Device{}).Where("id = ? AND version = ?", problem.ObjectID, expectedVersion+1).Update("retry_requested_at", nil).Error
if clearErr != nil {
return Problem{}, errors.Join(fmt.Errorf("execute device retry: %w", err), fmt.Errorf("clear device retry gate: %w", clearErr))
}
return Problem{}, fmt.Errorf("execute device retry: %w", err)
}
}
case ObjectMedia:
result := s.DB.Model(&media.Route{}).
Where("id = ? AND version = ? AND actual <> ?", problem.ObjectID, expectedVersion, "retry_pending").
Updates(map[string]any{"actual": "retry_pending", "next_retry_at": now, "detail": "已请求受控重试,等待视频服务执行", "version": expectedVersion + 1, "updated_at": now})
if result.Error != nil {
return Problem{}, fmt.Errorf("queue media retry: %w", result.Error)
}
if result.RowsAffected == 0 {
return Problem{}, s.retryConflict(ObjectMedia, problem.ObjectID, expectedVersion)
}
default:
return Problem{}, ErrRetryNotAllowed
}
updated, err := s.Get(id)
if errors.Is(err, ErrProblemNotFound) {
problem.Actual, problem.Difference, problem.NextAction = "converged", "重试完成,状态已经收敛", "无需处理"
problem.Retryable, problem.RetryInProgress, problem.Version, problem.UpdatedAt = false, false, expectedVersion+2, now
return problem, nil
}
return updated, err
}
func (s *Service) project() ([]Problem, Summary, error) {
var devices []deviceModels.Device
if err := s.DB.Order("updated_at DESC").Find(&devices).Error; err != nil {
return nil, Summary{}, fmt.Errorf("list operation devices: %w", err)
}
var admissions []admissionResult
if err := s.DB.Find(&admissions).Error; err != nil {
return nil, Summary{}, fmt.Errorf("list admission results: %w", err)
}
var routes []media.Route
if err := s.DB.Order("updated_at DESC").Find(&routes).Error; err != nil {
return nil, Summary{}, fmt.Errorf("list operation media routes: %w", err)
}
admissionByDevice := make(map[string]admissionResult, len(admissions))
for _, item := range admissions {
admissionByDevice[item.DeviceID] = item
}
deviceByID := make(map[string]deviceModels.Device, len(devices))
problems := make([]Problem, 0)
summary := Summary{ManagedCount: len(devices) + len(routes), LocalInference: "configured"}
for _, device := range devices {
deviceByID[device.ID] = device
if problem, ok := deviceProblem(device, admissionByDevice[device.ID]); ok {
problems = append(problems, problem)
summary.ActionableProblems++
} else {
summary.ConvergedCount++
}
}
for _, route := range routes {
if problem, ok := mediaProblem(route, deviceByID); ok {
problems = append(problems, problem)
summary.ActionableProblems++
} else {
summary.ConvergedCount++
}
}
if !s.LocalInferenceConfigured {
summary.LocalInference = "unavailable"
now := s.now().UTC()
problems = append(problems, Problem{
ID: "local_inference:adapter", ObjectType: ObjectLocalInference, ObjectID: "adapter",
ObjectName: "本地推理适配器", ProblemType: ProblemUnavailable, Severity: "info",
Expected: "可选", Actual: "unavailable", Difference: "未配置本地推理适配器;不影响设备和媒体运维",
NextAction: "无需处理", SafetyGate: "不读取 Brain 数据库,不阻断本页", OperationalOnly: true, UpdatedAt: now,
})
}
sort.SliceStable(problems, func(i, j int) bool {
rank := map[string]int{"high": 0, "medium": 1, "info": 2}
if rank[problems[i].Severity] != rank[problems[j].Severity] {
return rank[problems[i].Severity] < rank[problems[j].Severity]
}
return problems[i].UpdatedAt.After(problems[j].UpdatedAt)
})
return problems, summary, nil
}
func deviceProblem(device deviceModels.Device, admission admissionResult) (Problem, bool) {
if device.Status == deviceModels.StatusDisabled {
return Problem{}, false
}
actual, detail, updatedAt := device.AdapterStatus, "设备尚未完成接入验证", device.UpdatedAt
if admission.DeviceID != "" {
actual, detail, updatedAt = admission.Status, admission.Detail, admission.CheckedAt
}
if device.Status == deviceModels.StatusActive && (actual == "ready" || actual == "verified") {
return Problem{}, false
}
problemType, severity, nextAction := ProblemUnready, "medium", "检查设备和接入配置后重试"
text := strings.ToLower(actual + " " + detail)
if strings.Contains(text, "auth") || strings.Contains(text, "认证") || strings.Contains(text, "凭据") {
problemType, severity, nextAction = ProblemAuthentication, "high", "确认设备账号未变更后执行受控重试"
} else if strings.Contains(text, "clock") || strings.Contains(text, "time drift") || strings.Contains(text, "时间漂移") || strings.Contains(text, "时钟") {
problemType, nextAction = ProblemClockDrift, "检查设备 NTP 和时区后重新检测"
}
retrying := device.RetryRequestedAt != nil
retryable := !retrying && strings.TrimSpace(admission.Address) != ""
if admission.DeviceID == "" || strings.TrimSpace(admission.Address) == "" {
nextAction = "先到视频接入完成地址与凭据验证"
}
return Problem{
ID: "device:" + device.ID, ObjectType: ObjectDevice, ObjectID: device.ID, ObjectName: device.Name,
Location: device.Location, ProblemType: problemType, Severity: severity, Expected: "active / ready", Actual: actual,
Difference: detail, NextAction: nextAction, AttemptCount: boolInt(retrying), Retryable: retryable,
RetryInProgress: retrying, Version: device.Version, SafetyGate: "校验设备版本且同一设备仅允许一个在途重试",
OperationalOnly: true, UpdatedAt: updatedAt,
}, true
}
func mediaProblem(route media.Route, devices map[string]deviceModels.Device) (Problem, bool) {
device, exists := devices[route.DeviceID]
name, location := route.Path, ""
if exists {
name, location = device.Name+" / "+route.Path, device.Location
}
if !exists {
return Problem{
ID: "media:" + route.ID, ObjectType: ObjectMedia, ObjectID: route.ID, ObjectName: name,
ProblemType: ProblemOrphan, Severity: "medium", Expected: "路由关联有效设备", Actual: "已隔离待确认",
Difference: "媒体路由存在,但找不到有效设备归属", NextAction: "人工核对;不会自动删除", Version: route.Version,
SafetyGate: "孤儿资源仅隔离和提示,本接口没有删除动作", OperationalOnly: true, UpdatedAt: route.UpdatedAt,
}, true
}
converged := (route.Desired == media.DesiredRunning && (route.Actual == "ready" || route.Actual == "waiting")) || (route.Desired == media.DesiredStopped && route.Actual == "stopped")
if converged {
return Problem{}, false
}
problemType, nextAction := ProblemUnready, "检查视频服务状态后重试"
if route.NextRetryAt != nil {
problemType, nextAction = ProblemBackoff, "等待退避到期或执行受控提前重试"
}
retrying := route.Actual == "retry_pending"
return Problem{
ID: "media:" + route.ID, ObjectType: ObjectMedia, ObjectID: route.ID, ObjectName: name, Location: location,
ProblemType: problemType, Severity: "medium", Expected: route.Desired, Actual: route.Actual, Difference: route.Detail,
NextAction: nextAction, NextRetryAt: route.NextRetryAt, AttemptCount: route.FailureCount, Retryable: !retrying,
RetryInProgress: retrying, Version: route.Version, SafetyGate: "校验路由版本且只排队,不删除路径或修改凭据",
OperationalOnly: true, UpdatedAt: route.UpdatedAt,
}, true
}
func (s *Service) retryConflict(objectType, objectID string, expectedVersion int64) error {
var version int64
var inProgress bool
switch objectType {
case ObjectDevice:
var item deviceModels.Device
if err := s.DB.Select("version", "retry_requested_at").First(&item, "id = ?", objectID).Error; err != nil {
return ErrProblemNotFound
}
version, inProgress = item.Version, item.RetryRequestedAt != nil
case ObjectMedia:
var item media.Route
if err := s.DB.Select("version", "actual").First(&item, "id = ?", objectID).Error; err != nil {
return ErrProblemNotFound
}
version, inProgress = item.Version, item.Actual == "retry_pending"
}
if inProgress {
return ErrRetryInProgress
}
if version != expectedVersion {
return ErrVersionConflict
}
return ErrVersionConflict
}
func (s *Service) now() time.Time {
if s.Now != nil {
return s.Now()
}
return time.Now()
}
func validateFilter(request PageRequest) error {
valid := func(value string, values ...string) bool {
if value == "" {
return true
}
for _, candidate := range values {
if value == candidate {
return true
}
}
return false
}
if !valid(request.ObjectType, ObjectDevice, ObjectMedia, ObjectLocalInference) ||
!valid(request.ProblemType, ProblemAuthentication, ProblemBackoff, ProblemClockDrift, ProblemOrphan, ProblemUnavailable, ProblemUnready) ||
!valid(request.Severity, "high", "medium", "info") || len([]rune(request.Keyword)) > 128 {
return ErrInvalidFilter
}
return nil
}
func boolInt(value bool) int {
if value {
return 1
}
return 0
}
@@ -0,0 +1,187 @@
package operations
import (
"context"
"errors"
"testing"
"time"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
)
func operationsTestService(t *testing.T) (*Service, *gorm.DB, time.Time) {
t.Helper()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&deviceModels.Device{}, &admissionResult{}, &media.Route{}, &adminModels.SysOperaLog{}); err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 28, 2, 0, 0, 0, time.UTC)
service := NewService(db)
service.Now = func() time.Time { return now }
service.LocalInferenceConfigured = false
service.DeviceRetry = nil
return service, db, now
}
func TestProjectionCoversOperationsStatesAndIndependentBoundaries(t *testing.T) {
service, db, now := operationsTestService(t)
devices := []deviceModels.Device{
{ID: "ready", Name: "东门", Location: "一层", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1},
{ID: "auth", Name: "仓库", Location: "北区", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3},
{ID: "clock", Name: "南门", Location: "室外", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 4},
}
for index := range devices {
devices[index].CreatedAt, devices[index].UpdatedAt = now, now
}
if err := db.Create(&devices).Error; err != nil {
t.Fatal(err)
}
admissions := []admissionResult{
{DeviceID: "ready", Address: "http://192.0.2.10/onvif", Status: "ready", Detail: "接入验证完成", CheckedAt: now},
{DeviceID: "auth", Address: "http://192.0.2.11/onvif", Status: "authentication_failed", Detail: "ONVIF 认证失败", CheckedAt: now},
{DeviceID: "clock", Address: "http://192.0.2.12/onvif", Status: "clock_drift", Detail: "设备时间漂移 96 秒", CheckedAt: now},
}
if err := db.Create(&admissions).Error; err != nil {
t.Fatal(err)
}
next := now.Add(5 * time.Minute)
routes := []media.Route{
{ID: "auth:main", DeviceID: "auth", ProfileToken: "main", Path: "sense_auth", Desired: "running", Actual: "apply_failed", Detail: "媒体路径配置失败", FailureCount: 2, NextRetryAt: &next, Version: 5, UpdatedAt: now},
{ID: "missing:main", DeviceID: "missing", ProfileToken: "main", Path: "sense_orphan", Desired: "running", Actual: "ready", Detail: "上游拉流正常", Version: 1, UpdatedAt: now},
}
if err := db.Create(&routes).Error; err != nil {
t.Fatal(err)
}
response, err := service.List(PageRequest{})
if err != nil {
t.Fatal(err)
}
if response.Summary.ManagedCount != 5 || response.Summary.ConvergedCount != 1 || response.Summary.ActionableProblems != 4 || response.Summary.LocalInference != "unavailable" || response.Count != 5 {
t.Fatalf("unexpected summary: %#v count=%d", response.Summary, response.Count)
}
wanted := map[string]bool{ProblemAuthentication: false, ProblemClockDrift: false, ProblemBackoff: false, ProblemOrphan: false, ProblemUnavailable: false}
for _, item := range response.List {
wanted[item.ProblemType] = true
if !item.OperationalOnly {
t.Fatalf("problem can be mistaken for Bell alert: %#v", item)
}
}
for state, found := range wanted {
if !found {
t.Fatalf("missing problem state %s: %#v", state, response.List)
}
}
filtered, err := service.List(PageRequest{ObjectType: ObjectDevice, ProblemType: ProblemClockDrift, Keyword: "南门"})
if err != nil || filtered.Count != 1 || filtered.List[0].ObjectID != "clock" {
t.Fatalf("filtered=%#v err=%v", filtered, err)
}
}
func TestControlledRetryUsesVersionAndInProgressGate(t *testing.T) {
service, db, now := operationsTestService(t)
device := deviceModels.Device{ID: "auth", Name: "东门", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3}
device.CreatedAt, device.UpdatedAt = now, now
if err := db.Create(&device).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&admissionResult{DeviceID: "auth", Address: "http://192.0.2.11/onvif", Status: "authentication_failed", Detail: "认证失败", CheckedAt: now}).Error; err != nil {
t.Fatal(err)
}
item, err := service.Retry(context.Background(), "device:auth", 3, 7)
if err != nil || !item.RetryInProgress || item.Retryable || item.Version != 4 {
t.Fatalf("item=%#v err=%v", item, err)
}
if _, err = service.Retry(context.Background(), "device:auth", 3, 7); !errors.Is(err, ErrRetryInProgress) {
t.Fatalf("expected in-progress gate, got %v", err)
}
var stored deviceModels.Device
if err = db.First(&stored, "id = ?", "auth").Error; err != nil || stored.RetryRequestedAt == nil {
t.Fatalf("stored=%#v err=%v", stored, err)
}
}
func TestMediaRetryQueuesWithoutDeletingOrChangingCredentials(t *testing.T) {
service, db, now := operationsTestService(t)
device := deviceModels.Device{ID: "camera", Name: "仓库", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1}
device.CreatedAt, device.UpdatedAt = now, now
if err := db.Create(&device).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&admissionResult{DeviceID: "camera", Address: "http://192.0.2.20/onvif", Status: "ready", Detail: "接入完成", CheckedAt: now}).Error; err != nil {
t.Fatal(err)
}
next := now.Add(time.Minute)
route := media.Route{ID: "camera:main", DeviceID: "camera", ProfileToken: "main", Path: "sense_camera", Desired: "running", Actual: "apply_failed", FailureCount: 2, NextRetryAt: &next, Detail: "配置失败", Version: 6, UpdatedAt: now}
if err := db.Create(&route).Error; err != nil {
t.Fatal(err)
}
item, err := service.Retry(context.Background(), "media:camera:main", 6, 7)
if err != nil || item.Actual != "retry_pending" || item.Version != 7 || !item.RetryInProgress {
t.Fatalf("item=%#v err=%v", item, err)
}
var count int64
if err = db.Model(&media.Route{}).Where("id = ?", route.ID).Count(&count).Error; err != nil || count != 1 {
t.Fatalf("route was removed: count=%d err=%v", count, err)
}
}
func TestOrphanAndUnavailableCannotRetry(t *testing.T) {
service, db, now := operationsTestService(t)
route := media.Route{ID: "missing:main", DeviceID: "missing", ProfileToken: "main", Path: "orphan", Desired: "running", Actual: "ready", Version: 1, UpdatedAt: now}
if err := db.Create(&route).Error; err != nil {
t.Fatal(err)
}
if _, err := service.Retry(context.Background(), "media:missing:main", 1, 7); !errors.Is(err, ErrRetryNotAllowed) {
t.Fatalf("orphan retry should be rejected: %v", err)
}
if _, err := service.Retry(context.Background(), "local_inference:adapter", 1, 7); !errors.Is(err, ErrRetryNotAllowed) {
t.Fatalf("optional adapter retry should be rejected: %v", err)
}
}
func TestDeviceRetryUsesExistingAdmissionPortAndOperatorIdentity(t *testing.T) {
service, db, now := operationsTestService(t)
device := deviceModels.Device{ID: "clock", Name: "南门", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3}
device.CreatedAt, device.UpdatedAt = now, now
if err := db.Create(&device).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&admissionResult{DeviceID: "clock", Address: "http://192.0.2.12/onvif", Status: "clock_drift", Detail: "设备时间漂移", CheckedAt: now}).Error; err != nil {
t.Fatal(err)
}
called := false
service.DeviceRetry = func(_ context.Context, request admission.ProbeRequest) error {
called = true
if request.DeviceID != "clock" || request.Address != "http://192.0.2.12/onvif" || request.Version != 4 || request.UpdateBy != 9 {
t.Fatalf("unexpected retry request: %#v", request)
}
return db.Model(&deviceModels.Device{}).Where("id = ? AND version = ?", request.DeviceID, request.Version).Updates(map[string]any{"retry_requested_at": nil, "version": request.Version + 1}).Error
}
item, err := service.Retry(context.Background(), "device:clock", 3, 9)
if err != nil || !called || item.Version != 5 || item.RetryInProgress {
t.Fatalf("item=%#v called=%v err=%v", item, called, err)
}
}
func TestAuditIsMinimalAndDesensitized(t *testing.T) {
_, db, now := operationsTestService(t)
if err := WriteAudit(db, Audit{Action: "Retry", Method: "POST", Status: "1", Username: "operator", UserID: 7, ClientIP: "127.0.0.1", Route: "/api/v1/operations/:id/retry", Remark: "受控重试已排队 device:synthetic", At: now}); err != nil {
t.Fatal(err)
}
var stored adminModels.SysOperaLog
if err := db.First(&stored).Error; err != nil {
t.Fatal(err)
}
if stored.Title != "运维中心" || stored.RequestMethod != "POST" || stored.OperParam != "" || stored.JsonResult != "" || stored.CreateBy != 7 {
t.Fatalf("unexpected audit: %#v", stored)
}
}
@@ -12,6 +12,8 @@ import (
"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"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
)
type API struct{ api.Api }
@@ -137,6 +139,10 @@ func (e *API) writeError(err error) {
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, ErrBatchNotFound), errors.Is(err, ErrItemNotFound):
e.Error(http.StatusNotFound, err, err.Error())
case errors.Is(err, quota.ErrExceeded):
e.Error(http.StatusConflict, err, "当前配额已满,无法继续批量开通")
case errors.Is(err, quota.ErrUnavailable):
e.Error(http.StatusServiceUnavailable, err, "配额配置不可读取,已拒绝批量开通写入")
default:
e.Error(http.StatusInternalServerError, err, "批量开通操作失败")
}
+14 -20
View File
@@ -5,9 +5,7 @@ import (
"errors"
"fmt"
"net/url"
"os"
"sort"
"strconv"
"strings"
"time"
@@ -20,6 +18,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
deviceService "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service"
deviceDTO "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service/dto"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
)
var (
@@ -36,23 +35,11 @@ type Service struct {
Activator Activator
}
func QuotaFromEnvironment() int {
value := strings.TrimSpace(os.Getenv("SENSE_PROVISIONING_QUOTA"))
if value == "" {
return 16
}
parsed, err := strconv.Atoi(value)
if err != nil || parsed < 1 || parsed > 100000 {
return 16
}
return parsed
}
func (s *Service) quota() int {
func (s *Service) quotaLimit() (int, error) {
if s.Quota > 0 {
return s.Quota
return s.Quota, nil
}
return QuotaFromEnvironment()
return quota.ReadLimit(s.Orm)
}
func (s *Service) CreateBatch(request CreateBatchRequest) (BatchResponse, error) {
@@ -71,12 +58,15 @@ func (s *Service) CreateBatch(request CreateBatchRequest) (BatchResponse, error)
if err := s.Orm.Table("sense_devices").Where("status <> ?", "disabled").Count(&existingDevices).Error; err != nil {
return BatchResponse{}, fmt.Errorf("count provisioned devices: %w", err)
}
quota := s.quota()
available := quota - int(existingDevices)
quotaLimit, err := s.quotaLimit()
if err != nil {
return BatchResponse{}, err
}
available := quotaLimit - int(existingDevices)
if available < 0 {
available = 0
}
batch := Batch{ID: uuid.NewString(), IdempotencyKey: request.IdempotencyKey, Status: BatchReady, QuotaLimit: quota, ExistingCount: int(existingDevices), TotalCount: len(request.Rows)}
batch := Batch{ID: uuid.NewString(), IdempotencyKey: request.IdempotencyKey, Status: BatchReady, QuotaLimit: quotaLimit, ExistingCount: int(existingDevices), TotalCount: len(request.Rows)}
batch.CreateBy, batch.UpdateBy = request.CreateBy, request.CreateBy
seenLines := map[int]bool{}
seenAddresses := map[string]bool{}
@@ -351,6 +341,10 @@ func safeFailure(err error) (string, string) {
return "invalid_device", "设备信息或凭据不符合要求"
case errors.Is(err, deviceService.ErrVersionConflict), errors.Is(err, admission.ErrConflict):
return "version_conflict", "设备已被其他操作更新,请重试"
case errors.Is(err, quota.ErrExceeded):
return "quota_exceeded", "当前配额已满,请调整配额或停用其他设备后重试"
case errors.Is(err, quota.ErrUnavailable):
return "quota_unavailable", "配额配置不可读取,已拒绝新增或启用设备"
default:
return "activation_failed", "设备开通失败,请检查网络、地址和凭据后重试"
}
@@ -153,17 +153,6 @@ func TestExecuteRequiresCredentialsWithoutCallingActivator(t *testing.T) {
}
}
func TestQuotaFromEnvironment(t *testing.T) {
t.Setenv("SENSE_PROVISIONING_QUOTA", "24")
if got := QuotaFromEnvironment(); got != 24 {
t.Fatalf("quota=%d", got)
}
t.Setenv("SENSE_PROVISIONING_QUOTA", "invalid")
if got := QuotaFromEnvironment(); got != 16 {
t.Fatalf("fallback quota=%d", got)
}
}
func TestConcurrentCreateUsesOneIdempotentBatch(t *testing.T) {
service := provisioningService(t, 16)
sqlDB, err := service.Orm.DB()
+73
View File
@@ -0,0 +1,73 @@
package quota
import (
"errors"
"net/http"
"github.com/gin-gonic/gin"
"github.com/gin-gonic/gin/binding"
"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"
)
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 &Service{DB: base.Orm}, nil
}
func (e *API) Get(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := PageRequest{}
if err = e.MakeContext(c).Bind(&request).Errors; err != nil {
e.Error(http.StatusBadRequest, err, "查询条件格式不正确")
return
}
response, err := service.Overview(request)
if err != nil {
e.writeError(err)
return
}
e.OK(response, "查询成功")
}
func (e *API) Update(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := UpdateRequest{UpdateBy: user.GetUserId(c)}
if err = e.MakeContext(c).Bind(&request, binding.JSON).Errors; err != nil {
e.Error(http.StatusBadRequest, err, "请求内容格式不正确")
return
}
response, err := service.Update(request)
if err != nil {
e.writeError(err)
return
}
e.OK(response, "配额已更新")
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrInvalid):
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, ErrExceeded), errors.Is(err, ErrVersionConflict):
e.Error(http.StatusConflict, err, err.Error())
case errors.Is(err, ErrUnavailable):
e.Error(http.StatusServiceUnavailable, err, err.Error())
default:
e.Error(http.StatusInternalServerError, err, "容量与配额操作失败")
}
}
+60
View File
@@ -0,0 +1,60 @@
package quota
import (
"time"
commonDto "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
)
const (
ReadStatusReadable = "readable"
ReadStatusUnreadable = "unreadable"
)
type PageRequest struct {
commonDto.Pagination `search:"-"`
Keyword string `form:"keyword"`
Status string `form:"status"`
}
type UpdateRequest struct {
Limit int `json:"limit"`
Reason string `json:"reason"`
Version int64 `json:"version"`
UpdateBy int `json:"-"`
}
type Summary struct {
Limit int `json:"limit"`
Used int `json:"used"`
Remaining int `json:"remaining"`
ReadStatus string `json:"readStatus"`
Source string `json:"source"`
Version int64 `json:"version"`
LastReadAt time.Time `json:"lastReadAt"`
LastChanged time.Time `json:"lastChangedAt,omitempty"`
}
type DeviceOccupancy struct {
ID string `json:"id"`
Name string `json:"name"`
Location string `json:"location"`
Status string `json:"status"`
Occupied bool `json:"occupied"`
Version int64 `json:"version"`
UpdatedAt time.Time `json:"updatedAt"`
}
type Tier struct {
Limit int `json:"limit"`
Configured bool `json:"configured"`
Validation string `json:"validation"`
DeliveryStatus string `json:"deliveryStatus"`
}
type Overview struct {
Summary Summary `json:"summary"`
List []DeviceOccupancy `json:"list"`
Count int64 `json:"count"`
Tiers []Tier `json:"tiers"`
}
+31
View File
@@ -0,0 +1,31 @@
package quota
import (
"time"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
const SettingID uint = 1
type Setting struct {
ID uint `gorm:"primaryKey;autoIncrement:false" json:"id"`
Limit int `gorm:"not null" json:"limit"`
Source string `gorm:"size:32;not null;default:'database'" json:"source"`
Version int64 `gorm:"not null;default:1" json:"version"`
common.ControlBy
common.ModelTime
}
func (Setting) TableName() string { return "sense_quota_settings" }
type Change struct {
ID string `gorm:"size:36;primaryKey" json:"id"`
OldLimit int `gorm:"not null" json:"oldLimit"`
NewLimit int `gorm:"not null" json:"newLimit"`
Reason string `gorm:"size:512;not null" json:"reason"`
ChangedBy int `gorm:"not null;index" json:"changedBy"`
CreatedAt time.Time `gorm:"not null;index" json:"createdAt"`
}
func (Change) TableName() string { return "sense_quota_changes" }
+168
View File
@@ -0,0 +1,168 @@
package quota
import (
"errors"
"fmt"
"strings"
"time"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
var (
ErrUnavailable = errors.New("配额配置不可读取")
ErrExceeded = errors.New("当前配额已满")
ErrInvalid = errors.New("配额配置不符合要求")
ErrVersionConflict = errors.New("配额已被其他用户修改,请刷新后重试")
)
type Service struct{ DB *gorm.DB }
func (s Service) Overview(request PageRequest) (Overview, error) {
if s.DB == nil {
return Overview{}, ErrUnavailable
}
request.Keyword = strings.TrimSpace(request.Keyword)
if request.Status != "" && request.Status != "active" && request.Status != "pending" && request.Status != "disabled" {
return Overview{}, ErrInvalid
}
used, err := occupiedCount(s.DB)
if err != nil {
return Overview{}, fmt.Errorf("count quota occupancy: %w", err)
}
now := time.Now().UTC()
summary := Summary{Used: int(used), ReadStatus: ReadStatusUnreadable, LastReadAt: now}
var setting Setting
if err = s.DB.First(&setting, "id = ?", SettingID).Error; err == nil && validLimit(setting.Limit) {
summary.Limit = setting.Limit
summary.Remaining = setting.Limit - int(used)
if summary.Remaining < 0 {
summary.Remaining = 0
}
summary.ReadStatus = ReadStatusReadable
summary.Source = setting.Source
summary.Version = setting.Version
summary.LastChanged = setting.UpdatedAt
}
query := s.DB.Table("sense_devices").Select("id, name, location, status, version, updated_at")
if request.Keyword != "" {
pattern := "%" + strings.ToLower(request.Keyword) + "%"
query = query.Where("LOWER(name) LIKE ? OR LOWER(location) LIKE ?", pattern, pattern)
}
if request.Status != "" {
query = query.Where("status = ?", request.Status)
}
var count int64
if err = query.Count(&count).Error; err != nil {
return Overview{}, fmt.Errorf("count quota devices: %w", err)
}
pageSize := request.GetPageSize()
if pageSize > 100 {
pageSize = 100
}
var rows []DeviceOccupancy
if err = query.Order("updated_at DESC").Limit(pageSize).Offset((request.GetPageIndex() - 1) * pageSize).Scan(&rows).Error; err != nil {
return Overview{}, fmt.Errorf("list quota devices: %w", err)
}
for index := range rows {
rows[index].Occupied = rows[index].Status != "disabled"
}
return Overview{Summary: summary, List: rows, Count: count, Tiers: deliveryTiers(summary.Limit)}, nil
}
func (s Service) Update(request UpdateRequest) (Summary, error) {
request.Reason = strings.TrimSpace(request.Reason)
if s.DB == nil || !validLimit(request.Limit) || request.Version < 1 || request.Reason == "" || len([]rune(request.Reason)) > 512 {
return Summary{}, ErrInvalid
}
err := s.DB.Transaction(func(tx *gorm.DB) error {
setting, err := lockSetting(tx)
if err != nil {
return err
}
if setting.Version != request.Version {
return ErrVersionConflict
}
now := time.Now().UTC()
result := tx.Model(&Setting{}).Where("id = ? AND version = ?", SettingID, request.Version).Updates(map[string]any{
"limit": request.Limit, "source": "database", "version": request.Version + 1,
"update_by": request.UpdateBy, "updated_at": now,
})
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return ErrVersionConflict
}
return tx.Create(&Change{ID: uuid.NewString(), OldLimit: setting.Limit, NewLimit: request.Limit, Reason: request.Reason, ChangedBy: request.UpdateBy, CreatedAt: now}).Error
})
if err != nil {
return Summary{}, err
}
overview, err := s.Overview(PageRequest{})
return overview.Summary, err
}
func ReadLimit(db *gorm.DB) (int, error) {
if db == nil {
return 0, ErrUnavailable
}
var setting Setting
if err := db.First(&setting, "id = ?", SettingID).Error; err != nil || !validLimit(setting.Limit) {
return 0, ErrUnavailable
}
return setting.Limit, nil
}
func WithAvailableSlot(db *gorm.DB, write func(*gorm.DB) error) error {
if db == nil {
return ErrUnavailable
}
return db.Transaction(func(tx *gorm.DB) error {
setting, err := lockSetting(tx)
if err != nil {
return err
}
used, err := occupiedCount(tx)
if err != nil {
return fmt.Errorf("count quota occupancy: %w", err)
}
if used >= int64(setting.Limit) {
return ErrExceeded
}
return write(tx)
})
}
func lockSetting(tx *gorm.DB) (Setting, error) {
var setting Setting
err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&setting, "id = ?", SettingID).Error
if err != nil || !validLimit(setting.Limit) {
return Setting{}, ErrUnavailable
}
return setting, nil
}
func occupiedCount(db *gorm.DB) (int64, error) {
var count int64
err := db.Table("sense_devices").Where("status <> ?", "disabled").Count(&count).Error
return count, err
}
func validLimit(limit int) bool { return limit >= 1 && limit <= 100000 }
func deliveryTiers(configured int) []Tier {
tiers := []Tier{
{Limit: 16, Validation: "verified", DeliveryStatus: "默认学校试点交付档位"},
{Limit: 32, Validation: "unverified", DeliveryStatus: "需完成目标硬件压测后启用"},
{Limit: 64, Validation: "unverified", DeliveryStatus: "不作单机容量承诺"},
{Limit: 128, Validation: "unverified", DeliveryStatus: "当前不在交付承诺范围"},
}
for index := range tiers {
tiers[index].Configured = tiers[index].Limit == configured
}
return tiers
}
@@ -0,0 +1,130 @@
package quota
import (
"sync"
"testing"
"github.com/google/uuid"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
)
func testQuotaDB(t *testing.T, limit int) *gorm.DB {
t.Helper()
db, err := gorm.Open(sqlite.Open("file:"+uuid.NewString()+"?mode=memory&cache=shared"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatal(err)
}
sqlDB.SetMaxOpenConns(1)
if err = db.AutoMigrate(&deviceModels.Device{}, &Setting{}, &Change{}); err != nil {
t.Fatal(err)
}
if err = db.Create(&Setting{ID: SettingID, Limit: limit, Source: "test", Version: 1}).Error; err != nil {
t.Fatal(err)
}
return db
}
func TestOverviewShowsOccupancyPaginationAndDeliveryBoundaries(t *testing.T) {
db := testQuotaDB(t, 16)
for index, status := range []string{"active", "pending", "disabled"} {
device := deviceModels.Device{ID: uuid.NewString(), Name: "设备", Location: "位置", Modality: "video", CapabilitiesJSON: "[]", Status: status, AdapterStatus: "ready", Version: int64(index + 1)}
if err := db.Create(&device).Error; err != nil {
t.Fatal(err)
}
}
request := PageRequest{}
request.PageIndex, request.PageSize = 1, 64
response, err := (Service{DB: db}).Overview(request)
if err != nil {
t.Fatal(err)
}
if response.Summary.Limit != 16 || response.Summary.Used != 2 || response.Summary.Remaining != 14 || response.Count != 3 || len(response.List) != 3 {
t.Fatalf("unexpected overview: %#v", response)
}
if response.Tiers[0].Validation != "verified" || response.Tiers[1].Validation != "unverified" || !response.Tiers[0].Configured {
t.Fatalf("delivery tiers overpromise capacity: %#v", response.Tiers)
}
}
func TestUnreadableQuotaKeepsReadsAndRejectsWrites(t *testing.T) {
db := testQuotaDB(t, 16)
if err := db.Create(&deviceModels.Device{ID: uuid.NewString(), Name: "已有设备", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1}).Error; err != nil {
t.Fatal(err)
}
if err := db.Delete(&Setting{}, SettingID).Error; err != nil {
t.Fatal(err)
}
response, err := (Service{DB: db}).Overview(PageRequest{})
if err != nil {
t.Fatal(err)
}
if response.Summary.ReadStatus != ReadStatusUnreadable || response.Summary.Used != 1 || len(response.List) != 1 {
t.Fatalf("read-only fallback failed: %#v", response)
}
if err = WithAvailableSlot(db, func(*gorm.DB) error { return nil }); err != ErrUnavailable {
t.Fatalf("write with unreadable quota error=%v", err)
}
}
func TestUpdateAllowsLowerLimitWithoutStoppingExistingDevices(t *testing.T) {
db := testQuotaDB(t, 16)
for index := 0; index < 2; index++ {
if err := db.Create(&deviceModels.Device{ID: uuid.NewString(), Name: "设备", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1}).Error; err != nil {
t.Fatal(err)
}
}
summary, err := (Service{DB: db}).Update(UpdateRequest{Limit: 1, Reason: "测试降低配额", Version: 1, UpdateBy: 7})
if err != nil {
t.Fatal(err)
}
if summary.Limit != 1 || summary.Used != 2 || summary.Remaining != 0 {
t.Fatalf("unexpected lowered quota: %#v", summary)
}
var devices int64
if err = db.Model(&deviceModels.Device{}).Count(&devices).Error; err != nil || devices != 2 {
t.Fatalf("existing devices changed: count=%d err=%v", devices, err)
}
var changes int64
if err = db.Model(&Change{}).Count(&changes).Error; err != nil || changes != 1 {
t.Fatalf("audit changes=%d err=%v", changes, err)
}
}
func TestConcurrentSlotReservationDoesNotExceedQuota(t *testing.T) {
db := testQuotaDB(t, 4)
const workers = 12
var wait sync.WaitGroup
var accepted int
var lock sync.Mutex
for index := 0; index < workers; index++ {
wait.Add(1)
go func() {
defer wait.Done()
err := WithAvailableSlot(db, func(tx *gorm.DB) error {
return tx.Create(&deviceModels.Device{ID: uuid.NewString(), Name: "并发设备", Modality: "video", CapabilitiesJSON: "[]", Status: "pending", AdapterStatus: "ready", Version: 1}).Error
})
if err == nil {
lock.Lock()
accepted++
lock.Unlock()
} else if err != ErrExceeded {
t.Errorf("unexpected reservation error: %v", err)
}
}()
}
wait.Wait()
if accepted != 4 {
t.Fatalf("accepted=%d want=4", accepted)
}
var count int64
if err := db.Model(&deviceModels.Device{}).Where("status <> ?", "disabled").Count(&count).Error; err != nil || count != 4 {
t.Fatalf("occupied=%d err=%v", count, err)
}
}
@@ -0,0 +1,62 @@
package version
import (
"fmt"
"runtime"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"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), migrateSenseOperations)
}
func migrateSenseOperations(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
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: "SenseOperations", Title: "运维中心", Icon: "operation", Path: "operations", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:operations:list", ParentId: root.MenuId, Component: "/sense/operations/index", Sort: 7, Visible: "0", IsFrame: "1"})
if err != nil {
return err
}
detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOperationsDetail", Title: "查看运维详情", MenuType: "F", Action: "GET", Permission: "sense:operations: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
}
retry, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOperationsRetry", Title: "受控重试", MenuType: "F", Action: "POST", Permission: "sense:operations:retry", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 2, Visible: "1", IsFrame: "1"})
if err != nil {
return err
}
readPolicies := [][2]string{{"/api/v1/operations", "GET"}, {"/api/v1/operations/:id", "GET"}}
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, detail}); err != nil {
return err
}
for _, policy := range readPolicies {
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
}
}
}
for _, role := range []string{"implementation_operator", "site_admin"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{retry}); err != nil {
return err
}
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: "/api/v1/operations/:id/retry", V2: "POST"}).Error; err != nil {
return err
}
}
if err = rebuildSenseMenuPaths(tx, root.MenuId, "/0"); err != nil {
return err
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,85 @@
package version
import (
"os"
"testing"
"gorm.io/driver/postgres"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
func TestOperationsMigrationAddsLeastPrivilegeMenus(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 = "2026082811000_operations.go"
if err = migrateSenseOperations(db, version); err != nil {
t.Fatal(err)
}
var menus, reads, retries, viewerRetries, applied int64
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOperations%").Count(&menus)
db.Model(&deviceCasbinRule{}).Where("v1 IN ?", []string{"/api/v1/operations", "/api/v1/operations/:id"}).Count(&reads)
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/operations/:id/retry", "POST").Count(&retries)
db.Model(&deviceCasbinRule{}).Where("v0 = ? AND v1 = ?", "viewer", "/api/v1/operations/:id/retry").Count(&viewerRetries)
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
if menus != 3 || reads != 6 || retries != 2 || viewerRetries != 0 || applied != 1 {
t.Fatalf("menus=%d reads=%d retries=%d viewerRetries=%d applied=%d", menus, reads, retries, viewerRetries, applied)
}
}
func TestOperationsMigrationOnPostgres(t *testing.T) {
dsn := os.Getenv("SENSE_OPERATIONS_MIGRATION_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_OPERATIONS_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_operations_74_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 = migrateSenseOperations(db, "2026082811000_operations.go"); err != nil {
t.Fatal(err)
}
var menus, policies int64
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOperations%").Count(&menus)
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ?", "/api/v1/operations%").Count(&policies)
if menus != 3 || policies != 8 {
t.Fatalf("menus=%d policies=%d", menus, policies)
}
}
@@ -0,0 +1,100 @@
package version
import (
"fmt"
"os"
"runtime"
"strconv"
"strings"
"time"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/quota"
"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), migrateSenseQuota)
}
func migrateSenseQuota(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&quota.Setting{}, &quota.Change{}); err != nil {
return err
}
now := time.Now().UTC()
setting := quota.Setting{ID: quota.SettingID, Limit: initialQuotaLimit(), Source: "migration", Version: 1}
setting.CreatedAt, setting.UpdatedAt = now, now
if err := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&setting).Error; 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: "SenseQuota", Title: "容量与配额", Icon: "data-line", Path: "quota",
Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:quota:list",
ParentId: root.MenuId, Component: "/sense/quota/index", Sort: 8, Visible: "0", IsFrame: "1",
})
if err != nil {
return err
}
update, err := ensureDeviceMenu(tx, migrationModels.SysMenu{
MenuName: "SenseQuotaUpdate", Title: "调整配额", MenuType: "F", Action: "PUT",
Permission: "sense:quota:update", ParentId: page.MenuId,
Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 1, Visible: "1", IsFrame: "1",
})
if err != nil {
return err
}
var devicePage migrationModels.SysMenu
if err = tx.Where("menu_name = ?", "SenseDeviceManage").First(&devicePage).Error; err != nil {
return err
}
enable, err := ensureDeviceMenu(tx, migrationModels.SysMenu{
MenuName: "SenseDeviceEnable", Title: "启用设备", MenuType: "F", Action: "PUT",
Permission: "sense:device:enable", ParentId: devicePage.MenuId,
Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, devicePage.MenuId), Sort: 5, 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}); err != nil {
return err
}
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: "/api/v1/quota", V2: "GET"}).Error; err != nil {
return err
}
}
if err = attachDeviceRole(tx, "site_admin", []migrationModels.SysMenu{update, enable}); err != nil {
return err
}
for _, policy := range [][2]string{{"/api/v1/quota", "PUT"}, {"/api/v1/devices/:id/enable", "PUT"}} {
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: "site_admin", 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
})
}
func initialQuotaLimit() int {
value, err := strconv.Atoi(strings.TrimSpace(os.Getenv("SENSE_PROVISIONING_QUOTA")))
if err != nil || value < 1 || value > 100000 {
return 16
}
return value
}
@@ -0,0 +1,90 @@
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/quota"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
func prepareQuotaMigrationDB(t *testing.T, db *gorm.DB) {
t.Helper()
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 := db.Create(&migrationModels.SysMenu{MenuName: "SenseDeviceManage", Title: "设备管理", Path: "devices", MenuType: "C", Component: "/sense/device/index"}).Error; err != nil {
t.Fatal(err)
}
}
func assertQuotaMigration(t *testing.T, db *gorm.DB, version string) {
t.Helper()
var setting quota.Setting
if err := db.First(&setting, "id = ?", quota.SettingID).Error; err != nil {
t.Fatal(err)
}
var menus, reads, updates, enables, viewerWrites, applied int64
db.Model(&migrationModels.SysMenu{}).Where("menu_name IN ?", []string{"SenseQuota", "SenseQuotaUpdate", "SenseDeviceEnable"}).Count(&menus)
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/quota", "GET").Count(&reads)
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/quota", "PUT").Count(&updates)
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/devices/:id/enable", "PUT").Count(&enables)
db.Model(&deviceCasbinRule{}).Where("v0 = ? AND v2 = ?", "viewer", "PUT").Count(&viewerWrites)
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
if setting.Limit != 16 || menus != 3 || reads != 3 || updates != 1 || enables != 1 || viewerWrites != 0 || applied != 1 {
t.Fatalf("setting=%#v menus=%d reads=%d updates=%d enables=%d viewerWrites=%d applied=%d", setting, menus, reads, updates, enables, viewerWrites, applied)
}
}
func TestQuotaMigrationAddsDefaultAndLeastPrivilegeMenus(t *testing.T) {
t.Setenv("SENSE_PROVISIONING_QUOTA", "")
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
prepareQuotaMigrationDB(t, db)
const version = "2026082812000_quota.go"
if err = migrateSenseQuota(db, version); err != nil {
t.Fatal(err)
}
assertQuotaMigration(t, db, version)
}
func TestQuotaMigrationOnPostgres(t *testing.T) {
dsn := os.Getenv("SENSE_QUOTA_MIGRATION_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_QUOTA_MIGRATION_TEST_DATABASE_URL to run the PostgreSQL migration test")
}
t.Setenv("SENSE_PROVISIONING_QUOTA", "")
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
const schema = "sense_quota_75_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 })
if err = db.Exec("SET search_path TO " + schema).Error; err != nil {
t.Fatal(err)
}
prepareQuotaMigrationDB(t, db)
const version = "2026082812000_quota.go"
if err = migrateSenseQuota(db, version); err != nil {
t.Fatal(err)
}
assertQuotaMigration(t, db, version)
}
@@ -0,0 +1,53 @@
package version
import (
"fmt"
"runtime"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/edge_node"
"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), migrateSenseEdgeNode)
}
func migrateSenseEdgeNode(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&edge_node.Node{}, &edge_node.Event{}); 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: "SenseEdgeNode", Title: "边缘节点", Icon: "monitor", Path: "edge-node", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:edge-node:list", ParentId: root.MenuId, Component: "/sense/edge-node/index", Sort: 9, Visible: "0", IsFrame: "1"})
if err != nil {
return err
}
detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseEdgeNodeDetail", Title: "查看节点详情", MenuType: "F", Action: "GET", Permission: "sense:edge-node: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
}
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, detail}); err != nil {
return err
}
for _, policy := range [][2]string{{"/api/v1/edge-nodes", "GET"}, {"/api/v1/edge-nodes/:id", "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
})
}
@@ -0,0 +1,90 @@
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/edge_node"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
func TestEdgeNodeMigrationAddsTablesAndReadOnlyPermissions(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 = "2026082813000_edge_node.go"
if err = migrateSenseEdgeNode(db, version); err != nil {
t.Fatal(err)
}
if !db.Migrator().HasTable(&edge_node.Node{}) || !db.Migrator().HasTable(&edge_node.Event{}) {
t.Fatal("edge node tables missing")
}
var menus, reads, writes, applied int64
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseEdgeNode%").Count(&menus)
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 = ?", "/api/v1/edge-nodes%", "GET").Count(&reads)
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 <> ?", "/api/v1/edge-nodes%", "GET").Count(&writes)
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
if menus != 2 || reads != 6 || writes != 0 || applied != 1 {
t.Fatalf("menus=%d reads=%d writes=%d applied=%d", menus, reads, writes, applied)
}
var fixtures int64
db.Model(&edge_node.Node{}).Count(&fixtures)
if fixtures != 0 {
t.Fatalf("migration seeded %d fixtures", fixtures)
}
}
func TestEdgeNodeMigrationOnPostgres(t *testing.T) {
dsn := os.Getenv("SENSE_EDGE_NODE_MIGRATION_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_EDGE_NODE_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_edge_node_76_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 = migrateSenseEdgeNode(db, "2026082813000_edge_node.go"); err != nil {
t.Fatal(err)
}
if !db.Migrator().HasTable(&edge_node.Node{}) || !db.Migrator().HasTable(&edge_node.Event{}) {
t.Fatal("edge node tables missing")
}
}
@@ -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
})
}
@@ -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")
}
}
@@ -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=
+79 -1
View File
@@ -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)
}
+3
View File
@@ -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')
+4
View File
@@ -20,6 +20,10 @@ export function disableDevice(id, data) {
return request({ url: `/api/v1/devices/${id}/disable`, method: 'put', data })
}
export function enableDevice(id, data) {
return request({ url: `/api/v1/devices/${id}/enable`, method: 'put', data })
}
export function updateDeviceCredentials(id, data) {
return request({ url: `/api/v1/devices/${id}/credentials`, method: 'put', data })
}
+9
View File
@@ -0,0 +1,9 @@
import request from '@/utils/request'
export function listEdgeNodes(query) {
return request({ url: '/api/v1/edge-nodes', method: 'get', params: query })
}
export function getEdgeNode(id) {
return request({ url: `/api/v1/edge-nodes/${encodeURIComponent(id)}`, method: 'get' })
}
+13
View File
@@ -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' })
}
+13
View File
@@ -0,0 +1,13 @@
import request from '@/utils/request'
export function listOperationProblems(query) {
return request({ url: '/api/v1/operations', method: 'get', params: query })
}
export function getOperationProblem(id) {
return request({ url: `/api/v1/operations/${encodeURIComponent(id)}`, method: 'get' })
}
export function retryOperationProblem(id, expectedVersion) {
return request({ url: `/api/v1/operations/${encodeURIComponent(id)}/retry`, method: 'post', data: { expectedVersion }})
}
+9
View File
@@ -0,0 +1,9 @@
import request from '@/utils/request'
export function getQuotaOverview(query) {
return request({ url: '/api/v1/quota', method: 'get', params: query })
}
export function updateQuota(data) {
return request({ url: '/api/v1/quota', method: 'put', data })
}
+10
View File
@@ -69,6 +69,7 @@
<el-button v-permisaction="['sense:device:credential']" type="primary" link size="small" :icon="Key" @click="handleCredential(scope.row)">更新凭据</el-button>
<el-divider direction="vertical" />
<el-button v-if="scope.row.status !== 'disabled'" v-permisaction="['sense:device:disable']" type="danger" link size="small" @click="handleDisable(scope.row)">停用</el-button>
<el-button v-else v-permisaction="['sense:device:enable']" type="success" link size="small" @click="handleEnable(scope.row)">启用</el-button>
</template>
</el-table-column>
</el-table>
@@ -141,6 +142,7 @@ import { Edit, Key, Plus, Refresh, Search } from '@element-plus/icons-vue'
import {
addDevice,
disableDevice,
enableDevice,
getDevice,
listDevices,
updateDevice,
@@ -284,6 +286,14 @@ export default {
this.msgSuccess(response.msg)
this.getList()
}).catch(() => {})
},
handleEnable(row) {
this.$confirm(`启用“${row.name}”会占用一路配额,并要求重新完成视频接入验证。是否继续?`, '启用设备', {
confirmButtonText: '确认启用', cancelButtonText: '取消', type: 'warning'
}).then(() => enableDevice(row.id, disableDevicePayload(row.version))).then(response => {
this.msgSuccess(response.msg)
this.getList()
}).catch(() => {})
}
}
}
@@ -0,0 +1,36 @@
const stateLabels = {
online: '在线',
offline: '离线',
recovering: '恢复中',
ready: '正常',
unavailable: '不可用',
idle: '无待回填',
pending: '等待回填'
}
const eventLabels = {
heartbeat_received: '首次心跳',
heartbeat_timeout: '心跳超时',
node_recovered: '心跳恢复',
recovery_started: '开始恢复'
}
export function stateLabel(value) { return stateLabels[value] || value || '未知' }
export function eventLabel(value) { return eventLabels[value] || value || '状态变化' }
export function stateType(value) { return ({ online: 'success', offline: 'danger', recovering: 'warning', ready: 'success', unavailable: 'danger', idle: 'info', pending: 'warning' })[value] || 'info' }
export function freshnessText(item) {
if (!item?.lastHeartbeatAt) return '从未收到心跳'
const seconds = Math.max(0, Number(item.freshnessSeconds) || 0)
if (item.stale) return `已陈旧 ${durationText(seconds)}`
return `${durationText(seconds)}前`
}
export function durationText(seconds) {
const value = Math.max(0, Number(seconds) || 0)
if (value < 60) return `${Math.floor(value)} 秒`
if (value < 3600) return `${Math.floor(value / 60)} 分钟`
if (value < 86400) return `${Math.floor(value / 3600)} 小时`
return `${Math.floor(value / 86400)} 天`
}
export function buildEdgeNodeQuery(query) {
return { pageIndex: Number(query.pageIndex) || 1, pageSize: Number(query.pageSize) || 10 }
}
@@ -0,0 +1,88 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div><h3>边缘节点</h3><p>查看各现场节点的运行负载、连接状态和待回填数据。离线时保留最后一次已知状态。</p></div>
<el-button :icon="Refresh" :loading="loading" @click="load">刷新状态</el-button>
</div>
<el-alert v-if="summary.offline" type="warning" :closable="false" show-icon class="status-alert">
<template #title>{{ summary.offline }} 个节点离线</template>
离线阈值为 {{ summary.heartbeatTimeoutSeconds }} 秒。页面中的隧道、视频和回填数据是节点离线前的最后已知状态。
</el-alert>
<el-row :gutter="12" class="summary-row" aria-label="边缘节点概览">
<el-col v-for="item in summaryCards" :key="item.label" :xs="12" :sm="6">
<div class="summary-item"><span>{{ item.label }}</span><strong :class="item.className">{{ item.value }}</strong></div>
</el-col>
</el-row>
<el-table v-loading="loading" :data="nodes" border stripe empty-text="暂无边缘节点">
<el-table-column label="节点" min-width="180">
<template #default="scope"><strong>{{ scope.row.name }}</strong><div class="muted">{{ scope.row.id }} · {{ scope.row.location || '未填写位置' }}</div></template>
</el-table-column>
<el-table-column label="状态" width="106" align="center"><template #default="scope"><el-tag :type="stateType(scope.row.status)">{{ stateLabel(scope.row.status) }}</el-tag></template></el-table-column>
<el-table-column label="最近心跳" min-width="150"><template #default="scope"><div>{{ freshnessText(scope.row) }}</div><small>{{ formatTime(scope.row.lastHeartbeatAt) }}</small></template></el-table-column>
<el-table-column label="负载" min-width="112"><template #default="scope">{{ scope.row.loadUsed }} / {{ scope.row.loadCapacity }} 路</template></el-table-column>
<el-table-column label="控制隧道" min-width="142"><template #default="scope"><el-tag size="small" :type="stateType(scope.row.controlTunnelStatus)">{{ stateLabel(scope.row.controlTunnelStatus) }}</el-tag><div class="muted">{{ scope.row.controlLatencyMs ? `${scope.row.controlLatencyMs} ms` : scope.row.controlTunnelDetail }}</div></template></el-table-column>
<el-table-column label="视频数据面" min-width="150"><template #default="scope"><el-tag size="small" :type="stateType(scope.row.videoPlaneStatus)">{{ stateLabel(scope.row.videoPlaneStatus) }}</el-tag><div class="muted">{{ scope.row.activeStreamCount }} 路视频 · {{ scope.row.videoPlaneDetail }}</div></template></el-table-column>
<el-table-column label="回填队列" min-width="128"><template #default="scope"><el-tag size="small" :type="stateType(scope.row.backfillStatus)">{{ stateLabel(scope.row.backfillStatus) }}</el-tag><div class="muted">{{ scope.row.backfillQueueDepth }} 条待回填</div></template></el-table-column>
<el-table-column label="版本" min-width="130"><template #default="scope"><div>{{ scope.row.runtimeVersion || '—' }}</div><small>{{ scope.row.modelVersion || '未部署模型' }}</small></template></el-table-column>
<el-table-column label="操作" width="92" fixed="right"><template #default="scope"><el-button v-permisaction="['sense:edge-node:detail']" type="primary" link @click="openDetail(scope.row.id)">详情</el-button></template></el-table-column>
</el-table>
<pagination v-show="total > 0" v-model:current-page="query.pageIndex" v-model:page-size="query.pageSize" :total="total" @pagination="load" />
</el-card>
<el-dialog v-model="detailOpen" title="边缘节点详情" width="min(800px, calc(100vw - 24px))" :close-on-click-modal="false">
<template v-if="selected">
<el-alert v-if="selected.stale" title="当前显示最后已知状态" type="warning" :closable="false" show-icon class="detail-alert">最近心跳距今 {{ durationText(selected.freshnessSeconds) }};连接恢复前,下列通道数据不会继续更新。</el-alert>
<el-descriptions :column="2" border>
<el-descriptions-item label="节点">{{ selected.name }}({{ selected.id }})</el-descriptions-item><el-descriptions-item label="位置">{{ selected.location || '未填写' }}</el-descriptions-item>
<el-descriptions-item label="运行版本">{{ selected.runtimeVersion || '—' }}</el-descriptions-item><el-descriptions-item label="模型版本">{{ selected.modelVersion || '未部署' }}</el-descriptions-item>
<el-descriptions-item label="运行时长">{{ durationText(selected.uptimeSeconds) }}</el-descriptions-item><el-descriptions-item label="最近心跳">{{ formatTime(selected.lastHeartbeatAt) }}</el-descriptions-item>
<el-descriptions-item label="控制隧道">{{ stateLabel(selected.controlTunnelStatus) }} · {{ selected.controlTunnelDetail }}</el-descriptions-item><el-descriptions-item label="视频数据面">{{ stateLabel(selected.videoPlaneStatus) }} · {{ selected.videoPlaneDetail }}</el-descriptions-item>
<el-descriptions-item label="回填队列">{{ stateLabel(selected.backfillStatus) }} · {{ selected.backfillQueueDepth }} 条</el-descriptions-item><el-descriptions-item label="恢复阶段">{{ recoveryText(selected) }}</el-descriptions-item>
</el-descriptions>
<h4 class="audit-title">状态审计</h4>
<el-timeline v-if="selected.events?.length">
<el-timeline-item v-for="event in selected.events" :key="event.id" :timestamp="formatTime(event.occurredAt)" placement="top"><strong>{{ eventLabel(event.eventType) }}</strong><div>{{ stateLabel(event.fromStatus) }} → {{ stateLabel(event.toStatus) }}</div><small>{{ event.detail }}</small></el-timeline-item>
</el-timeline>
<el-empty v-else description="暂无状态变化记录" :image-size="72" />
</template>
<template #footer><el-button @click="detailOpen = false">关闭</el-button></template>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { computed, onMounted, reactive, ref } from 'vue'
import { Refresh } from '@element-plus/icons-vue'
import { ElMessage } from 'element-plus'
import { getEdgeNode, listEdgeNodes } from '@/api/sense/edge-node'
import { buildEdgeNodeQuery, durationText, eventLabel, freshnessText, stateLabel, stateType } from './edgeNodeState'
defineOptions({ name: 'SenseEdgeNode' })
const loading = ref(false)
const nodes = ref([])
const total = ref(0)
const detailOpen = ref(false)
const selected = ref(null)
const query = reactive({ pageIndex: 1, pageSize: 10 })
const summary = reactive({ total: 0, online: 0, offline: 0, recovering: 0, loadUsed: 0, loadCapacity: 0, backfillQueueDepth: 0, heartbeatTimeoutSeconds: 90 })
const summaryCards = computed(() => [
{ label: '节点总数', value: summary.total }, { label: '在线 / 恢复中', value: `${summary.online} / ${summary.recovering}`, className: 'online-number' },
{ label: '总负载', value: `${summary.loadUsed} / ${summary.loadCapacity} 路` }, { label: '待回填', value: `${summary.backfillQueueDepth} 条`, className: summary.backfillQueueDepth ? 'warning-number' : '' }
])
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function formatTime(value) { return value ? new Date(value).toLocaleString('zh-CN', { hour12: false }) : '—' }
function recoveryText(item) { return item.status === 'recovering' ? (item.recoveryPhase || '通道收敛中') : '已收敛' }
async function load() { loading.value = true; try { const payload = unwrap(await listEdgeNodes(buildEdgeNodeQuery(query))) || {}; nodes.value = payload.list || []; total.value = payload.count || 0; Object.assign(summary, payload.summary || {}) } catch (error) { ElMessage.error(error.message || '边缘节点加载失败') } finally { loading.value = false } }
async function openDetail(id) { try { selected.value = unwrap(await getEdgeNode(id)); detailOpen.value = true } catch (error) { ElMessage.error(error.message || '节点详情加载失败') } }
onMounted(load)
</script>
<style scoped>
.page-header{display:flex;align-items:flex-start;justify-content:space-between;gap:16px}.page-header h3{margin:0 0 6px}.page-header p{margin:0;color:#909399}.status-alert{margin:16px 0}.summary-row{margin:16px 0 20px}.summary-item{display:flex;align-items:center;justify-content:space-between;min-height:72px;padding:12px 16px;border:1px solid #ebeef5;border-radius:4px}.summary-item span{color:#606266}.summary-item strong{font-size:20px}.online-number{color:#67c23a}.warning-number{color:#e6a23c}.muted,small{color:#909399;font-size:12px}.detail-alert{margin-bottom:16px}.audit-title{margin:20px 0 14px}@media(max-width:768px){.page-header{align-items:stretch;flex-direction:column}.summary-item{margin-bottom:8px;min-height:64px}.page-header .el-button{min-height:44px}}
</style>
@@ -0,0 +1,65 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div><h3>媒体分片</h3><p>查看 MediaMTX 分片容量、路径归属与故障影响。故障不会自动迁移路径。</p></div>
<el-button :icon="Refresh" :loading="loading" @click="load">刷新状态</el-button>
</div>
<el-alert v-if="summary.failed" type="error" :closable="false" show-icon class="status-alert">
<template #title>{{ summary.failed }} 个媒体分片异常,影响 {{ summary.impactedPaths }} 条路径</template>
路径仍保持原分片归属,便于定位故障范围;请先查看详情,不要直接变更生产配置。
</el-alert>
<el-row :gutter="12" class="summary-row" aria-label="媒体分片概览">
<el-col v-for="item in summaryCards" :key="item.label" :xs="12" :sm="6"><div class="summary-item"><span>{{ item.label }}</span><strong :class="item.className">{{ item.value }}</strong></div></el-col>
</el-row>
<el-table v-loading="loading" :data="shards" border stripe empty-text="暂无媒体分片运行记录">
<el-table-column label="分片" min-width="180"><template #default="scope"><strong>{{ scope.row.name }}</strong><div class="muted">{{ scope.row.id }} · {{ stateLabel(scope.row.mode) }}</div></template></el-table-column>
<el-table-column label="状态" width="112" align="center"><template #default="scope"><el-tag :type="stateType(scope.row.status)">{{ stateLabel(scope.row.status) }}</el-tag><div class="muted">{{ freshnessText(scope.row) }}</div></template></el-table-column>
<el-table-column label="已分配 / 容量" min-width="190"><template #default="scope"><div>{{ scope.row.assignedPaths }} / {{ scope.row.capacity }} 条路径</div><el-progress :percentage="usagePercent(scope.row)" :status="usagePercent(scope.row) >= 90 ? 'warning' : ''" /></template></el-table-column>
<el-table-column label="剩余容量" width="110" align="center"><template #default="scope">{{ scope.row.remaining }}</template></el-table-column>
<el-table-column prop="detail" label="运行说明" min-width="210" show-overflow-tooltip />
<el-table-column label="最近探测" min-width="170"><template #default="scope">{{ formatTime(scope.row.lastProbeAt) }}</template></el-table-column>
<el-table-column label="操作" width="156" fixed="right"><template #default="scope"><el-button v-permisaction="['sense:media-shard:detail']" type="primary" link @click="openDetail(scope.row.id)">影响详情</el-button><el-button v-permisaction="['sense:media-shard:preflight']" type="primary" link @click="openPreflight(scope.row.id)">迁移预检</el-button></template></el-table-column>
</el-table>
</el-card>
<el-dialog v-model="detailOpen" title="分片故障影响" width="min(880px, calc(100vw - 24px))" :close-on-click-modal="false">
<template v-if="selected"><el-descriptions :column="2" border><el-descriptions-item label="分片">{{ selected.name }}({{ selected.id }})</el-descriptions-item><el-descriptions-item label="状态"><el-tag :type="stateType(selected.status)">{{ stateLabel(selected.status) }}</el-tag></el-descriptions-item><el-descriptions-item label="路径占用">{{ selected.assignedPaths }} / {{ selected.capacity }}</el-descriptions-item><el-descriptions-item label="运行说明">{{ selected.detail }}</el-descriptions-item></el-descriptions>
<h4>受影响设备与 Profile</h4><el-table :data="selected.impact || []" border empty-text="此分片暂无路径归属"><el-table-column prop="deviceName" label="设备" min-width="150"><template #default="scope">{{ scope.row.deviceName || scope.row.deviceId }}<div class="muted">{{ scope.row.location || '未填写位置' }}</div></template></el-table-column><el-table-column prop="profileToken" label="Profile" min-width="130" /><el-table-column prop="path" label="媒体路径" min-width="180" /><el-table-column prop="actual" label="最近状态" width="120" /></el-table>
</template><template #footer><el-button @click="detailOpen = false">关闭</el-button></template>
</el-dialog>
<el-dialog v-model="preflightOpen" title="跨分片迁移预检(只读)" width="min(720px, calc(100vw - 24px))" :close-on-click-modal="false">
<template v-if="preflight"><el-alert :type="preflight.ready ? 'warning' : 'error'" :title="preflightConclusion(preflight)" :closable="false" show-icon /><el-descriptions :column="2" border class="preflight-summary"><el-descriptions-item label="源分片">{{ preflight.sourceShardId }}</el-descriptions-item><el-descriptions-item label="候选目标">{{ preflight.targetShardName || '无可用目标' }}</el-descriptions-item><el-descriptions-item label="影响路径">{{ preflight.impactedPaths }}</el-descriptions-item><el-descriptions-item label="执行授权">未授权</el-descriptions-item></el-descriptions><el-table :data="preflight.checks || []" border><el-table-column label="结果" width="90"><template #default="scope"><el-tag :type="scope.row.passed ? 'success' : 'danger'">{{ scope.row.passed ? '通过' : '未通过' }}</el-tag></template></el-table-column><el-table-column prop="name" label="检查项" width="130" /><el-table-column prop="detail" label="说明" min-width="260" /></el-table></template>
<template #footer><span class="read-only-note">本页面不会修改任何路径归属。</span><el-button @click="preflightOpen = false">关闭</el-button></template>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { computed, onMounted, reactive, ref } from 'vue'
import { Refresh } from '@element-plus/icons-vue'
import { ElMessage } from 'element-plus'
import { getMediaShard, listMediaShards, preflightMediaShard } from '@/api/sense/media-shard'
import { freshnessText, preflightConclusion, stateLabel, stateType, usagePercent } from './mediaShardState'
defineOptions({ name: 'SenseMediaShard' })
const loading = ref(false); const shards = ref([]); const selected = ref(null); const preflight = ref(null); const detailOpen = ref(false); const preflightOpen = ref(false)
const summary = reactive({ total: 0, available: 0, failed: 0, configuredCapacity: 0, assignedPaths: 0, impactedPaths: 0 })
const summaryCards = computed(() => [{ label: '分片总数', value: summary.total }, { label: '可用 / 异常', value: `${summary.available} / ${summary.failed}`, className: summary.failed ? 'danger-number' : 'success-number' }, { label: '已分配 / 总容量', value: `${summary.assignedPaths} / ${summary.configuredCapacity}` }, { label: '故障影响', value: `${summary.impactedPaths} 条路径`, className: summary.impactedPaths ? 'danger-number' : '' }])
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function formatTime(value) { return value ? new Date(value).toLocaleString('zh-CN', { hour12: false }) : '—' }
async function load() { loading.value = true; try { const payload = unwrap(await listMediaShards()) || {}; shards.value = payload.list || []; Object.assign(summary, payload.summary || {}) } catch (error) { ElMessage.error(error.message || '媒体分片加载失败') } finally { loading.value = false } }
async function openDetail(id) { try { selected.value = unwrap(await getMediaShard(id)); detailOpen.value = true } catch (error) { ElMessage.error(error.message || '分片详情加载失败') } }
async function openPreflight(id) { try { preflight.value = unwrap(await preflightMediaShard(id)); preflightOpen.value = true } catch (error) { ElMessage.error(error.message || '迁移预检失败') } }
onMounted(load)
</script>
<style scoped>
.page-header{display:flex;align-items:flex-start;justify-content:space-between;gap:16px}.page-header h3{margin:0 0 6px}.page-header p{margin:0;color:#909399}.status-alert{margin:16px 0}.summary-row{margin:16px 0 20px}.summary-item{display:flex;align-items:center;justify-content:space-between;min-height:72px;padding:12px 16px;border:1px solid #ebeef5;border-radius:4px}.summary-item span{color:#606266}.summary-item strong{font-size:20px}.success-number{color:#67c23a}.danger-number{color:#f56c6c}.muted{color:#909399;font-size:12px;margin-top:4px}.preflight-summary{margin:16px 0}.read-only-note{margin-right:16px;color:#909399}@media(max-width:768px){.page-header{align-items:stretch;flex-direction:column}.summary-item{margin-bottom:8px;min-height:64px}.page-header .el-button{min-height:44px}}
</style>
@@ -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 '预检通过,但仍需单独建立高风险工单并人工确认' }
@@ -0,0 +1,170 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div>
<h3>运维中心</h3>
<p>集中查看设备、媒体和可选本地推理的期望态、实际态与未收敛问题。</p>
</div>
<el-button :icon="Refresh" :loading="loading" @click="load">刷新状态</el-button>
</div>
<el-alert type="info" :closable="false" show-icon class="boundary-alert">
<template #title>本地推理是可选能力</template>
Brain 未安装时显示“不可用(未安装)”,不会导致页面失败,也不会阻断设备和媒体运维。运维异常不会发送到 Bell 作为业务预警。
</el-alert>
<el-row :gutter="12" class="summary-row" aria-label="运维概览">
<el-col :xs="24" :sm="8">
<div class="summary-item"><span>设备与媒体已收敛</span><strong>{{ summary.convergedCount }} / {{ summary.managedCount }}</strong></div>
</el-col>
<el-col :xs="24" :sm="8">
<div class="summary-item"><span>待处理问题</span><strong class="warning-number">{{ summary.actionableProblems }}</strong></div>
</el-col>
<el-col :xs="24" :sm="8">
<div class="summary-item"><span>本地推理适配器</span><strong class="adapter-state">{{ inferenceState }}</strong></div>
</el-col>
</el-row>
<div class="section-header">
<div><h4>问题队列</h4><p>默认只读诊断;只有明确可重试的项目才显示受控重试。</p></div>
</div>
<el-form :model="query" label-width="76px" class="filter-form" @submit.prevent="search">
<el-form-item label="对象类型">
<el-select v-model="query.objectType" clearable placeholder="全部">
<el-option label="设备" value="device" /><el-option label="媒体" value="media" /><el-option label="本地推理" value="local_inference" />
</el-select>
</el-form-item>
<el-form-item label="问题类型">
<el-select v-model="query.problemType" clearable placeholder="全部">
<el-option label="认证失败" value="authentication_failed" /><el-option label="退避等待" value="backoff_wait" />
<el-option label="时间漂移" value="clock_drift" /><el-option label="孤儿安全闸" value="orphan_safety_gate" />
<el-option label="能力未安装" value="capability_unavailable" /><el-option label="尚未收敛" value="not_converged" />
</el-select>
</el-form-item>
<el-form-item label="级别">
<el-select v-model="query.severity" clearable placeholder="全部">
<el-option label="高" value="high" /><el-option label="中" value="medium" /><el-option label="提示" value="info" />
</el-select>
</el-form-item>
<el-form-item label="关键词"><el-input v-model="query.keyword" clearable placeholder="设备、位置或差异" @keyup.enter="search" /></el-form-item>
<el-form-item class="filter-actions">
<el-button type="primary" :icon="Search" native-type="submit">查询</el-button>
<el-button :icon="RefreshLeft" @click="reset">重置</el-button>
</el-form-item>
</el-form>
<el-table v-loading="loading" :data="problems" border stripe empty-text="当前筛选条件下没有未收敛问题">
<el-table-column label="对象" min-width="170">
<template #default="scope"><strong>{{ scope.row.objectName }}</strong><div class="muted">{{ objectLabel(scope.row.objectType) }} · {{ scope.row.location || scope.row.objectId }}</div></template>
</el-table-column>
<el-table-column label="问题" min-width="138"><template #default="scope"><div>{{ problemLabel(scope.row.problemType) }}</div><small>{{ scope.row.id }}</small></template></el-table-column>
<el-table-column label="期望态" min-width="110"><template #default="scope">{{ stateLabel(scope.row.expected) }}</template></el-table-column>
<el-table-column label="实际态" min-width="120"><template #default="scope"><el-tag size="small" :type="severityType(scope.row.severity)">{{ stateLabel(scope.row.actual) }}</el-tag></template></el-table-column>
<el-table-column prop="difference" label="未收敛差异" min-width="210" show-overflow-tooltip />
<el-table-column label="下次动作" min-width="190"><template #default="scope"><div>{{ scope.row.nextAction }}</div><small v-if="scope.row.nextRetryAt">自动重试:{{ formatTime(scope.row.nextRetryAt) }}</small></template></el-table-column>
<el-table-column label="级别" width="78" align="center"><template #default="scope"><el-tag size="small" :type="severityType(scope.row.severity)">{{ severityLabel(scope.row.severity) }}</el-tag></template></el-table-column>
<el-table-column label="操作" width="156" fixed="right">
<template #default="scope">
<el-button v-permisaction="['sense:operations:detail']" type="primary" link @click="openDetail(scope.row.id)">详情</el-button>
<el-button v-if="canRetry(scope.row)" v-permisaction="['sense:operations:retry']" type="warning" link @click="confirmRetry(scope.row)">受控重试</el-button>
<span v-else-if="scope.row.retryInProgress" class="muted">重试已排队</span>
</template>
</el-table-column>
</el-table>
<pagination v-show="total > 0" v-model:current-page="query.pageIndex" v-model:page-size="query.pageSize" :total="total" @pagination="load" />
<p class="safety-help">受控重试会校验权限、当前版本和并发占用,并写入操作审计。孤儿资源仅隔离和提示,本页没有删除操作。</p>
</el-card>
<el-dialog v-model="detailOpen" title="状态详情" width="min(760px, calc(100vw - 24px))" :close-on-click-modal="false">
<el-descriptions v-if="selected" :column="2" border>
<el-descriptions-item label="问题编号">{{ selected.id }}</el-descriptions-item>
<el-descriptions-item label="对象">{{ selected.objectName }}</el-descriptions-item>
<el-descriptions-item label="期望态">{{ stateLabel(selected.expected) }}</el-descriptions-item>
<el-descriptions-item label="实际态">{{ stateLabel(selected.actual) }}</el-descriptions-item>
<el-descriptions-item label="未收敛差异" :span="2">{{ selected.difference }}</el-descriptions-item>
<el-descriptions-item label="退避 / 下次重试">{{ selected.nextRetryAt ? formatTime(selected.nextRetryAt) : '无退避' }}</el-descriptions-item>
<el-descriptions-item label="尝试次数">{{ selected.attemptCount }}</el-descriptions-item>
<el-descriptions-item label="建议处理" :span="2">{{ selected.nextAction }}</el-descriptions-item>
</el-descriptions>
<el-alert v-if="selected" title="安全闸已启用" type="warning" :closable="false" show-icon class="safety-alert">{{ selected.safetyGate }}</el-alert>
<p class="boundary-text">Brain/Bell 状态不会参与本地设备和媒体重试判定。</p>
<template #footer>
<el-button @click="detailOpen = false">关闭</el-button>
<el-button v-if="canRetry(selected)" v-permisaction="['sense:operations:retry']" type="warning" @click="confirmRetry(selected)">受控重试</el-button>
</template>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { onMounted, reactive, ref } from 'vue'
import { Refresh, RefreshLeft, Search } from '@element-plus/icons-vue'
import { ElMessage, ElMessageBox } from 'element-plus'
import { getOperationProblem, listOperationProblems, retryOperationProblem } from '@/api/sense/operations'
import { buildOperationsQuery, canRetry, objectLabel, problemLabel, severityLabel, severityType, stateLabel } from './operationsState'
defineOptions({ name: 'SenseOperations' })
const loading = ref(false)
const problems = ref([])
const total = ref(0)
const summary = reactive({ managedCount: 0, convergedCount: 0, actionableProblems: 0, localInference: 'unavailable' })
const query = reactive({ pageIndex: 1, pageSize: 10, objectType: '', problemType: '', severity: '', keyword: '' })
const detailOpen = ref(false)
const selected = ref(null)
const inferenceState = ref('不可用(未安装)')
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function formatTime(value) { return value ? new Date(value).toLocaleString('zh-CN', { hour12: false }) : '—' }
async function load() {
loading.value = true
try {
const payload = unwrap(await listOperationProblems(buildOperationsQuery(query))) || {}
problems.value = payload.list || []
total.value = payload.count || 0
Object.assign(summary, payload.summary || { managedCount: 0, convergedCount: 0, actionableProblems: 0 })
inferenceState.value = summary.localInference === 'configured' ? '已配置' : '不可用(未安装)'
} catch (error) {
ElMessage.error(error.message || '运维状态加载失败')
} finally {
loading.value = false
}
}
function search() { query.pageIndex = 1; load() }
function reset() { Object.assign(query, { pageIndex: 1, objectType: '', problemType: '', severity: '', keyword: '' }); load() }
async function openDetail(id) {
try {
selected.value = unwrap(await getOperationProblem(id))
detailOpen.value = true
} catch (error) {
ElMessage.error(error.message || '状态详情加载失败')
}
}
async function confirmRetry(item) {
try {
await ElMessageBox.confirm(`将对“${item.objectName}”发起一次重试。此操作不会修改凭据,也不会删除媒体资源。`, '确认受控重试', {
type: 'warning', confirmButtonText: '确认重试', cancelButtonText: '取消'
})
await retryOperationProblem(item.id, item.version)
ElMessage.success('重试任务已排队')
detailOpen.value = false
await load()
} catch (error) {
if (error !== 'cancel' && error !== 'close') ElMessage.warning(error.message || '重试未排队,请刷新状态后再试')
}
}
onMounted(load)
</script>
<style scoped>
.page-header,.section-header{display:flex;align-items:flex-start;justify-content:space-between;gap:16px}.page-header h3,.section-header h4{margin:0 0 6px}.page-header p,.section-header p{margin:0;color:#909399}.boundary-alert{margin:16px 0}.summary-row{margin-bottom:20px}.summary-item{display:flex;align-items:center;justify-content:space-between;min-height:72px;padding:12px 16px;border:1px solid #ebeef5;border-radius:4px}.summary-item span{color:#606266}.summary-item strong{font-size:22px}.summary-item .warning-number{color:#e6a23c}.summary-item .adapter-state{font-size:15px;color:#909399}.filter-form{display:flex;align-items:flex-end;flex-wrap:wrap;gap:0 12px;margin:16px 0 2px}.filter-form .el-form-item{width:210px}.filter-form .filter-actions{width:auto}.filter-form :deep(.el-select){width:100%}.muted,small{color:#909399;font-size:12px}.safety-help,.boundary-text{margin:12px 0 0;color:#909399;font-size:12px}.safety-alert{margin-top:16px}@media(max-width:768px){.filter-form .el-form-item{width:100%}.page-header{align-items:stretch;flex-direction:column}.summary-item{margin-bottom:8px}}
</style>
@@ -0,0 +1,50 @@
const objectLabels = {
device: '设备',
media: '媒体',
local_inference: '本地推理'
}
const problemLabels = {
authentication_failed: '认证失败',
backoff_wait: '退避等待',
clock_drift: '时间漂移',
orphan_safety_gate: '孤儿安全闸',
capability_unavailable: '能力未安装',
not_converged: '尚未收敛'
}
const stateLabels = {
active: '已启用',
ready: '就绪',
running: '运行',
stopped: '已停止',
waiting: '等待拉流',
pending: '等待处理',
retry_pending: '重试已排队',
authentication_failed: '认证失败',
clock_drift: '时间漂移',
apply_failed: '配置失败',
process_unavailable: '进程不可用',
status_unavailable: '状态不可用',
unavailable: '不可用(未安装)'
}
export function objectLabel(value) { return objectLabels[value] || value || '未知' }
export function problemLabel(value) { return problemLabels[value] || value || '未知问题' }
export function stateLabel(value) {
return String(value || '').split(' / ').map(item => stateLabels[item] || item).join(' / ') || '未知'
}
export function severityLabel(value) { return ({ high: '高', medium: '中', info: '提示' })[value] || value || '未知' }
export function severityType(value) { return ({ high: 'danger', medium: 'warning', info: 'info' })[value] || 'info' }
export function canRetry(item) { return Boolean(item?.retryable && !item?.retryInProgress && item?.version > 0) }
export function buildOperationsQuery(query) {
return {
pageIndex: Number(query.pageIndex) || 1,
pageSize: Number(query.pageSize) || 10,
objectType: String(query.objectType || '').trim(),
problemType: String(query.problemType || '').trim(),
severity: String(query.severity || '').trim(),
keyword: String(query.keyword || '').trim()
}
}
+177
View File
@@ -0,0 +1,177 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div>
<h3>容量与配额</h3>
<p>查看当前占用与安全闸状态。配额是交付配置,不等于单机性能承诺。</p>
</div>
<div class="header-actions">
<el-button v-permisaction="['sense:device:add']" :icon="Plus" :disabled="!gate.writable" @click="router.push('/sense/devices')">新增设备</el-button>
<el-button
v-permisaction="['sense:quota:update']"
type="primary"
:icon="Setting"
:disabled="summary.readStatus !== 'readable'"
@click="openQuotaDialog"
>调整配额</el-button>
</div>
</div>
<el-alert :title="gate.title" :type="gate.type" :closable="false" show-icon class="gate-alert">
{{ gate.detail }}
</el-alert>
<el-row :gutter="12" class="summary-row" aria-label="容量摘要">
<el-col :xs="24" :sm="12" :lg="6"><div class="summary-item"><span>当前配置配额</span><strong>{{ readableLimit }}</strong><small>来源:{{ sourceLabel }}</small></div></el-col>
<el-col :xs="24" :sm="12" :lg="6"><div class="summary-item"><span>已占用设备</span><strong>{{ summary.used }} 路</strong><small>已接入和待接入均占用</small></div></el-col>
<el-col :xs="24" :sm="12" :lg="6"><div class="summary-item"><span>剩余可用</span><strong>{{ readableRemaining }}</strong><small>写入时再次原子校验</small></div></el-col>
<el-col :xs="24" :sm="12" :lg="6"><div class="summary-item"><span>配额读取状态</span><strong class="read-state"><el-tag :type="summary.readStatus === 'readable' ? 'success' : 'danger'">{{ summary.readStatus === 'readable' ? '可读取' : '读取失败' }}</el-tag></strong><small>最后读取:{{ formatTime(summary.lastReadAt) }}</small></div></el-col>
</el-row>
<section class="usage-section">
<div class="section-header"><h4>占用情况</h4><span>{{ usageCaption }}</span></div>
<el-progress :percentage="percent" :status="gate.type === 'error' ? 'exception' : gate.type === 'warning' ? 'warning' : ''" />
</section>
<div class="section-header"><div><h4>设备占用明细</h4><p>停用设备不占用配额;重新启用时会执行安全闸校验。</p></div></div>
<el-form :model="query" class="filter-form" label-width="76px" @submit.prevent="search">
<el-form-item label="状态">
<el-select v-model="query.status" clearable placeholder="全部状态">
<el-option label="已接入" value="active" /><el-option label="待接入" value="pending" /><el-option label="已停用" value="disabled" />
</el-select>
</el-form-item>
<el-form-item label="设备名称"><el-input v-model="query.keyword" clearable placeholder="请输入设备名称或位置" @keyup.enter="search" /></el-form-item>
<el-form-item class="filter-actions"><el-button type="primary" :icon="Search" native-type="submit">查询</el-button><el-button :icon="RefreshLeft" @click="reset">重置</el-button></el-form-item>
</el-form>
<el-table v-loading="loading" :data="devices" border stripe empty-text="当前没有设备">
<el-table-column prop="name" label="设备名称" min-width="160" />
<el-table-column prop="location" label="安装位置" min-width="180" show-overflow-tooltip />
<el-table-column label="接入状态" width="110" align="center"><template #default="scope"><el-tag :type="scope.row.status === 'active' ? 'success' : scope.row.status === 'disabled' ? 'info' : 'warning'">{{ deviceStatusLabel(scope.row.status) }}</el-tag></template></el-table-column>
<el-table-column label="配额占用" width="110" align="center"><template #default="scope">{{ scope.row.occupied ? '占用 1 路' : '不占用' }}</template></el-table-column>
<el-table-column label="更新时间" width="180" align="center"><template #default="scope">{{ formatTime(scope.row.updatedAt) }}</template></el-table-column>
<el-table-column label="操作" width="100" fixed="right">
<template #default="scope">
<el-button
v-if="scope.row.status === 'disabled'"
v-permisaction="['sense:device:enable']"
type="success"
link
:disabled="!gate.writable"
@click="enable(scope.row)"
>启用</el-button>
<el-button v-else type="primary" link @click="router.push('/sense/devices')">查看</el-button>
</template>
</el-table-column>
</el-table>
<pagination v-show="total > 0" v-model:current-page="query.pageIndex" v-model:page-size="query.pageSize" :total="total" @pagination="load" />
<div class="section-header tier-header"><div><h4>交付档位与验证状态</h4><p>“未验证”不代表系统承诺可稳定承载,必须先完成目标硬件性能测试。</p></div></div>
<el-table :data="tiers" border>
<el-table-column label="档位" width="100"><template #default="scope">{{ scope.row.limit }} 路</template></el-table-column>
<el-table-column label="当前配置" width="120"><template #default="scope"><el-tag v-if="scope.row.configured" type="success">当前配置</el-tag><span v-else>未配置</span></template></el-table-column>
<el-table-column label="性能测试" width="120"><template #default="scope"><el-tag :type="scope.row.validation === 'verified' ? 'success' : 'warning'">{{ tierValidationLabel(scope.row.validation) }}</el-tag></template></el-table-column>
<el-table-column prop="deliveryStatus" label="说明" min-width="220" />
</el-table>
</el-card>
<el-dialog v-model="dialogOpen" title="调整配额" width="min(520px, calc(100vw - 24px))" :close-on-click-modal="false" @closed="resetDialog">
<el-alert title="调整前请确认目标硬件已经完成容量测试" type="warning" :closable="false" show-icon class="dialog-alert">
降低到当前占用以下不会中断已有视频,但会阻止新增与启用,直到占用回到配额内。
</el-alert>
<el-form ref="quotaFormRef" :model="quotaForm" :rules="quotaRules" label-width="110px">
<el-form-item label="新配额(路)" prop="limit"><el-input-number v-model="quotaForm.limit" :min="1" :max="100000" controls-position="right" /></el-form-item>
<el-form-item label="变更原因" prop="reason"><el-input v-model.trim="quotaForm.reason" maxlength="512" show-word-limit placeholder="请输入变更原因" /></el-form-item>
</el-form>
<template #footer><el-button @click="dialogOpen = false">取消</el-button><el-button type="primary" :loading="saving" @click="saveQuota">确认调整</el-button></template>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { computed, onMounted, reactive, ref } from 'vue'
import { useRouter } from 'vue-router'
import { Plus, RefreshLeft, Search, Setting } from '@element-plus/icons-vue'
import { ElMessage, ElMessageBox } from 'element-plus'
import { getQuotaOverview, updateQuota } from '@/api/sense/quota'
import { enableDevice } from '@/api/sense/device'
import { buildQuotaQuery, deviceStatusLabel, quotaGate, tierValidationLabel, usagePercent } from './quotaState'
defineOptions({ name: 'SenseQuota' })
const loading = ref(false)
const router = useRouter()
const saving = ref(false)
const devices = ref([])
const tiers = ref([])
const total = ref(0)
const summary = reactive({ limit: 0, used: 0, remaining: 0, readStatus: 'unreadable', source: '', version: 0 })
const query = reactive({ pageIndex: 1, pageSize: 10, keyword: '', status: '' })
const dialogOpen = ref(false)
const quotaFormRef = ref()
const quotaForm = reactive({ limit: 16, reason: '' })
const quotaRules = {
limit: [{ required: true, type: 'number', min: 1, max: 100000, message: '请输入 1 至 100000 的整数配额', trigger: 'change' }],
reason: [{ required: true, message: '请输入变更原因,以便审计追溯', trigger: 'blur' }]
}
const gate = computed(() => quotaGate(summary))
const percent = computed(() => usagePercent(summary))
const readableLimit = computed(() => summary.readStatus === 'readable' ? `${summary.limit} 路` : '—')
const readableRemaining = computed(() => summary.readStatus === 'readable' ? `${summary.remaining} 路` : '—')
const sourceLabel = computed(() => summary.readStatus === 'readable' ? (summary.source === 'migration' ? '初始化配置' : '数据库配置') : '不可读取')
const usageCaption = computed(() => summary.readStatus === 'readable' ? `已占用 ${summary.used} / ${summary.limit} 路(${percent.value}%)` : `已占用 ${summary.used} 路;配额上限不可读取`)
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function formatTime(value) { return value ? new Date(value).toLocaleString('zh-CN', { hour12: false }) : '—' }
async function load() {
loading.value = true
try {
const payload = unwrap(await getQuotaOverview(buildQuotaQuery(query))) || {}
Object.assign(summary, payload.summary || {})
devices.value = payload.list || []
tiers.value = payload.tiers || []
total.value = payload.count || 0
} catch (error) {
ElMessage.error(error.message || '容量与配额加载失败')
} finally {
loading.value = false
}
}
function search() { query.pageIndex = 1; load() }
function reset() { Object.assign(query, { pageIndex: 1, keyword: '', status: '' }); load() }
function openQuotaDialog() { quotaForm.limit = summary.limit; quotaForm.reason = ''; dialogOpen.value = true }
function resetDialog() { quotaForm.reason = ''; quotaFormRef.value?.clearValidate() }
async function enable(row) {
try {
await ElMessageBox.confirm(`启用“${row.name}”会占用一路配额,并要求重新完成视频接入验证。是否继续?`, '启用设备', {
type: 'warning', confirmButtonText: '确认启用', cancelButtonText: '取消'
})
await enableDevice(row.id, { version: row.version })
ElMessage.success('设备已启用,请重新完成视频接入验证')
await load()
} catch (error) {
if (error !== 'cancel' && error !== 'close') ElMessage.warning(error.message || '设备启用失败')
}
}
async function saveQuota() {
if (!await quotaFormRef.value.validate().catch(() => false)) return
saving.value = true
try {
await updateQuota({ limit: quotaForm.limit, reason: quotaForm.reason, version: summary.version })
ElMessage.success('配额已更新')
dialogOpen.value = false
await load()
} catch (error) {
ElMessage.warning(error.message || '配额更新失败,请刷新后重试')
} finally {
saving.value = false
}
}
onMounted(load)
</script>
<style scoped>
.page-header,.section-header{display:flex;align-items:flex-start;justify-content:space-between;gap:16px}.page-header h3,.section-header h4{margin:0 0 6px}.page-header p,.section-header p{margin:0;color:#909399}.header-actions{display:flex;gap:8px;flex-wrap:wrap}.gate-alert{margin:16px 0}.summary-row{margin-bottom:18px}.summary-item{display:flex;flex-direction:column;min-height:108px;padding:14px 16px;border:1px solid #ebeef5;border-radius:4px}.summary-item span,.summary-item small{color:#606266}.summary-item strong{margin:8px 0 4px;font-size:24px}.summary-item .read-state{font-size:16px}.usage-section{margin:4px 0 22px;padding:16px;border:1px solid #ebeef5}.filter-form{display:flex;align-items:flex-end;flex-wrap:wrap;gap:0 12px;margin:16px 0 2px}.filter-form .el-form-item{width:250px}.filter-form .filter-actions{width:auto}.filter-form :deep(.el-select){width:100%}.tier-header{margin-top:24px}.dialog-alert{margin-bottom:18px}@media(max-width:768px){.page-header{align-items:stretch;flex-direction:column}.filter-form .el-form-item{width:100%}.summary-item{margin-bottom:8px}}
</style>
@@ -0,0 +1,48 @@
export function buildQuotaQuery(query) {
return {
pageIndex: Math.max(1, Number(query.pageIndex) || 1),
pageSize: Math.min(100, Math.max(1, Number(query.pageSize) || 10)),
keyword: String(query.keyword || '').trim(),
status: String(query.status || '').trim()
}
}
export function quotaGate(summary = {}) {
if (summary.readStatus !== 'readable') {
return {
type: 'error',
title: '配额配置不可读取',
detail: '为避免超配额,新增和启用请求已关闭;已有设备、视频流与查询不受影响。请检查配置来源后重试。',
writable: false
}
}
if ((Number(summary.remaining) || 0) <= 0) {
return {
type: 'warning',
title: '当前配额已满',
detail: '新增和启用请求将被拒绝;已有设备、视频流与查询继续正常运行。',
writable: false
}
}
return {
type: 'success',
title: '安全闸正常',
detail: '新增和启用设备可继续;系统会在写入时再次原子校验剩余配额。',
writable: true
}
}
export function usagePercent(summary = {}) {
const limit = Number(summary.limit) || 0
const used = Number(summary.used) || 0
if (limit <= 0) return 0
return Math.min(100, Math.max(0, Math.round((used / limit) * 100)))
}
export function deviceStatusLabel(status) {
return ({ active: '已接入', pending: '待接入', disabled: '已停用' })[status] || status || '未知'
}
export function tierValidationLabel(value) {
return value === 'verified' ? '已验证' : '未验证'
}
@@ -0,0 +1,14 @@
import request from '@/utils/request'
import { getEdgeNode, listEdgeNodes } from '@/api/sense/edge-node'
jest.mock('@/utils/request', () => jest.fn())
describe('Sense edge node API', () => {
beforeEach(() => request.mockReset())
test('uses read-only endpoints and encodes identifiers', () => {
listEdgeNodes({ pageIndex: 1, pageSize: 10 })
getEdgeNode('edge/site 01')
expect(request).toHaveBeenNthCalledWith(1, { url: '/api/v1/edge-nodes', method: 'get', params: { pageIndex: 1, pageSize: 10 } })
expect(request).toHaveBeenNthCalledWith(2, { url: '/api/v1/edge-nodes/edge%2Fsite%2001', method: 'get' })
})
})
@@ -0,0 +1,18 @@
import { buildEdgeNodeQuery, durationText, eventLabel, freshnessText, stateLabel, stateType } from '@/views/sense/edge-node/edgeNodeState'
describe('Sense edge node presentation state', () => {
test('uses operator-facing connection and channel labels', () => {
expect(stateLabel('offline')).toBe('离线')
expect(stateLabel('recovering')).toBe('恢复中')
expect(stateType('offline')).toBe('danger')
expect(eventLabel('node_recovered')).toBe('心跳恢复')
})
test('makes stale last-known data explicit', () => {
expect(freshnessText({ stale: true, freshnessSeconds: 245, lastHeartbeatAt: '2026-08-28T00:00:00Z' })).toBe('已陈旧 4 分钟')
expect(freshnessText({ stale: false, freshnessSeconds: 18, lastHeartbeatAt: '2026-08-28T00:00:00Z' })).toBe('18 秒前')
expect(durationText(172800)).toBe('2 天')
})
test('sends only allowlisted pagination fields', () => {
expect(buildEdgeNodeQuery({ pageIndex: '2', pageSize: 20, ignored: 'not-sent' })).toEqual({ pageIndex: 2, pageSize: 20 })
})
})
@@ -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) })
})
@@ -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('不满足') })
})
@@ -0,0 +1,17 @@
import request from '@/utils/request'
import { getOperationProblem, listOperationProblems, retryOperationProblem } from '@/api/sense/operations'
jest.mock('@/utils/request', () => jest.fn())
describe('Sense operations API', () => {
beforeEach(() => request.mockReset())
test('encodes identifiers and sends only the expected version for controlled retry', () => {
listOperationProblems({ pageIndex: 1, pageSize: 10 })
getOperationProblem('media:camera/unsafe')
retryOperationProblem('media:camera/unsafe', 7)
expect(request).toHaveBeenNthCalledWith(1, { url: '/api/v1/operations', method: 'get', params: { pageIndex: 1, pageSize: 10 }})
expect(request).toHaveBeenNthCalledWith(2, { url: '/api/v1/operations/media%3Acamera%2Funsafe', method: 'get' })
expect(request).toHaveBeenNthCalledWith(3, { url: '/api/v1/operations/media%3Acamera%2Funsafe/retry', method: 'post', data: { expectedVersion: 7 }})
})
})
@@ -0,0 +1,23 @@
import { buildOperationsQuery, canRetry, objectLabel, problemLabel, severityType, stateLabel } from '@/views/sense/operations/operationsState'
describe('Sense operations presentation state', () => {
test('uses operator-facing state labels', () => {
expect(objectLabel('local_inference')).toBe('本地推理')
expect(problemLabel('orphan_safety_gate')).toBe('孤儿安全闸')
expect(stateLabel('unavailable')).toBe('不可用(未安装)')
expect(stateLabel('active / ready')).toBe('已启用 / 就绪')
expect(severityType('high')).toBe('danger')
})
test('allows retry only when the server projection says it is safe', () => {
expect(canRetry({ retryable: true, retryInProgress: false, version: 2 })).toBe(true)
expect(canRetry({ retryable: true, retryInProgress: true, version: 2 })).toBe(false)
expect(canRetry({ retryable: false, retryInProgress: false, version: 2 })).toBe(false)
})
test('builds an allowlisted query', () => {
expect(buildOperationsQuery({ pageIndex: 2, pageSize: 20, objectType: ' device ', problemType: 'clock_drift', severity: ' medium ', keyword: ' 南门 ', ignored: 'not-sent' })).toEqual({
pageIndex: 2, pageSize: 20, objectType: 'device', problemType: 'clock_drift', severity: 'medium', keyword: '南门'
})
})
})
@@ -0,0 +1,21 @@
import { buildQuotaQuery, quotaGate, usagePercent } from '@/views/sense/quota/quotaState'
describe('Sense quota state', () => {
test('keeps pagination independent from the configured quota', () => {
expect(buildQuotaQuery({ pageIndex: 2, pageSize: 64, keyword: ' 东门 ', status: 'active' })).toEqual({
pageIndex: 2, pageSize: 64, keyword: '东门', status: 'active'
})
})
test('closes writes when full or unreadable while preserving an actionable explanation', () => {
expect(quotaGate({ readStatus: 'readable', remaining: 4 }).writable).toBe(true)
expect(quotaGate({ readStatus: 'readable', remaining: 0 })).toMatchObject({ writable: false, title: '当前配额已满' })
expect(quotaGate({ readStatus: 'unreadable', remaining: 4 })).toMatchObject({ writable: false, title: '配额配置不可读取' })
})
test('bounds the occupancy percentage without treating quota as a list limit', () => {
expect(usagePercent({ used: 12, limit: 16 })).toBe(75)
expect(usagePercent({ used: 20, limit: 16 })).toBe(100)
expect(usagePercent({ used: 12, limit: 0 })).toBe(0)
})
})
+64 -2
View File
@@ -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: 3c95336ab96aa1e245f67c598299f2904f0e4437
synchronized_at: 2026-08-27T11:07:58Z
wiki_revision: 578ddbaae3d7e846037d958085e40609cd398bef
synchronized_at: 2026-08-28T08:02:32Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -174,3 +174,65 @@ Bell 使用独立 PostgreSQL、JWT realm、token key 和首次迁移管理员环
- `dev_scripts/harness.py sync` 只从 Gitea Wiki 导出 `wiki-docs.json` 映射的核心镜像;`sync --check` 只检查一致性。
- `archive`、`export` 和 `export --all` 只在人工明确要求时处理可选任务快照,任务快照不进入核心映射。
- `dev_scripts/wiki_docs.py` 负责 Gitea Wiki 读取、revision、镜像头、脏文件保护和安全路径校验。
<!-- sense-provisioning:start -->
## Sense 批量开通代码入口
工单 #72 的领域代码位于 `Sense/server/app/sense/provisioning/`,GoAdmin JWT/Casbin/PermissionAction 路由注册在 `Sense/server/app/admin/router/sense_provisioning.go`,迁移位于 `Sense/server/cmd/migrate/migration/version/2026082809000_provisioning.go`。前端页面为 `Sense/ui/src/views/sense/provisioning/index.vue`,API 封装为 `Sense/ui/src/api/sense/provisioning.js`;页面复用 BasicLayout、动态菜单、Axios、Element Plus Steps/Upload/Form/Table/Pagination/Dialog/Tag/Alert/Button 和权限指令。
PostgreSQL 表 `sense_provisioning_batches` 保存幂等键、配额快照和汇总状态,`sense_provisioning_items` 保存行号、设备信息、逐项状态、失败原因、尝试次数和已建立的设备引用;两表都没有凭据字段。执行链按条目条件领取 ready/failed 状态,调用既有 Device、Credential Vault 与 Admission 服务,成功项不整体回滚。API 根路径为 `/api/v1/provisioning/batches`,覆盖创建/预校验、列表、详情、执行、失败项重试、单项重试和无秘密 CSV 导出。
<!-- sense-provisioning:end -->
<!-- sense-local-events:start -->
## Sense 本地事件只读链路
- 内部模型、DTO、服务、API 与合成夹具:`Sense/server/app/sense/local_event/`。
- GoAdmin 路由:`Sense/server/app/admin/router/sense_local_event.go`,仅开放 `GET /api/v1/local-events` 和 `GET /api/v1/local-events/:id`。
- PostgreSQL 表:`sense_local_event_candidates`;迁移同时建立“本地事件”菜单、详情功能权限和 implementation_operator/site_admin/viewer 的只读 Casbin 策略。
- go-admin-ui 页面:`Sense/ui/src/views/sense/local-event/index.vue`;复用 BasicLayout、Element Plus 表单、表格、分页、Dialog、Tag、权限指令与 Axios 封装。
- API 响应只包含 Sense 内部候选、匿名源引用、规则引用、证据状态和保留信息,不包含 Bell/Brain schema 或送达字段。页面上的 Bell 送达状态固定解释为“不适用”,避免把本地候选冒充外部预警。
- 列表和详情成功/失败读取会同步写入脱敏的 `sys_opera_log`,不依赖可关闭的全局数据库日志开关;记录只含动作、路由、操作者、结果和候选 ID,不保存筛选值、响应或证据内容。拒绝访问继续由认证/RBAC 身份审计记录。
<!-- sense-local-events:end -->
<!-- sense-operations:start -->
## Sense 运维中心代码路径
- 后端入口:`Sense/server/app/admin/router/sense_operations.go`。
- 领域投影:`Sense/server/app/sense/operations/`,只读聚合 `sense_devices`、`sense_admission_results` 和 `sense_media_routes`,不建立第二套运维状态表。
- API:`GET /api/v1/operations`、`GET /api/v1/operations/:id`、`POST /api/v1/operations/:id/retry`;均复用 GoAdmin JWT、Casbin、PermissionAction 和响应封装。
- 设备重试先用版本 CAS 取得在途所有权,再调用既有 admission Probe;失败时释放在途闸。媒体重试以 CAS 写入 `retry_pending` 和到期时间,由既有 MediaMTX 对账循环执行。
- 前端:`Sense/ui/src/views/sense/operations/` 与 `Sense/ui/src/api/sense/operations.js`,复用 BasicLayout、Element Plus 表单、表格、分页、Dialog、Tag 和权限按钮。
- 本模块不导入 Brain/Bell 模型,不访问其数据库,不定义共享契约。
<!-- sense-operations:end -->
<!-- sense-quota:start -->
## Sense 容量与配额代码路径
- 领域模型、DTO、服务、API 与测试:`Sense/server/app/sense/quota/`。
- GoAdmin 路由:`Sense/server/app/admin/router/sense_quota.go`;API 为 `GET /api/v1/quota` 和 `PUT /api/v1/quota`,复用 JWT、Casbin、PermissionAction 和 GoAdmin 响应封装。
- PostgreSQL 表:`sense_quota_settings` 保存单行配置与版本,`sense_quota_changes` 保存不可变变更记录。迁移 `2026082812000_quota.go` 初始化默认 16、动态菜单、最小权限及设备启用权限。
- 设备新增和重新启用由 `quota.WithAvailableSlot` 在事务内锁定配置行、计算非停用设备占用并执行写入;设备服务与批量开通服务共用该事实源。
- 前端:`Sense/ui/src/views/sense/quota/` 与 `Sense/ui/src/api/sense/quota.js`;复用 BasicLayout、Axios、Element Plus Alert/Row/Progress/Form/Table/Pagination/Dialog/Tag/Button 和权限指令。
- 本模块不导入 Brain/Bell 模型,不访问其数据库,也不定义跨项目配额契约。
<!-- sense-quota:end -->
<!-- sense-edge-nodes:start -->
## Sense 边缘节点代码路径
- 领域投影、DTO、服务、审计、合成夹具与测试:`Sense/server/app/sense/edge_node/`。
- GoAdmin 路由:`Sense/server/app/admin/router/sense_edge_node.go`;只开放 `GET /api/v1/edge-nodes` 与 `GET /api/v1/edge-nodes/:id`,复用 JWT、Casbin、PermissionAction 和统一响应封装。
- PostgreSQL 表:`sense_edge_nodes` 保存最后已知投影,`sense_edge_node_events` 保存不可变状态转换;迁移 `2026082813000_edge_node.go` 建立动态菜单及 implementation_operator/site_admin/viewer 的只读权限,不写入夹具。
- 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-edge-nodes:end -->
<!-- sense-media-shards:start -->
## 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 建立只读权限;生产迁移和启动不写入合成分片。
<!-- sense-media-shards:end -->
+58 -2
View File
@@ -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: 96fca4d2aa411725282bd3911c133efc75f7146a
synchronized_at: 2026-08-27T09:04:59Z
wiki_revision: 999eb1aee3558ff75cfa929acac77841af20b742
synchronized_at: 2026-08-28T08:02:42Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -154,3 +154,59 @@ synchronized_at: 2026-08-27T09:04:59Z
- Gitea Wiki 只保存长期有效的项目事实;核心页面由 `wiki-docs.json` 显式映射。
- `docs/task/` 是按人工明确要求形成的专项或历史兼容快照,可能不完整或不是最新状态,不得替代工单。
- 默认不创建、导出或更新任务快照;导出过程不自动删除本地历史文件。
<!-- sense-local-events:start -->
## Sense 本地事件候选规则
- “本地事件候选”是 Sense 内部匿名记录,不是 Bell Event 或 Alert;任何跨项目输出必须等待版本化协调契约。
- 候选状态与证据状态是两组独立状态,不能用证据成功推断事件已确认,也不能用事件已确认推断证据成功。
- 列表和详情只读;新增、确认、删除、重试或送达不属于本页面。
- 每条记录保存明确的 `retain_until`,页面展示实际到期时间;保留时长由服务端生产者/策略决定,页面不虚构全局固定天数。
- 合成夹具只由测试显式装载,生产启动和迁移都不会自动写入假事件。
- GET 列表和详情继续经过 Sense JWT、Casbin RBAC、数据权限中间件及系统操作审计。
<!-- sense-local-events:end -->
<!-- sense-operations:start -->
## Sense 运维状态规则
- 运维问题是设备、媒体或本地推理适配器的本地状态,不是 Event 或 Bell Alert;不得进入 Bell 业务预警队列。
- 问题必须同时保留期望态、实际态、差异、建议动作、对象版本和安全闸;退避问题还要展示下次重试时间。
- 设备认证失败和时间漂移复用既有视频接入探测;媒体失败复用既有 MediaMTX 对账循环。
- 受控重试必须校验对象版本并防止同一对象并发执行;实施/运维和站点管理员可执行,viewer 只读。
- 孤儿媒体路由只能显示“隔离待确认”,不得从运维中心自动删除。
- Brain 未配置或未安装时显示本地推理 unavailable,不读取 Brain 数据库,也不阻断设备/媒体运维。
<!-- sense-operations:end -->
<!-- sense-quota:start -->
## Sense 容量与配额规则
- 配额是 Sense 单产品的交付配置,不是单机性能承诺;默认初始化为 16 路,数据库、循环、分页和列表容量不得以 16 为硬上限。
- 配额占用按非停用设备计算:已接入和待接入各占用一路,停用设备不占用;重新启用必须重新经过配额安全闸并回到待接入状态。
- 设备新增、重新启用及批量开通写入必须使用同一 PostgreSQL 配额事实源。写事务锁定单行配额配置、读取当前占用、检查剩余量后再写设备,防止并发超配。
- 配额缺失、非法或读取失败时,新增、启用和批量相关写入必须拒绝;已有设备、视频流和只读查询继续可用。
- 降低配额到当前占用以下不会自动停用设备或中断视频;剩余量按 0 显示,后续新增/启用持续拒绝,直到占用回到配额内。
- 32/64/128 只展示当前配置和目标硬件性能测试状态。未验证不得解释为支持或稳定承载承诺。
- 配额调整必须记录旧值、新值、操作人、时间和原因;viewer 与实施/运维角色只读,只有站点管理员可调整配额和重新启用设备。
<!-- sense-quota:end -->
<!-- sense-edge-nodes:start -->
## Sense 边缘节点状态规则
- **边缘节点投影**:Sense 保存的节点最后已知运行状态,不是 Brain/Bell 的共享机器身份或跨产品事实源。
- **在线**:最近心跳距读取时刻不超过 90 秒,且没有未完成的恢复阶段。
- **离线**:最近心跳距读取时刻超过 90 秒。离线是读取时推导状态,不得用“未知”覆盖最后已知的控制隧道、视频数据面、负载和回填队列;所有这些值必须同时标明陈旧。
- **恢复中**:心跳已经恢复,但控制隧道、视频数据面或回填通道仍在独立收敛。只有恢复阶段明确收敛后才显示在线。
- 状态转换写入不可变事件;列表和详情读取写入脱敏操作审计,不记录凭据或原始心跳载荷。
- 合成夹具必须显式调用,生产迁移和启动不得自动写入节点。
<!-- sense-edge-nodes:end -->
<!-- sense-media-shards:start -->
## 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 不停止外部实例。
<!-- sense-media-shards:end -->
+118 -2
View File
@@ -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: e1755deff68c95188e6bf2f759c7ec93293c8715
synchronized_at: 2026-08-27T15:22:48Z
wiki_revision: 05435d53528271a866c525655486689afc198762
synchronized_at: 2026-08-28T08:02:52Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
@@ -399,3 +399,119 @@ D:\supervisor\supervisord.exe ctl /c D:\supervisor\supervisord.conf reload
当前 Go Supervisor 的 `reload` 会重新读取独立配置;实际受影响实例必须以命令输出和 reload 前后 PID 为准。切换托管前先停止占用 Sense 端口的非 Supervisor 实例,防止自动启动进入 Backoff。验证至少包含 Supervisor 状态为 Running、`http://127.0.0.1:18080/health` 与首页返回 200,以及受控重启后 Sense 和受管 MediaMTX PID 均更新。
<!-- sense-supervisor:end -->
<!-- sense-local-events:start -->
### Sense 本地事件验证
合成数据通过 `local_event.SeedSyntheticFixture` 在测试中显式装载,不存在生产自动开关,也不要求 Brain 或 Bell 进程。定向及回归命令:
```powershell
cd Sense/server
go test -race ./app/sense/local_event
go test ./...
go vet ./...
go build ./...
cd ../ui
corepack pnpm@9.15.1 lint
corepack pnpm@9.15.1 test:unit
corepack pnpm@9.15.1 build:prod
```
PostgreSQL 迁移测试使用独立 schema;连接值只放当前进程环境,不写入仓库或日志:
```powershell
$env:SENSE_LOCAL_EVENT_MIGRATION_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
go test ./cmd/migrate/migration/version -run TestLocalEventMigrationOnPostgres -count=1 -v
```
验证应覆盖:candidate/confirmed 与 pending/success/failed 的独立组合、分页和时间/状态/关键词筛选、详情及保留提示、viewer/实施/站点管理员只读权限、操作审计,以及 Brain/Bell 均未启动时的合成候选 smoke。没有可用 PostgreSQL 测试连接时必须明确记录为未验证,不得以 SQLite 单测替代 PostgreSQL 结论。
<!-- sense-local-events:end -->
<!-- sense-operations:start -->
### Sense 运维中心验证
定向验证:
```powershell
cd Sense/server
go test -race ./app/sense/operations
go test ./app/admin/router ./cmd/migrate/migration/version
go test ./...
go vet ./...
go build ./...
cd ../ui
corepack pnpm@9.15.1 lint
corepack pnpm@9.15.1 test:unit
corepack pnpm@9.15.1 build:prod
```
合成 smoke 应覆盖设备认证失败、设备时间漂移、媒体退避、孤儿路由、本地推理 unavailable、筛选分页、详情、设备同步探测重试、媒体排队重试、版本冲突、在途防重和脱敏审计。Brain/Bell 不应启动或成为测试依赖。真实摄像头与 MediaMTX 故障恢复仍须在获准环境验证,日志不得记录设备地址、Stream URI 或凭据。
<!-- sense-operations:end -->
<!-- sense-edge-nodes:start -->
### Sense 边缘节点验证
定向验证:
```powershell
cd Sense/server
go test -race ./app/sense/edge_node
go test ./cmd/migrate/migration/version -run EdgeNode
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
```
PostgreSQL 迁移测试只在提供隔离连接时执行:
```powershell
$env:SENSE_EDGE_NODE_MIGRATION_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
go test ./cmd/migrate/migration/version -run TestEdgeNodeMigrationOnPostgres -count=1 -v
```
验证应覆盖在线、超过 90 秒的离线、心跳恢复但通道未收敛、最后已知状态与陈旧时长、负载和回填队列、不可变状态事件、只读权限、脱敏读取审计,以及 Brain/Bell 不运行时的独立性。合成节点只能由测试或显式开发调用;没有隔离 PostgreSQL 连接时必须记录迁移真库测试未执行。
<!-- sense-edge-nodes:end -->
<!-- sense-media-shards:start -->
### 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 = '<mediamtx.exe 的绝对路径>'
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 连接时必须记录真库迁移测试未执行。
<!-- sense-media-shards:end -->
+44 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Product-Requirements
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Product-Requirements.-
wiki_revision: e912a1ca1410e01680a0f11f6199ccb42cd8fe8f
synchronized_at: 2026-08-27T09:06:08Z
wiki_revision: dcbdcf563017a1749fa76ad5f78c74a2c3cc6be1
synchronized_at: 2026-08-28T06:16:00Z
<!-- gitea-wiki-mirror:end -->
# 产品需求
@@ -190,3 +190,45 @@ YoVision 首个可交付目标是在民办寄宿学校以默认 16 路高风险
- Sense 账号不能登录 Bell,Bell 账号不能登录 Sense。
- 三个项目独立构建、测试和版本化;Sense/Bell 独立迁移、备份、恢复和打包。
- 集成断开不阻断各自核心能力,恢复后按 Outbox 和幂等收据继续。
<!-- sense-provisioning:start -->
## Sense 批量开通长期规则
SEN-004 已在工单 #72 实现、合入 `dev` 并通过用户验收。Sense 可独立导入 CSV/XLSX 摄像头清单,服务端预校验地址、重复项和剩余配额,再按条目执行接入。批次保留创建时的配额与占用快照,但新的批次与设备写入必须读取 SEN-011 的统一配额事实源;16 不是数据库、循环、分页或单机容量硬上限。
导入清单禁止账号和密码字段。ONVIF/RTSP 凭据只在受控表单与单次执行请求中短暂存在,随后进入既有加密 Vault;批次、条目、响应、日志和结果导出都不得包含凭据。批次允许部分成功,成功设备不因其他条目失败而回滚;批量重试与单项重试只领取失败条目,并通过批次幂等键、持久化设备引用和条件状态更新避免重复创建设备。
<!-- sense-provisioning:end -->
<!-- sense-local-events:start -->
## Sense 本地事件候选
SEN-008 在 Sense 内提供只读的本地事件候选查询。网管或非技术人员可按发生时间、候选状态、证据状态、规则引用和关键词筛选,查看候选详情、证据处理结果与逐条保留到期时间。
候选状态只表示 Sense 内部的 `candidate`(候选)或 `confirmed`(已确认事件);证据状态独立表示 `pending`(处理中)、`success`(成功)或 `failed`(失败)。本能力在 Brain、Bell 均未启动时仍可使用。本地候选不是跨项目事件契约,不等同于 Bell Alert,也不表示已经向 Bell 送达。
<!-- sense-local-events:end -->
<!-- sense-operations:start -->
## Sense 运维中心
SEN-009 由 Sense 独立提供设备、媒体和可选本地推理的运维概览。页面以“期望态—实际态—未收敛差异—下一步动作”展示认证失败、退避等待、时间漂移、孤儿安全闸和能力未安装,面向网管或非技术运维人员,不暴露凭据或内部数据库结构。
运维中心在 Brain、Bell 均未启动时可用。本地推理适配器尚未建立协调契约时显示 unavailable;该状态不导致页面失败。所有问题均为 Sense 运维事实,不等同于 Bell 业务 Alert,也不会从本页面发送给 Bell。
读取默认只读。受控重试仅对明确可重试的设备/媒体问题开放,必须通过 RBAC、对象版本和在途状态校验并写入脱敏审计。孤儿媒体资源只隔离和提示,不提供自动删除。
<!-- sense-operations:end -->
<!-- sense-quota:start -->
## Sense 容量与配额长期规则
SEN-011 将默认 16 路实现为可配置的 Sense 本地交付配额。统一配置保存在 PostgreSQL;新增设备、重新启用和批量开通在写事务内原子校验。配额不可读时拒绝这些写入,但已有设备、视频流和读取保持可用。降低配额不会自动停用现有设备。
容量页面只展示 16/32/64/128 的当前配置和性能测试状态。除已经验证的 16 路学校试点基线外,其余档位在完成目标硬件压测前都不得解释为单机承载承诺。
<!-- sense-quota:end -->
<!-- sense-edge-nodes:start -->
## Sense 边缘节点资产与离线状态
Sense 为网管和非技术运维人员提供只读的“边缘节点”页面,集中查看节点名称、位置、运行版本、模型版本、运行时长、承载路数、控制隧道、视频数据面和回填队列。节点最近心跳超过 90 秒即在读取时显示为离线;离线不会清空最后一次采集到的通道和队列状态,页面必须同时标明陈旧时长。心跳恢复后,连接状态与控制、视频、回填通道分别收敛,未全部收敛时显示“恢复中”。列表、详情及状态变化均可审计。
本能力只管理 Sense 自有投影,不建立 Brain/Bell 共享身份、控制协议或跨项目回填执行;Brain、Bell 未运行时仍可独立查看。合成节点仅供显式开发和测试,不由生产启动或迁移自动写入。
<!-- sense-edge-nodes:end -->
+34 -2
View File
@@ -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: 367c6aed6eaeece793622431c7d1b11b1b9dbec5
synchronized_at: 2026-08-27T09:07:20Z
wiki_revision: ce246054849b5b797dc4da0a55116238a4c00d7e
synchronized_at: 2026-08-28T08:04:52Z
<!-- gitea-wiki-mirror:end -->
# YoVision 部署与运维
@@ -80,3 +80,35 @@ Sense\start_sense.bat
- MediaMTX 启用时进程、路径和播放链路状态可定位。
- Supervisor 不会与手工进程重复占用端口。
- 日志和文档没有秘密;未验证的 Brain、Bell、真机或生产行为明确标注。
<!-- sense-provisioning:start -->
## Sense 批量开通与容量配额配置排错
数据库迁移 `2026082812000_quota.go` 首次创建统一配额配置:读取当次迁移进程的可选 `SENSE_PROVISIONING_QUOTA` 正整数作为初值,未设置或无效时使用 16。迁移完成后,运行期配额以 PostgreSQL `sense_quota_settings` 为事实源,并由“容量与配额”页面受控调整;后续修改环境变量不会覆盖数据库值。每个批量开通批次仍记录创建时的配额与占用快照。
升级后看不到“容量与配额”时,确认迁移成功并重新登录刷新动态菜单。implementation_operator、site_admin、viewer 可读取容量;只有 site_admin 可调整配额和重新启用设备。调整必须填写原因。降低配额不会停止已有流;当占用达到或超过配额时,新增和启用返回冲突。
页面显示“配额配置不可读取”时,先确认 `sense_quota_settings` 的 ID 1 记录存在、limit 为 1–100000 的整数且数据库可读;不要通过手工插入设备绕过安全闸。配置不可读时已有流和查询应继续,新设备、重新启用和批量开通写入会返回服务不可用。32/64/128 的“未验证”状态不能作为容量承诺。
批量导入被拒绝时还应检查模板只含 line_number/name/location/address,地址必须为不带账号、查询参数或片段的 HTTP(S) ONVIF 地址。条目失败时按页面原因检查网络、获准网段、凭据和接入状态;仅重试失败项,不删除已成功设备。日志、导出和问题记录不得粘贴摄像头凭据。
<!-- sense-provisioning:end -->
<!-- sense-edge-nodes:start -->
## Sense 边缘节点状态排错
升级后看不到“边缘节点”菜单时,先确认数据库迁移 `2026082813000_edge_node.go` 已成功,再重新登录或刷新动态菜单。implementation_operator、site_admin、viewer 均只有列表和详情读取权限,本模块没有网页写入、删除或远程控制接口。
节点超过 90 秒没有心跳即显示离线。离线时页面继续显示最后已知的隧道、视频、负载和回填状态,并标明陈旧时长,这不表示通道仍实时可用。心跳恢复但通道未收敛时显示“恢复中”,应分别检查控制隧道、视频数据面和回填队列。当前版本只提供 Sense 内部投影入口,不包含外部心跳接入协议或跨节点回填执行;没有节点数据时不会自动生成演示节点。
<!-- sense-edge-nodes:end -->
<!-- sense-media-shards:start -->
## 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/路径,不要手工改数据库归属。迁移预检不会执行迁移;任何实际跨分片迁移都必须另建高风险工单和回退方案。
<!-- sense-media-shards:end -->