Compare commits

..
Author SHA1 Message Date
QiuSW 02a5af5e3b feat: 建立 Sense 可靠投递状态机 (#78) 2026-08-28 17:22:44 +08:00
ila e06904272a Merge PR #121: [SEN] MediaMTX 分片分配与故障范围可观察
实现工单 #77;保持工单待人工验收。
2026-08-28 16:10:21 +08:00
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
51 changed files with 2868 additions and 18 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'
@@ -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,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/outbox"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseOutboxRouter) }
func registerSenseOutboxRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &outbox.API{}
r := v1.Group("/outbox").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("", api.List)
r.GET("/:id", api.Get)
r.POST("/:id/requeue", api.Requeue)
}
@@ -0,0 +1,26 @@
package router
import (
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"net/http"
"testing"
)
func TestSenseOutboxRoutes(t *testing.T) {
gin.SetMode(gin.TestMode)
engine := gin.New()
registerSenseOutboxRouter(engine.Group("/api/v1"), &jwt.GinJWTMiddleware{})
wanted := map[string]bool{http.MethodGet + " /api/v1/outbox": false, http.MethodGet + " /api/v1/outbox/:id": false, http.MethodPost + " /api/v1/outbox/:id/requeue": 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,29 @@
package local_event
import (
"context"
"encoding/json"
"fmt"
"time"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
)
// CreateWithOutbox commits the local candidate and its internal delivery
// record atomically. The payload remains Sense-internal and is not a Bell or
// Brain contract.
func CreateWithOutbox(ctx context.Context, db *gorm.DB, candidate EventCandidate, payload map[string]interface{}, now time.Time) error {
encoded, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("encode local event outbox payload: %w", err)
}
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&candidate).Error; err != nil {
return fmt.Errorf("create local event candidate: %w", err)
}
_, err = outbox.Enqueue(tx, outbox.EnqueueInput{InternalType: "local_event_candidate", BusinessRef: candidate.ID, IdempotencyKey: "local-event:" + candidate.ID + ":v1", PayloadJSON: encoded}, now)
return err
})
}
@@ -0,0 +1,46 @@
package local_event
import (
"context"
"testing"
"time"
"github.com/google/uuid"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
)
func TestCreateWithOutboxCommitsAndRollsBackAtomically(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file:"+uuid.NewString()+"?mode=memory&cache=shared"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
sqlDB, _ := db.DB()
sqlDB.SetMaxOpenConns(1)
if err = db.AutoMigrate(&EventCandidate{}, &outbox.Message{}, &outbox.DeliveryRecord{}, &outbox.Attempt{}); err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 28, 10, 0, 0, 0, time.UTC)
candidate := EventCandidate{ID: uuid.NewString(), OccurredAt: now, SourceRef: "SEN-CAM-01", RuleRef: "rule-1", RuleName: "区域闯入", CandidateState: CandidateStateCandidate, EvidenceState: EvidenceStatePending, RetainUntil: now.Add(24 * time.Hour)}
if err = CreateWithOutbox(context.Background(), db, candidate, map[string]interface{}{"eventId": candidate.ID}, now); err != nil {
t.Fatal(err)
}
var candidates, messages int64
db.Model(&EventCandidate{}).Count(&candidates)
db.Model(&outbox.Message{}).Count(&messages)
if candidates != 1 || messages != 1 {
t.Fatalf("candidates=%d messages=%d", candidates, messages)
}
duplicate := candidate
duplicate.ID = candidate.ID
if err = CreateWithOutbox(context.Background(), db, duplicate, map[string]interface{}{"eventId": duplicate.ID}, now); err == nil {
t.Fatal("expected duplicate transaction failure")
}
db.Model(&EventCandidate{}).Count(&candidates)
db.Model(&outbox.Message{}).Count(&messages)
if candidates != 1 || messages != 1 {
t.Fatalf("atomic rollback failed candidates=%d messages=%d", candidates, messages)
}
}
+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)
}
}
+103
View File
@@ -0,0 +1,103 @@
package outbox
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"
)
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
}
response, 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, "读取可靠投递详情 "+response.Message.ID)
e.OK(response, "查询成功")
}
func (e *API) Requeue(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
request := RequeueRequest{}
if err = e.MakeContext(c).Bind(&request, binding.JSON).Errors; err != nil {
e.audit(c, service, "Requeue", auditFailure, "死信重新排队请求格式不正确")
e.Error(http.StatusBadRequest, err, "请求格式不正确")
return
}
message, err := service.Requeue(c.Param("id"), request, user.GetUserId(c))
if err != nil {
e.audit(c, service, "Requeue", auditFailure, "死信重新排队被拒绝")
e.writeError(err)
return
}
e.audit(c, service, "Requeue", auditSuccess, "死信已重新排队 "+message.ID)
e.OK(message, "已重新排队")
}
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("outbox audit failed: %s", err.Error())
}
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrInvalidInput):
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, ErrNotFound):
e.Error(http.StatusNotFound, err, err.Error())
case errors.Is(err, ErrNotDead), errors.Is(err, ErrVersionConflict):
e.Error(http.StatusConflict, err, err.Error())
default:
e.Error(http.StatusInternalServerError, err, "可靠投递操作失败")
}
}
+23
View File
@@ -0,0 +1,23 @@
package outbox
import (
"time"
"gorm.io/gorm"
adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
)
const auditSuccess, auditFailure = "1", "2"
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: "outbox.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
}
+41
View File
@@ -0,0 +1,41 @@
package outbox
import commonDTO "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
type PageRequest struct {
commonDTO.Pagination `search:"-"`
State string `form:"state"`
InternalType string `form:"internalType"`
Keyword string `form:"keyword"`
}
type Summary struct {
Pending int64 `json:"pending"`
Retry int64 `json:"retry"`
Processing int64 `json:"processing"`
Dead int64 `json:"dead"`
}
type PageResponse struct {
List []Message `json:"list"`
Count int64 `json:"count"`
Summary Summary `json:"summary"`
}
type DetailResponse struct {
Message Message `json:"message"`
Attempts []Attempt `json:"attempts"`
}
type RequeueRequest struct {
ExpectedVersion int64 `json:"expectedVersion" binding:"required,min=1"`
Reason string `json:"reason" binding:"required"`
}
type EnqueueInput struct {
InternalType string
BusinessRef string
IdempotencyKey string
PayloadJSON []byte
MaxAttempts int
}
+58
View File
@@ -0,0 +1,58 @@
package outbox
import "time"
const (
StatePending = "pending"
StateProcessing = "processing"
StateRetry = "retry"
StateDead = "dead"
StateDelivered = "delivered"
)
// Message is a Sense-internal delivery record. PayloadJSON is deliberately
// excluded from management APIs and is not a cross-product contract.
type Message struct {
ID string `gorm:"size:36;primaryKey" json:"id"`
InternalType string `gorm:"size:64;not null;index" json:"internalType"`
BusinessRef string `gorm:"size:128;not null;index" json:"businessRef"`
IdempotencyKey string `gorm:"size:191;not null;uniqueIndex" json:"idempotencyKey"`
PayloadJSON string `gorm:"column:payload;type:jsonb;not null" json:"-"`
State string `gorm:"size:24;not null;index" json:"state"`
AttemptCount int `gorm:"not null;default:0" json:"attemptCount"`
MaxAttempts int `gorm:"not null;default:12" json:"maxAttempts"`
AvailableAt time.Time `gorm:"not null;index" json:"availableAt"`
LeaseOwner string `gorm:"size:128;not null;default:''" json:"leaseOwner,omitempty"`
LeaseUntil *time.Time `gorm:"index" json:"leaseUntil,omitempty"`
LastError string `gorm:"size:512;not null;default:''" json:"lastError,omitempty"`
DeliveredAt *time.Time `json:"deliveredAt,omitempty"`
Version int64 `gorm:"not null;default:1" json:"version"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
func (Message) TableName() string { return "sense_outbox_messages" }
// DeliveryRecord is permanent idempotency evidence. It is never deleted by
// queue cleanup and prevents a delivered key from being processed again.
type DeliveryRecord struct {
IdempotencyKey string `gorm:"size:191;primaryKey" json:"idempotencyKey"`
MessageID string `gorm:"size:36;not null;uniqueIndex" json:"messageId"`
DeliveredAt time.Time `gorm:"not null" json:"deliveredAt"`
CreatedAt time.Time `json:"createdAt"`
}
func (DeliveryRecord) TableName() string { return "sense_outbox_deliveries" }
type Attempt struct {
ID uint `gorm:"primaryKey;autoIncrement" json:"id"`
MessageID string `gorm:"size:36;not null;index" json:"messageId"`
Number int `gorm:"not null" json:"number"`
Outcome string `gorm:"size:32;not null" json:"outcome"`
Detail string `gorm:"size:512;not null;default:''" json:"detail"`
Worker string `gorm:"size:128;not null;default:''" json:"worker,omitempty"`
ActorUserID int `gorm:"not null;default:0" json:"actorUserId,omitempty"`
CreatedAt time.Time `json:"createdAt"`
}
func (Attempt) TableName() string { return "sense_outbox_attempts" }
@@ -0,0 +1,92 @@
package outbox
import (
"fmt"
"net/url"
"os"
"strings"
"sync"
"testing"
"time"
"gorm.io/driver/postgres"
"gorm.io/gorm"
)
func TestPostgresConcurrentWorkersDoNotClaimSameMessage(t *testing.T) {
dsn := os.Getenv("SENSE_OUTBOX_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_OUTBOX_TEST_DATABASE_URL to run the PostgreSQL multi-worker test")
}
base, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
schema := fmt.Sprintf("sense_outbox_78_%d", time.Now().UnixNano())
if err = base.Exec("CREATE SCHEMA " + schema).Error; err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = base.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error })
scoped, err := gorm.Open(postgres.Open(withSearchPath(dsn, schema)), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = scoped.AutoMigrate(&Message{}, &DeliveryRecord{}, &Attempt{}); err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 28, 12, 0, 0, 0, time.UTC)
for index := 0; index < 20; index++ {
if _, err = Enqueue(scoped, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: fmt.Sprintf("event-%d", index), IdempotencyKey: fmt.Sprintf("event:%d:v1", index), PayloadJSON: []byte(`{}`)}, now); err != nil {
t.Fatal(err)
}
}
workers := []string{"worker-a", "worker-b"}
results := make(chan []Message, len(workers))
errorsCh := make(chan error, len(workers))
var group sync.WaitGroup
for _, worker := range workers {
group.Add(1)
go func(name string) {
defer group.Done()
relay := NewRelay(scoped)
relay.Now = func() time.Time { return now }
items, claimErr := relay.Claim(name, 20)
if claimErr != nil {
errorsCh <- claimErr
return
}
results <- items
}(worker)
}
group.Wait()
close(results)
close(errorsCh)
for claimErr := range errorsCh {
t.Fatal(claimErr)
}
seen := map[string]string{}
for batch := range results {
for _, item := range batch {
if owner, exists := seen[item.ID]; exists {
t.Fatalf("message %s claimed by %s and %s", item.ID, owner, item.LeaseOwner)
}
seen[item.ID] = item.LeaseOwner
}
}
if len(seen) != 20 {
t.Fatalf("claimed=%d want=20", len(seen))
}
}
func withSearchPath(dsn, schema string) string {
if strings.Contains(dsn, "://") {
parsed, err := url.Parse(dsn)
if err == nil {
query := parsed.Query()
query.Set("search_path", schema)
parsed.RawQuery = query.Encode()
return parsed.String()
}
}
return strings.TrimSpace(dsn) + " search_path=" + schema
}
+177
View File
@@ -0,0 +1,177 @@
package outbox
import (
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
var (
ErrDuplicateIdempotency = errors.New("可靠投递幂等键已存在")
ErrLeaseLost = errors.New("可靠投递租约已失效")
)
func Enqueue(tx *gorm.DB, input EnqueueInput, now time.Time) (Message, error) {
input.InternalType, input.BusinessRef, input.IdempotencyKey = strings.TrimSpace(input.InternalType), strings.TrimSpace(input.BusinessRef), strings.TrimSpace(input.IdempotencyKey)
if input.InternalType == "" || len(input.InternalType) > 64 || input.BusinessRef == "" || len(input.BusinessRef) > 128 || input.IdempotencyKey == "" || len(input.IdempotencyKey) > 191 || !json.Valid(input.PayloadJSON) {
return Message{}, ErrInvalidInput
}
if input.MaxAttempts == 0 {
input.MaxAttempts = 12
}
if input.MaxAttempts < 1 || input.MaxAttempts > 100 {
return Message{}, ErrInvalidInput
}
now = now.UTC()
message := Message{ID: uuid.NewString(), InternalType: input.InternalType, BusinessRef: input.BusinessRef, IdempotencyKey: input.IdempotencyKey, PayloadJSON: string(input.PayloadJSON), State: StatePending, MaxAttempts: input.MaxAttempts, AvailableAt: now, Version: 1, CreatedAt: now, UpdatedAt: now}
var existing int64
if err := tx.Model(&Message{}).Where("idempotency_key = ?", input.IdempotencyKey).Count(&existing).Error; err != nil {
return Message{}, fmt.Errorf("check outbox idempotency: %w", err)
}
if existing > 0 {
return Message{}, ErrDuplicateIdempotency
}
if err := tx.Create(&message).Error; err != nil {
return Message{}, fmt.Errorf("enqueue outbox message: %w", err)
}
return message, nil
}
type Relay struct {
DB *gorm.DB
Now func() time.Time
LeaseDuration time.Duration
Backoff func(int) time.Duration
}
func NewRelay(db *gorm.DB) *Relay {
return &Relay{DB: db, Now: time.Now, LeaseDuration: 30 * time.Second, Backoff: defaultBackoff}
}
func (r *Relay) Claim(worker string, limit int) ([]Message, error) {
worker = strings.TrimSpace(worker)
if worker == "" || len(worker) > 128 || limit < 1 || limit > 100 {
return nil, ErrInvalidInput
}
now, leaseUntil := r.now(), r.now().Add(r.leaseDuration())
claimed := make([]Message, 0, limit)
err := r.DB.Transaction(func(tx *gorm.DB) error {
var candidates []Message
query := tx.Where("((state IN ?) AND available_at <= ?) OR (state = ? AND lease_until < ?)", []string{StatePending, StateRetry}, now, StateProcessing, now).Order("available_at ASC, created_at ASC").Limit(limit)
if tx.Dialector.Name() == "postgres" {
query = query.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"})
}
if err := query.Find(&candidates).Error; err != nil {
return err
}
for _, item := range candidates {
result := tx.Model(&Message{}).Where("id = ? AND version = ?", item.ID, item.Version).Updates(map[string]interface{}{"state": StateProcessing, "lease_owner": worker, "lease_until": leaseUntil, "version": gorm.Expr("version + 1"), "updated_at": now})
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 1 {
item.State, item.LeaseOwner, item.LeaseUntil, item.Version = StateProcessing, worker, &leaseUntil, item.Version+1
claimed = append(claimed, item)
}
}
return nil
})
if err != nil {
return nil, fmt.Errorf("claim outbox messages: %w", err)
}
return claimed, nil
}
func (r *Relay) MarkSuccess(id, worker string) error {
now := r.now()
return r.DB.Transaction(func(tx *gorm.DB) error {
var message Message
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&message, "id = ?", id).Error; err != nil {
return err
}
var delivered int64
if err := tx.Model(&DeliveryRecord{}).Where("idempotency_key = ?", message.IdempotencyKey).Count(&delivered).Error; err != nil {
return err
}
if delivered > 0 {
return nil
}
if message.State != StateProcessing || message.LeaseOwner != worker || message.LeaseUntil == nil || !message.LeaseUntil.After(now) {
return ErrLeaseLost
}
if err := tx.Create(&DeliveryRecord{IdempotencyKey: message.IdempotencyKey, MessageID: message.ID, DeliveredAt: now, CreatedAt: now}).Error; err != nil {
return err
}
number := message.AttemptCount + 1
if err := tx.Create(&Attempt{MessageID: message.ID, Number: number, Outcome: StateDelivered, Detail: "投递成功", Worker: worker, CreatedAt: now}).Error; err != nil {
return err
}
return tx.Model(&message).Updates(map[string]interface{}{"state": StateDelivered, "attempt_count": number, "delivered_at": now, "lease_owner": "", "lease_until": nil, "last_error": "", "version": gorm.Expr("version + 1"), "updated_at": now}).Error
})
}
func (r *Relay) MarkFailure(id, worker, detail string) error {
now := r.now()
detail = strings.TrimSpace(detail)
if len(detail) > 512 {
detail = detail[:512]
}
if detail == "" {
detail = "投递失败"
}
return r.DB.Transaction(func(tx *gorm.DB) error {
var message Message
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&message, "id = ?", id).Error; err != nil {
return err
}
if message.State != StateProcessing || message.LeaseOwner != worker || message.LeaseUntil == nil || !message.LeaseUntil.After(now) {
return ErrLeaseLost
}
nextAttempt := message.AttemptCount + 1
state := StateRetry
available := now.Add(r.backoff(nextAttempt))
if nextAttempt >= message.MaxAttempts {
state = StateDead
available = now
}
if err := tx.Model(&message).Updates(map[string]interface{}{"state": state, "attempt_count": nextAttempt, "available_at": available, "lease_owner": "", "lease_until": nil, "last_error": detail, "version": gorm.Expr("version + 1"), "updated_at": now}).Error; err != nil {
return err
}
return tx.Create(&Attempt{MessageID: message.ID, Number: nextAttempt, Outcome: state, Detail: detail, Worker: worker, CreatedAt: now}).Error
})
}
func (r *Relay) now() time.Time {
if r.Now != nil {
return r.Now().UTC()
}
return time.Now().UTC()
}
func (r *Relay) leaseDuration() time.Duration {
if r.LeaseDuration <= 0 {
return 30 * time.Second
}
return r.LeaseDuration
}
func (r *Relay) backoff(attempt int) time.Duration {
if r.Backoff != nil {
return r.Backoff(attempt)
}
return defaultBackoff(attempt)
}
func defaultBackoff(attempt int) time.Duration {
if attempt < 1 {
attempt = 1
}
delay := time.Second * time.Duration(1<<min(attempt-1, 8))
if delay > 5*time.Minute {
return 5 * time.Minute
}
return delay
}
+122
View File
@@ -0,0 +1,122 @@
package outbox
import (
"errors"
"testing"
"time"
"github.com/google/uuid"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
func testDB(t *testing.T) *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, _ := db.DB()
sqlDB.SetMaxOpenConns(1)
if err = db.AutoMigrate(&Message{}, &DeliveryRecord{}, &Attempt{}); err != nil {
t.Fatal(err)
}
return db
}
func TestEnqueueClaimLeaseRecoveryBackoffAndIdempotentSuccess(t *testing.T) {
db := testDB(t)
now := time.Date(2026, 8, 28, 8, 0, 0, 0, time.UTC)
message, err := Enqueue(db, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: "event-1", IdempotencyKey: "local-event:event-1:v1", PayloadJSON: []byte(`{"event":"event-1"}`), MaxAttempts: 3}, now)
if err != nil {
t.Fatal(err)
}
if _, err = Enqueue(db, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: "event-1", IdempotencyKey: message.IdempotencyKey, PayloadJSON: []byte(`{}`)}, now); !errors.Is(err, ErrDuplicateIdempotency) {
t.Fatalf("expected duplicate error, got %v", err)
}
relay := NewRelay(db)
relay.Now = func() time.Time { return now }
relay.LeaseDuration = 10 * time.Second
relay.Backoff = func(int) time.Duration { return 5 * time.Second }
first, err := relay.Claim("worker-1", 1)
if err != nil || len(first) != 1 {
t.Fatalf("first claim=%#v err=%v", first, err)
}
second, err := relay.Claim("worker-2", 1)
if err != nil || len(second) != 0 {
t.Fatalf("concurrent claim=%#v err=%v", second, err)
}
now = now.Add(11 * time.Second)
recovered, err := relay.Claim("worker-2", 1)
if err != nil || len(recovered) != 1 || recovered[0].LeaseOwner != "worker-2" {
t.Fatalf("recovered=%#v err=%v", recovered, err)
}
if err = relay.MarkFailure(message.ID, "worker-2", "temporary outage"); err != nil {
t.Fatal(err)
}
now = now.Add(4 * time.Second)
waiting, _ := relay.Claim("worker-3", 1)
if len(waiting) != 0 {
t.Fatalf("claimed before backoff elapsed: %#v", waiting)
}
now = now.Add(2 * time.Second)
retry, err := relay.Claim("worker-3", 1)
if err != nil || len(retry) != 1 {
t.Fatalf("retry=%#v err=%v", retry, err)
}
if err = relay.MarkSuccess(message.ID, "worker-3"); err != nil {
t.Fatal(err)
}
if err = relay.MarkSuccess(message.ID, "worker-3"); err != nil {
t.Fatalf("idempotent success failed: %v", err)
}
var deliveries int64
db.Model(&DeliveryRecord{}).Where("idempotency_key = ?", message.IdempotencyKey).Count(&deliveries)
if deliveries != 1 {
t.Fatalf("deliveries=%d", deliveries)
}
}
func TestFailureBecomesDeadAndManualRequeuePreservesHistory(t *testing.T) {
db := testDB(t)
now := time.Date(2026, 8, 28, 9, 0, 0, 0, time.UTC)
message, err := Enqueue(db, EnqueueInput{InternalType: "audit_projection", BusinessRef: "audit-1", IdempotencyKey: "audit:audit-1:v1", PayloadJSON: []byte(`{}`), MaxAttempts: 1}, now)
if err != nil {
t.Fatal(err)
}
relay := NewRelay(db)
relay.Now = func() time.Time { return now }
claimed, _ := relay.Claim("worker", 1)
if len(claimed) != 1 {
t.Fatal("message not claimed")
}
if err = relay.MarkFailure(message.ID, "worker", "permanent failure"); err != nil {
t.Fatal(err)
}
service := NewService(db)
service.Now = func() time.Time { return now.Add(time.Minute) }
detail, err := service.Get(message.ID)
if err != nil || detail.Message.State != StateDead || len(detail.Attempts) != 1 {
t.Fatalf("detail=%#v err=%v", detail, err)
}
requeued, err := service.Requeue(message.ID, RequeueRequest{ExpectedVersion: detail.Message.Version, Reason: "出口故障已排除"}, 7)
if err != nil {
t.Fatal(err)
}
if requeued.State != StatePending {
t.Fatalf("state=%s", requeued.State)
}
detail, _ = service.Get(message.ID)
if len(detail.Attempts) != 2 || detail.Attempts[0].Outcome != "manual_requeue" || detail.Attempts[0].ActorUserID != 7 {
t.Fatalf("attempt history=%#v", detail.Attempts)
}
}
func TestProductionCannotCreateTestSink(t *testing.T) {
if _, err := NewTestSink("prod"); !errors.Is(err, ErrTestSinkForbidden) {
t.Fatalf("err=%v", err)
}
if _, err := NewTestSink("test"); err != nil {
t.Fatal(err)
}
}
+124
View File
@@ -0,0 +1,124 @@
package outbox
import (
"errors"
"fmt"
"strings"
"time"
"unicode/utf8"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
var (
ErrInvalidInput = errors.New("可靠投递请求不符合要求")
ErrNotFound = errors.New("可靠投递记录不存在")
ErrNotDead = errors.New("仅死信记录可以重新排队")
ErrVersionConflict = 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) {
request.State = strings.TrimSpace(request.State)
request.InternalType = strings.TrimSpace(request.InternalType)
request.Keyword = strings.TrimSpace(request.Keyword)
if request.GetPageSize() > 100 || utf8.RuneCountInString(request.Keyword) > 128 || len(request.InternalType) > 64 || (request.State != "" && !validState(request.State)) {
return PageResponse{}, ErrInvalidInput
}
query := s.Orm.Model(&Message{})
if request.State != "" {
query = query.Where("state = ?", request.State)
}
if request.InternalType != "" {
query = query.Where("internal_type = ?", request.InternalType)
}
if request.Keyword != "" {
pattern := "%" + strings.ToLower(request.Keyword) + "%"
query = query.Where("LOWER(id) LIKE ? OR LOWER(idempotency_key) LIKE ? OR LOWER(business_ref) LIKE ?", pattern, pattern, pattern)
}
var response PageResponse
if err := query.Count(&response.Count).Error; err != nil {
return PageResponse{}, fmt.Errorf("count outbox messages: %w", err)
}
if err := query.Order("created_at DESC, id DESC").Limit(request.GetPageSize()).Offset((request.GetPageIndex() - 1) * request.GetPageSize()).Find(&response.List).Error; err != nil {
return PageResponse{}, fmt.Errorf("list outbox messages: %w", err)
}
for state, target := range map[string]*int64{StatePending: &response.Summary.Pending, StateRetry: &response.Summary.Retry, StateProcessing: &response.Summary.Processing, StateDead: &response.Summary.Dead} {
if err := s.Orm.Model(&Message{}).Where("state = ?", state).Count(target).Error; err != nil {
return PageResponse{}, fmt.Errorf("summarize outbox: %w", err)
}
}
return response, nil
}
func (s *Service) Get(id string) (DetailResponse, error) {
id = strings.TrimSpace(id)
if id == "" || len(id) > 64 {
return DetailResponse{}, ErrInvalidInput
}
var message Message
if err := s.Orm.First(&message, "id = ?", id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return DetailResponse{}, ErrNotFound
}
return DetailResponse{}, fmt.Errorf("get outbox message: %w", err)
}
var attempts []Attempt
if err := s.Orm.Where("message_id = ?", id).Order("created_at DESC, id DESC").Limit(50).Find(&attempts).Error; err != nil {
return DetailResponse{}, fmt.Errorf("list outbox attempts: %w", err)
}
return DetailResponse{Message: message, Attempts: attempts}, nil
}
func (s *Service) Requeue(id string, request RequeueRequest, actorUserID int) (Message, error) {
reason := strings.TrimSpace(request.Reason)
if utf8.RuneCountInString(reason) < 6 || utf8.RuneCountInString(reason) > 256 || request.ExpectedVersion < 1 {
return Message{}, ErrInvalidInput
}
now := s.now()
var result Message
err := s.Orm.Transaction(func(tx *gorm.DB) error {
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&result, "id = ?", strings.TrimSpace(id)).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return ErrNotFound
}
return err
}
if result.Version != request.ExpectedVersion {
return ErrVersionConflict
}
if result.State != StateDead {
return ErrNotDead
}
result.State, result.AvailableAt, result.LastError = StatePending, now, ""
result.LeaseOwner, result.LeaseUntil, result.Version = "", nil, result.Version+1
if err := tx.Save(&result).Error; err != nil {
return err
}
return tx.Create(&Attempt{MessageID: result.ID, Number: result.AttemptCount, Outcome: "manual_requeue", Detail: reason, ActorUserID: actorUserID, CreatedAt: now}).Error
})
if err != nil {
return Message{}, err
}
return result, nil
}
func (s *Service) now() time.Time {
if s.Now != nil {
return s.Now().UTC()
}
return time.Now().UTC()
}
func validState(value string) bool {
return value == StatePending || value == StateProcessing || value == StateRetry || value == StateDead || value == StateDelivered
}
+24
View File
@@ -0,0 +1,24 @@
package outbox
import (
"errors"
"strings"
)
var ErrTestSinkForbidden = errors.New("production 模式禁止启用测试接收器")
type Sink interface{ Deliver(Message) error }
type TestSink struct{ Delivered []string }
func (s *TestSink) Deliver(message Message) error {
s.Delivered = append(s.Delivered, message.IdempotencyKey)
return nil
}
func NewTestSink(applicationMode string) (*TestSink, error) {
if strings.EqualFold(strings.TrimSpace(applicationMode), "prod") || strings.EqualFold(strings.TrimSpace(applicationMode), "production") {
return nil, ErrTestSinkForbidden
}
return &TestSink{}, nil
}
@@ -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")
}
}
@@ -0,0 +1,65 @@
package version
import (
"fmt"
"runtime"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
"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), migrateSenseOutbox)
}
func migrateSenseOutbox(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&outbox.Message{}, &outbox.DeliveryRecord{}, &outbox.Attempt{}); 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: "SenseOutbox", Title: "可靠投递", Icon: "connection", Path: "outbox", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:outbox:list", ParentId: root.MenuId, Component: "/sense/outbox/index", Sort: 11, Visible: "0", IsFrame: "1"})
if err != nil {
return err
}
detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOutboxDetail", Title: "查看投递详情", MenuType: "F", Action: "GET", Permission: "sense:outbox: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
}
requeue, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOutboxRequeue", Title: "死信重新排队", MenuType: "F", Action: "POST", Permission: "sense:outbox:requeue", 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}); err != nil {
return err
}
for _, policy := range [][2]string{{"/api/v1/outbox", "GET"}, {"/api/v1/outbox/: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
}
}
}
for _, role := range []string{"implementation_operator", "site_admin"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{requeue}); err != nil {
return err
}
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: "/api/v1/outbox/:id/requeue", 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,41 @@
package version
import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"testing"
)
func TestOutboxMigrationAddsRBACWithoutFixtures(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 = "2026082815000_outbox.go"
if err = migrateSenseOutbox(db, version); err != nil {
t.Fatal(err)
}
if !db.Migrator().HasTable(&outbox.Message{}) || !db.Migrator().HasTable(&outbox.DeliveryRecord{}) || !db.Migrator().HasTable(&outbox.Attempt{}) {
t.Fatal("outbox tables missing")
}
var menus, reads, writes, fixtures, applied int64
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOutbox%").Count(&menus)
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 = ?", "/api/v1/outbox%", "GET").Count(&reads)
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/outbox/:id/requeue", "POST").Count(&writes)
db.Model(&outbox.Message{}).Count(&fixtures)
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
if menus != 3 || reads != 6 || writes != 2 || fixtures != 0 || applied != 1 {
t.Fatalf("menus=%d reads=%d writes=%d fixtures=%d applied=%d", menus, reads, writes, fixtures, applied)
}
}
@@ -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')
+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' })
}
+20
View File
@@ -0,0 +1,20 @@
import request from '@/utils/request'
export function listOutbox(query) {
return request({ url: '/api/v1/outbox', method: 'get', params: query })
}
export function getOutbox(id) {
return request({
url: `/api/v1/outbox/${encodeURIComponent(id)}`,
method: 'get'
})
}
export function requeueOutbox(id, expectedVersion, reason) {
return request({
url: `/api/v1/outbox/${encodeURIComponent(id)}/requeue`,
method: 'post',
data: { expectedVersion, reason }
})
}
@@ -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 '预检通过,但仍需单独建立高风险工单并人工确认' }
+447
View File
@@ -0,0 +1,447 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div>
<h3>可靠投递</h3>
<p>查看 Sense 内部待投递记录、重试进度和人工恢复。</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>未配置外部投递出口不影响 Sense 核心功能</template>
内部记录会安全保留,正式 connector 由后续协调工单提供;测试接收器在
production 模式不可启用。
</el-alert>
<el-row :gutter="12" class="summary-row" aria-label="可靠投递队列概览">
<el-col
v-for="card in summaryCards"
:key="card.key"
:xs="12"
:sm="6"
><div class="summary-item">
<span>{{ card.label }}</span><strong :class="card.className">{{ summary[card.key] }}</strong>
</div></el-col>
</el-row>
<el-form
:model="query"
label-width="76px"
class="filter-form"
@submit.prevent="search"
>
<el-form-item label="关键词"><el-input
v-model="query.keyword"
clearable
placeholder="内部记录、业务引用或幂等键"
@keyup.enter="search"
/></el-form-item>
<el-form-item label="状态"><el-select
v-model="query.state"
clearable
placeholder="全部状态"
><el-option
v-for="(label, value) in stateLabels"
:key="value"
:label="label"
:value="value"
/></el-select></el-form-item>
<el-form-item label="内部类型"><el-select
v-model="query.internalType"
clearable
placeholder="全部类型"
><el-option
v-for="(label, value) in typeLabels"
:key="value"
:label="label"
:value="value"
/></el-select></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="items"
border
stripe
empty-text="当前筛选条件下没有可靠投递记录"
>
<el-table-column
label="内部记录"
min-width="255"
><template #default="scope"><strong>{{ typeLabel(scope.row.internalType) }}</strong>
<div class="code-text">{{ scope.row.id }}</div>
<small>幂等键:{{ scope.row.idempotencyKey }}</small></template></el-table-column>
<el-table-column
label="状态"
width="112"
align="center"
><template #default="scope"><el-tag :type="stateType(scope.row.state)">{{
stateLabel(scope.row.state)
}}</el-tag></template></el-table-column>
<el-table-column
label="尝试"
width="74"
align="center"
><template #default="scope">{{ scope.row.attemptCount }} 次</template></el-table-column>
<el-table-column
label="下次动作"
min-width="180"
><template #default="scope">{{
nextAction(scope.row)
}}</template></el-table-column>
<el-table-column
label="租约"
min-width="142"
><template #default="scope"><div>{{ scope.row.leaseOwner || "未领取" }}</div>
<small v-if="scope.row.leaseUntil">{{
formatTime(scope.row.leaseUntil)
}}</small></template></el-table-column>
<el-table-column
label="最近结果"
min-width="210"
><template #default="scope">{{
scope.row.lastError ||
(scope.row.state === "delivered" ? "投递成功" : "尚未失败")
}}</template></el-table-column>
<el-table-column
label="操作"
width="150"
fixed="right"
><template #default="scope"><el-button
v-permisaction="['sense:outbox:detail']"
type="primary"
link
@click="openDetail(scope.row.id)"
>详情</el-button><el-button
v-if="scope.row.state === 'dead'"
v-permisaction="['sense:outbox:requeue']"
type="warning"
link
@click="openRequeue(scope.row)"
>重新排队</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"
/>
<p class="boundary-text">
页面不展示内部
payload、外部凭据或机器身份;人工恢复保留原业务记录、幂等键和失败历史。
</p>
</el-card>
<el-dialog
v-model="detailOpen"
title="投递记录详情"
width="min(780px, calc(100vw - 24px))"
:close-on-click-modal="false"
>
<el-descriptions v-if="selected" :column="2" border>
<el-descriptions-item label="内部记录">{{
selected.message.id
}}</el-descriptions-item><el-descriptions-item label="内部类型">{{
typeLabel(selected.message.internalType)
}}</el-descriptions-item>
<el-descriptions-item
label="幂等键"
:span="2"
><span class="code-text">{{
selected.message.idempotencyKey
}}</span></el-descriptions-item><el-descriptions-item label="当前状态"><el-tag :type="stateType(selected.message.state)">{{
stateLabel(selected.message.state)
}}</el-tag></el-descriptions-item><el-descriptions-item label="业务引用">{{
selected.message.businessRef
}}</el-descriptions-item>
<el-descriptions-item label="最近结果" :span="2">{{
selected.message.lastError || "尚未失败"
}}</el-descriptions-item>
</el-descriptions>
<h4 class="timeline-title">处理时间线</h4>
<el-timeline v-if="selected"><el-timeline-item
v-for="attempt in selected.attempts"
:key="attempt.id"
:timestamp="formatTime(attempt.createdAt)"
placement="top"
><strong>第 {{ attempt.number }} 次 · {{ attempt.outcome }}</strong>
<div>{{ attempt.detail }}</div>
<small v-if="attempt.actorUserId">操作人编号:{{ attempt.actorUserId }}</small></el-timeline-item><el-timeline-item
v-if="!selected.attempts.length"
:timestamp="formatTime(selected.message.createdAt)"
>与业务记录在同一事务中创建</el-timeline-item></el-timeline>
<template #footer><el-button @click="detailOpen = false">关闭</el-button></template>
</el-dialog>
<el-dialog
v-model="requeueOpen"
title="将死信重新排队"
width="min(560px, calc(100vw - 24px))"
:close-on-click-modal="false"
@closed="resetRequeue"
>
<el-alert
title="这是受控恢复操作"
type="warning"
:closable="false"
show-icon
>只创建新的投递尝试,不修改原业务记录或删除失败历史。执行前请先排除失败原因。</el-alert>
<el-form
ref="requeueFormRef"
:model="requeueForm"
:rules="requeueRules"
label-position="top"
class="requeue-form"
><el-form-item
label="恢复原因"
prop="reason"
><el-input
v-model="requeueForm.reason"
type="textarea"
:rows="3"
maxlength="256"
show-word-limit
placeholder="例如:出口配置已恢复,已核对幂等键和目标状态"
/></el-form-item></el-form>
<p class="boundary-text">
原因将与操作人、对象版本和时间一起写入脱敏审计。
</p>
<template #footer><el-button @click="requeueOpen = false">取消</el-button><el-button
type="warning"
:loading="requeueLoading"
@click="confirmRequeue"
>确认重新排队</el-button></template>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { computed, onMounted, reactive, ref } from 'vue'
import { Refresh, RefreshLeft, Search } from '@element-plus/icons-vue'
import { ElMessage } from 'element-plus'
import { getOutbox, listOutbox, requeueOutbox } from '@/api/sense/outbox'
import {
buildOutboxQuery,
formatTime,
nextAction,
stateLabel,
stateLabels,
stateType,
typeLabel,
typeLabels
} from './outboxState'
defineOptions({ name: 'SenseOutbox' })
const loading = ref(false)
const items = ref([])
const total = ref(0)
const detailOpen = ref(false)
const selected = ref(null)
const requeueOpen = ref(false)
const requeueLoading = ref(false)
const requeueTarget = ref(null)
const requeueFormRef = ref(null)
const summary = reactive({ pending: 0, retry: 0, processing: 0, dead: 0 })
const query = reactive({
pageIndex: 1,
pageSize: 10,
state: '',
internalType: '',
keyword: ''
})
const requeueForm = reactive({ reason: '' })
const requeueRules = {
reason: [
{ required: true, message: '请填写恢复原因', trigger: 'blur' },
{ min: 6, max: 256, message: '请填写 6 至 256 个字符', trigger: 'blur' }
]
}
const summaryCards = computed(() => [
{ key: 'pending', label: '等待投递' },
{ key: 'retry', label: '重试等待', className: 'warning-number' },
{ key: 'processing', label: '处理中 / 租约' },
{ key: 'dead', label: '死信待处理', className: 'danger-number' }
])
function unwrap(response) {
return response?.data?.data ?? response?.data ?? response
}
async function load() {
loading.value = true
try {
const payload = unwrap(await listOutbox(buildOutboxQuery(query))) || {}
items.value = payload.list || []
total.value = payload.count || 0
Object.assign(
summary,
payload.summary || { pending: 0, retry: 0, processing: 0, dead: 0 }
)
} catch (error) {
ElMessage.error(error.message || '可靠投递状态加载失败')
} finally {
loading.value = false
}
}
function search() {
query.pageIndex = 1
load()
}
function reset() {
Object.assign(query, {
pageIndex: 1,
state: '',
internalType: '',
keyword: ''
})
load()
}
async function openDetail(id) {
try {
selected.value = unwrap(await getOutbox(id))
detailOpen.value = true
} catch (error) {
ElMessage.error(error.message || '投递详情加载失败')
}
}
function openRequeue(item) {
requeueTarget.value = item
requeueOpen.value = true
}
function resetRequeue() {
requeueForm.reason = ''
requeueTarget.value = null
requeueFormRef.value?.clearValidate()
}
async function confirmRequeue() {
if (!requeueFormRef.value || !requeueTarget.value) return
const valid = await requeueFormRef.value.validate().catch(() => false)
if (!valid) return
requeueLoading.value = true
try {
await requeueOutbox(
requeueTarget.value.id,
requeueTarget.value.version,
requeueForm.reason.trim()
)
ElMessage.success('死信已重新排队,原失败历史已保留')
requeueOpen.value = false
detailOpen.value = false
await load()
} catch (error) {
ElMessage.warning(error.message || '重新排队失败,请刷新状态后重试')
} finally {
requeueLoading.value = false
}
}
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;
}
.boundary-alert {
margin: 16px 0;
}
.summary-row {
margin-bottom: 18px;
}
.summary-item {
display: flex;
align-items: center;
justify-content: space-between;
min-height: 68px;
padding: 12px 16px;
border: 1px solid #ebeef5;
border-radius: 4px;
}
.summary-item span {
color: #606266;
}
.summary-item strong {
font-size: 22px;
}
.warning-number {
color: #e6a23c;
}
.danger-number {
color: #f56c6c;
}
.filter-form {
display: flex;
align-items: flex-end;
flex-wrap: wrap;
gap: 0 12px;
margin-bottom: 2px;
}
.filter-form .el-form-item {
width: 260px;
}
.filter-form .filter-actions {
width: auto;
}
.filter-form :deep(.el-select) {
width: 100%;
}
.code-text {
font-family: ui-monospace, SFMono-Regular, Consolas, monospace;
overflow-wrap: anywhere;
}
.boundary-text,
small {
color: #909399;
font-size: 12px;
}
.boundary-text {
margin: 12px 0 0;
}
.timeline-title {
margin: 20px 0 14px;
}
.requeue-form {
margin-top: 16px;
}
@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,54 @@
export const stateLabels = {
pending: '等待投递',
processing: '处理中',
retry: '重试等待',
dead: '死信',
delivered: '已投递'
}
export const typeLabels = {
local_event_candidate: '本地事件候选',
audit_projection: '审计投影'
}
export function stateLabel(value) {
return stateLabels[value] || value || '—'
}
export function stateType(value) {
return (
{
pending: 'info',
processing: 'primary',
retry: 'warning',
dead: 'danger',
delivered: 'success'
}[value] || 'info'
)
}
export function typeLabel(value) {
return typeLabels[value] || value || '—'
}
export function buildOutboxQuery(query) {
return {
pageIndex: query.pageIndex,
pageSize: query.pageSize,
state: query.state || undefined,
internalType: query.internalType || undefined,
keyword: String(query.keyword || '').trim() || undefined
}
}
export function nextAction(item, now = Date.now()) {
if (item.state === 'dead') return '等待人工处理'
if (item.state === 'processing') { return item.leaseUntil ? `租约至 ${formatTime(item.leaseUntil)}` : '处理中' }
if (item.state === 'retry') {
return item.availableAt && new Date(item.availableAt).getTime() > now
? `自动重试 ${formatTime(item.availableAt)}`
: '等待重试领取'
}
if (item.state === 'delivered') return '已完成'
return '等待 worker 领取'
}
export function formatTime(value) {
return value
? new Date(value).toLocaleString('zh-CN', { hour12: false })
: '—'
}
@@ -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,27 @@
import request from '@/utils/request'
import { getOutbox, listOutbox, requeueOutbox } from '@/api/sense/outbox'
jest.mock('@/utils/request', () => jest.fn())
describe('Sense outbox API', () => {
beforeEach(() => request.mockReset())
test('encodes identifiers and sends only the recovery fields', () => {
listOutbox({ pageIndex: 1, pageSize: 10 })
getOutbox('outbox/id unsafe')
requeueOutbox('outbox/id unsafe', 4, '出口故障已排除')
expect(request).toHaveBeenNthCalledWith(1, {
url: '/api/v1/outbox',
method: 'get',
params: { pageIndex: 1, pageSize: 10 }
})
expect(request).toHaveBeenNthCalledWith(2, {
url: '/api/v1/outbox/outbox%2Fid%20unsafe',
method: 'get'
})
expect(request).toHaveBeenNthCalledWith(3, {
url: '/api/v1/outbox/outbox%2Fid%20unsafe/requeue',
method: 'post',
data: { expectedVersion: 4, reason: '出口故障已排除' }
})
})
})
@@ -0,0 +1,37 @@
import {
buildOutboxQuery,
nextAction,
stateLabel,
stateType,
typeLabel
} from '@/views/sense/outbox/outboxState'
describe('Sense outbox presentation state', () => {
test('uses operator-facing labels and non-color-only states', () => {
expect(stateLabel('dead')).toBe('死信')
expect(stateType('dead')).toBe('danger')
expect(typeLabel('local_event_candidate')).toBe('本地事件候选')
})
test('builds an allowlisted trimmed query', () => {
expect(
buildOutboxQuery({
pageIndex: 2,
pageSize: 20,
state: 'retry',
internalType: 'local_event_candidate',
keyword: ' key ',
ignored: 'no'
})
).toEqual({
pageIndex: 2,
pageSize: 20,
state: 'retry',
internalType: 'local_event_candidate',
keyword: 'key'
})
})
test('explains the next recovery action', () => {
expect(nextAction({ state: 'dead' })).toBe('等待人工处理')
expect(nextAction({ state: 'pending' })).toBe('等待 worker 领取')
})
})
+13 -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: 5c2b5df080b160c1174be3e935fffba420885bff
synchronized_at: 2026-08-28T06:14:33Z
wiki_revision: 578ddbaae3d7e846037d958085e40609cd398bef
synchronized_at: 2026-08-28T08:02:32Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -225,3 +225,14 @@ PostgreSQL 表 `sense_provisioning_batches` 保存幂等键、配额快照和汇
- go-admin-ui 页面:`Sense/ui/src/views/sense/edge-node/`;API 封装:`Sense/ui/src/api/sense/edge-node.js`。页面复用 BasicLayout、Element Plus 表格/分页/Dialog/Tag/Alert 和权限指令。
- 内部 `ApplyHeartbeat` 是未来 Sense 适配器的投影写入口,本工单不把它暴露为 HTTP API,也不定义跨项目契约。
<!-- sense-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 -->
+13 -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: 8f81b925af671209341f1282e128c7f048bca8ec
synchronized_at: 2026-08-28T06:14:44Z
wiki_revision: 999eb1aee3558ff75cfa929acac77841af20b742
synchronized_at: 2026-08-28T08:02:42Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -199,3 +199,14 @@ synchronized_at: 2026-08-28T06:14:44Z
- 状态转换写入不可变事件;列表和详情读取写入脱敏操作审计,不记录凭据或原始心跳载荷。
- 合成夹具必须显式调用,生产迁移和启动不得自动写入节点。
<!-- 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 -->
+40 -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: cfec18809b3cb0d938003102c9caa5e6d6f4e818
synchronized_at: 2026-08-28T06:25:12Z
wiki_revision: 05435d53528271a866c525655486689afc198762
synchronized_at: 2026-08-28T08:02:52Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
@@ -477,3 +477,41 @@ go test ./cmd/migrate/migration/version -run TestEdgeNodeMigrationOnPostgres -co
验证应覆盖在线、超过 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 -->
+14 -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: 34b85f1c3fc2c3e0e7fa0de3df056940fcb71333
synchronized_at: 2026-08-28T06:17:09Z
wiki_revision: ce246054849b5b797dc4da0a55116238a4c00d7e
synchronized_at: 2026-08-28T08:04:52Z
<!-- gitea-wiki-mirror:end -->
# YoVision 部署与运维
@@ -100,3 +100,15 @@ Sense\start_sense.bat
节点超过 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 -->