diff --git a/Sense/server/app/admin/router/sense_operations.go b/Sense/server/app/admin/router/sense_operations.go new file mode 100644 index 0000000..2b9b04e --- /dev/null +++ b/Sense/server/app/admin/router/sense_operations.go @@ -0,0 +1,20 @@ +package router + +import ( + "github.com/gin-gonic/gin" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/operations" + "git.ilapage.cn/ila/yovision/Sense/server/common/actions" + "git.ilapage.cn/ila/yovision/Sense/server/common/middleware" +) + +func init() { routerCheckRole = append(routerCheckRole, registerSenseOperationsRouter) } + +func registerSenseOperationsRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) { + api := &operations.API{} + r := v1.Group("/operations").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction()) + r.GET("", api.List) + r.GET("/:id", api.Get) + r.POST("/:id/retry", api.Retry) +} diff --git a/Sense/server/app/admin/router/sense_operations_test.go b/Sense/server/app/admin/router/sense_operations_test.go new file mode 100644 index 0000000..f221a72 --- /dev/null +++ b/Sense/server/app/admin/router/sense_operations_test.go @@ -0,0 +1,31 @@ +package router + +import ( + "net/http" + "testing" + + "github.com/gin-gonic/gin" + jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth" +) + +func TestSenseOperationsRoutes(t *testing.T) { + gin.SetMode(gin.TestMode) + engine := gin.New() + registerSenseOperationsRouter(engine.Group("/api/v1"), &jwt.GinJWTMiddleware{}) + wanted := map[string]bool{ + http.MethodGet + " /api/v1/operations": false, + http.MethodGet + " /api/v1/operations/:id": false, + http.MethodPost + " /api/v1/operations/:id/retry": false, + } + for _, route := range engine.Routes() { + key := route.Method + " " + route.Path + if _, ok := wanted[key]; ok { + wanted[key] = true + } + } + for route, found := range wanted { + if !found { + t.Fatalf("route not registered: %s", route) + } + } +} diff --git a/Sense/server/app/sense/operations/apis.go b/Sense/server/app/sense/operations/apis.go new file mode 100644 index 0000000..2a19481 --- /dev/null +++ b/Sense/server/app/sense/operations/apis.go @@ -0,0 +1,107 @@ +package operations + +import ( + "errors" + "net/http" + "time" + + "github.com/gin-gonic/gin" + "github.com/gin-gonic/gin/binding" + "github.com/go-admin-team/go-admin-core/sdk/api" + "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user" + coreService "github.com/go-admin-team/go-admin-core/sdk/service" + + "git.ilapage.cn/ila/yovision/Sense/server/common" +) + +const auditSuccess, auditFailure = "1", "2" + +type API struct{ api.Api } + +func (e *API) service(c *gin.Context) (*Service, error) { + base := coreService.Service{} + if err := e.MakeContext(c).MakeOrm().MakeService(&base).Errors; err != nil { + return nil, err + } + return NewService(base.Orm), nil +} + +func (e *API) List(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + request := PageRequest{} + if err = e.MakeContext(c).Bind(&request).Errors; err != nil { + e.audit(c, service, "List", auditFailure, "运维中心查询条件格式不正确") + e.Error(http.StatusBadRequest, err, "查询条件格式不正确") + return + } + response, err := service.List(request) + if err != nil { + e.audit(c, service, "List", auditFailure, "运维中心列表查询失败") + e.writeError(err) + return + } + e.audit(c, service, "List", auditSuccess, "读取运维中心问题列表") + e.OK(response, "查询成功") +} + +func (e *API) Get(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + item, err := service.Get(c.Param("id")) + if err != nil { + e.audit(c, service, "Get", auditFailure, "运维问题详情查询失败") + e.writeError(err) + return + } + e.audit(c, service, "Get", auditSuccess, "读取运维问题详情 "+item.ID) + e.OK(item, "查询成功") +} + +func (e *API) Retry(c *gin.Context) { + service, err := e.service(c) + if err != nil { + e.writeError(err) + return + } + request := RetryRequest{} + if err = e.MakeContext(c).Bind(&request, binding.JSON).Errors; err != nil { + e.audit(c, service, "Retry", auditFailure, "受控重试请求格式不正确") + e.Error(http.StatusBadRequest, err, "请求格式不正确") + return + } + item, err := service.Retry(c.Request.Context(), c.Param("id"), request.ExpectedVersion, user.GetUserId(c)) + if err != nil { + e.audit(c, service, "Retry", auditFailure, "受控重试被拒绝") + e.writeError(err) + return + } + e.audit(c, service, "Retry", auditSuccess, "受控重试已排队 "+item.ID) + e.OK(item, "重试任务已排队") +} + +func (e *API) audit(c *gin.Context, service *Service, action, status, remark string) { + err := WriteAudit(service.DB, Audit{Action: action, Method: c.Request.Method, Status: status, Username: user.GetUserName(c), UserID: user.GetUserId(c), ClientIP: common.GetClientIP(c), Route: c.FullPath(), Remark: remark, At: time.Now()}) + if err != nil { + api.GetRequestLogger(c).Errorf("operations audit failed: %s", err.Error()) + } +} + +func (e *API) writeError(err error) { + switch { + case errors.Is(err, ErrInvalidFilter): + e.Error(http.StatusBadRequest, err, err.Error()) + case errors.Is(err, ErrVersionConflict), errors.Is(err, ErrRetryInProgress), errors.Is(err, ErrRetryNotAllowed): + e.Error(http.StatusConflict, err, err.Error()) + case errors.Is(err, ErrProblemNotFound): + e.Error(http.StatusNotFound, err, err.Error()) + default: + e.Error(http.StatusInternalServerError, err, "运维中心操作失败") + } +} diff --git a/Sense/server/app/sense/operations/audit.go b/Sense/server/app/sense/operations/audit.go new file mode 100644 index 0000000..76b4da4 --- /dev/null +++ b/Sense/server/app/sense/operations/audit.go @@ -0,0 +1,26 @@ +package operations + +import ( + "time" + + "gorm.io/gorm" + + adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models" +) + +type Audit struct { + Action, Method, Status, Username, ClientIP, Route, Remark string + UserID int + At time.Time +} + +func WriteAudit(db *gorm.DB, input Audit) error { + model := adminModels.SysOperaLog{ + Title: "运维中心", BusinessType: "other", Method: "operations.API." + input.Action, + RequestMethod: input.Method, OperatorType: "1", OperName: input.Username, OperUrl: input.Route, + OperIp: input.ClientIP, Status: input.Status, OperTime: input.At.UTC(), Remark: input.Remark, + CreatedAt: input.At.UTC(), UpdatedAt: input.At.UTC(), + } + model.CreateBy, model.UpdateBy = input.UserID, input.UserID + return db.Create(&model).Error +} diff --git a/Sense/server/app/sense/operations/dto.go b/Sense/server/app/sense/operations/dto.go new file mode 100644 index 0000000..aef82af --- /dev/null +++ b/Sense/server/app/sense/operations/dto.go @@ -0,0 +1,68 @@ +package operations + +import ( + "time" + + commonDto "git.ilapage.cn/ila/yovision/Sense/server/common/dto" +) + +const ( + ObjectDevice = "device" + ObjectMedia = "media" + ObjectLocalInference = "local_inference" + + ProblemAuthentication = "authentication_failed" + ProblemBackoff = "backoff_wait" + ProblemClockDrift = "clock_drift" + ProblemOrphan = "orphan_safety_gate" + ProblemUnavailable = "capability_unavailable" + ProblemUnready = "not_converged" +) + +type PageRequest struct { + commonDto.Pagination `search:"-"` + ObjectType string `form:"objectType"` + ProblemType string `form:"problemType"` + Severity string `form:"severity"` + Keyword string `form:"keyword"` +} + +type RetryRequest struct { + ExpectedVersion int64 `json:"expectedVersion"` +} + +type Summary struct { + ManagedCount int `json:"managedCount"` + ConvergedCount int `json:"convergedCount"` + ActionableProblems int `json:"actionableProblems"` + LocalInference string `json:"localInference"` +} + +type Problem struct { + ID string `json:"id"` + ObjectType string `json:"objectType"` + ObjectID string `json:"objectId"` + ObjectName string `json:"objectName"` + Location string `json:"location,omitempty"` + ProblemType string `json:"problemType"` + Severity string `json:"severity"` + Expected string `json:"expected"` + Actual string `json:"actual"` + Difference string `json:"difference"` + NextAction string `json:"nextAction"` + NextRetryAt *time.Time `json:"nextRetryAt,omitempty"` + LastSuccessAt *time.Time `json:"lastSuccessAt,omitempty"` + AttemptCount int `json:"attemptCount"` + Retryable bool `json:"retryable"` + RetryInProgress bool `json:"retryInProgress"` + Version int64 `json:"version"` + SafetyGate string `json:"safetyGate"` + OperationalOnly bool `json:"operationalOnly"` + UpdatedAt time.Time `json:"updatedAt"` +} + +type PageResponse struct { + Summary Summary `json:"summary"` + List []Problem `json:"list"` + Count int64 `json:"count"` +} diff --git a/Sense/server/app/sense/operations/service.go b/Sense/server/app/sense/operations/service.go new file mode 100644 index 0000000..acfcdf9 --- /dev/null +++ b/Sense/server/app/sense/operations/service.go @@ -0,0 +1,359 @@ +package operations + +import ( + "context" + "errors" + "fmt" + "sort" + "strings" + "time" + + "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission" + deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models" + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media" + coreService "github.com/go-admin-team/go-admin-core/sdk/service" +) + +var ( + ErrInvalidFilter = errors.New("运维中心查询条件不符合要求") + ErrProblemNotFound = errors.New("运维问题不存在或已经收敛") + ErrVersionConflict = errors.New("状态已经变化,请刷新后重试") + ErrRetryInProgress = errors.New("该对象已有重试任务") + ErrRetryNotAllowed = errors.New("该问题不允许重试") +) + +type Service struct { + DB *gorm.DB + Now func() time.Time + LocalInferenceConfigured bool + DeviceRetry func(context.Context, admission.ProbeRequest) error +} + +type admissionResult struct { + DeviceID string + Address string + Status string + Detail string + CheckedAt time.Time +} + +func (admissionResult) TableName() string { return "sense_admission_results" } + +func NewService(db *gorm.DB) *Service { + service := &Service{DB: db, Now: time.Now, LocalInferenceConfigured: false} + service.DeviceRetry = func(ctx context.Context, request admission.ProbeRequest) error { + base := coreService.Service{} + base.Orm = db + runtimeService, err := admission.NewRuntime(base) + if err != nil { + return err + } + _, err = runtimeService.Probe(ctx, request) + if err == nil { + // Keep the existing admission boundary: route intent is best-effort and + // must not turn a successful device probe into a MediaMTX failure. + _ = media.EnsureDeviceRoutes(ctx, db, request.DeviceID) + } + return err + } + return service +} + +func (s *Service) List(request PageRequest) (PageResponse, error) { + if err := validateFilter(request); err != nil { + return PageResponse{}, err + } + problems, summary, err := s.project() + if err != nil { + return PageResponse{}, err + } + filtered := make([]Problem, 0, len(problems)) + keyword := strings.ToLower(strings.TrimSpace(request.Keyword)) + for _, item := range problems { + if request.ObjectType != "" && item.ObjectType != request.ObjectType { + continue + } + if request.ProblemType != "" && item.ProblemType != request.ProblemType { + continue + } + if request.Severity != "" && item.Severity != request.Severity { + continue + } + searchable := strings.ToLower(strings.Join([]string{item.ObjectID, item.ObjectName, item.Location, item.Difference}, " ")) + if keyword != "" && !strings.Contains(searchable, keyword) { + continue + } + filtered = append(filtered, item) + } + pageIndex, pageSize := request.GetPageIndex(), request.GetPageSize() + if pageSize > 100 { + pageSize = 100 + } + start := (pageIndex - 1) * pageSize + if start > len(filtered) { + start = len(filtered) + } + end := start + pageSize + if end > len(filtered) { + end = len(filtered) + } + return PageResponse{Summary: summary, List: filtered[start:end], Count: int64(len(filtered))}, nil +} + +func (s *Service) Get(id string) (Problem, error) { + problems, _, err := s.project() + if err != nil { + return Problem{}, err + } + for _, item := range problems { + if item.ID == id { + return item, nil + } + } + return Problem{}, ErrProblemNotFound +} + +func (s *Service) Retry(ctx context.Context, id string, expectedVersion int64, userID int) (Problem, error) { + if expectedVersion < 1 { + return Problem{}, ErrVersionConflict + } + problem, err := s.Get(id) + if err != nil { + return Problem{}, err + } + if !problem.Retryable { + if problem.RetryInProgress { + return Problem{}, ErrRetryInProgress + } + return Problem{}, ErrRetryNotAllowed + } + now := s.now().UTC() + switch problem.ObjectType { + case ObjectDevice: + var admissionState admissionResult + if err := s.DB.First(&admissionState, "device_id = ?", problem.ObjectID).Error; err != nil || strings.TrimSpace(admissionState.Address) == "" { + return Problem{}, ErrRetryNotAllowed + } + result := s.DB.Model(&deviceModels.Device{}). + Where("id = ? AND version = ? AND retry_requested_at IS NULL", problem.ObjectID, expectedVersion). + Updates(map[string]any{"retry_requested_at": now, "version": expectedVersion + 1, "updated_at": now}) + if result.Error != nil { + return Problem{}, fmt.Errorf("queue device retry: %w", result.Error) + } + if result.RowsAffected == 0 { + return Problem{}, s.retryConflict(ObjectDevice, problem.ObjectID, expectedVersion) + } + if s.DeviceRetry != nil { + err = s.DeviceRetry(ctx, admission.ProbeRequest{DeviceID: problem.ObjectID, Address: admissionState.Address, Version: expectedVersion + 1, UpdateBy: userID}) + if err != nil { + clearErr := s.DB.Model(&deviceModels.Device{}).Where("id = ? AND version = ?", problem.ObjectID, expectedVersion+1).Update("retry_requested_at", nil).Error + if clearErr != nil { + return Problem{}, errors.Join(fmt.Errorf("execute device retry: %w", err), fmt.Errorf("clear device retry gate: %w", clearErr)) + } + return Problem{}, fmt.Errorf("execute device retry: %w", err) + } + } + case ObjectMedia: + result := s.DB.Model(&media.Route{}). + Where("id = ? AND version = ? AND actual <> ?", problem.ObjectID, expectedVersion, "retry_pending"). + Updates(map[string]any{"actual": "retry_pending", "next_retry_at": now, "detail": "已请求受控重试,等待视频服务执行", "version": expectedVersion + 1, "updated_at": now}) + if result.Error != nil { + return Problem{}, fmt.Errorf("queue media retry: %w", result.Error) + } + if result.RowsAffected == 0 { + return Problem{}, s.retryConflict(ObjectMedia, problem.ObjectID, expectedVersion) + } + default: + return Problem{}, ErrRetryNotAllowed + } + updated, err := s.Get(id) + if errors.Is(err, ErrProblemNotFound) { + problem.Actual, problem.Difference, problem.NextAction = "converged", "重试完成,状态已经收敛", "无需处理" + problem.Retryable, problem.RetryInProgress, problem.Version, problem.UpdatedAt = false, false, expectedVersion+2, now + return problem, nil + } + return updated, err +} + +func (s *Service) project() ([]Problem, Summary, error) { + var devices []deviceModels.Device + if err := s.DB.Order("updated_at DESC").Find(&devices).Error; err != nil { + return nil, Summary{}, fmt.Errorf("list operation devices: %w", err) + } + var admissions []admissionResult + if err := s.DB.Find(&admissions).Error; err != nil { + return nil, Summary{}, fmt.Errorf("list admission results: %w", err) + } + var routes []media.Route + if err := s.DB.Order("updated_at DESC").Find(&routes).Error; err != nil { + return nil, Summary{}, fmt.Errorf("list operation media routes: %w", err) + } + admissionByDevice := make(map[string]admissionResult, len(admissions)) + for _, item := range admissions { + admissionByDevice[item.DeviceID] = item + } + deviceByID := make(map[string]deviceModels.Device, len(devices)) + problems := make([]Problem, 0) + summary := Summary{ManagedCount: len(devices) + len(routes), LocalInference: "configured"} + for _, device := range devices { + deviceByID[device.ID] = device + if problem, ok := deviceProblem(device, admissionByDevice[device.ID]); ok { + problems = append(problems, problem) + summary.ActionableProblems++ + } else { + summary.ConvergedCount++ + } + } + for _, route := range routes { + if problem, ok := mediaProblem(route, deviceByID); ok { + problems = append(problems, problem) + summary.ActionableProblems++ + } else { + summary.ConvergedCount++ + } + } + if !s.LocalInferenceConfigured { + summary.LocalInference = "unavailable" + now := s.now().UTC() + problems = append(problems, Problem{ + ID: "local_inference:adapter", ObjectType: ObjectLocalInference, ObjectID: "adapter", + ObjectName: "本地推理适配器", ProblemType: ProblemUnavailable, Severity: "info", + Expected: "可选", Actual: "unavailable", Difference: "未配置本地推理适配器;不影响设备和媒体运维", + NextAction: "无需处理", SafetyGate: "不读取 Brain 数据库,不阻断本页", OperationalOnly: true, UpdatedAt: now, + }) + } + sort.SliceStable(problems, func(i, j int) bool { + rank := map[string]int{"high": 0, "medium": 1, "info": 2} + if rank[problems[i].Severity] != rank[problems[j].Severity] { + return rank[problems[i].Severity] < rank[problems[j].Severity] + } + return problems[i].UpdatedAt.After(problems[j].UpdatedAt) + }) + return problems, summary, nil +} + +func deviceProblem(device deviceModels.Device, admission admissionResult) (Problem, bool) { + if device.Status == deviceModels.StatusDisabled { + return Problem{}, false + } + actual, detail, updatedAt := device.AdapterStatus, "设备尚未完成接入验证", device.UpdatedAt + if admission.DeviceID != "" { + actual, detail, updatedAt = admission.Status, admission.Detail, admission.CheckedAt + } + if device.Status == deviceModels.StatusActive && (actual == "ready" || actual == "verified") { + return Problem{}, false + } + problemType, severity, nextAction := ProblemUnready, "medium", "检查设备和接入配置后重试" + text := strings.ToLower(actual + " " + detail) + if strings.Contains(text, "auth") || strings.Contains(text, "认证") || strings.Contains(text, "凭据") { + problemType, severity, nextAction = ProblemAuthentication, "high", "确认设备账号未变更后执行受控重试" + } else if strings.Contains(text, "clock") || strings.Contains(text, "time drift") || strings.Contains(text, "时间漂移") || strings.Contains(text, "时钟") { + problemType, nextAction = ProblemClockDrift, "检查设备 NTP 和时区后重新检测" + } + retrying := device.RetryRequestedAt != nil + retryable := !retrying && strings.TrimSpace(admission.Address) != "" + if admission.DeviceID == "" || strings.TrimSpace(admission.Address) == "" { + nextAction = "先到视频接入完成地址与凭据验证" + } + return Problem{ + ID: "device:" + device.ID, ObjectType: ObjectDevice, ObjectID: device.ID, ObjectName: device.Name, + Location: device.Location, ProblemType: problemType, Severity: severity, Expected: "active / ready", Actual: actual, + Difference: detail, NextAction: nextAction, AttemptCount: boolInt(retrying), Retryable: retryable, + RetryInProgress: retrying, Version: device.Version, SafetyGate: "校验设备版本且同一设备仅允许一个在途重试", + OperationalOnly: true, UpdatedAt: updatedAt, + }, true +} + +func mediaProblem(route media.Route, devices map[string]deviceModels.Device) (Problem, bool) { + device, exists := devices[route.DeviceID] + name, location := route.Path, "" + if exists { + name, location = device.Name+" / "+route.Path, device.Location + } + if !exists { + return Problem{ + ID: "media:" + route.ID, ObjectType: ObjectMedia, ObjectID: route.ID, ObjectName: name, + ProblemType: ProblemOrphan, Severity: "medium", Expected: "路由关联有效设备", Actual: "已隔离待确认", + Difference: "媒体路由存在,但找不到有效设备归属", NextAction: "人工核对;不会自动删除", Version: route.Version, + SafetyGate: "孤儿资源仅隔离和提示,本接口没有删除动作", OperationalOnly: true, UpdatedAt: route.UpdatedAt, + }, true + } + converged := (route.Desired == media.DesiredRunning && (route.Actual == "ready" || route.Actual == "waiting")) || (route.Desired == media.DesiredStopped && route.Actual == "stopped") + if converged { + return Problem{}, false + } + problemType, nextAction := ProblemUnready, "检查视频服务状态后重试" + if route.NextRetryAt != nil { + problemType, nextAction = ProblemBackoff, "等待退避到期或执行受控提前重试" + } + retrying := route.Actual == "retry_pending" + return Problem{ + ID: "media:" + route.ID, ObjectType: ObjectMedia, ObjectID: route.ID, ObjectName: name, Location: location, + ProblemType: problemType, Severity: "medium", Expected: route.Desired, Actual: route.Actual, Difference: route.Detail, + NextAction: nextAction, NextRetryAt: route.NextRetryAt, AttemptCount: route.FailureCount, Retryable: !retrying, + RetryInProgress: retrying, Version: route.Version, SafetyGate: "校验路由版本且只排队,不删除路径或修改凭据", + OperationalOnly: true, UpdatedAt: route.UpdatedAt, + }, true +} + +func (s *Service) retryConflict(objectType, objectID string, expectedVersion int64) error { + var version int64 + var inProgress bool + switch objectType { + case ObjectDevice: + var item deviceModels.Device + if err := s.DB.Select("version", "retry_requested_at").First(&item, "id = ?", objectID).Error; err != nil { + return ErrProblemNotFound + } + version, inProgress = item.Version, item.RetryRequestedAt != nil + case ObjectMedia: + var item media.Route + if err := s.DB.Select("version", "actual").First(&item, "id = ?", objectID).Error; err != nil { + return ErrProblemNotFound + } + version, inProgress = item.Version, item.Actual == "retry_pending" + } + if inProgress { + return ErrRetryInProgress + } + if version != expectedVersion { + return ErrVersionConflict + } + return ErrVersionConflict +} + +func (s *Service) now() time.Time { + if s.Now != nil { + return s.Now() + } + return time.Now() +} + +func validateFilter(request PageRequest) error { + valid := func(value string, values ...string) bool { + if value == "" { + return true + } + for _, candidate := range values { + if value == candidate { + return true + } + } + return false + } + if !valid(request.ObjectType, ObjectDevice, ObjectMedia, ObjectLocalInference) || + !valid(request.ProblemType, ProblemAuthentication, ProblemBackoff, ProblemClockDrift, ProblemOrphan, ProblemUnavailable, ProblemUnready) || + !valid(request.Severity, "high", "medium", "info") || len([]rune(request.Keyword)) > 128 { + return ErrInvalidFilter + } + return nil +} + +func boolInt(value bool) int { + if value { + return 1 + } + return 0 +} diff --git a/Sense/server/app/sense/operations/service_test.go b/Sense/server/app/sense/operations/service_test.go new file mode 100644 index 0000000..aa19a14 --- /dev/null +++ b/Sense/server/app/sense/operations/service_test.go @@ -0,0 +1,187 @@ +package operations + +import ( + "context" + "errors" + "testing" + "time" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models" + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission" + deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models" + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/media" +) + +func operationsTestService(t *testing.T) (*Service, *gorm.DB, time.Time) { + t.Helper() + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&deviceModels.Device{}, &admissionResult{}, &media.Route{}, &adminModels.SysOperaLog{}); err != nil { + t.Fatal(err) + } + now := time.Date(2026, 8, 28, 2, 0, 0, 0, time.UTC) + service := NewService(db) + service.Now = func() time.Time { return now } + service.LocalInferenceConfigured = false + service.DeviceRetry = nil + return service, db, now +} + +func TestProjectionCoversOperationsStatesAndIndependentBoundaries(t *testing.T) { + service, db, now := operationsTestService(t) + devices := []deviceModels.Device{ + {ID: "ready", Name: "东门", Location: "一层", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1}, + {ID: "auth", Name: "仓库", Location: "北区", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3}, + {ID: "clock", Name: "南门", Location: "室外", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 4}, + } + for index := range devices { + devices[index].CreatedAt, devices[index].UpdatedAt = now, now + } + if err := db.Create(&devices).Error; err != nil { + t.Fatal(err) + } + admissions := []admissionResult{ + {DeviceID: "ready", Address: "http://192.0.2.10/onvif", Status: "ready", Detail: "接入验证完成", CheckedAt: now}, + {DeviceID: "auth", Address: "http://192.0.2.11/onvif", Status: "authentication_failed", Detail: "ONVIF 认证失败", CheckedAt: now}, + {DeviceID: "clock", Address: "http://192.0.2.12/onvif", Status: "clock_drift", Detail: "设备时间漂移 96 秒", CheckedAt: now}, + } + if err := db.Create(&admissions).Error; err != nil { + t.Fatal(err) + } + next := now.Add(5 * time.Minute) + routes := []media.Route{ + {ID: "auth:main", DeviceID: "auth", ProfileToken: "main", Path: "sense_auth", Desired: "running", Actual: "apply_failed", Detail: "媒体路径配置失败", FailureCount: 2, NextRetryAt: &next, Version: 5, UpdatedAt: now}, + {ID: "missing:main", DeviceID: "missing", ProfileToken: "main", Path: "sense_orphan", Desired: "running", Actual: "ready", Detail: "上游拉流正常", Version: 1, UpdatedAt: now}, + } + if err := db.Create(&routes).Error; err != nil { + t.Fatal(err) + } + response, err := service.List(PageRequest{}) + if err != nil { + t.Fatal(err) + } + if response.Summary.ManagedCount != 5 || response.Summary.ConvergedCount != 1 || response.Summary.ActionableProblems != 4 || response.Summary.LocalInference != "unavailable" || response.Count != 5 { + t.Fatalf("unexpected summary: %#v count=%d", response.Summary, response.Count) + } + wanted := map[string]bool{ProblemAuthentication: false, ProblemClockDrift: false, ProblemBackoff: false, ProblemOrphan: false, ProblemUnavailable: false} + for _, item := range response.List { + wanted[item.ProblemType] = true + if !item.OperationalOnly { + t.Fatalf("problem can be mistaken for Bell alert: %#v", item) + } + } + for state, found := range wanted { + if !found { + t.Fatalf("missing problem state %s: %#v", state, response.List) + } + } + filtered, err := service.List(PageRequest{ObjectType: ObjectDevice, ProblemType: ProblemClockDrift, Keyword: "南门"}) + if err != nil || filtered.Count != 1 || filtered.List[0].ObjectID != "clock" { + t.Fatalf("filtered=%#v err=%v", filtered, err) + } +} + +func TestControlledRetryUsesVersionAndInProgressGate(t *testing.T) { + service, db, now := operationsTestService(t) + device := deviceModels.Device{ID: "auth", Name: "东门", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3} + device.CreatedAt, device.UpdatedAt = now, now + if err := db.Create(&device).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&admissionResult{DeviceID: "auth", Address: "http://192.0.2.11/onvif", Status: "authentication_failed", Detail: "认证失败", CheckedAt: now}).Error; err != nil { + t.Fatal(err) + } + item, err := service.Retry(context.Background(), "device:auth", 3, 7) + if err != nil || !item.RetryInProgress || item.Retryable || item.Version != 4 { + t.Fatalf("item=%#v err=%v", item, err) + } + if _, err = service.Retry(context.Background(), "device:auth", 3, 7); !errors.Is(err, ErrRetryInProgress) { + t.Fatalf("expected in-progress gate, got %v", err) + } + var stored deviceModels.Device + if err = db.First(&stored, "id = ?", "auth").Error; err != nil || stored.RetryRequestedAt == nil { + t.Fatalf("stored=%#v err=%v", stored, err) + } +} + +func TestMediaRetryQueuesWithoutDeletingOrChangingCredentials(t *testing.T) { + service, db, now := operationsTestService(t) + device := deviceModels.Device{ID: "camera", Name: "仓库", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "ready", Version: 1} + device.CreatedAt, device.UpdatedAt = now, now + if err := db.Create(&device).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&admissionResult{DeviceID: "camera", Address: "http://192.0.2.20/onvif", Status: "ready", Detail: "接入完成", CheckedAt: now}).Error; err != nil { + t.Fatal(err) + } + next := now.Add(time.Minute) + route := media.Route{ID: "camera:main", DeviceID: "camera", ProfileToken: "main", Path: "sense_camera", Desired: "running", Actual: "apply_failed", FailureCount: 2, NextRetryAt: &next, Detail: "配置失败", Version: 6, UpdatedAt: now} + if err := db.Create(&route).Error; err != nil { + t.Fatal(err) + } + item, err := service.Retry(context.Background(), "media:camera:main", 6, 7) + if err != nil || item.Actual != "retry_pending" || item.Version != 7 || !item.RetryInProgress { + t.Fatalf("item=%#v err=%v", item, err) + } + var count int64 + if err = db.Model(&media.Route{}).Where("id = ?", route.ID).Count(&count).Error; err != nil || count != 1 { + t.Fatalf("route was removed: count=%d err=%v", count, err) + } +} + +func TestOrphanAndUnavailableCannotRetry(t *testing.T) { + service, db, now := operationsTestService(t) + route := media.Route{ID: "missing:main", DeviceID: "missing", ProfileToken: "main", Path: "orphan", Desired: "running", Actual: "ready", Version: 1, UpdatedAt: now} + if err := db.Create(&route).Error; err != nil { + t.Fatal(err) + } + if _, err := service.Retry(context.Background(), "media:missing:main", 1, 7); !errors.Is(err, ErrRetryNotAllowed) { + t.Fatalf("orphan retry should be rejected: %v", err) + } + if _, err := service.Retry(context.Background(), "local_inference:adapter", 1, 7); !errors.Is(err, ErrRetryNotAllowed) { + t.Fatalf("optional adapter retry should be rejected: %v", err) + } +} + +func TestDeviceRetryUsesExistingAdmissionPortAndOperatorIdentity(t *testing.T) { + service, db, now := operationsTestService(t) + device := deviceModels.Device{ID: "clock", Name: "南门", Modality: "video", CapabilitiesJSON: "[]", Status: "active", AdapterStatus: "verification_failed", Version: 3} + device.CreatedAt, device.UpdatedAt = now, now + if err := db.Create(&device).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&admissionResult{DeviceID: "clock", Address: "http://192.0.2.12/onvif", Status: "clock_drift", Detail: "设备时间漂移", CheckedAt: now}).Error; err != nil { + t.Fatal(err) + } + called := false + service.DeviceRetry = func(_ context.Context, request admission.ProbeRequest) error { + called = true + if request.DeviceID != "clock" || request.Address != "http://192.0.2.12/onvif" || request.Version != 4 || request.UpdateBy != 9 { + t.Fatalf("unexpected retry request: %#v", request) + } + return db.Model(&deviceModels.Device{}).Where("id = ? AND version = ?", request.DeviceID, request.Version).Updates(map[string]any{"retry_requested_at": nil, "version": request.Version + 1}).Error + } + item, err := service.Retry(context.Background(), "device:clock", 3, 9) + if err != nil || !called || item.Version != 5 || item.RetryInProgress { + t.Fatalf("item=%#v called=%v err=%v", item, called, err) + } +} + +func TestAuditIsMinimalAndDesensitized(t *testing.T) { + _, db, now := operationsTestService(t) + if err := WriteAudit(db, Audit{Action: "Retry", Method: "POST", Status: "1", Username: "operator", UserID: 7, ClientIP: "127.0.0.1", Route: "/api/v1/operations/:id/retry", Remark: "受控重试已排队 device:synthetic", At: now}); err != nil { + t.Fatal(err) + } + var stored adminModels.SysOperaLog + if err := db.First(&stored).Error; err != nil { + t.Fatal(err) + } + if stored.Title != "运维中心" || stored.RequestMethod != "POST" || stored.OperParam != "" || stored.JsonResult != "" || stored.CreateBy != 7 { + t.Fatalf("unexpected audit: %#v", stored) + } +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082811000_operations.go b/Sense/server/cmd/migrate/migration/version/2026082811000_operations.go new file mode 100644 index 0000000..441237e --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082811000_operations.go @@ -0,0 +1,62 @@ +package version + +import ( + "fmt" + "runtime" + + "gorm.io/gorm" + "gorm.io/gorm/clause" + + "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration" + migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models" + common "git.ilapage.cn/ila/yovision/Sense/server/common/models" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSenseOperations) +} + +func migrateSenseOperations(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + root, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: senseLayoutMenuName, Title: "视频感知", Icon: "video-camera", Path: "/sense", MenuType: "M", Component: "Layout", Sort: 5, Visible: "0", IsFrame: "1"}) + if err != nil { + return err + } + page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOperations", Title: "运维中心", Icon: "operation", Path: "operations", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:operations:list", ParentId: root.MenuId, Component: "/sense/operations/index", Sort: 7, Visible: "0", IsFrame: "1"}) + if err != nil { + return err + } + detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOperationsDetail", Title: "查看运维详情", MenuType: "F", Action: "GET", Permission: "sense:operations:detail", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 1, Visible: "1", IsFrame: "1"}) + if err != nil { + return err + } + retry, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOperationsRetry", Title: "受控重试", MenuType: "F", Action: "POST", Permission: "sense:operations:retry", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 2, Visible: "1", IsFrame: "1"}) + if err != nil { + return err + } + readPolicies := [][2]string{{"/api/v1/operations", "GET"}, {"/api/v1/operations/:id", "GET"}} + for _, role := range []string{"implementation_operator", "site_admin", "viewer"} { + if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, detail}); err != nil { + return err + } + for _, policy := range readPolicies { + if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}).Error; err != nil { + return err + } + } + } + for _, role := range []string{"implementation_operator", "site_admin"} { + if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{retry}); err != nil { + return err + } + if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: "/api/v1/operations/:id/retry", V2: "POST"}).Error; err != nil { + return err + } + } + if err = rebuildSenseMenuPaths(tx, root.MenuId, "/0"); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082811000_operations_test.go b/Sense/server/cmd/migrate/migration/version/2026082811000_operations_test.go new file mode 100644 index 0000000..4bfb1e2 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082811000_operations_test.go @@ -0,0 +1,85 @@ +package version + +import ( + "os" + "testing" + + "gorm.io/driver/postgres" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models" + common "git.ilapage.cn/ila/yovision/Sense/server/common/models" +) + +func TestOperationsMigrationAddsLeastPrivilegeMenus(t *testing.T) { + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil { + t.Fatal(err) + } + for _, role := range []string{"implementation_operator", "site_admin", "viewer"} { + if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil { + t.Fatal(err) + } + } + const version = "2026082811000_operations.go" + if err = migrateSenseOperations(db, version); err != nil { + t.Fatal(err) + } + var menus, reads, retries, viewerRetries, applied int64 + db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOperations%").Count(&menus) + db.Model(&deviceCasbinRule{}).Where("v1 IN ?", []string{"/api/v1/operations", "/api/v1/operations/:id"}).Count(&reads) + db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/operations/:id/retry", "POST").Count(&retries) + db.Model(&deviceCasbinRule{}).Where("v0 = ? AND v1 = ?", "viewer", "/api/v1/operations/:id/retry").Count(&viewerRetries) + db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied) + if menus != 3 || reads != 6 || retries != 2 || viewerRetries != 0 || applied != 1 { + t.Fatalf("menus=%d reads=%d retries=%d viewerRetries=%d applied=%d", menus, reads, retries, viewerRetries, applied) + } +} + +func TestOperationsMigrationOnPostgres(t *testing.T) { + dsn := os.Getenv("SENSE_OPERATIONS_MIGRATION_TEST_DATABASE_URL") + if dsn == "" { + t.Skip("set SENSE_OPERATIONS_MIGRATION_TEST_DATABASE_URL to run the PostgreSQL migration test") + } + db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + const schema = "sense_operations_74_test" + if err = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error; err != nil { + t.Fatal(err) + } + if err = db.Exec("CREATE SCHEMA " + schema).Error; err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error }) + sqlDB, err := db.DB() + if err != nil { + t.Fatal(err) + } + sqlDB.SetMaxOpenConns(1) + if err = db.Exec("SET search_path TO " + schema).Error; err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil { + t.Fatal(err) + } + for _, role := range []string{"implementation_operator", "site_admin", "viewer"} { + if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil { + t.Fatal(err) + } + } + if err = migrateSenseOperations(db, "2026082811000_operations.go"); err != nil { + t.Fatal(err) + } + var menus, policies int64 + db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOperations%").Count(&menus) + db.Model(&deviceCasbinRule{}).Where("v1 LIKE ?", "/api/v1/operations%").Count(&policies) + if menus != 3 || policies != 8 { + t.Fatalf("menus=%d policies=%d", menus, policies) + } +} diff --git a/Sense/ui/src/api/sense/operations.js b/Sense/ui/src/api/sense/operations.js new file mode 100644 index 0000000..577596c --- /dev/null +++ b/Sense/ui/src/api/sense/operations.js @@ -0,0 +1,13 @@ +import request from '@/utils/request' + +export function listOperationProblems(query) { + return request({ url: '/api/v1/operations', method: 'get', params: query }) +} + +export function getOperationProblem(id) { + return request({ url: `/api/v1/operations/${encodeURIComponent(id)}`, method: 'get' }) +} + +export function retryOperationProblem(id, expectedVersion) { + return request({ url: `/api/v1/operations/${encodeURIComponent(id)}/retry`, method: 'post', data: { expectedVersion }}) +} diff --git a/Sense/ui/src/views/sense/operations/index.vue b/Sense/ui/src/views/sense/operations/index.vue new file mode 100644 index 0000000..e4614d3 --- /dev/null +++ b/Sense/ui/src/views/sense/operations/index.vue @@ -0,0 +1,170 @@ + + + + + diff --git a/Sense/ui/src/views/sense/operations/operationsState.js b/Sense/ui/src/views/sense/operations/operationsState.js new file mode 100644 index 0000000..243119f --- /dev/null +++ b/Sense/ui/src/views/sense/operations/operationsState.js @@ -0,0 +1,50 @@ +const objectLabels = { + device: '设备', + media: '媒体', + local_inference: '本地推理' +} + +const problemLabels = { + authentication_failed: '认证失败', + backoff_wait: '退避等待', + clock_drift: '时间漂移', + orphan_safety_gate: '孤儿安全闸', + capability_unavailable: '能力未安装', + not_converged: '尚未收敛' +} + +const stateLabels = { + active: '已启用', + ready: '就绪', + running: '运行', + stopped: '已停止', + waiting: '等待拉流', + pending: '等待处理', + retry_pending: '重试已排队', + authentication_failed: '认证失败', + clock_drift: '时间漂移', + apply_failed: '配置失败', + process_unavailable: '进程不可用', + status_unavailable: '状态不可用', + unavailable: '不可用(未安装)' +} + +export function objectLabel(value) { return objectLabels[value] || value || '未知' } +export function problemLabel(value) { return problemLabels[value] || value || '未知问题' } +export function stateLabel(value) { + return String(value || '').split(' / ').map(item => stateLabels[item] || item).join(' / ') || '未知' +} +export function severityLabel(value) { return ({ high: '高', medium: '中', info: '提示' })[value] || value || '未知' } +export function severityType(value) { return ({ high: 'danger', medium: 'warning', info: 'info' })[value] || 'info' } +export function canRetry(item) { return Boolean(item?.retryable && !item?.retryInProgress && item?.version > 0) } + +export function buildOperationsQuery(query) { + return { + pageIndex: Number(query.pageIndex) || 1, + pageSize: Number(query.pageSize) || 10, + objectType: String(query.objectType || '').trim(), + problemType: String(query.problemType || '').trim(), + severity: String(query.severity || '').trim(), + keyword: String(query.keyword || '').trim() + } +} diff --git a/Sense/ui/tests/unit/sense/operationsApi.spec.js b/Sense/ui/tests/unit/sense/operationsApi.spec.js new file mode 100644 index 0000000..88ad98c --- /dev/null +++ b/Sense/ui/tests/unit/sense/operationsApi.spec.js @@ -0,0 +1,17 @@ +import request from '@/utils/request' +import { getOperationProblem, listOperationProblems, retryOperationProblem } from '@/api/sense/operations' + +jest.mock('@/utils/request', () => jest.fn()) + +describe('Sense operations API', () => { + beforeEach(() => request.mockReset()) + + test('encodes identifiers and sends only the expected version for controlled retry', () => { + listOperationProblems({ pageIndex: 1, pageSize: 10 }) + getOperationProblem('media:camera/unsafe') + retryOperationProblem('media:camera/unsafe', 7) + expect(request).toHaveBeenNthCalledWith(1, { url: '/api/v1/operations', method: 'get', params: { pageIndex: 1, pageSize: 10 }}) + expect(request).toHaveBeenNthCalledWith(2, { url: '/api/v1/operations/media%3Acamera%2Funsafe', method: 'get' }) + expect(request).toHaveBeenNthCalledWith(3, { url: '/api/v1/operations/media%3Acamera%2Funsafe/retry', method: 'post', data: { expectedVersion: 7 }}) + }) +}) diff --git a/Sense/ui/tests/unit/sense/operationsState.spec.js b/Sense/ui/tests/unit/sense/operationsState.spec.js new file mode 100644 index 0000000..0471b75 --- /dev/null +++ b/Sense/ui/tests/unit/sense/operationsState.spec.js @@ -0,0 +1,23 @@ +import { buildOperationsQuery, canRetry, objectLabel, problemLabel, severityType, stateLabel } from '@/views/sense/operations/operationsState' + +describe('Sense operations presentation state', () => { + test('uses operator-facing state labels', () => { + expect(objectLabel('local_inference')).toBe('本地推理') + expect(problemLabel('orphan_safety_gate')).toBe('孤儿安全闸') + expect(stateLabel('unavailable')).toBe('不可用(未安装)') + expect(stateLabel('active / ready')).toBe('已启用 / 就绪') + expect(severityType('high')).toBe('danger') + }) + + test('allows retry only when the server projection says it is safe', () => { + expect(canRetry({ retryable: true, retryInProgress: false, version: 2 })).toBe(true) + expect(canRetry({ retryable: true, retryInProgress: true, version: 2 })).toBe(false) + expect(canRetry({ retryable: false, retryInProgress: false, version: 2 })).toBe(false) + }) + + test('builds an allowlisted query', () => { + expect(buildOperationsQuery({ pageIndex: 2, pageSize: 20, objectType: ' device ', problemType: 'clock_drift', severity: ' medium ', keyword: ' 南门 ', ignored: 'not-sent' })).toEqual({ + pageIndex: 2, pageSize: 20, objectType: 'device', problemType: 'clock_drift', severity: 'medium', keyword: '南门' + }) + }) +})