Compare commits

...
Author SHA1 Message Date
QiuSW 35b9c7b8a4 docs: 记录工单 #67 验收通过 2026-08-14 18:16:15 +08:00
QiuSW 263a68c98d docs: 归档工单 #67 待验收证据 2026-08-14 18:00:44 +08:00
QiuSW 4149cc4426 docs: 记录 Sense 视频服务生命周期 (#67) 2026-08-14 17:51:16 +08:00
QiuSW 19f9bfa1a0 feat: 重建 MediaMTX 生命周期与状态对账 (#67) 2026-08-14 17:42:26 +08:00
ila 24cf4bbd65 Merge pull request '#84' from feature/66-sense-onvif-admission into dev
[SEN] 在 GoAdmin 基线上重建 ONVIF 发现、手工接入与 RTSP Profile (#66)
2026-08-14 17:18:27 +08:00
QiuSW ab54235819 docs: 记录工单 #66 验收通过 2026-08-14 17:18:13 +08:00
QiuSW d74e1aee85 docs: 归档工单 #66 待验收证据 2026-08-14 17:09:40 +08:00
QiuSW 26e2e632d5 docs: 记录 Sense 视频接入安全边界 (#66) 2026-08-14 16:58:13 +08:00
QiuSW 2bb1614a3c feat: 重建 ONVIF 与 RTSP 视频接入 (#66) 2026-08-14 16:50:00 +08:00
ila 5856b76de9 Merge pull request '#83' from feature/65-sense-device-credential into dev
[SEN] 在 GoAdmin 基线上重建设备台账与凭据边界 (#65)
2026-08-14 16:27:59 +08:00
56 changed files with 3350 additions and 11 deletions
+4
View File
@@ -42,4 +42,8 @@ go run . server -c C:\secure-path\sense-settings.yml
设备台账本身可以在不配置摄像头凭据的情况下使用。创建或更新 ONVIF/RTSP 凭据前,还必须在启动进程环境中设置 `SENSE_CREDENTIAL_KEY`:该值是随机 32 字节密钥的 Base64 编码,仅保存在仓库外。变量名模板见 `server/config/credential.env.example`;不要把真实值写入配置、脚本、日志或工单。密钥缺失或格式不正确时,Sense 会拒绝凭据写入,不会降级为明文存储。
使用 ONVIF 发现或手工接入前,还必须设置 `SENSE_ONVIF_DISCOVERY_IP` 和 `SENSE_ONVIF_ALLOWED_CIDRS`。前者只能是获准用于 WS-Discovery 的本机网卡地址;后者是获准访问的摄像头网段(多个 CIDR 用逗号分隔)。未配置时系统会给出可行动提示且不会扫描任意网卡;手工地址、Media XAddr 和 Stream URI 同样受该网段限制,并拒绝重定向或 URL 内凭据。
MediaMTX 保持独立二进制。配置 `SENSE_MEDIAMTX_BINARY`、`SENSE_MEDIAMTX_CONFIG` 和只允许回环地址的 `SENSE_MEDIAMTX_API`。Sense 只生成无摄像头凭据的基础配置;路径和凭据在运行时通过回环 Control API 下发。模板见 `server/config/mediamtx/mediamtx.yml.example`。
仓库不提供默认账号、默认密码或可用密钥。首位管理员通过受仓库外 `SENSE_BOOTSTRAP_TOKEN` 保护的一次性初始化接口创建,详细步骤以项目 Wiki 的本地开发与验证页为准。
@@ -0,0 +1,18 @@
package router
import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseAdmissionRouter) }
func registerSenseAdmissionRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &admission.API{}
r := v1.Group("/admission").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("/discover", api.Discover)
r.GET("/devices/:id", api.Get)
r.POST("/devices/:id/probe", api.Probe)
}
@@ -0,0 +1,21 @@
package router
import (
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
)
func init() { routerCheckRole = append(routerCheckRole, registerSenseMediaRouter) }
func registerSenseMediaRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
api := &media.API{}
r := v1.Group("/media").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
r.GET("/routes", api.List)
r.GET("/process", api.Process)
r.POST("/reconcile", api.ReconcileAll)
r.POST("/routes/:id/reconcile", api.Reconcile)
r.POST("/routes/:id/stop", api.Stop)
}
+109
View File
@@ -0,0 +1,109 @@
package admission
import (
"encoding/json"
"errors"
"io"
"net/http"
"strings"
"github.com/gin-gonic/gin"
"github.com/go-admin-team/go-admin-core/sdk/api"
"github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
)
type API struct{ api.Api }
func (e *API) runtime(c *gin.Context) (*Service, error) {
base := Service{}
if err := e.MakeContext(c).MakeOrm().MakeService(&base.Service).Errors; err != nil {
return nil, err
}
return NewRuntime(base.Service)
}
func (e *API) Discover(c *gin.Context) {
service, err := e.runtime(c)
if err != nil {
e.writeError(err)
return
}
result, err := service.Discover(c.Request.Context())
if err != nil {
e.writeError(err)
return
}
e.OK(result, "发现完成")
}
func (e *API) Probe(c *gin.Context) {
service, err := e.runtime(c)
if err != nil {
e.writeError(err)
return
}
request := ProbeRequest{DeviceID: c.Param("id"), UpdateBy: user.GetUserId(c)}
if err = bindStrict(c, &request); err != nil {
e.Error(http.StatusBadRequest, err, "请求内容格式不正确")
return
}
result, err := service.Probe(c.Request.Context(), request)
if err != nil {
e.writeError(err)
return
}
// Media route intent is best-effort and credential-free. A MediaMTX
// failure must never roll back the verified device/Profile transaction.
_ = media.EnsureDeviceRoutes(c.Request.Context(), service.Orm, request.DeviceID)
e.OK(result, "探测完成")
}
func (e *API) Get(c *gin.Context) {
service := &Service{}
if err := e.MakeContext(c).MakeOrm().MakeService(&service.Service).Errors; err != nil {
e.writeError(err)
return
}
result, err := service.Get(c.Param("id"))
if err != nil {
e.writeError(err)
return
}
e.OK(result, "查询成功")
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrInvalid):
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, ErrNotFound):
e.Error(http.StatusNotFound, err, err.Error())
case errors.Is(err, ErrConflict):
e.Error(http.StatusConflict, err, err.Error())
case errors.Is(err, onvif.ErrDiscoveryNotConfigured), errors.Is(err, onvif.ErrDiscoveryInterface), errors.Is(err, onvif.ErrTargetNotAllowed):
e.Error(http.StatusBadRequest, err, err.Error())
case errors.Is(err, credential.ErrCredentialNotConfigured):
e.Error(http.StatusConflict, err, "请先在设备管理中配置摄像头凭据")
case errors.Is(err, credential.ErrKeyUnavailable):
e.Error(http.StatusServiceUnavailable, err, "摄像头凭据安全配置不可用")
default:
e.Error(http.StatusInternalServerError, err, "视频接入操作失败")
}
}
func bindStrict(c *gin.Context, target any) error {
if !strings.HasPrefix(strings.ToLower(strings.TrimSpace(c.GetHeader("Content-Type"))), "application/json") {
return errors.New("content type must be application/json")
}
decoder := json.NewDecoder(http.MaxBytesReader(c.Writer, c.Request.Body, 64<<10))
decoder.DisallowUnknownFields()
if err := decoder.Decode(target); err != nil {
return err
}
if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
if err == nil {
return errors.New("request body must contain one JSON object")
}
return err
}
return nil
}
@@ -0,0 +1,21 @@
package admission
import (
"net/http/httptest"
"strings"
"testing"
"github.com/gin-gonic/gin"
)
func TestProbePayloadRejectsCredentialFields(t *testing.T) {
gin.SetMode(gin.TestMode)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Request = httptest.NewRequest("POST", "/api/v1/admission/devices/device/probe", strings.NewReader(`{"address":"http://192.0.2.10/onvif","version":1,"password":"must-not-be-accepted"}`))
ctx.Request.Header.Set("Content-Type", "application/json")
var request ProbeRequest
if err := bindStrict(ctx, &request); err == nil {
t.Fatal("credential-like unknown field accepted")
}
}
+34
View File
@@ -0,0 +1,34 @@
package admission
import "time"
type ProbeRequest struct {
DeviceID string `json:"-"`
Address string `json:"address"`
Version int64 `json:"version"`
UpdateBy int `json:"-"`
}
type ProfileResponse struct {
Token string `json:"token"`
Name string `json:"name"`
Width int `json:"width"`
Height int `json:"height"`
Encoding string `json:"encoding"`
StreamURI string `json:"streamUri"`
Kind string `json:"kind"`
VerificationStatus string `json:"verificationStatus"`
VerificationLatencyMS int64 `json:"verificationLatencyMs"`
VerificationDetail string `json:"verificationDetail"`
}
type ResultResponse struct {
DeviceID string `json:"deviceId"`
Address string `json:"address"`
Status string `json:"status"`
Detail string `json:"detail"`
CheckedAt time.Time `json:"checkedAt"`
Profiles []ProfileResponse `json:"profiles"`
}
type DiscoveryResponse struct {
Addresses []string `json:"addresses"`
Interface string `json:"interface"`
}
@@ -0,0 +1,32 @@
package admission
import "time"
type Result struct {
DeviceID string `gorm:"size:36;primaryKey"`
Address string `gorm:"size:1024;not null"`
Status string `gorm:"size:32;not null;index"`
Detail string `gorm:"size:512;not null"`
CheckedAt time.Time `gorm:"not null"`
UpdatedAt time.Time
Profiles []Profile `gorm:"foreignKey:DeviceID;references:DeviceID;constraint:OnDelete:CASCADE"`
}
func (Result) TableName() string { return "sense_admission_results" }
type Profile struct {
DeviceID string `gorm:"size:36;primaryKey"`
Token string `gorm:"size:255;primaryKey"`
Name string `gorm:"size:255;not null"`
Width int `gorm:"not null"`
Height int `gorm:"not null"`
Encoding string `gorm:"size:32;not null"`
StreamURI string `gorm:"size:2048;not null"`
Kind string `gorm:"size:16;not null"`
VerificationStatus string `gorm:"size:32;not null"`
VerificationLatencyMS int64 `gorm:"not null"`
VerificationDetail string `gorm:"size:512;not null"`
UpdatedAt time.Time
}
func (Profile) TableName() string { return "sense_admission_profiles" }
+171
View File
@@ -0,0 +1,171 @@
package admission
import (
"context"
"errors"
"fmt"
"os"
"sort"
"strings"
"time"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/rtsp"
)
var (
ErrNotFound = errors.New("尚无该设备的接入结果")
ErrInvalid = errors.New("接入请求不符合要求")
ErrConflict = errors.New("设备已被其他用户更新,请刷新后重试")
)
type Service struct {
coreService.Service
ONVIF onvif.Client
RTSP rtsp.Verifier
Policy onvif.Policy
DiscoveryIP string
Timeout time.Duration
}
func NewRuntime(service coreService.Service) (*Service, error) {
policy, err := onvif.ParsePolicy(os.Getenv("SENSE_ONVIF_ALLOWED_CIDRS"))
if err != nil {
return nil, err
}
return &Service{Service: service, Policy: policy, DiscoveryIP: strings.TrimSpace(os.Getenv("SENSE_ONVIF_DISCOVERY_IP")), ONVIF: onvif.NewHTTPClient(8*time.Second, policy), RTSP: rtsp.NetVerifier{Timeout: 5 * time.Second, Policy: policy}}, nil
}
func (s *Service) Discover(ctx context.Context) (DiscoveryResponse, error) {
addresses, err := onvif.Discover(ctx, s.DiscoveryIP, 3*time.Second, s.Policy)
return DiscoveryResponse{Addresses: addresses, Interface: s.DiscoveryIP}, err
}
func (s *Service) Probe(ctx context.Context, request ProbeRequest) (ResultResponse, error) {
if request.Version < 1 || strings.TrimSpace(request.Address) == "" || len(request.Address) > 1024 {
return ResultResponse{}, ErrInvalid
}
var device deviceModels.Device
if err := s.Orm.Select("id", "modality", "version").First(&device, "id = ?", request.DeviceID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return ResultResponse{}, ErrNotFound
}
return ResultResponse{}, err
}
if device.Modality != deviceModels.ModalityVideo {
return ResultResponse{}, ErrInvalid
}
if device.Version != request.Version {
return ResultResponse{}, ErrConflict
}
onvifValue, err := credential.Read(s.Orm, request.DeviceID, credential.PurposeONVIF)
if err != nil {
return ResultResponse{}, err
}
rtspValue, err := credential.Read(s.Orm, request.DeviceID, credential.PurposeRTSP)
if err != nil {
return ResultResponse{}, err
}
profiles, probeErr := s.ONVIF.Profiles(ctx, request.Address, onvif.Credential{Username: onvifValue.Username, Password: onvifValue.Password})
now := time.Now().UTC()
result := Result{DeviceID: request.DeviceID, Address: strings.TrimSpace(request.Address), CheckedAt: now, UpdatedAt: now}
if probeErr != nil {
result.Status, result.Detail = classify(probeErr)
if err = s.save(result, request, false); err != nil {
return ResultResponse{}, err
}
return s.Get(request.DeviceID)
}
for _, profile := range profiles {
verification, verifyErr := s.RTSP.Verify(ctx, profile.StreamURI, rtsp.Credential{Username: rtspValue.Username, Password: rtspValue.Password})
if verifyErr != nil {
verification = rtsp.Result{Status: "failed", Detail: "视频地址未通过安全检查"}
}
result.Profiles = append(result.Profiles, Profile{DeviceID: request.DeviceID, Token: profile.Token, Name: profile.Name, Width: profile.Width, Height: profile.Height, Encoding: profile.Encoding, StreamURI: profile.StreamURI, Kind: "other", VerificationStatus: verification.Status, VerificationLatencyMS: verification.LatencyMS, VerificationDetail: verification.Detail, UpdatedAt: now})
}
sort.Slice(result.Profiles, func(i, j int) bool {
return result.Profiles[i].Width*result.Profiles[i].Height > result.Profiles[j].Width*result.Profiles[j].Height
})
if len(result.Profiles) > 0 {
result.Profiles[0].Kind = "main"
}
if len(result.Profiles) > 1 {
result.Profiles[len(result.Profiles)-1].Kind = "sub"
}
result.Status = "ready"
result.Detail = "设备与视频 Profile 已验证"
for _, profile := range result.Profiles {
if profile.VerificationStatus != "ready" {
result.Status = "profile_failed"
result.Detail = "部分视频 Profile 验证失败"
}
}
if err = s.save(result, request, true); err != nil {
return ResultResponse{}, err
}
return response(result), nil
}
func (s *Service) save(result Result, request ProbeRequest, replaceProfiles bool) error {
return s.Orm.Transaction(func(tx *gorm.DB) error {
updates := map[string]any{"version": request.Version + 1, "update_by": request.UpdateBy, "updated_at": result.UpdatedAt, "retry_requested_at": nil}
if replaceProfiles {
updates["status"] = map[bool]string{true: deviceModels.StatusActive, false: deviceModels.StatusPending}[result.Status == "ready"]
updates["adapter_status"] = map[bool]string{true: deviceModels.AdapterReady, false: deviceModels.AdapterFailed}[result.Status == "ready"]
}
update := tx.Model(&deviceModels.Device{}).Where("id = ? AND version = ?", request.DeviceID, request.Version).Updates(updates)
if update.Error != nil {
return update.Error
}
if update.RowsAffected == 0 {
return ErrConflict
}
if err := tx.Omit("Profiles").Save(&result).Error; err != nil {
return err
}
if replaceProfiles {
if err := tx.Where("device_id = ?", request.DeviceID).Delete(&Profile{}).Error; err != nil {
return err
}
}
if replaceProfiles && len(result.Profiles) > 0 {
return tx.Create(&result.Profiles).Error
}
return nil
})
}
func (s *Service) Get(deviceID string) (ResultResponse, error) {
var result Result
if err := s.Orm.Preload("Profiles", func(db *gorm.DB) *gorm.DB { return db.Order("width * height DESC") }).First(&result, "device_id = ?", deviceID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return ResultResponse{}, ErrNotFound
}
return ResultResponse{}, fmt.Errorf("read admission: %w", err)
}
return response(result), nil
}
func response(result Result) ResultResponse {
out := ResultResponse{DeviceID: result.DeviceID, Address: result.Address, Status: result.Status, Detail: result.Detail, CheckedAt: result.CheckedAt, Profiles: make([]ProfileResponse, 0, len(result.Profiles))}
for _, p := range result.Profiles {
out.Profiles = append(out.Profiles, ProfileResponse{Token: p.Token, Name: p.Name, Width: p.Width, Height: p.Height, Encoding: p.Encoding, StreamURI: p.StreamURI, Kind: p.Kind, VerificationStatus: p.VerificationStatus, VerificationLatencyMS: p.VerificationLatencyMS, VerificationDetail: p.VerificationDetail})
}
return out
}
func classify(err error) (string, string) {
switch {
case errors.Is(err, onvif.ErrAuthentication):
return "authentication_failed", "设备拒绝了当前凭据,请更新后重试"
case errors.Is(err, onvif.ErrTargetNotAllowed):
return "target_not_allowed", "设备地址不在获准网段内"
case errors.Is(err, onvif.ErrRedirect):
return "redirect_rejected", "设备返回了不允许的重定向"
case errors.Is(err, context.DeadlineExceeded) || strings.Contains(strings.ToLower(err.Error()), "timeout"):
return "timeout", "设备响应超时"
case strings.Contains(strings.ToLower(err.Error()), "time") || strings.Contains(strings.ToLower(err.Error()), "clock"):
return "clock_skew", "设备时间可能不准确,请校时后重试"
default:
return "unreachable", "无法读取设备信息,请检查地址和网络"
}
}
@@ -0,0 +1,107 @@
package admission
import (
"context"
"encoding/base64"
"errors"
"testing"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/rtsp"
)
type fakeONVIF struct{ err error }
func (f fakeONVIF) Profiles(context.Context, string, onvif.Credential) ([]onvif.Profile, error) {
if f.err != nil {
return nil, f.err
}
return []onvif.Profile{{Token: "main", Name: "主码流", Width: 1920, Height: 1080, Encoding: "H264", StreamURI: "rtsp://192.0.2.10/main"}, {Token: "sub", Name: "子码流", Width: 640, Height: 360, Encoding: "H264", StreamURI: "rtsp://192.0.2.10/sub"}}, nil
}
type fakeRTSP struct{}
func (fakeRTSP) Verify(context.Context, string, rtsp.Credential) (rtsp.Result, error) {
return rtsp.Result{Status: "ready", Detail: "码流可访问"}, nil
}
func admissionService(t *testing.T) *Service {
t.Helper()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&deviceModels.Device{}, &credential.DeviceCredential{}, &Result{}, &Profile{}); err != nil {
t.Fatal(err)
}
key := []byte("0123456789abcdef0123456789abcdef")
t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
vault, _ := credential.NewVault(key)
device := deviceModels.Device{ID: "device-1", Name: "东门摄像机", Modality: deviceModels.ModalityVideo, Version: 1, Status: "pending", AdapterStatus: "ready"}
if err = db.Create(&device).Error; err != nil {
t.Fatal(err)
}
for _, purpose := range []string{credential.PurposeONVIF, credential.PurposeRTSP} {
cipher, _ := vault.Encrypt(device.ID, purpose, "synthetic-user", "synthetic-password")
if err = db.Create(&credential.DeviceCredential{DeviceID: device.ID, Purpose: purpose, Ciphertext: cipher, KeyVersion: credential.Version()}).Error; err != nil {
t.Fatal(err)
}
}
return &Service{Service: coreService.Service{Orm: db}, ONVIF: fakeONVIF{}, RTSP: fakeRTSP{}}
}
func TestProbePersistsProfilesWithoutReturningCredentials(t *testing.T) {
service := admissionService(t)
result, err := service.Probe(context.Background(), ProbeRequest{DeviceID: "device-1", Address: "http://192.0.2.10/onvif", Version: 1})
if err != nil {
t.Fatal(err)
}
if result.Status != "ready" || len(result.Profiles) != 2 || result.Profiles[0].Kind != "main" || result.Profiles[1].Kind != "sub" {
t.Fatalf("result=%#v", result)
}
if result.Profiles[0].StreamURI == "" {
t.Fatal("stream URI missing")
}
saved, err := service.Get("device-1")
if err != nil || len(saved.Profiles) != 2 {
t.Fatalf("saved=%#v err=%v", saved, err)
}
}
func TestFailedReprobePreservesLastVerifiedProfiles(t *testing.T) {
service := admissionService(t)
if _, err := service.Probe(context.Background(), ProbeRequest{DeviceID: "device-1", Address: "http://192.0.2.10/onvif", Version: 1}); err != nil {
t.Fatal(err)
}
service.ONVIF = fakeONVIF{err: onvif.ErrAuthentication}
result, err := service.Probe(context.Background(), ProbeRequest{DeviceID: "device-1", Address: "http://192.0.2.10/onvif", Version: 2})
if err != nil {
t.Fatal(err)
}
if result.Status != "authentication_failed" || len(result.Profiles) != 2 {
t.Fatalf("last verified profiles lost: %#v", result)
}
if _, err = service.Probe(context.Background(), ProbeRequest{DeviceID: "device-1", Address: "http://192.0.2.10/onvif", Version: 2}); !errors.Is(err, ErrConflict) {
t.Fatalf("stale error=%v", err)
}
}
func TestProbeErrorsHaveActionableStates(t *testing.T) {
for _, test := range []struct {
err error
status string
}{
{onvif.ErrAuthentication, "authentication_failed"},
{onvif.ErrTargetNotAllowed, "target_not_allowed"},
{errors.New("device clock time fault"), "clock_skew"},
{context.DeadlineExceeded, "timeout"},
} {
status, detail := classify(test.err)
if status != test.status || detail == "" {
t.Fatalf("error=%v status=%s detail=%s", test.err, status, detail)
}
}
}
@@ -0,0 +1,35 @@
package credential
import (
"errors"
"fmt"
"gorm.io/gorm"
)
var ErrCredentialNotConfigured = errors.New("摄像头凭据尚未配置")
type Value struct {
Username string
Password string
}
// Read is an internal adapter port. HTTP handlers must never expose Value.
func Read(db *gorm.DB, deviceID, purpose string) (Value, error) {
var row DeviceCredential
if err := db.First(&row, "device_id = ? AND purpose = ?", deviceID, purpose).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return Value{}, ErrCredentialNotConfigured
}
return Value{}, fmt.Errorf("read device credential: %w", err)
}
vault, err := NewVaultFromEnvironment()
if err != nil {
return Value{}, err
}
username, password, err := vault.Decrypt(deviceID, purpose, row.Ciphertext)
if err != nil {
return Value{}, err
}
return Value{Username: username, Password: password}, nil
}
@@ -9,8 +9,10 @@ import (
const (
ModalityVideo = "video"
StatusPending = "pending"
StatusActive = "active"
StatusDisabled = "disabled"
AdapterReady = "ready"
AdapterFailed = "verification_failed"
AdapterNotReady = "adapter_not_ready"
)
+98
View File
@@ -0,0 +1,98 @@
package media
import (
"errors"
"net/http"
"github.com/gin-gonic/gin"
"github.com/go-admin-team/go-admin-core/sdk/api"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
)
type API struct{ api.Api }
func (e *API) service(c *gin.Context) (*Service, error) {
base := coreService.Service{}
if err := e.MakeContext(c).MakeOrm().MakeService(&base).Errors; err != nil {
return nil, err
}
return serviceFor(base.Orm), nil
}
func (e *API) List(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
items, err := service.List(c.Request.Context())
if err != nil {
e.writeError(err)
return
}
e.OK(gin.H{"list": items, "total": len(items)}, "查询成功")
}
func (e *API) Process(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
e.OK(service.ProcessState(), "查询成功")
}
func (e *API) ReconcileAll(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
if err = service.EnsureAllVerified(c.Request.Context()); err == nil {
err = service.ReconcileDue(c.Request.Context())
}
if err != nil {
e.writeError(err)
return
}
e.OK(nil, "对账完成")
}
func (e *API) Reconcile(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
item, err := service.Reconcile(c.Request.Context(), c.Param("id"))
if err != nil && item.ID == "" {
e.writeError(err)
return
}
e.OK(item, "对账完成")
}
func (e *API) Stop(c *gin.Context) {
service, err := e.service(c)
if err != nil {
e.writeError(err)
return
}
item, err := service.StopRoute(c.Request.Context(), c.Param("id"))
if err != nil {
e.writeError(err)
return
}
e.OK(item, "媒体路径已停止")
}
func (e *API) writeError(err error) {
switch {
case errors.Is(err, ErrNotFound):
e.Error(http.StatusNotFound, err, err.Error())
case errors.Is(err, ErrRuntimeUnavailable), errors.Is(err, ErrUnsafeControlAPI):
e.Error(http.StatusServiceUnavailable, err, "视频服务运行配置不可用")
default:
e.Error(http.StatusInternalServerError, err, "视频服务操作失败")
}
}
+91
View File
@@ -0,0 +1,91 @@
package media
import (
"errors"
"fmt"
"io"
"net"
"net/url"
"os"
"path/filepath"
"strings"
"time"
)
var ErrUnsafeControlAPI = errors.New("MediaMTX Control API 必须使用本机回环地址")
type RuntimeConfig struct {
Binary string
ConfigPath string
APIBase string
PollInterval time.Duration
StartTimeout time.Duration
}
func ConfigFromEnvironment() (RuntimeConfig, error) {
c := RuntimeConfig{
Binary: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_BINARY")), ConfigPath: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_CONFIG")),
APIBase: strings.TrimSpace(os.Getenv("SENSE_MEDIAMTX_API")), PollInterval: 5 * time.Second, StartTimeout: 8 * time.Second,
}
if c.APIBase == "" {
c.APIBase = "http://127.0.0.1:9997"
}
if err := validateControlAPI(c.APIBase); err != nil {
return RuntimeConfig{}, err
}
if c.Binary != "" && c.ConfigPath == "" {
return RuntimeConfig{}, errors.New("SENSE_MEDIAMTX_CONFIG is required when SENSE_MEDIAMTX_BINARY is configured")
}
return c, nil
}
func validateControlAPI(value string) error {
u, err := url.Parse(value)
if err != nil || u.Scheme != "http" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || u.Path != "" {
return ErrUnsafeControlAPI
}
host := u.Hostname()
ip := net.ParseIP(host)
if host != "localhost" && (ip == nil || !ip.IsLoopback()) {
return ErrUnsafeControlAPI
}
return nil
}
func RenderBaseConfig(w io.Writer, apiBase string) error {
if err := validateControlAPI(apiBase); err != nil {
return err
}
u, _ := url.Parse(apiBase)
_, err := fmt.Fprintf(w, "logLevel: info\napi: true\napiAddress: %s\nmetrics: false\npaths: {}\n", u.Host)
return err
}
func EnsureBaseConfig(path, apiBase string) error {
if path == "" {
return errors.New("MediaMTX config path is empty")
}
if _, err := os.Stat(path); err == nil {
return nil
} else if !errors.Is(err, os.ErrNotExist) {
return err
}
if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
return err
}
temporary := path + ".tmp"
f, err := os.OpenFile(temporary, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
if err != nil {
return err
}
writeErr := RenderBaseConfig(f, apiBase)
closeErr := f.Close()
if writeErr != nil || closeErr != nil {
_ = os.Remove(temporary)
return errors.Join(writeErr, closeErr)
}
if err = os.Rename(temporary, path); err != nil {
_ = os.Remove(temporary)
}
return err
}
+168
View File
@@ -0,0 +1,168 @@
package media
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"regexp"
"strings"
"time"
)
var validPath = regexp.MustCompile(`^[A-Za-z0-9_-]{1,96}$`)
type Source struct {
Path, URI, Username, Password string
}
type PathStatus struct {
Exists, Ready bool
Readers int
}
type Controller interface {
Health(context.Context) error
Apply(context.Context, Source) error
Delete(context.Context, string) error
Status(context.Context, string) (PathStatus, error)
}
type HTTPController struct {
base string
client *http.Client
}
func NewHTTPController(base string) (*HTTPController, error) {
if err := validateControlAPI(base); err != nil {
return nil, err
}
transport := http.DefaultTransport.(*http.Transport).Clone()
transport.Proxy = nil
return &HTTPController{base: strings.TrimRight(base, "/"), client: &http.Client{
Timeout: 5 * time.Second, Transport: transport,
CheckRedirect: func(*http.Request, []*http.Request) error {
return errors.New("MediaMTX Control API redirect rejected")
},
}}, nil
}
func (c *HTTPController) Health(ctx context.Context) error {
return c.request(ctx, http.MethodGet, "/v3/config/global/get", nil, nil)
}
func (c *HTTPController) Apply(ctx context.Context, source Source) error {
if !validPath.MatchString(source.Path) {
return errors.New("invalid MediaMTX path")
}
u, err := url.Parse(source.URI)
if err != nil || u.User != nil || u.Host == "" || (u.Scheme != "rtsp" && u.Scheme != "rtsps") {
return errors.New("invalid credential-free RTSP source")
}
payload := map[string]any{"source": u.String(), "sourceOnDemand": true, "rtspTransport": "tcp"}
if source.Username != "" {
// MediaMTX v1.19.3 has no sourceUser/sourcePass fields. Credentials
// are assembled only for this loopback request and are never stored,
// logged, or returned by a Sense endpoint.
u.User = url.UserPassword(source.Username, source.Password)
payload["source"] = u.String()
}
data, err := json.Marshal(payload)
if err != nil {
return err
}
configured, err := c.configured(ctx, source.Path)
if err != nil {
return err
}
method, action := http.MethodPost, "add"
if configured {
method, action = http.MethodPatch, "patch"
}
return c.request(ctx, method, "/v3/config/paths/"+action+"/"+url.PathEscape(source.Path), bytes.NewReader(data), nil)
}
func (c *HTTPController) configured(ctx context.Context, path string) (bool, error) {
status := 0
err := c.request(ctx, http.MethodGet, "/v3/config/paths/get/"+url.PathEscape(path), nil, &status)
if status == http.StatusNotFound {
return false, nil
}
return err == nil, err
}
func (c *HTTPController) Delete(ctx context.Context, path string) error {
if !validPath.MatchString(path) {
return errors.New("invalid MediaMTX path")
}
status := 0
err := c.request(ctx, http.MethodDelete, "/v3/config/paths/delete/"+url.PathEscape(path), nil, &status)
if status == http.StatusNotFound {
return nil
}
return err
}
func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, error) {
statusCode := 0
var raw struct {
Ready bool `json:"ready"`
Readers []any `json:"readers"`
}
err := c.requestJSON(ctx, http.MethodGet, "/v3/paths/get/"+url.PathEscape(path), &statusCode, &raw)
if statusCode == http.StatusNotFound {
return PathStatus{}, nil
}
if err != nil {
return PathStatus{}, err
}
return PathStatus{Exists: true, Ready: raw.Ready, Readers: len(raw.Readers)}, nil
}
func (c *HTTPController) request(ctx context.Context, method, path string, body io.Reader, statusOut *int) error {
return c.requestJSON(ctx, method, path, body, statusOut, nil)
}
func (c *HTTPController) requestJSON(ctx context.Context, method, path string, args ...any) error {
var body io.Reader
var statusOut *int
var target any
for _, arg := range args {
switch value := arg.(type) {
case io.Reader:
body = value
case *int:
statusOut = value
default:
target = value
}
}
req, err := http.NewRequestWithContext(ctx, method, c.base+path, body)
if err != nil {
return err
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
res, err := c.client.Do(req)
if err != nil {
return fmt.Errorf("MediaMTX Control API unavailable: %w", err)
}
defer res.Body.Close()
if statusOut != nil {
*statusOut = res.StatusCode
}
if res.StatusCode < 200 || res.StatusCode >= 300 {
_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
return fmt.Errorf("MediaMTX Control API returned %d", res.StatusCode)
}
if target == nil {
_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
return nil
}
return json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(target)
}
@@ -0,0 +1,67 @@
package media
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
func TestControllerUsesCredentialsOnlyAtLoopbackBoundary(t *testing.T) {
var payload map[string]any
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.URL.Path == "/v3/config/paths/get/sense_test":
http.NotFound(w, r)
case r.URL.Path == "/v3/config/paths/add/sense_test":
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
t.Error(err)
}
w.WriteHeader(http.StatusOK)
default:
http.NotFound(w, r)
}
}))
defer server.Close()
controller, err := NewHTTPController(server.URL)
if err != nil {
t.Fatal(err)
}
source := Source{Path: "sense_test", URI: "rtsp://192.0.2.1/live", Username: "synthetic-user", Password: "synthetic-pass"}
if err = controller.Apply(context.Background(), source); err != nil {
t.Fatal(err)
}
if payload["source"] != "rtsp://synthetic-user:synthetic-pass@192.0.2.1/live" {
t.Fatalf("unexpected MediaMTX payload: %#v", payload)
}
if source.URI != "rtsp://192.0.2.1/live" || strings.Contains(source.URI, "synthetic") {
t.Fatalf("caller source was mutated: %#v", source)
}
}
func TestControlAPIAndGeneratedConfigAreLoopbackOnly(t *testing.T) {
if _, err := NewHTTPController("http://192.0.2.1:9997"); err != ErrUnsafeControlAPI {
t.Fatalf("error=%v", err)
}
var output strings.Builder
if err := RenderBaseConfig(&output, "http://127.0.0.1:9997"); err != nil {
t.Fatal(err)
}
if strings.Contains(strings.ToLower(output.String()), "password") || !strings.Contains(output.String(), "paths: {}") {
t.Fatalf("unsafe config: %s", output.String())
}
}
func TestExternalProcessUsesOrphanSafetyGate(t *testing.T) {
supervisor := NewSupervisor("", "")
supervisor.MarkExternal()
if err := supervisor.Stop(context.Background()); err != nil {
t.Fatal(err)
}
state := supervisor.State()
if !state.External || state.Owned || state.Phase != "running" {
t.Fatalf("unexpected state: %#v", state)
}
}
+24
View File
@@ -0,0 +1,24 @@
package media
import "time"
type RouteResponse struct {
ID string `json:"id"`
DeviceID string `json:"deviceId"`
ProfileToken string `json:"profileToken"`
Path string `json:"path"`
Desired string `json:"desired"`
Actual string `json:"actual"`
Readers int `json:"readers"`
SourceReady bool `json:"sourceReady"`
FailureCount int `json:"failureCount"`
NextRetryAt *time.Time `json:"nextRetryAt,omitempty"`
LastErrorCode string `json:"lastErrorCode,omitempty"`
Detail string `json:"detail"`
Version int64 `json:"version"`
UpdatedAt time.Time `json:"updatedAt"`
}
func routeResponse(r Route) RouteResponse {
return RouteResponse{ID: r.ID, DeviceID: r.DeviceID, ProfileToken: r.ProfileToken, Path: r.Path, Desired: r.Desired, Actual: r.Actual, Readers: r.Readers, SourceReady: r.SourceReady, FailureCount: r.FailureCount, NextRetryAt: r.NextRetryAt, LastErrorCode: r.LastErrorCode, Detail: r.Detail, Version: r.Version, UpdatedAt: r.UpdatedAt}
}
+37
View File
@@ -0,0 +1,37 @@
package media
import "time"
const (
DesiredRunning = "running"
DesiredStopped = "stopped"
)
type Route struct {
ID string `gorm:"size:96;primaryKey"`
DeviceID string `gorm:"size:36;not null;uniqueIndex:media_device_profile"`
ProfileToken string `gorm:"size:255;not null;uniqueIndex:media_device_profile"`
Path string `gorm:"size:96;not null;uniqueIndex"`
Desired string `gorm:"size:16;not null;index"`
Actual string `gorm:"size:32;not null;index"`
Readers int `gorm:"not null"`
SourceReady bool `gorm:"not null"`
FailureCount int `gorm:"not null"`
NextRetryAt *time.Time `gorm:"index"`
LastErrorCode string `gorm:"size:64;not null"`
Detail string `gorm:"size:512;not null"`
Version int64 `gorm:"not null"`
UpdatedAt time.Time
}
func (Route) TableName() string { return "sense_media_routes" }
// admissionProfile intentionally maps only non-secret fields from #66.
type admissionProfile struct {
DeviceID string `gorm:"column:device_id"`
Token string `gorm:"column:token"`
StreamURI string `gorm:"column:stream_uri"`
VerificationStatus string `gorm:"column:verification_status"`
}
func (admissionProfile) TableName() string { return "sense_admission_profiles" }
+70
View File
@@ -0,0 +1,70 @@
package media
import (
"context"
"errors"
"sync"
"time"
"gorm.io/gorm"
)
var runtimeState struct {
sync.RWMutex
service *Service
cancel context.CancelFunc
}
func StartRuntime(parent context.Context, db *gorm.DB) error {
if db == nil {
return errors.New("Sense database is unavailable for MediaMTX runtime")
}
config, err := ConfigFromEnvironment()
if err != nil {
return err
}
controller, err := NewHTTPController(config.APIBase)
if err != nil {
return err
}
service := NewService(db, controller, NewSupervisor(config.Binary, config.ConfigPath), config)
ctx, cancel := context.WithCancel(parent)
runtimeState.Lock()
if runtimeState.cancel != nil {
runtimeState.cancel()
}
runtimeState.service, runtimeState.cancel = service, cancel
runtimeState.Unlock()
go service.Run(ctx)
return nil
}
func ShutdownRuntime(ctx context.Context) error {
runtimeState.Lock()
service, cancel := runtimeState.service, runtimeState.cancel
runtimeState.service, runtimeState.cancel = nil, nil
runtimeState.Unlock()
if cancel != nil {
cancel()
}
if service != nil && service.process != nil {
return service.process.Stop(ctx)
}
return nil
}
func serviceFor(db *gorm.DB) *Service {
runtimeState.RLock()
service := runtimeState.service
runtimeState.RUnlock()
if service != nil {
return service
}
return NewService(db, nil, nil, RuntimeConfig{PollInterval: 5 * time.Second})
}
// EnsureDeviceRoutes is an internal post-admission port. It persists only
// credential-free route intent and deliberately does not fail device probing.
func EnsureDeviceRoutes(ctx context.Context, db *gorm.DB, deviceID string) error {
return serviceFor(db).EnsureDevice(ctx, deviceID)
}
+305
View File
@@ -0,0 +1,305 @@
package media
import (
"context"
"crypto/sha256"
"errors"
"fmt"
"strings"
"sync"
"time"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/reconcile"
)
var (
ErrNotFound = errors.New("媒体路由不存在")
ErrRuntimeUnavailable = errors.New("视频服务运行配置不可用")
)
type Service struct {
db *gorm.DB
controller Controller
process *Supervisor
config RuntimeConfig
now func() time.Time
mu sync.Mutex
}
func NewService(db *gorm.DB, controller Controller, process *Supervisor, config RuntimeConfig) *Service {
return &Service{db: db, controller: controller, process: process, config: config, now: time.Now}
}
func (s *Service) EnsureAllVerified(ctx context.Context) error {
var ids []string
if err := s.db.WithContext(ctx).Model(&admissionProfile{}).Where("verification_status = ?", "ready").Distinct().Pluck("device_id", &ids).Error; err != nil {
return err
}
for _, id := range ids {
if err := s.ensureDevice(ctx, id, false); err != nil {
return err
}
}
return nil
}
func (s *Service) EnsureDevice(ctx context.Context, deviceID string) error {
return s.ensureDevice(ctx, deviceID, true)
}
func (s *Service) ensureDevice(ctx context.Context, deviceID string, reactivate bool) error {
if strings.TrimSpace(deviceID) == "" {
return errors.New("device id is required")
}
var profiles []admissionProfile
if err := s.db.WithContext(ctx).Where("device_id = ? AND verification_status = ?", deviceID, "ready").Find(&profiles).Error; err != nil {
return err
}
now := s.now().UTC()
for _, profile := range profiles {
id := profile.DeviceID + ":" + profile.Token
digest := sha256.Sum256([]byte(id))
route := Route{ID: id, DeviceID: profile.DeviceID, ProfileToken: profile.Token, Path: fmt.Sprintf("sense_%x", digest[:12]), Desired: DesiredRunning, Actual: "pending", Detail: "等待视频服务对账", Version: 1, UpdatedAt: now}
if err := s.db.WithContext(ctx).Clauses(clause.OnConflict{Columns: []clause.Column{{Name: "id"}}, DoNothing: true}).Create(&route).Error; err != nil {
return err
}
var existing Route
if err := s.db.WithContext(ctx).First(&existing, "id = ?", id).Error; err != nil {
return err
}
if reactivate && existing.Desired != DesiredRunning {
if err := s.db.WithContext(ctx).Model(&existing).Updates(map[string]any{"desired": DesiredRunning, "actual": "pending", "detail": "等待视频服务对账", "next_retry_at": nil, "version": existing.Version + 1, "updated_at": now}).Error; err != nil {
return err
}
}
}
return nil
}
func (s *Service) Run(ctx context.Context) {
_ = s.EnsureAllVerified(ctx)
_ = s.ReconcileDue(ctx)
interval := s.config.PollInterval
if interval <= 0 {
interval = 5 * time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
_ = s.ReconcileDue(ctx)
}
}
}
func (s *Service) ReconcileDue(ctx context.Context) error {
now := s.now().UTC()
var routes []Route
if err := s.db.WithContext(ctx).Where("desired = ? AND (next_retry_at IS NULL OR next_retry_at <= ?)", DesiredRunning, now).Order("updated_at ASC").Limit(128).Find(&routes).Error; err != nil {
return err
}
var result error
for _, route := range routes {
var err error
if route.FailureCount == 0 && (route.Actual == "waiting" || route.Actual == "ready") {
_, err = s.Refresh(ctx, route.ID)
} else {
_, err = s.Reconcile(ctx, route.ID)
}
if err != nil {
result = errors.Join(result, err)
}
}
return result
}
func (s *Service) Refresh(ctx context.Context, id string) (RouteResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
var route Route
if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
return RouteResponse{}, ErrNotFound
}
if err := s.ensureControl(ctx); err != nil {
return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err)
}
status, err := s.controller.Status(ctx, route.Path)
if err != nil {
return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err)
}
if !status.Exists {
// A cold MediaMTX start begins with paths: {}; re-apply the route
// after releasing the service lock.
s.mu.Unlock()
response, reconcileErr := s.Reconcile(ctx, id)
s.mu.Lock()
return response, reconcileErr
}
if status.Ready {
return s.saveSuccess(ctx, route, "ready", true, status.Readers, "上游拉流正常")
}
return s.saveSuccess(ctx, route, "waiting", false, status.Readers, "等待播放器连接并按需拉流")
}
func (s *Service) Reconcile(ctx context.Context, id string) (RouteResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
var route Route
if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return RouteResponse{}, ErrNotFound
}
return RouteResponse{}, err
}
if route.Desired == DesiredStopped {
if s.controller != nil {
_ = s.controller.Delete(ctx, route.Path)
}
return s.saveSuccess(ctx, route, "stopped", false, 0, "媒体路径已停止")
}
if err := s.ensureControl(ctx); err != nil {
return s.saveFailure(ctx, route, "process_unavailable", "MediaMTX 未启动或 Control API 未就绪", err)
}
var profile admissionProfile
if err := s.db.WithContext(ctx).Where("device_id = ? AND token = ? AND verification_status = ?", route.DeviceID, route.ProfileToken, "ready").First(&profile).Error; err != nil {
return s.saveFailure(ctx, route, "profile_unavailable", "已验证 Profile 不可用", err)
}
value, err := credential.Read(s.db.WithContext(ctx), route.DeviceID, credential.PurposeRTSP)
if err != nil {
return s.saveFailure(ctx, route, "credential_unavailable", "RTSP 凭据不可用", err)
}
if err = s.controller.Apply(ctx, Source{Path: route.Path, URI: profile.StreamURI, Username: value.Username, Password: value.Password}); err != nil {
return s.saveFailure(ctx, route, "apply_failed", "媒体路径配置失败", err)
}
status, err := s.controller.Status(ctx, route.Path)
if err != nil {
return s.saveFailure(ctx, route, "status_unavailable", "尚未取得媒体路径状态", err)
}
if !status.Exists {
return s.saveFailure(ctx, route, "path_missing", "媒体路径不存在", errors.New("MediaMTX path missing after apply"))
}
if status.Ready {
return s.saveSuccess(ctx, route, "ready", true, status.Readers, "上游拉流正常")
}
return s.saveSuccess(ctx, route, "waiting", false, status.Readers, "等待播放器连接并按需拉流")
}
func (s *Service) ensureControl(ctx context.Context) error {
if s.controller == nil {
return ErrRuntimeUnavailable
}
if err := s.controller.Health(ctx); err == nil {
if s.process != nil && !s.process.State().Owned {
s.process.MarkExternal()
}
return nil
}
if s.process == nil {
return ErrRuntimeUnavailable
}
if err := EnsureBaseConfig(s.config.ConfigPath, s.config.APIBase); err != nil {
s.process.MarkFailed("MediaMTX 基础配置不可用")
return err
}
if err := s.process.Start(ctx); err != nil {
return err
}
deadline := s.now().Add(s.config.StartTimeout)
for s.now().Before(deadline) {
if err := s.controller.Health(ctx); err == nil {
s.process.MarkReady()
return nil
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(100 * time.Millisecond):
}
}
s.process.MarkFailed("MediaMTX Control API 就绪超时")
return errors.New("MediaMTX readiness timeout")
}
func (s *Service) saveSuccess(ctx context.Context, route Route, actual string, sourceReady bool, readers int, detail string) (RouteResponse, error) {
if route.Actual == actual && route.SourceReady == sourceReady && route.Readers == readers && route.FailureCount == 0 && route.NextRetryAt == nil && route.LastErrorCode == "" && route.Detail == detail {
return routeResponse(route), nil
}
now := s.now().UTC()
updates := map[string]any{"actual": actual, "source_ready": sourceReady, "readers": readers, "failure_count": 0, "next_retry_at": nil, "last_error_code": "", "detail": detail, "version": route.Version + 1, "updated_at": now}
if err := s.db.WithContext(ctx).Model(&route).Updates(updates).Error; err != nil {
return RouteResponse{}, err
}
for key, value := range updates {
switch key {
case "actual":
route.Actual = value.(string)
case "source_ready":
route.SourceReady = value.(bool)
case "readers":
route.Readers = value.(int)
case "failure_count":
route.FailureCount = value.(int)
case "last_error_code":
route.LastErrorCode = value.(string)
case "detail":
route.Detail = value.(string)
case "version":
route.Version = value.(int64)
case "updated_at":
route.UpdatedAt = value.(time.Time)
}
}
route.NextRetryAt = nil
return routeResponse(route), nil
}
func (s *Service) saveFailure(ctx context.Context, route Route, code, detail string, cause error) (RouteResponse, error) {
now := s.now().UTC()
failures := route.FailureCount + 1
next := now.Add(reconcile.Backoff(failures))
updates := map[string]any{"actual": code, "source_ready": false, "readers": 0, "failure_count": failures, "next_retry_at": &next, "last_error_code": code, "detail": detail, "version": route.Version + 1, "updated_at": now}
if err := s.db.WithContext(ctx).Model(&route).Updates(updates).Error; err != nil {
return RouteResponse{}, errors.Join(cause, err)
}
route.Actual, route.SourceReady, route.Readers, route.FailureCount, route.NextRetryAt, route.LastErrorCode, route.Detail = code, false, 0, failures, &next, code, detail
route.Version, route.UpdatedAt = route.Version+1, now
return routeResponse(route), cause
}
func (s *Service) StopRoute(ctx context.Context, id string) (RouteResponse, error) {
var route Route
if err := s.db.WithContext(ctx).First(&route, "id = ?", id).Error; err != nil {
return RouteResponse{}, ErrNotFound
}
route.Desired = DesiredStopped
if err := s.db.WithContext(ctx).Model(&route).Updates(map[string]any{"desired": DesiredStopped, "next_retry_at": nil, "version": route.Version + 1, "updated_at": s.now().UTC()}).Error; err != nil {
return RouteResponse{}, err
}
return s.Reconcile(ctx, id)
}
func (s *Service) List(ctx context.Context) ([]RouteResponse, error) {
var routes []Route
if err := s.db.WithContext(ctx).Order("updated_at DESC").Find(&routes).Error; err != nil {
return nil, err
}
items := make([]RouteResponse, 0, len(routes))
for _, route := range routes {
items = append(items, routeResponse(route))
}
return items, nil
}
func (s *Service) ProcessState() ProcessState {
if s.process == nil {
return ProcessState{Phase: "configuration_failed", Detail: "视频服务运行配置不可用", LastChanged: s.now().UTC()}
}
return s.process.State()
}
@@ -0,0 +1,128 @@
package media
import (
"context"
"crypto/rand"
"encoding/base64"
"errors"
"os"
"testing"
"time"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
)
type fakeController struct {
healthErr, applyErr, statusErr error
applies int
status PathStatus
}
func (f *fakeController) Health(context.Context) error { return f.healthErr }
func (f *fakeController) Apply(context.Context, Source) error { f.applies++; return f.applyErr }
func (f *fakeController) Delete(context.Context, string) error { return nil }
func (f *fakeController) Status(context.Context, string) (PathStatus, error) {
return f.status, f.statusErr
}
func mediaTestService(t *testing.T, controller Controller) (*Service, *gorm.DB) {
t.Helper()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&Route{}, &admissionProfile{}, &credential.DeviceCredential{}); err != nil {
t.Fatal(err)
}
key := make([]byte, 32)
if _, err = rand.Read(key); err != nil {
t.Fatal(err)
}
t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
vault, _ := credential.NewVault(key)
ciphertext, err := vault.Encrypt("device-1", credential.PurposeRTSP, "synthetic-user", "synthetic-password")
if err != nil {
t.Fatal(err)
}
if err = db.Create(&credential.DeviceCredential{DeviceID: "device-1", Purpose: credential.PurposeRTSP, Ciphertext: ciphertext, KeyVersion: credential.Version()}).Error; err != nil {
t.Fatal(err)
}
if err = db.Create(&admissionProfile{DeviceID: "device-1", Token: "main", StreamURI: "rtsp://192.0.2.1/live", VerificationStatus: "ready"}).Error; err != nil {
t.Fatal(err)
}
service := NewService(db, controller, nil, RuntimeConfig{})
service.now = func() time.Time { return time.Date(2026, 8, 14, 0, 0, 0, 0, time.UTC) }
return service, db
}
func TestEnsureAndReconcileAreIdempotent(t *testing.T) {
controller := &fakeController{status: PathStatus{Exists: true, Ready: true, Readers: 2}}
service, db := mediaTestService(t, controller)
if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
t.Fatal(err)
}
if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
t.Fatal(err)
}
var count int64
if err := db.Model(&Route{}).Count(&count).Error; err != nil || count != 1 {
t.Fatalf("count=%d err=%v", count, err)
}
item, err := service.Reconcile(context.Background(), "device-1:main")
if err != nil {
t.Fatal(err)
}
if item.Actual != "ready" || item.Readers != 2 || controller.applies != 1 {
t.Fatalf("item=%#v applies=%d", item, controller.applies)
}
if err = service.ReconcileDue(context.Background()); err != nil {
t.Fatal(err)
}
items, err := service.List(context.Background())
if err != nil || controller.applies != 1 || items[0].Version != item.Version {
t.Fatalf("steady route was rewritten: items=%#v applies=%d err=%v", items, controller.applies, err)
}
}
func TestColdStartDoesNotReactivateStoppedRoute(t *testing.T) {
controller := &fakeController{status: PathStatus{Exists: true}}
service, _ := mediaTestService(t, controller)
if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
t.Fatal(err)
}
if _, err := service.StopRoute(context.Background(), "device-1:main"); err != nil {
t.Fatal(err)
}
if err := service.EnsureAllVerified(context.Background()); err != nil {
t.Fatal(err)
}
items, err := service.List(context.Background())
if err != nil || len(items) != 1 || items[0].Desired != DesiredStopped {
t.Fatalf("stopped route was reactivated: %#v err=%v", items, err)
}
}
func TestFailurePersistsBackoffWithoutChangingProfile(t *testing.T) {
controller := &fakeController{applyErr: errors.New("synthetic apply failure")}
service, db := mediaTestService(t, controller)
if err := service.EnsureDevice(context.Background(), "device-1"); err != nil {
t.Fatal(err)
}
item, err := service.Reconcile(context.Background(), "device-1:main")
if err == nil || item.Actual != "apply_failed" || item.FailureCount != 1 || item.NextRetryAt == nil {
t.Fatalf("item=%#v err=%v", item, err)
}
var profile admissionProfile
if err = db.First(&profile, "device_id = ? AND token = ?", "device-1", "main").Error; err != nil {
t.Fatal(err)
}
if profile.VerificationStatus != "ready" {
t.Fatalf("profile changed: %#v", profile)
}
if value := os.Getenv(credential.EnvironmentKey); value == "" {
t.Fatal("test key unexpectedly missing")
}
}
+149
View File
@@ -0,0 +1,149 @@
package media
import (
"context"
"errors"
"os"
"os/exec"
"path/filepath"
"sync"
"time"
)
type ProcessState struct {
Phase string `json:"phase"`
PID int `json:"pid,omitempty"`
Owned bool `json:"owned"`
External bool `json:"external"`
Restarts int `json:"restarts"`
Detail string `json:"detail"`
LastChanged time.Time `json:"lastChanged"`
}
type Process interface {
Start(context.Context) error
Stop(context.Context) error
State() ProcessState
}
type Supervisor struct {
binary, config string
mu sync.Mutex
command *exec.Cmd
state ProcessState
}
func NewSupervisor(binary, config string) *Supervisor {
phase, detail := "stopped", "MediaMTX 尚未启动"
if binary == "" {
phase, detail = "not_configured", "未配置 MediaMTX 二进制;可连接外部已启动实例"
}
return &Supervisor{binary: binary, config: config, state: ProcessState{Phase: phase, Detail: detail, LastChanged: time.Now().UTC()}}
}
func (s *Supervisor) MarkExternal() {
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Owned {
return
}
s.state.Phase, s.state.External, s.state.Detail, s.state.LastChanged = "running", true, "检测到外部 MediaMTX;孤儿安全闸禁止 Sense 停止该进程", time.Now().UTC()
}
func (s *Supervisor) Start(_ context.Context) error {
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Owned && s.command != nil {
return nil
}
if s.binary == "" {
return errors.New("SENSE_MEDIAMTX_BINARY is not configured")
}
args := []string{}
if s.config != "" {
args = append(args, s.config)
}
cmd := exec.Command(s.binary, args...)
if s.config != "" {
cmd.Dir = filepath.Dir(s.config)
}
if err := cmd.Start(); err != nil {
s.state.Phase, s.state.Detail, s.state.LastChanged = "failed", "MediaMTX 进程启动失败", time.Now().UTC()
return err
}
if s.state.Phase == "failed" {
s.state.Restarts++
}
s.command = cmd
s.state.Phase, s.state.PID, s.state.Owned, s.state.External = "starting", cmd.Process.Pid, true, false
s.state.Detail, s.state.LastChanged = "等待 MediaMTX Control API 就绪", time.Now().UTC()
go s.wait(cmd)
return nil
}
func (s *Supervisor) MarkReady() {
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Owned {
s.state.Phase, s.state.Detail, s.state.LastChanged = "running", "MediaMTX Control API 已就绪", time.Now().UTC()
}
}
func (s *Supervisor) MarkFailed(detail string) {
s.mu.Lock()
defer s.mu.Unlock()
s.state.Phase, s.state.Detail, s.state.LastChanged = "failed", detail, time.Now().UTC()
}
func (s *Supervisor) wait(cmd *exec.Cmd) {
err := cmd.Wait()
s.mu.Lock()
defer s.mu.Unlock()
if s.command != cmd {
return
}
s.command = nil
s.state.Phase, s.state.PID, s.state.Owned, s.state.External = "failed", 0, false, false
s.state.Detail = "MediaMTX 进程已退出"
if err == nil {
s.state.Phase, s.state.Detail = "stopped", "MediaMTX 进程已停止"
}
s.state.LastChanged = time.Now().UTC()
}
func (s *Supervisor) Stop(ctx context.Context) error {
s.mu.Lock()
cmd := s.command
owned := s.state.Owned
s.mu.Unlock()
if !owned || cmd == nil {
return nil
}
if err := cmd.Process.Signal(os.Interrupt); err != nil {
if killErr := cmd.Process.Kill(); killErr != nil {
return errors.Join(err, killErr)
}
}
ticker := time.NewTicker(50 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
_ = cmd.Process.Kill()
return ctx.Err()
case <-ticker.C:
s.mu.Lock()
finished := s.command != cmd
s.mu.Unlock()
if finished {
return nil
}
}
}
}
func (s *Supervisor) State() ProcessState {
s.mu.Lock()
defer s.mu.Unlock()
return s.state
}
+296
View File
@@ -0,0 +1,296 @@
package onvif
import (
"bytes"
"context"
"crypto/md5"
"crypto/rand"
"crypto/sha256"
"encoding/xml"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"strconv"
"strings"
"time"
)
var (
ErrAuthentication = errors.New("设备拒绝了当前凭据")
ErrRedirect = errors.New("设备返回了不允许的重定向")
)
type Credential struct{ Username, Password string }
type Profile struct {
Token string `json:"token"`
Name string `json:"name"`
Width int `json:"width"`
Height int `json:"height"`
Encoding string `json:"encoding"`
StreamURI string `json:"streamUri"`
}
type Client interface {
Profiles(context.Context, string, Credential) ([]Profile, error)
}
type HTTPClient struct {
client *http.Client
policy Policy
}
func NewHTTPClient(timeout time.Duration, policy Policy) *HTTPClient {
if timeout <= 0 {
timeout = 8 * time.Second
}
dialer := net.Dialer{Timeout: timeout}
transport := &http.Transport{Proxy: nil, DialContext: func(ctx context.Context, network, address string) (net.Conn, error) {
host, port, err := net.SplitHostPort(address)
if err != nil {
return nil, ErrAddressInvalid
}
_, ip, err := policy.ValidateURL(ctx, "http://"+net.JoinHostPort(host, port), "http")
if err != nil {
return nil, err
}
return dialer.DialContext(ctx, network, net.JoinHostPort(ip.String(), port))
}}
return &HTTPClient{policy: policy, client: &http.Client{Timeout: timeout, Transport: transport, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}}
}
func (c *HTTPClient) Profiles(ctx context.Context, address string, credential Credential) ([]Profile, error) {
device, _, err := c.policy.ValidateURL(ctx, address, "http", "https")
if err != nil {
return nil, err
}
capabilities, err := c.soap(ctx, device.String(), credential, `<?xml version="1.0"?><s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope"><s:Body><GetCapabilities xmlns="http://www.onvif.org/ver10/device/wsdl"><Category>All</Category></GetCapabilities></s:Body></s:Envelope>`)
if err != nil {
return nil, err
}
mediaRaw, err := parseElement(capabilities, "Media", "XAddr")
if err != nil {
return nil, err
}
media, err := c.normalizeService(ctx, device, mediaRaw)
if err != nil {
return nil, err
}
data, err := c.soap(ctx, media.String(), credential, `<?xml version="1.0"?><s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope"><s:Body><GetProfiles xmlns="http://www.onvif.org/ver10/media/wsdl"/></s:Body></s:Envelope>`)
if err != nil {
return nil, err
}
profiles, err := parseProfiles(data)
if err != nil {
return nil, err
}
for index := range profiles {
body := fmt.Sprintf(`<?xml version="1.0"?><s:Envelope xmlns:s="http://www.w3.org/2003/05/soap-envelope"><s:Body><GetStreamUri xmlns="http://www.onvif.org/ver10/media/wsdl"><StreamSetup><Stream xmlns="http://www.onvif.org/ver10/schema">RTP-Unicast</Stream><Transport xmlns="http://www.onvif.org/ver10/schema"><Protocol>RTSP</Protocol></Transport></StreamSetup><ProfileToken>%s</ProfileToken></GetStreamUri></s:Body></s:Envelope>`, xmlEscape(profiles[index].Token))
response, requestErr := c.soap(ctx, media.String(), credential, body)
if requestErr != nil {
return nil, requestErr
}
raw, parseErr := parseElement(response, "", "Uri")
if parseErr != nil {
return nil, parseErr
}
profiles[index].StreamURI, err = c.normalizeStream(ctx, device, raw)
if err != nil {
return nil, err
}
}
return profiles, nil
}
func (c *HTTPClient) normalizeService(ctx context.Context, device *url.URL, raw string) (*url.URL, error) {
advertised, err := url.Parse(strings.TrimSpace(raw))
if err != nil || advertised.Hostname() == "" || advertised.User != nil || (advertised.Scheme != "http" && advertised.Scheme != "https") {
return nil, ErrAddressInvalid
}
if !strings.EqualFold(advertised.Hostname(), device.Hostname()) {
advertised.Scheme = device.Scheme
advertised.Host = device.Host
}
if _, _, err = c.policy.ValidateURL(ctx, advertised.String(), "http", "https"); err != nil {
return nil, err
}
return advertised, nil
}
func (c *HTTPClient) normalizeStream(ctx context.Context, device *url.URL, raw string) (string, error) {
stream, err := url.Parse(strings.TrimSpace(raw))
if err != nil || len(raw) > 2048 || stream.Hostname() == "" || stream.User != nil || stream.Scheme != "rtsp" {
return "", ErrAddressInvalid
}
if !strings.EqualFold(stream.Hostname(), device.Hostname()) {
port := stream.Port()
stream.Host = device.Hostname()
if port != "" {
stream.Host = net.JoinHostPort(device.Hostname(), port)
}
}
if _, _, err = c.policy.ValidateURL(ctx, stream.String(), "rtsp"); err != nil {
return "", err
}
return stream.String(), nil
}
func (c *HTTPClient) soap(ctx context.Context, endpoint string, credential Credential, body string) ([]byte, error) {
return c.soapAttempt(ctx, endpoint, credential, body, "")
}
func (c *HTTPClient) soapAttempt(ctx context.Context, endpoint string, credential Credential, body, authorization string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewBufferString(body))
if err != nil {
return nil, err
}
req.Header.Set("Content-Type", "application/soap+xml; charset=utf-8")
if authorization != "" {
req.Header.Set("Authorization", authorization)
} else if credential.Username != "" {
req.SetBasicAuth(credential.Username, credential.Password)
}
res, err := c.client.Do(req)
if err != nil {
return nil, fmt.Errorf("onvif request: %w", err)
}
defer res.Body.Close()
data, err := io.ReadAll(io.LimitReader(res.Body, 2<<20))
if err != nil {
return nil, err
}
if res.StatusCode >= 300 && res.StatusCode < 400 {
return nil, ErrRedirect
}
if res.StatusCode == http.StatusUnauthorized {
if authorization == "" && credential.Username != "" {
challenge, challengeErr := parseDigestChallenge(res.Header.Values("WWW-Authenticate"))
if challengeErr == nil {
digest, digestErr := digestAuthorization(http.MethodPost, req.URL.RequestURI(), credential, challenge)
if digestErr != nil {
return nil, digestErr
}
return c.soapAttempt(ctx, endpoint, credential, body, digest)
}
}
return nil, ErrAuthentication
}
if res.StatusCode < 200 || res.StatusCode >= 300 {
return nil, fmt.Errorf("onvif http status %d", res.StatusCode)
}
return data, nil
}
type digestChallenge struct{ realm, nonce, opaque, algorithm, qop string }
func parseDigestChallenge(values []string) (digestChallenge, error) {
for _, value := range values {
parts := strings.SplitN(strings.TrimSpace(value), " ", 2)
if len(parts) != 2 || !strings.EqualFold(parts[0], "Digest") {
continue
}
params, err := parseAuthParameters(parts[1])
if err != nil {
return digestChallenge{}, err
}
c := digestChallenge{realm: params["realm"], nonce: params["nonce"], opaque: params["opaque"], algorithm: strings.ToUpper(params["algorithm"])}
if c.realm == "" || c.nonce == "" {
return digestChallenge{}, ErrAuthentication
}
if c.algorithm == "" {
c.algorithm = "MD5"
}
if c.algorithm != "MD5" && c.algorithm != "SHA-256" {
return digestChallenge{}, ErrAuthentication
}
for _, q := range strings.Split(params["qop"], ",") {
if strings.EqualFold(strings.TrimSpace(q), "auth") {
c.qop = "auth"
}
}
if params["qop"] != "" && c.qop == "" {
return digestChallenge{}, ErrAuthentication
}
return c, nil
}
return digestChallenge{}, ErrAuthentication
}
func parseAuthParameters(value string) (map[string]string, error) {
result := map[string]string{}
for position := 0; position < len(value); {
for position < len(value) && (value[position] == ' ' || value[position] == ',') {
position++
}
start := position
for position < len(value) && value[position] != '=' && value[position] != ',' {
position++
}
if position == start || position >= len(value) || value[position] != '=' {
return nil, ErrAuthentication
}
name := strings.ToLower(strings.TrimSpace(value[start:position]))
position++
var parameter string
if position < len(value) && value[position] == '"' {
position++
var builder strings.Builder
closed := false
for position < len(value) {
if value[position] == '"' {
position++
closed = true
break
}
if value[position] == '\\' && position+1 < len(value) {
position++
}
builder.WriteByte(value[position])
position++
}
if !closed {
return nil, ErrAuthentication
}
parameter = builder.String()
} else {
start = position
for position < len(value) && value[position] != ',' {
position++
}
parameter = strings.TrimSpace(value[start:position])
}
result[name] = parameter
}
return result, nil
}
func digestAuthorization(method, uri string, credential Credential, c digestChallenge) (string, error) {
random := make([]byte, 16)
if _, err := rand.Read(random); err != nil {
return "", err
}
cnonce := fmt.Sprintf("%x", random)
hash := func(value string) string {
if c.algorithm == "SHA-256" {
sum := sha256.Sum256([]byte(value))
return fmt.Sprintf("%x", sum)
}
sum := md5.Sum([]byte(value))
return fmt.Sprintf("%x", sum)
}
ha1 := hash(credential.Username + ":" + c.realm + ":" + credential.Password)
ha2 := hash(method + ":" + uri)
nc := "00000001"
response := hash(ha1 + ":" + c.nonce + ":" + ha2)
if c.qop != "" {
response = hash(ha1 + ":" + c.nonce + ":" + nc + ":" + cnonce + ":" + c.qop + ":" + ha2)
}
values := []string{`username=` + strconv.Quote(credential.Username), `realm=` + strconv.Quote(c.realm), `nonce=` + strconv.Quote(c.nonce), `uri=` + strconv.Quote(uri), `response=` + strconv.Quote(response), `algorithm=` + c.algorithm}
if c.opaque != "" {
values = append(values, `opaque=`+strconv.Quote(c.opaque))
}
if c.qop != "" {
values = append(values, `qop=`+c.qop, `nc=`+nc, `cnonce=`+strconv.Quote(cnonce))
}
return "Digest " + strings.Join(values, ", "), nil
}
func xmlEscape(value string) string {
var b strings.Builder
_ = xml.EscapeText(&b, []byte(value))
return b.String()
}
@@ -0,0 +1,81 @@
package onvif
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
)
func loopbackPolicy(t *testing.T) Policy {
t.Helper()
policy, err := ParsePolicy("127.0.0.0/8")
if err != nil {
t.Fatal(err)
}
return policy
}
func TestProfilesSupportsDigestNormalizesAdvertisedHostsAndRejectsCredentials(t *testing.T) {
digestSeen := false
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body := make([]byte, r.ContentLength)
_, _ = r.Body.Read(body)
value := string(body)
switch {
case strings.Contains(value, "GetCapabilities"):
fmt.Fprint(w, `<Envelope><Body><GetCapabilitiesResponse><Capabilities><Media><XAddr>http://unusable.invalid/onvif/media</XAddr></Media></Capabilities></GetCapabilitiesResponse></Body></Envelope>`)
case !strings.HasPrefix(r.Header.Get("Authorization"), "Digest "):
w.Header().Set("WWW-Authenticate", `Digest realm="camera", nonce="n", algorithm=MD5, qop="auth"`)
w.WriteHeader(http.StatusUnauthorized)
case strings.Contains(value, "GetProfiles"):
digestSeen = strings.HasPrefix(r.Header.Get("Authorization"), "Digest ")
fmt.Fprint(w, `<Envelope><Body><GetProfilesResponse><Profiles token="main"><Name>主码流</Name><VideoEncoderConfiguration><Encoding>H264</Encoding><Resolution><Width>1920</Width><Height>1080</Height></Resolution></VideoEncoderConfiguration></Profiles></GetProfilesResponse></Body></Envelope>`)
default:
fmt.Fprint(w, `<Envelope><Body><GetStreamUriResponse><MediaUri><Uri>rtsp://unusable.invalid:8554/live</Uri></MediaUri></GetStreamUriResponse></Body></Envelope>`)
}
}))
defer server.Close()
profiles, err := NewHTTPClient(2*time.Second, loopbackPolicy(t)).Profiles(context.Background(), server.URL+"/onvif/device", Credential{Username: "synthetic", Password: "synthetic"})
if err != nil {
t.Fatal(err)
}
if !digestSeen || len(profiles) != 1 || !strings.Contains(profiles[0].StreamURI, "127.0.0.1:8554") {
t.Fatalf("profiles=%#v digest=%v", profiles, digestSeen)
}
if _, _, err = loopbackPolicy(t).ValidateURL(context.Background(), "http://user:pass@127.0.0.1/onvif", "http"); err == nil {
t.Fatal("credential URL accepted")
}
if _, _, err = loopbackPolicy(t).ValidateURL(context.Background(), "http://127.0.0.1/onvif?access_token=synthetic", "http"); err == nil {
t.Fatal("credential-like query accepted")
}
}
func TestPolicyRequiresExplicitCIDRAndRejectsOutsideTarget(t *testing.T) {
if _, err := ParsePolicy(""); err == nil {
t.Fatal("empty policy accepted")
}
policy := loopbackPolicy(t)
if _, _, err := policy.ValidateURL(context.Background(), "http://192.0.2.1/onvif", "http"); err == nil {
t.Fatal("outside target accepted")
}
}
func TestDiscoveryRequiresApprovedInterface(t *testing.T) {
_, err := Discover(context.Background(), "", time.Millisecond, loopbackPolicy(t))
if err != ErrDiscoveryNotConfigured {
t.Fatalf("error=%v", err)
}
}
func TestClientRejectsRedirect(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Location", "http://127.0.0.1/other")
w.WriteHeader(http.StatusFound)
}))
defer server.Close()
_, err := NewHTTPClient(time.Second, loopbackPolicy(t)).Profiles(context.Background(), server.URL+"/onvif", Credential{})
if err != ErrRedirect {
t.Fatalf("error=%v", err)
}
}
+96
View File
@@ -0,0 +1,96 @@
package onvif
import (
"context"
"encoding/xml"
"errors"
"io"
"net"
"sort"
"strings"
"time"
)
var (
ErrDiscoveryNotConfigured = errors.New("未配置获准的发现网卡")
ErrDiscoveryInterface = errors.New("配置的发现地址不是本机网卡")
)
const discoveryProbe = `<?xml version="1.0"?><e:Envelope xmlns:e="http://www.w3.org/2003/05/soap-envelope" xmlns:w="http://schemas.xmlsoap.org/ws/2004/08/addressing" xmlns:d="http://schemas.xmlsoap.org/ws/2005/04/discovery" xmlns:dn="http://www.onvif.org/ver10/network/wsdl"><e:Header><w:MessageID>uuid:sense-controlled-discovery</w:MessageID><w:To e:mustUnderstand="true">urn:schemas-xmlsoap-org:ws:2005:04:discovery</w:To><w:Action e:mustUnderstand="true">http://schemas.xmlsoap.org/ws/2005/04/discovery/Probe</w:Action></e:Header><e:Body><d:Probe><d:Types>dn:NetworkVideoTransmitter</d:Types></d:Probe></e:Body></e:Envelope>`
func Discover(ctx context.Context, localIP string, timeout time.Duration, policy Policy) ([]string, error) {
ip := net.ParseIP(strings.TrimSpace(localIP))
if ip == nil {
return nil, ErrDiscoveryNotConfigured
}
approved := false
interfaces, _ := net.Interfaces()
for _, iface := range interfaces {
addresses, _ := iface.Addrs()
for _, address := range addresses {
host, _, _ := net.ParseCIDR(address.String())
if host != nil && host.Equal(ip) {
approved = true
}
}
}
if !approved {
return nil, ErrDiscoveryInterface
}
connection, err := net.ListenUDP("udp4", &net.UDPAddr{IP: ip})
if err != nil {
return nil, err
}
defer connection.Close()
if timeout <= 0 {
timeout = 3 * time.Second
}
_ = connection.SetDeadline(time.Now().Add(timeout))
if _, err = connection.WriteToUDP([]byte(discoveryProbe), &net.UDPAddr{IP: net.ParseIP("239.255.255.250"), Port: 3702}); err != nil {
return nil, err
}
found := map[string]bool{}
buffer := make([]byte, 65535)
for {
n, _, readErr := connection.ReadFromUDP(buffer)
if readErr != nil {
if e, ok := readErr.(net.Error); ok && e.Timeout() {
break
}
return nil, readErr
}
for _, candidate := range extractXAddrs(string(buffer[:n])) {
if _, _, validErr := policy.ValidateURL(ctx, candidate, "http", "https"); validErr == nil {
found[candidate] = true
}
}
}
result := make([]string, 0, len(found))
for value := range found {
result = append(result, value)
}
sort.Strings(result)
return result, nil
}
func extractXAddrs(value string) []string {
decoder := xml.NewDecoder(strings.NewReader(value))
for {
token, err := decoder.Token()
if errors.Is(err, io.EOF) {
break
}
if err != nil {
return nil
}
start, ok := token.(xml.StartElement)
if !ok || start.Name.Local != "XAddrs" {
continue
}
var addresses string
if decoder.DecodeElement(&addresses, &start) == nil {
return strings.Fields(addresses)
}
}
return nil
}
+75
View File
@@ -0,0 +1,75 @@
package onvif
import (
"encoding/xml"
"errors"
"io"
"strings"
)
var ErrInvalidResponse = errors.New("ONVIF 返回内容无法识别")
type profileEnvelope struct {
Profiles []struct {
Token string `xml:"token,attr"`
Name string `xml:"Name"`
Encoder struct {
Encoding string `xml:"Encoding"`
Resolution struct {
Width int `xml:"Width"`
Height int `xml:"Height"`
} `xml:"Resolution"`
} `xml:"VideoEncoderConfiguration"`
} `xml:"Body>GetProfilesResponse>Profiles"`
}
func parseProfiles(data []byte) ([]Profile, error) {
var envelope profileEnvelope
if err := xml.Unmarshal(data, &envelope); err != nil || len(envelope.Profiles) == 0 || len(envelope.Profiles) > 128 {
return nil, ErrInvalidResponse
}
result := make([]Profile, 0, len(envelope.Profiles))
for _, p := range envelope.Profiles {
if strings.TrimSpace(p.Token) == "" || len(p.Token) > 255 || len([]rune(p.Name)) > 255 || p.Encoder.Resolution.Width <= 0 || p.Encoder.Resolution.Width > 32768 || p.Encoder.Resolution.Height <= 0 || p.Encoder.Resolution.Height > 32768 || len(p.Encoder.Encoding) > 32 {
continue
}
result = append(result, Profile{Token: p.Token, Name: p.Name, Width: p.Encoder.Resolution.Width, Height: p.Encoder.Resolution.Height, Encoding: p.Encoder.Encoding})
}
if len(result) == 0 {
return nil, ErrInvalidResponse
}
return result, nil
}
func parseElement(data []byte, parent, name string) (string, error) {
decoder := xml.NewDecoder(strings.NewReader(string(data)))
depth := 0
for {
token, err := decoder.Token()
if err != nil {
if errors.Is(err, io.EOF) {
return "", ErrInvalidResponse
}
return "", ErrInvalidResponse
}
switch value := token.(type) {
case xml.StartElement:
if parent == "" || value.Name.Local == parent {
if value.Name.Local == parent {
depth++
}
}
if (parent == "" || depth > 0) && value.Name.Local == name {
var text string
if err := decoder.DecodeElement(&text, &value); err != nil {
return "", ErrInvalidResponse
}
return strings.TrimSpace(text), nil
}
case xml.EndElement:
if value.Name.Local == parent && depth > 0 {
depth--
}
}
}
}
+76
View File
@@ -0,0 +1,76 @@
package onvif
import (
"context"
"errors"
"net"
"net/url"
"strings"
)
var (
ErrTargetNotAllowed = errors.New("目标地址不在获准网段内")
ErrAddressInvalid = errors.New("设备地址格式不正确")
)
type Policy struct{ Networks []*net.IPNet }
func ParsePolicy(value string) (Policy, error) {
var policy Policy
for _, item := range strings.Split(value, ",") {
item = strings.TrimSpace(item)
if item == "" {
continue
}
_, network, err := net.ParseCIDR(item)
if err != nil {
return Policy{}, ErrTargetNotAllowed
}
policy.Networks = append(policy.Networks, network)
}
if len(policy.Networks) == 0 {
return Policy{}, ErrTargetNotAllowed
}
return policy, nil
}
func (p Policy) ValidateURL(ctx context.Context, raw string, schemes ...string) (*url.URL, net.IP, error) {
u, err := url.Parse(strings.TrimSpace(raw))
if err != nil || u.Hostname() == "" || u.User != nil || u.Fragment != "" {
return nil, nil, ErrAddressInvalid
}
for key := range u.Query() {
lower := strings.ToLower(key)
for _, sensitive := range []string{"user", "password", "passwd", "token", "auth", "credential", "secret", "key"} {
if strings.Contains(lower, sensitive) {
return nil, nil, ErrAddressInvalid
}
}
}
allowedScheme := false
for _, scheme := range schemes {
if strings.EqualFold(u.Scheme, scheme) {
allowedScheme = true
}
}
if !allowedScheme {
return nil, nil, ErrAddressInvalid
}
addresses, err := net.DefaultResolver.LookupIP(ctx, "ip", u.Hostname())
if err != nil || len(addresses) == 0 {
return nil, nil, ErrTargetNotAllowed
}
for _, address := range addresses {
ok := false
for _, network := range p.Networks {
if network.Contains(address) {
ok = true
break
}
}
if !ok {
return nil, nil, ErrTargetNotAllowed
}
}
return u, addresses[0], nil
}
@@ -0,0 +1,17 @@
package reconcile
import "time"
func Backoff(failures int) time.Duration {
if failures <= 1 {
return time.Second
}
d := time.Second
for i := 1; i < failures; i++ {
if d >= 30*time.Second {
return time.Minute
}
d *= 2
}
return d
}
@@ -0,0 +1,12 @@
package reconcile
import (
"testing"
"time"
)
func TestBackoffIsBounded(t *testing.T) {
if Backoff(1) != time.Second || Backoff(4) != 8*time.Second || Backoff(20) != time.Minute {
t.Fatal("unexpected retry backoff")
}
}
+70
View File
@@ -0,0 +1,70 @@
package rtsp
import (
"bufio"
"context"
"encoding/base64"
"fmt"
"net"
"strings"
"time"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
)
type Credential struct{ Username, Password string }
type Result struct {
Status string `json:"status"`
LatencyMS int64 `json:"latencyMs"`
Detail string `json:"detail"`
}
type Verifier interface {
Verify(context.Context, string, Credential) (Result, error)
}
type NetVerifier struct {
Timeout time.Duration
Policy onvif.Policy
}
func (v NetVerifier) Verify(ctx context.Context, raw string, credential Credential) (Result, error) {
parsed, ip, err := v.Policy.ValidateURL(ctx, raw, "rtsp")
if err != nil {
return Result{}, err
}
port := parsed.Port()
if port == "" {
port = "554"
}
timeout := v.Timeout
if timeout <= 0 {
timeout = 5 * time.Second
}
started := time.Now()
connection, err := (&net.Dialer{Timeout: timeout}).DialContext(ctx, "tcp", net.JoinHostPort(ip.String(), port))
if err != nil {
return Result{Status: "unreachable", Detail: "无法连接视频端口"}, nil
}
defer connection.Close()
_ = connection.SetDeadline(time.Now().Add(timeout))
authorization := ""
if credential.Username != "" {
authorization = "Authorization: Basic " + base64.StdEncoding.EncodeToString([]byte(credential.Username+":"+credential.Password)) + "\r\n"
}
request := fmt.Sprintf("OPTIONS %s RTSP/1.0\r\nCSeq: 1\r\nUser-Agent: YoVision-Sense\r\n%s\r\n", parsed.String(), authorization)
if _, err = connection.Write([]byte(request)); err != nil {
return Result{}, err
}
line, err := bufio.NewReader(connection).ReadString('\n')
if err != nil {
return Result{Status: "timeout", Detail: "等待视频响应超时"}, nil
}
result := Result{Status: "ready", Detail: "码流可访问", LatencyMS: time.Since(started).Milliseconds()}
if strings.Contains(line, " 401 ") {
result.Status = "authentication_failed"
result.Detail = "设备拒绝了当前 RTSP 凭据"
} else if !strings.Contains(line, " 200 ") {
result.Status = "failed"
result.Detail = "设备返回非成功状态"
}
return result, nil
}
@@ -0,0 +1,55 @@
package rtsp
import (
"bufio"
"context"
"net"
"strings"
"testing"
"time"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
)
func TestVerifierUsesCredentialHeaderWithoutPuttingItInURI(t *testing.T) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer listener.Close()
requestChannel := make(chan string, 1)
go func() {
connection, _ := listener.Accept()
if connection == nil {
return
}
defer connection.Close()
reader := bufio.NewReader(connection)
request, _ := reader.ReadString('\n')
headers := request
for {
line, _ := reader.ReadString('\n')
headers += line
if line == "\r\n" || line == "" {
break
}
}
requestChannel <- headers
_, _ = connection.Write([]byte("RTSP/1.0 200 OK\r\nCSeq: 1\r\n\r\n"))
}()
policy, _ := onvif.ParsePolicy("127.0.0.0/8")
result, err := (NetVerifier{Timeout: time.Second, Policy: policy}).Verify(context.Background(), "rtsp://"+listener.Addr().String()+"/live", Credential{Username: "synthetic", Password: "synthetic"})
if err != nil || result.Status != "ready" {
t.Fatalf("result=%#v err=%v", result, err)
}
request := <-requestChannel
if strings.Contains(strings.Split(request, "\r\n")[0], "synthetic") {
t.Fatal("credentials leaked into request URI")
}
if !strings.Contains(request, "Authorization: Basic ") {
t.Fatal("authorization header missing")
}
if _, err = (NetVerifier{Policy: policy}).Verify(context.Background(), "rtsp://user:pass@"+listener.Addr().String()+"/live", Credential{}); err == nil {
t.Fatal("credential-bearing URI accepted")
}
}
+18
View File
@@ -20,6 +20,7 @@ import (
"git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/admin/router"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/common/database"
"git.ilapage.cn/ila/yovision/Sense/server/common/global"
common "git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
@@ -85,6 +86,19 @@ func run() error {
for _, f := range AppRouters {
f()
}
runtimeCtx, runtimeCancel := context.WithCancel(context.Background())
defer runtimeCancel()
var runtimeDBFound bool
for _, db := range sdk.Runtime.GetDb() {
runtimeDBFound = true
if err := media.StartRuntime(runtimeCtx, db); err != nil {
log.Errorf("MediaMTX runtime unavailable: %v", err)
}
break
}
if !runtimeDBFound {
log.Error("MediaMTX runtime unavailable: Sense database is not initialized")
}
srv := &http.Server{
Addr: fmt.Sprintf("%s:%d", config.ApplicationConfig.Host, config.ApplicationConfig.Port),
@@ -139,6 +153,10 @@ func run() error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
log.Info("Shutdown Server ... ")
runtimeCancel()
if err := media.ShutdownRuntime(ctx); err != nil && !errors.Is(err, context.Canceled) {
log.Errorf("Shutdown MediaMTX runtime: %v", err)
}
if err := srv.Shutdown(ctx); err != nil {
log.Fatal("Server Shutdown:", err)
@@ -0,0 +1,60 @@
package version
import (
"runtime"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
"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), migrateSenseAdmission)
}
func migrateSenseAdmission(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&admission.Result{}, &admission.Profile{}); err != nil {
return err
}
page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseAdmission", Title: "视频接入", Icon: "video-camera", Path: "/sense/admission", MenuType: "C", Permission: "sense:admission:list", Component: "/sense/admission/index", Sort: 6, Visible: "0", IsFrame: "1"})
if err != nil {
return err
}
discover, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseAdmissionDiscover", Title: "发现设备", MenuType: "F", Action: "GET", Permission: "sense:admission:discover", ParentId: page.MenuId, Sort: 1, Visible: "1", IsFrame: "1"})
if err != nil {
return err
}
probe, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseAdmissionProbe", Title: "验证接入", MenuType: "F", Action: "POST", Permission: "sense:admission:probe", ParentId: page.MenuId, Sort: 2, Visible: "1", IsFrame: "1"})
if err != nil {
return err
}
for _, role := range []string{"implementation_operator", "site_admin"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, discover, probe}); err != nil {
return err
}
}
if err = attachDeviceRole(tx, "viewer", []migrationModels.SysMenu{page}); err != nil {
return err
}
read := [][2]string{{"/api/v1/admission/devices/:id", "GET"}}
write := [][2]string{{"/api/v1/admission/discover", "GET"}, {"/api/v1/admission/devices/:id/probe", "POST"}}
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
policies := append([][2]string{}, read...)
if role != "viewer" {
policies = append(policies, write...)
}
for _, policy := range policies {
rule := deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&rule).Error; err != nil {
return err
}
}
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,60 @@
package version
import (
"runtime"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
"git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSenseMedia)
}
func migrateSenseMedia(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(&media.Route{}); err != nil {
return err
}
page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMedia", Title: "视频服务", Icon: "video-play", Path: "/sense/media", MenuType: "C", Permission: "sense:media:list", Component: "/sense/media/index", Sort: 7, Visible: "0", IsFrame: "1"})
if err != nil {
return err
}
reconcile, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaReconcile", Title: "立即对账", MenuType: "F", Action: "POST", Permission: "sense:media:reconcile", ParentId: page.MenuId, Sort: 1, Visible: "1", IsFrame: "1"})
if err != nil {
return err
}
stop, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseMediaStop", Title: "停止路径", MenuType: "F", Action: "POST", Permission: "sense:media:stop", ParentId: page.MenuId, Sort: 2, Visible: "1", IsFrame: "1"})
if err != nil {
return err
}
for _, role := range []string{"implementation_operator", "site_admin"} {
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, reconcile, stop}); err != nil {
return err
}
}
if err = attachDeviceRole(tx, "viewer", []migrationModels.SysMenu{page}); err != nil {
return err
}
read := [][2]string{{"/api/v1/media/routes", "GET"}, {"/api/v1/media/process", "GET"}}
write := [][2]string{{"/api/v1/media/reconcile", "POST"}, {"/api/v1/media/routes/:id/reconcile", "POST"}, {"/api/v1/media/routes/:id/stop", "POST"}}
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
policies := append([][2]string{}, read...)
if role != "viewer" {
policies = append(policies, write...)
}
for _, policy := range policies {
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}).Error; err != nil {
return err
}
}
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,54 @@
package version
import (
"os"
"testing"
"gorm.io/driver/postgres"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
)
func TestMediaMigrationOnPostgres(t *testing.T) {
dsn := os.Getenv("SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL to run the PostgreSQL migration test")
}
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
db.Exec("DROP TABLE IF EXISTS sense_media_routes, sys_role_menu, sys_menu, sys_role, casbin_rule, sys_migration CASCADE")
})
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil {
t.Fatal(err)
}
}
if err = migrateSenseMedia(db, "2026081419000_media.go"); err != nil {
t.Fatal(err)
}
var routes, menus, policies, applied int64
if err = db.Model(&media.Route{}).Count(&routes).Error; err != nil {
t.Fatal(err)
}
if err = db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseMedia%").Count(&menus).Error; err != nil {
t.Fatal(err)
}
if err = db.Model(&deviceCasbinRule{}).Where("v1 LIKE ?", "/api/v1/media%").Count(&policies).Error; err != nil {
t.Fatal(err)
}
if err = db.Model(&common.Migration{}).Where("version = ?", "2026081419000_media.go").Count(&applied).Error; err != nil {
t.Fatal(err)
}
if routes != 0 || menus != 3 || policies != 12 || applied != 1 {
t.Fatalf("routes=%d menus=%d policies=%d applied=%d", routes, menus, policies, applied)
}
}
@@ -1,3 +1,11 @@
# Required only when creating or updating camera credentials.
# Set outside the repository to a Base64-encoded random 32-byte key.
SENSE_CREDENTIAL_KEY=
# Required before ONVIF discovery or manual probing. Use only explicitly
# approved local interface/IP ranges; comma-separate multiple CIDRs.
SENSE_ONVIF_DISCOVERY_IP=
SENSE_ONVIF_ALLOWED_CIDRS=
SENSE_MEDIAMTX_BINARY=
SENSE_MEDIAMTX_CONFIG=
SENSE_MEDIAMTX_API=http://127.0.0.1:9997
@@ -0,0 +1,7 @@
# Sense generates only this credential-free base configuration.
# The Control API must stay on loopback; camera paths are applied at runtime.
logLevel: info
api: true
apiAddress: 127.0.0.1:9997
metrics: false
paths: {}
@@ -0,0 +1,85 @@
package admission_test
import (
"context"
"crypto/rand"
"encoding/base64"
"os"
"testing"
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
"gorm.io/driver/postgres"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/onvif"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/rtsp"
)
type onvifFixture struct{}
func (onvifFixture) Profiles(context.Context, string, onvif.Credential) ([]onvif.Profile, error) {
return []onvif.Profile{{Token: "main", Name: "主码流", Width: 1920, Height: 1080, Encoding: "H264", StreamURI: "rtsp://192.0.2.10/main"}, {Token: "sub", Name: "子码流", Width: 640, Height: 360, Encoding: "H264", StreamURI: "rtsp://192.0.2.10/sub"}}, nil
}
type rtspFixture struct{}
func (rtspFixture) Verify(context.Context, string, rtsp.Credential) (rtsp.Result, error) {
return rtsp.Result{Status: "ready", Detail: "合成 RTSP 可用"}, nil
}
func TestProfilesSurvivePostgreSQLReopen(t *testing.T) {
dsn := os.Getenv("SENSE_ADMISSION_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_ADMISSION_TEST_DATABASE_URL to an isolated PostgreSQL database")
}
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&deviceModels.Device{}, &credential.DeviceCredential{}, &admission.Result{}, &admission.Profile{}); err != nil {
t.Fatal(err)
}
deviceID := "issue66-postgres-device"
db.Where("device_id = ?", deviceID).Delete(&admission.Profile{})
db.Where("device_id = ?", deviceID).Delete(&admission.Result{})
db.Where("device_id = ?", deviceID).Delete(&credential.DeviceCredential{})
db.Where("id = ?", deviceID).Delete(&deviceModels.Device{})
t.Cleanup(func() {
db.Where("device_id = ?", deviceID).Delete(&admission.Profile{})
db.Where("device_id = ?", deviceID).Delete(&admission.Result{})
db.Where("device_id = ?", deviceID).Delete(&credential.DeviceCredential{})
db.Where("id = ?", deviceID).Delete(&deviceModels.Device{})
})
key := make([]byte, 32)
if _, err = rand.Read(key); err != nil {
t.Fatal(err)
}
t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
vault, _ := credential.NewVault(key)
device := deviceModels.Device{ID: deviceID, Name: "重启持久化摄像机", Modality: deviceModels.ModalityVideo, Status: deviceModels.StatusPending, AdapterStatus: deviceModels.AdapterReady, Version: 1}
if err = db.Create(&device).Error; err != nil {
t.Fatal(err)
}
for _, purpose := range []string{credential.PurposeONVIF, credential.PurposeRTSP} {
cipher, _ := vault.Encrypt(deviceID, purpose, "synthetic-user", "synthetic-password")
if err = db.Create(&credential.DeviceCredential{DeviceID: deviceID, Purpose: purpose, Ciphertext: cipher, KeyVersion: credential.Version()}).Error; err != nil {
t.Fatal(err)
}
}
service := admission.Service{Service: coreService.Service{Orm: db}, ONVIF: onvifFixture{}, RTSP: rtspFixture{}}
if _, err = service.Probe(context.Background(), admission.ProbeRequest{DeviceID: deviceID, Address: "http://192.0.2.10/onvif", Version: 1}); err != nil {
t.Fatal(err)
}
reopened, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
restarted := admission.Service{Service: coreService.Service{Orm: reopened}}
saved, err := restarted.Get(deviceID)
if err != nil || saved.Status != "ready" || len(saved.Profiles) != 2 {
t.Fatalf("saved=%#v err=%v", saved, err)
}
}
@@ -0,0 +1,73 @@
package media_test
import (
"context"
"fmt"
"net"
"os"
"path/filepath"
"testing"
"time"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
)
func freeAddress(t *testing.T) string {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
address := listener.Addr().String()
if err = listener.Close(); err != nil {
t.Fatal(err)
}
return address
}
func TestRealMediaMTXControlLifecycle(t *testing.T) {
binary := os.Getenv("SENSE_MEDIAMTX_TEST_BINARY")
if binary == "" {
t.Skip("set SENSE_MEDIAMTX_TEST_BINARY to run the real MediaMTX integration")
}
apiAddress, rtspAddress := freeAddress(t), freeAddress(t)
configPath := filepath.Join(t.TempDir(), "mediamtx.yml")
config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress)
if err := os.WriteFile(configPath, []byte(config), 0o600); err != nil {
t.Fatal(err)
}
controller, err := media.NewHTTPController("http://" + apiAddress)
if err != nil {
t.Fatal(err)
}
supervisor := media.NewSupervisor(binary, configPath)
if err = supervisor.Start(context.Background()); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = supervisor.Stop(ctx)
})
deadline := time.Now().Add(8 * time.Second)
for {
err = controller.Health(context.Background())
if err == nil {
break
}
if time.Now().After(deadline) {
t.Fatalf("MediaMTX did not become ready: %v", err)
}
time.Sleep(100 * time.Millisecond)
}
source := media.Source{Path: "sense_integration", URI: "rtsp://127.0.0.1:65530/test", Username: "synthetic-user", Password: "synthetic-password"}
if err = controller.Apply(context.Background(), source); err != nil {
t.Fatal(err)
}
if err = controller.Apply(context.Background(), source); err != nil {
t.Fatalf("replace must be idempotent: %v", err)
}
if err = controller.Delete(context.Background(), source.Path); err != nil {
t.Fatal(err)
}
}
+70
View File
@@ -0,0 +1,70 @@
package media_test
import (
"context"
"crypto/rand"
"encoding/base64"
"os"
"testing"
"gorm.io/driver/postgres"
"gorm.io/gorm"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
)
type readyController struct{}
func (readyController) Health(context.Context) error { return nil }
func (readyController) Apply(context.Context, media.Source) error { return nil }
func (readyController) Delete(context.Context, string) error { return nil }
func (readyController) Status(context.Context, string) (media.PathStatus, error) {
return media.PathStatus{Exists: true, Ready: true}, nil
}
func TestPostgresColdStartRestoresDesiredRoute(t *testing.T) {
dsn := os.Getenv("SENSE_MEDIA_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set SENSE_MEDIA_TEST_DATABASE_URL to run PostgreSQL media recovery")
}
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&media.Route{}, &admission.Profile{}, &credential.DeviceCredential{}); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
db.Exec("DROP TABLE IF EXISTS sense_media_routes, sense_admission_profiles, sense_device_credentials")
})
key := make([]byte, 32)
if _, err = rand.Read(key); err != nil {
t.Fatal(err)
}
t.Setenv(credential.EnvironmentKey, base64.StdEncoding.EncodeToString(key))
vault, _ := credential.NewVault(key)
ciphertext, err := vault.Encrypt("device-pg", credential.PurposeRTSP, "synthetic-user", "synthetic-password")
if err != nil {
t.Fatal(err)
}
if err = db.Create(&credential.DeviceCredential{DeviceID: "device-pg", Purpose: credential.PurposeRTSP, Ciphertext: ciphertext, KeyVersion: credential.Version()}).Error; err != nil {
t.Fatal(err)
}
if err = db.Create(&admission.Profile{DeviceID: "device-pg", Token: "main", Name: "Main", StreamURI: "rtsp://192.0.2.1/live", VerificationStatus: "ready", VerificationDetail: "synthetic"}).Error; err != nil {
t.Fatal(err)
}
first := media.NewService(db, readyController{}, nil, media.RuntimeConfig{})
if err = first.EnsureDevice(context.Background(), "device-pg"); err != nil {
t.Fatal(err)
}
second := media.NewService(db.Session(&gorm.Session{NewDB: true}), readyController{}, nil, media.RuntimeConfig{})
if err = second.ReconcileDue(context.Background()); err != nil {
t.Fatal(err)
}
items, err := second.List(context.Background())
if err != nil || len(items) != 1 || items[0].Actual != "ready" {
t.Fatalf("items=%#v err=%v", items, err)
}
}
+5
View File
@@ -0,0 +1,5 @@
import request from '@/utils/request'
export function discoverDevices() { return request({ url: '/api/v1/admission/discover', method: 'get' }) }
export function getAdmission(deviceId) { return request({ url: `/api/v1/admission/devices/${deviceId}`, method: 'get' }) }
export function probeDevice(deviceId, data) { return request({ url: `/api/v1/admission/devices/${deviceId}/probe`, method: 'post', data }) }
+7
View File
@@ -0,0 +1,7 @@
import request from '@/utils/request'
export function listMediaRoutes() { return request({ url: '/api/v1/media/routes', method: 'get' }) }
export function getMediaProcess() { return request({ url: '/api/v1/media/process', method: 'get' }) }
export function reconcileAllMedia() { return request({ url: '/api/v1/media/reconcile', method: 'post' }) }
export function reconcileMediaRoute(id) { return request({ url: `/api/v1/media/routes/${encodeURIComponent(id)}/reconcile`, method: 'post' }) }
export function stopMediaRoute(id) { return request({ url: `/api/v1/media/routes/${encodeURIComponent(id)}/stop`, method: 'post' }) }
@@ -0,0 +1,7 @@
export function buildProbePayload(form) {
return { address: String(form.address || '').trim(), version: Number(form.version) }
}
export function addressHasCredentials(value) {
try { return Boolean(new URL(value).username || new URL(value).password) } catch (_) { return false }
}
@@ -0,0 +1,80 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div><h3>视频接入</h3><p>选择已登记的视频设备,发现或填写 ONVIF 地址,然后验证主、子码流。</p></div>
<el-button v-permisaction="['sense:admission:discover']" type="primary" plain :loading="discovering" @click="handleDiscover">发现设备</el-button>
</div>
<el-alert title="发现只使用服务端配置的获准网卡;手工地址也只能访问获准网段。" type="info" :closable="false" show-icon />
<el-form ref="probeFormRef" :model="form" :rules="rules" label-width="120px" class="probe-form">
<el-form-item label="设备" prop="deviceId">
<el-select v-model="form.deviceId" filterable placeholder="请选择已登记的视频设备" @change="selectDevice">
<el-option v-for="device in devices" :key="device.id" :label="`${device.name} · ${device.location || '未填写位置'}`" :value="device.id" :disabled="device.status === 'disabled' || !device.onvifCredentialConfigured" />
</el-select>
<span class="field-hint">未配置凭据或已停用的设备不可探测。</span>
</el-form-item>
<el-form-item label="ONVIF 地址" prop="address">
<el-input v-model="form.address" placeholder="例如:http://设备地址/onvif/device_service" />
<span class="field-hint">地址中不能包含用户名或密码;凭据从设备管理安全读取。</span>
</el-form-item>
<el-form-item>
<el-button v-permisaction="['sense:admission:probe']" type="primary" :loading="probing" @click="handleProbe">验证接入</el-button>
<el-button @click="loadSaved">查看上次结果</el-button>
</el-form-item>
</el-form>
<el-divider />
<el-empty v-if="!result" description="尚无接入验证结果" />
<template v-else>
<el-descriptions :column="2" border>
<el-descriptions-item label="设备状态"><el-tag :type="result.status === 'ready' ? 'success' : 'warning'">{{ statusLabel(result.status) }}</el-tag></el-descriptions-item>
<el-descriptions-item label="检查时间">{{ parseTime(result.checkedAt) }}</el-descriptions-item>
<el-descriptions-item label="结果说明" :span="2">{{ result.detail }}</el-descriptions-item>
</el-descriptions>
<el-table :data="result.profiles || []" border stripe class="profile-table">
<el-table-column prop="name" label="Profile" min-width="130" />
<el-table-column label="用途" width="90"><template #default="scope"><el-tag size="small">{{ scope.row.kind === 'main' ? '主码流' : scope.row.kind === 'sub' ? '子码流' : '其他' }}</el-tag></template></el-table-column>
<el-table-column label="分辨率" width="110"><template #default="scope">{{ scope.row.width }} × {{ scope.row.height }}</template></el-table-column>
<el-table-column prop="encoding" label="编码" width="90" />
<el-table-column label="验证状态" width="130"><template #default="scope"><el-tag :type="scope.row.verificationStatus === 'ready' ? 'success' : 'danger'">{{ statusLabel(scope.row.verificationStatus) }}</el-tag></template></el-table-column>
<el-table-column prop="verificationDetail" label="说明" min-width="180" />
</el-table>
</template>
</el-card>
<el-dialog v-model="discoveryOpen" title="发现结果" width="720px">
<el-empty v-if="!discovered.length" description="获准网卡内未发现设备" />
<el-table v-else :data="discovered.map(address => ({ address }))" border>
<el-table-column prop="address" label="ONVIF 地址" show-overflow-tooltip />
<el-table-column label="操作" width="100"><template #default="scope"><el-button type="primary" link @click="useAddress(scope.row.address)">使用</el-button></template></el-table-column>
</el-table>
</el-dialog>
</template>
</BasicLayout>
</template>
<script setup>
import { onMounted, reactive, ref } from 'vue'
import { ElMessage } from 'element-plus'
import { listDevices } from '@/api/sense/device'
import { discoverDevices, getAdmission, probeDevice } from '@/api/sense/admission'
import { addressHasCredentials, buildProbePayload } from './admissionPayload'
const devices = ref([]); const result = ref(null); const discovered = ref([])
const discovering = ref(false); const probing = ref(false); const discoveryOpen = ref(false); const probeFormRef = ref()
const form = reactive({ deviceId: '', address: '', version: 0 })
const rules = { deviceId: [{ required: true, message: '请选择设备', trigger: 'change' }], address: [{ required: true, message: '请输入 ONVIF 地址', trigger: 'blur' }, { validator: (_r, value, done) => addressHasCredentials(value) ? done(new Error('地址中不能包含用户名或密码')) : done(), trigger: 'blur' }] }
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function selectDevice(id) { const device = devices.value.find(item => item.id === id); form.version = device?.version || 0; result.value = null }
async function loadDevices() { const response = await listDevices({ pageIndex: 1, pageSize: 100, modality: 'video' }); const payload = unwrap(response); devices.value = payload?.list || payload?.data || payload || [] }
async function handleDiscover() { discovering.value = true; try { const response = await discoverDevices(); const payload = unwrap(response); discovered.value = payload?.addresses || []; discoveryOpen.value = true } catch (error) { ElMessage.error(error.message || '发现失败,请检查获准网卡配置') } finally { discovering.value = false } }
function useAddress(address) { form.address = address; discoveryOpen.value = false }
async function handleProbe() { const valid = await probeFormRef.value?.validate().catch(() => false); if (!valid) return; probing.value = true; try { const response = await probeDevice(form.deviceId, buildProbePayload(form)); result.value = unwrap(response); const device = devices.value.find(item => item.id === form.deviceId); if (device) { device.version += 1; form.version = device.version } ElMessage.success('接入验证完成') } catch (error) { ElMessage.error(error.message || '接入验证失败') } finally { probing.value = false } }
async function loadSaved() { if (!form.deviceId) return ElMessage.warning('请先选择设备'); try { result.value = unwrap(await getAdmission(form.deviceId)) } catch (error) { ElMessage.warning(error.message || '尚无验证结果') } }
function statusLabel(status) { return ({ ready: '可用', profile_failed: '部分码流失败', authentication_failed: '认证失败', target_not_allowed: '目标未获准', redirect_rejected: '重定向已拒绝', timeout: '响应超时', clock_skew: '设备时间异常', unreachable: '无法连接', failed: '验证失败' })[status] || status || '未知' }
onMounted(loadDevices)
</script>
<style scoped>
.page-header{display:flex;justify-content:space-between;align-items:flex-start;margin-bottom:16px}.page-header h3{margin:0 0 6px}.page-header p{margin:0;color:#909399}.probe-form{max-width:820px;margin-top:22px}.probe-form .el-select{width:100%}.field-hint{display:block;color:#909399;font-size:12px;line-height:20px}.profile-table{margin-top:18px}
</style>
+2 -1
View File
@@ -14,6 +14,7 @@
<el-form-item label="状态" prop="status">
<el-select v-model="queryParams.status" placeholder="全部状态" clearable size="small">
<el-option label="待接入" value="pending" />
<el-option label="已接入" value="active" />
<el-option label="已停用" value="disabled" />
</el-select>
</el-form-item>
@@ -54,7 +55,7 @@
</el-table-column>
<el-table-column label="状态" width="100" align="center">
<template #default="scope">
<el-tag :type="scope.row.status === 'disabled' ? 'info' : 'success'">{{ scope.row.status === 'disabled' ? '已停用' : '待接入' }}</el-tag>
<el-tag :type="scope.row.status === 'active' ? 'success' : scope.row.status === 'disabled' ? 'info' : 'warning'">{{ scope.row.status === 'active' ? '已接入' : scope.row.status === 'disabled' ? '已停用' : '待接入' }}</el-tag>
</template>
</el-table-column>
<el-table-column label="版本" width="80" align="center" prop="version" />
+55
View File
@@ -0,0 +1,55 @@
<template>
<BasicLayout>
<template #wrapper>
<el-card class="box-card">
<div class="page-header">
<div><h3>视频服务</h3><p>查看 MediaMTX 进程、摄像头拉流路径和自动重试状态。</p></div>
<el-button v-permisaction="['sense:media:reconcile']" type="primary" :loading="reconciling" @click="handleReconcileAll">立即对账</el-button>
</div>
<el-alert title="视频服务故障不会删除已验证的设备和码流;外部进程受孤儿安全闸保护,不会被 Sense 停止。" type="info" :closable="false" show-icon />
<el-descriptions :column="3" border class="process-state">
<el-descriptions-item label="进程状态"><el-tag :type="mediaStatusType(process.phase)">{{ processLabel(process.phase) }}</el-tag></el-descriptions-item>
<el-descriptions-item label="进程归属">{{ process.owned ? 'Sense 启动' : process.external ? '外部启动(受保护)' : '未运行' }}</el-descriptions-item>
<el-descriptions-item label="进程号">{{ process.pid || '—' }}</el-descriptions-item>
<el-descriptions-item label="说明" :span="3">{{ process.detail || '尚未取得状态' }}</el-descriptions-item>
</el-descriptions>
<el-table v-loading="loading" :data="routes" border stripe class="route-table">
<el-table-column prop="path" label="媒体路径" min-width="190" show-overflow-tooltip />
<el-table-column prop="profileToken" label="Profile" min-width="120" show-overflow-tooltip />
<el-table-column label="期望状态" width="100"><template #default="scope"><el-tag size="small" type="info">{{ scope.row.desired === 'running' ? '运行' : '停止' }}</el-tag></template></el-table-column>
<el-table-column label="实际状态" width="120"><template #default="scope"><el-tag :type="mediaStatusType(scope.row.actual)" size="small">{{ mediaStatusLabel(scope.row.actual) }}</el-tag></template></el-table-column>
<el-table-column prop="readers" label="观看数" width="90" />
<el-table-column label="下次重试" min-width="165"><template #default="scope">{{ scope.row.nextRetryAt ? parseTime(scope.row.nextRetryAt) : '—' }}</template></el-table-column>
<el-table-column prop="detail" label="说明" min-width="210" show-overflow-tooltip />
<el-table-column label="操作" width="150" fixed="right">
<template #default="scope">
<el-button v-permisaction="['sense:media:reconcile']" type="primary" link @click="handleRoute(scope.row)">对账</el-button>
<el-button v-if="scope.row.desired === 'running'" v-permisaction="['sense:media:stop']" type="danger" link @click="handleStop(scope.row)">停止</el-button>
</template>
</el-table-column>
</el-table>
<el-empty v-if="!loading && routes.length === 0" description="暂无媒体路径;请先完成摄像头接入验证" />
</el-card>
</template>
</BasicLayout>
</template>
<script setup>
import { onMounted, ref } from 'vue'
import { ElMessage, ElMessageBox } from 'element-plus'
import { getMediaProcess, listMediaRoutes, reconcileAllMedia, reconcileMediaRoute, stopMediaRoute } from '@/api/sense/media'
import { mediaStatusLabel, mediaStatusType } from './mediaStatus'
const routes = ref([]); const process = ref({}); const loading = ref(false); const reconciling = ref(false)
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
function processLabel(value) { return ({ running: '运行中', starting: '启动中', stopped: '已停止', not_configured: '未配置', configuration_failed: '配置错误', failed: '启动失败' })[value] || value || '未知' }
async function load() { loading.value = true; try { const [routeResponse, processResponse] = await Promise.all([listMediaRoutes(), getMediaProcess()]); const routePayload = unwrap(routeResponse); routes.value = routePayload?.list || []; process.value = unwrap(processResponse) || {} } catch (error) { ElMessage.error(error.message || '视频服务状态加载失败') } finally { loading.value = false } }
async function handleReconcileAll() { reconciling.value = true; try { await reconcileAllMedia(); ElMessage.success('对账完成'); await load() } catch (error) { ElMessage.warning(error.message || '对账未完成,请查看路径状态') } finally { reconciling.value = false } }
async function handleRoute(route) { try { await reconcileMediaRoute(route.id); await load() } catch (error) { ElMessage.warning(error.message || '该路径对账失败') } }
async function handleStop(route) { try { await ElMessageBox.confirm('停止后该路径将不再拉取摄像头视频,设备和 Profile 不会被删除。', '停止媒体路径', { type: 'warning' }); await stopMediaRoute(route.id); ElMessage.success('路径已停止'); await load() } catch (error) { if (error !== 'cancel' && error !== 'close') ElMessage.warning(error.message || '停止失败') } }
onMounted(load)
</script>
<style scoped>
.page-header{display:flex;justify-content:space-between;align-items:flex-start;margin-bottom:16px}.page-header h3{margin:0 0 6px}.page-header p{margin:0;color:#909399}.process-state{margin-top:18px}.route-table{margin-top:18px}
</style>
@@ -0,0 +1,10 @@
export function mediaStatusLabel(value) {
return ({ ready: '拉流正常', waiting: '等待拉流', pending: '等待对账', stopped: '已停止', process_unavailable: '进程不可用', apply_failed: '配置失败', status_unavailable: '状态未知', path_missing: '路径缺失', profile_unavailable: 'Profile 不可用', credential_unavailable: '凭据不可用' })[value] || value || '未知'
}
export function mediaStatusType(value) {
if (value === 'ready' || value === 'running') return 'success'
if (value === 'waiting' || value === 'pending' || value === 'starting') return 'warning'
if (value === 'stopped' || value === 'not_configured') return 'info'
return 'danger'
}
@@ -0,0 +1,11 @@
import { addressHasCredentials, buildProbePayload } from '@/views/sense/admission/admissionPayload'
describe('Sense admission payload', () => {
it('only sends the approved probe fields', () => {
expect(buildProbePayload({ address: ' http://192.0.2.1/onvif ', version: '3', password: 'never-send' })).toEqual({ address: 'http://192.0.2.1/onvif', version: 3 })
})
it('detects credentials embedded in a URL', () => {
expect(addressHasCredentials('http://user:secret@192.0.2.1/onvif')).toBe(true)
expect(addressHasCredentials('http://192.0.2.1/onvif')).toBe(false)
})
})
@@ -0,0 +1,10 @@
import { mediaStatusLabel, mediaStatusType } from '@/views/sense/media/mediaStatus'
describe('Sense media status presentation', () => {
it('uses actionable labels for expected lifecycle states', () => {
expect(mediaStatusLabel('waiting')).toBe('等待拉流')
expect(mediaStatusLabel('process_unavailable')).toBe('进程不可用')
expect(mediaStatusType('ready')).toBe('success')
expect(mediaStatusType('apply_failed')).toBe('danger')
})
})
+10 -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: 6880a2918c46be27738467361d8d3f923f33fc69
synchronized_at: 2026-08-14T07:02:35Z
wiki_revision: c5482806a40428a436c202760eefdb6403406584
synchronized_at: 2026-08-14T09:46:37Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -92,6 +92,14 @@ Sense JWT realm 固定为 `Sense`;浏览器令牌 Cookie 为 `Sense-Admin-Toke
工单 #65 新增设备台账入口:后端按 `models → dto → service → api → router` 分层位于 `Sense/server/app/sense/device/`,管理路由在 `Sense/server/app/admin/router/sense_device.go`,前端页面位于 `Sense/ui/src/views/sense/device/index.vue`。设备凭据由 `Sense/server/app/sense/credential/` 独立存储和 AES-256-GCM 加密,HTTP 只返回是否已配置,不提供凭据读取接口。
设备写入采用版本号乐观并发控制;视频设备适配器状态为可接入,雷达、门磁、按钮、穿戴和其他类型明确显示“适配器未就绪”,不得伪装成已接入。`admin`、`implementation_operator`、`site_admin` 可维护设备,`viewer` 只读;停用替代物理删除。
工单 #66 在 Sense/server/app/sense/onvif/、rtsp/ 与 admission/ 建立视频接入边界:WS-Discovery 只能绑定 SENSE_ONVIF_DISCOVERY_IP 指定的本机网卡,所有 ONVIF、Media XAddr 与 RTSP Stream URI 都必须落在 SENSE_ONVIF_ALLOWED_CIDRS 明确授权的网段。HTTP 客户端禁止代理和重定向,并在每次连接时重新解析、校验和固定目标 IP,防止 DNS 重绑定;URL 用户信息及敏感查询参数被拒绝。
ONVIF 支持 Basic 与 MD5/SHA-256 Digest challenge,Profile 与无凭据 Stream URI 持久化到 PostgreSQL。接入失败会记录可行动状态但保留最后一次已验证 Profile;成功接入清除凭据更新触发的重试标记。前端继续复用 GoAdmin 动态菜单、权限链、BasicLayout 和 Element Plus 表单、Dialog、Table、Tag。
工单 #67 在 `Sense/server/app/sense/media/` 与 `reconcile/` 建立 MediaMTX 管理面:`cmd/api/server.go` 随 Sense 生命周期启动后台对账并只停止本实例拥有的子进程;检测到外部实例时设置孤儿安全闸,不发送停止信号。MediaMTX Control API 只允许 loopback HTTP,禁用代理与重定向。
媒体路由只保存设备/Profile 引用、无秘密路径名、期望态、实际态、reader、退避和下次重试;摄像头凭据从内部端口按需解密,仅在 loopback Control API 请求内临时组装,不写入路由表、基础配置、日志或 Sense 响应。MediaMTX 故障和退避不改变 #66 的设备/Profile 验证状态。
<!-- sense-runtime:end -->
<!-- sense-mvp:start -->
+23 -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: 013865e0b48d156bb984515d6fabd42139829459
synchronized_at: 2026-08-14T07:02:39Z
wiki_revision: ea2c540a9ac3f8ccf40123f8389eb392fb5e172b
synchronized_at: 2026-08-14T09:46:41Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -76,6 +76,17 @@ synchronized_at: 2026-08-14T07:02:39Z
- 登录成功/失败、登出、密码变更和鉴权拒绝必须留下身份审计;密码、令牌、Cookie、验证码、数据库连接和摄像头凭据不得进入审计正文。
- `site_admin` 可在账户维护流程中读取角色、部门、岗位和字典等必要支撑数据,但不能修改角色、菜单或系统配置;`implementation_operator` 与 `viewer` 不具备账户管理权限。未注册的配置和接口管理路由对所有角色返回 404。
<!-- sense-media:start -->
## Sense 视频服务规则
- MediaMTX 始终是独立二进制;Sense 管理配置、进程生命周期、路径期望态和状态对账,不把媒体内核放入 GoAdmin handler 或 GORM model。
- Control API 只能绑定回环地址。Sense 可启动配置的 MediaMTX,也可连接已由外部启动的实例;外部实例标记为非本实例所有,孤儿安全闸禁止 Sense 停止它。
- 已验证 Profile 幂等形成媒体路径;数据库不保存带凭据 Stream URI。摄像头凭据只在 loopback Control API 请求边界临时使用,不进入基础配置、日志或 Sense API。
- 路径状态区分 pending、waiting、ready、process_unavailable、apply_failed、status_unavailable、path_missing、stopped,并保存失败次数和有上限的下次重试时间。
- 冷启动恢复 desired=running 路径;用户明确停止的路径保持 stopped,不因启动扫描自动重新启用。稳定路径只刷新状态,不重复下发配置或无意义增加版本。
- MediaMTX 失败不得删除或降级设备台账与最后一次已验证 Profile。
<!-- sense-media:end -->
<!-- sense-mvp:start -->
## Sense 旧 MVP 规则状态
@@ -105,3 +116,13 @@ synchronized_at: 2026-08-14T07:02:39Z
- 凭据只写不可读:HTTP 和页面仅显示“已配置/未配置”,不得回填用户名、密码或密文;更新凭据后只记录状态并请求后续接入流程重试。
- `admin`、`implementation_operator`、`site_admin` 可维护设备与凭据,`viewer` 仅可查看设备台账。
<!-- sense-device-ledger:end -->
<!-- sense-admission:start -->
## Sense 视频接入规则
- “获准网卡”和“获准目标网段”都是部署人员显式配置的授权边界;私网地址不自动代表已授权。未配置发现网卡时不发送 WS-Discovery,手工地址也必须通过目标 CIDR 检查。
- ONVIF 设备地址、Media XAddr 和 RTSP Stream URI 禁止 URL 用户信息、敏感认证查询参数、HTTP 重定向和超出授权网段的目标。摄像机返回不可用主机名时,只能归一化为已验证设备主机并重新执行授权检查。
- ONVIF 支持 Basic、MD5 Digest 和 SHA-256 Digest 的 auth;不支持的算法或 qop 必须拒绝,不静默降级。
- Profile 保存 token、名称、分辨率、编码、用途、无凭据 Stream URI 和逐 Profile 验证状态;主码流默认取分辨率最高项,子码流取最低项。
- 认证失败、超时、时间异常、目标未授权和重定向拒绝必须给出不同状态。失败重探不得删除最后一次已验证 Profile;凭据更新后可重新探测。
<!-- sense-admission:end -->
+29 -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: f68bbd8c7a99689e65453ae00c7bbd2c3ee00357
synchronized_at: 2026-08-14T07:02:43Z
wiki_revision: e084a4ae3f5038c4ec7dae6028a2e3d7527680a8
synchronized_at: 2026-08-14T09:46:47Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
@@ -163,6 +163,33 @@ $env:SENSE_CREDENTIAL_KEY = [Convert]::ToBase64String($keyBytes)
```
缺少或格式错误的密钥时,普通设备台账仍可读写,但凭据更新返回服务不可用且不得产生部分写入。设备回归至少覆盖:中文名称与位置、未知 JSON 字段拒绝、版本冲突返回 409、非视频设备显示适配器未就绪、viewer 只读、凭据响应/操作日志不含明文,以及 PostgreSQL 迁移重复执行不增加菜单或权限记录。
视频接入还需在仓库外配置 SENSE_ONVIF_DISCOVERY_IP(获准的本机网卡 IP)和 SENSE_ONVIF_ALLOWED_CIDRS(逗号分隔的获准摄像头网段)。不要使用 0.0.0.0/0 代替授权清单。
协议回归位于 app/sense/onvif、app/sense/rtsp、app/sense/admission;隔离 PostgreSQL 重启恢复测试通过 SENSE_ADMISSION_TEST_DATABASE_URL 显式启用。验证至少覆盖 Digest/Basic、无配置发现提示、URL 凭据和敏感查询拒绝、目标网段、重定向、Media/Stream 主机归一化、主子码流、失败重探保留已验证 Profile,以及 viewer 只读权限。
MediaMTX 保持仓库外独立二进制。运行前在进程环境设置:
```powershell
$env:SENSE_MEDIAMTX_BINARY = '<MediaMTX 可执行文件>'
$env:SENSE_MEDIAMTX_CONFIG = '<仓库外 mediamtx.yml>'
$env:SENSE_MEDIAMTX_API = 'http://127.0.0.1:9997'
```
配置文件不存在时 Sense 只生成 loopback API 和空 `paths: {}` 的无凭据基础配置;模板位于 `Sense/server/config/mediamtx/mediamtx.yml.example`。Control API 不允许非回环地址。真实集成验证使用:
```powershell
$env:SENSE_MEDIAMTX_TEST_BINARY = '<MediaMTX 可执行文件>'
go test ./tests/media -run TestRealMediaMTXControlLifecycle -v
$env:SENSE_MEDIA_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
go test ./tests/media -run TestPostgresColdStartRestoresDesiredRoute -v
$env:SENSE_MEDIA_MIGRATION_TEST_DATABASE_URL = '<隔离 PostgreSQL 连接>'
go test ./cmd/migrate/migration/version -run TestMediaMigrationOnPostgres -v
```
测试必须使用隔离端口和数据库;结束后停止测试进程。不得输出连接串或摄像头凭据。
<!-- sense-runtime:end -->
+30 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Troubleshooting
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Troubleshooting
wiki_revision: 2be5d803b954df64e31085bcf097ee27f8611f9d
synchronized_at: 2026-08-14T01:17:28Z
wiki_revision: c615cd5235557a3888260beb5a682c6ed4e5b5dc
synchronized_at: 2026-08-14T09:46:56Z
<!-- gitea-wiki-mirror:end -->
# 故障排查
@@ -60,3 +60,31 @@ synchronized_at: 2026-08-14T01:17:28Z
|---|---|
| 在 `dev` 找不到 Sense/Bell 可运行代码 | 这是重建空基线的预期状态;旧实现位于 `explore`,新代码必须由 GoAdmin 源码派生工单建立。 |
| 新骨架只有相似页面、没有 GoAdmin 启动链或权限模块 | 不符合二次开发门禁;停止验收,对照 `goadmin-baseline.json`、上游源码和 go-admin-doc 重新实施。 |
<!-- sense-admission:start -->
## Sense 视频接入排错
| 现象 | 原因与处理 |
|---|---|
| 未配置获准的发现网卡 | 在服务进程环境设置本机实际网卡 IP SENSE_ONVIF_DISCOVERY_IP;不要填写摄像机 IP。 |
| 配置的发现地址不是本机网卡 | 网卡地址已变化或填写错误;用 Get-NetIPAddress 核对后重启服务。 |
| 目标地址不在获准网段内 | 核对摄像机实际地址与 SENSE_ONVIF_ALLOWED_CIDRS;只追加已审批的最小 CIDR,不使用全网放行。 |
| 认证失败 | 在设备管理重新填写 ONVIF/RTSP 凭据,再返回视频接入重新验证;页面不会回显旧凭据。 |
| 设备时间异常 | 在摄像机管理页或受控 NTP 环境校时后重新探测;Sense 不自动修改设备时间。 |
| 部分码流失败 | 查看逐 Profile 状态、设备 RTSP 权限和端口;最后一次已验证 Profile 会保留。 |
| 重定向已拒绝 | ONVIF 服务返回了 3xx;修正为摄像机最终服务地址,不允许 Sense 跟随到未知目标。 |
<!-- sense-media:start -->
## Sense 视频服务排错
| 现象 | 原因与处理 |
|---|---|
| 进程状态“未配置” | 未设置 SENSE_MEDIAMTX_BINARY;如由外部服务管理,先确认 loopback Control API 已就绪,否则配置二进制和仓库外配置路径。 |
| 进程启动失败 | 核对二进制存在、配置目录可写、MediaMTX 配置可解析,以及 RTSP/API 端口未被其他进程占用。 |
| 等待拉流 | 路径已建立但 sourceOnDemand 尚无 reader;打开实时监看后再观察,不等同于接入失败。 |
| 配置失败或路径缺失 | 在“视频服务”点击对账;检查 Control API 仍为 loopback、Profile 仍已验证、RTSP 凭据可用。 |
| 显示外部启动(受保护) | Sense 检测到不是本实例启动的 MediaMTX;孤儿安全闸生效,Sense 关闭时不会停止它。 |
| 持续自动重试 | 查看失败码、失败次数和下次重试时间;修正二进制、端口、凭据或上游后等待退避到期,或由有权限用户立即对账。 |
| Sense 重启后路径未恢复 | 确认数据库 route 的 desired 为 running、迁移已执行、Control API 可达;明确停止的路径不会自动恢复。 |
<!-- sense-media:end -->
<!-- sense-admission:end -->
+20 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Delivery-Documentation-Guide
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Delivery-Documentation-Guide.-
wiki_revision: 93cda004538c5114c5c4ca8a03b19e3e975b503d
synchronized_at: 2026-08-14T07:03:33Z
wiki_revision: 45c86f2a0e4d3252e8042df5ee725e633dd497c1
synchronized_at: 2026-08-14T09:47:31Z
<!-- gitea-wiki-mirror:end -->
# 交付文档指南
@@ -112,3 +112,21 @@ Sense 面向网管、实施人员和非技术现场人员,菜单按日常任
凭据窗口每次均为空,不会回显已保存用户名或密码;“已配置”标签只表示服务器保存了密文。部署人员必须在仓库外为服务进程配置 Base64 编码的随机 32 字节 `SENSE_CREDENTIAL_KEY`,丢失或更换该密钥会使旧凭据不可用,因此应纳入受控秘密备份。停用设备不会物理删除台账。
<!-- sense-device-ledger:end -->
<!-- sense-admission:start -->
### Sense 视频接入交付说明
部署人员必须先确认获准摄像头网段,再把本机对应网卡 IP 配置为 SENSE_ONVIF_DISCOVERY_IP,把获准网段配置为逗号分隔的 SENSE_ONVIF_ALLOWED_CIDRS。不得为了省事填写全网段。现场人员在“视频接入”选择已登记且已配置凭据的视频设备,可使用发现结果或手工填写不含账号密码的 ONVIF 地址。
验证结果区分可用、部分码流失败、认证失败、目标未获准、重定向拒绝、响应超时、设备时间异常和无法连接,并显示主/子码流及逐 Profile 状态。失败重试不会删除上次已验证 Profile;修改凭据后应重新验证。真实摄像机兼容性、网络 ACL 和设备校时仍需在客户授权环境完成。
<!-- sense-media:start -->
### Sense 视频服务交付说明
“视频服务”面向实施、运维和站点管理员展示 MediaMTX 进程归属、媒体路径、拉流状态、观看数、失败原因和下次重试。只读用户只能查看;有操作权限的人员可立即对账或停止单一路径。
交付时 MediaMTX 二进制和真实配置放在仓库外受控目录,Control API 只能绑定回环地址。Sense 关闭时只停止自己启动的 MediaMTX;外部启动进程会显示“外部启动(受保护)”。“等待拉流”表示按需路径尚无观看者,不等同于故障。停止路径不会删除设备或 Profile,后续重新接入可恢复期望态。
任何客户文档、截图和日志都不得包含 Control API 请求体、摄像头凭据或带凭据 URI。现场至少验证启动失败、端口冲突、外部进程保护、重复对账、冷启动恢复和 reader 状态。
<!-- sense-media:end -->
<!-- sense-admission:end -->
@@ -0,0 +1,67 @@
<!-- gitea-wiki-mirror:start -->
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Task-66-Sense视频接入与Profile
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Task-66-Sense%E8%A7%86%E9%A2%91%E6%8E%A5%E5%85%A5%E4%B8%8EProfile.-
wiki_revision: 322feb232fc03c3a9ba22f65504cf3e151fb0e57
synchronized_at: 2026-08-14T09:15:26Z
<!-- gitea-wiki-mirror:end -->
# 66 Sense视频接入与Profile
- 类型:需求
- 所属 Epic:#7
- 所属 MVP / 版本:#8
- 状态:已完成
- 日期:2026-08-14
- 验收:用户于 2026-08-14 明确验收通过
- 工单:https://git.ilapage.cn/ila/yovision/issues/66
- Pull Request:https://git.ilapage.cn/ila/yovision/pulls/84
- 主项目:Sense
## 背景与目标
在 #65 设备与凭据边界上重建获准网卡 ONVIF 发现、手工接入、认证、RTSP 验证和 Profile 持久化。Brain、Bell 不启动时可独立验收。
## 最终方案
- WS-Discovery 只绑定 SENSE_ONVIF_DISCOVERY_IP 指定的本机地址;无配置或非本机地址给出明确提示。
- ONVIF、Media XAddr 和 RTSP Stream URI 均受 SENSE_ONVIF_ALLOWED_CIDRS 限制;每次连接重新解析并固定目标 IP,禁用代理和重定向,拒绝 URL 用户信息及敏感查询参数。
- ONVIF 支持 Basic、MD5/SHA-256 Digest auth;不支持的算法或 qop 拒绝。
- 内部读取 #65 的 ONVIF/RTSP 分用途密文,HTTP 不提供凭据读取。
- Profile、主/子码流、无凭据 Stream URI 和逐项验证状态持久化到 PostgreSQL;失败重探保留最后一次已验证 Profile。
- 复用 GoAdmin JWT/Casbin/操作审计、迁移和 go-admin-ui BasicLayout、Element Plus Form/Dialog/Table/Tag、动态菜单与权限按钮。
- implementation_operator、site_admin 可发现和探测,viewer 只读保存结果。
## 验收结果
| 标准 | 结果 |
|---|---|
| 未配置获准网卡不扫描 | 单元测试及错误映射通过 |
| Digest/Basic 与 RTSP 验证 | 合成协议服务通过 |
| ONVIF/RTSP 分离凭据且不回显 | 内部端口与严格 DTO 通过 |
| SSRF/重定向/凭据泄漏防护 | CIDR、DNS、URL、重定向测试通过 |
| Profile 重启恢复 | PostgreSQL 17 连接重开测试通过 |
| 认证失败、时间异常、重探 | 状态分类和失败保留 Profile 测试通过 |
## 测试
- go test ./...:通过。
- go vet ./...:通过。
- go build ./...:通过。
- 前端 lint:0 error,31 条上游继承 warning。
- 前端单测:9 suites、34 tests 通过。
- 前端生产构建:通过,6 条继承 warning。
- PostgreSQL 17:迁移与重复迁移通过;migration=1、tables=2、menus=3、policies=7。
- SENSE_ADMISSION_TEST_DATABASE_URL 隔离测试:写入 Profile、重开连接、读取主子码流通过。
- Wiki 镜像检查:通过。
- 未验证部分:未连接客户或实验室真实摄像机;真实厂商 Digest/RTSP 兼容、网络 ACL、设备校时和目标浏览器留待授权现场验收,不声称已通过。
## 回退
回退 PR #84 的实现和文档提交可移除入口;已产生的接入表保留以避免破坏追溯数据。失败探测本身不删除最后一次已验证 Profile。
## 相关提交
- 2bb1614 feat: 重建 ONVIF 与 RTSP 视频接入 (#66)
- 26e2e63 docs: 记录 Sense 视频接入安全边界 (#66)
- d74e1ae docs: 归档工单 #66 待验收证据
@@ -0,0 +1,72 @@
<!-- gitea-wiki-mirror:start -->
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Task-67-Sense视频服务生命周期与状态对账
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Task-67-Sense%E8%A7%86%E9%A2%91%E6%9C%8D%E5%8A%A1%E7%94%9F%E5%91%BD%E5%91%A8%E6%9C%9F%E4%B8%8E%E7%8A%B6%E6%80%81%E5%AF%B9%E8%B4%A6.-
wiki_revision: 4642924705d7a3406874ef532579fb2d6a25a88b
synchronized_at: 2026-08-14T10:13:23Z
<!-- gitea-wiki-mirror:end -->
# 67 Sense视频服务生命周期与状态对账
- 类型:需求
- 所属 Epic:#7
- 所属 MVP / 版本:#8
- 状态:已完成
- 日期:2026-08-14
- 工单:https://git.ilapage.cn/ila/yovision/issues/67
- Pull Request:https://git.ilapage.cn/ila/yovision/pulls/85
- 主项目:Sense
## 背景与目标
在 #66 已验证设备/Profile 上重建 MediaMTX 安全配置、独立进程生命周期、按需路径和持久化状态对账。Brain、Bell 不启动时可独立运行和验收。
## 最终方案
- MediaMTX 保持独立二进制;Sense 在 GoAdmin server 生命周期内启动后台对账,只停止自己启动的子进程。
- Control API 只允许 loopback HTTP,禁用代理和重定向;检测到外部实例时标记“外部启动(受保护)”,孤儿安全闸禁止停止。
- 路由只持久化 Device/Profile 引用、哈希路径、desired/actual、reader、失败次数、退避和下次重试,不保存带凭据 URI。
- 摄像头凭据按需从 #65 内部端口读取,仅在 MediaMTX v1.19.3 loopback 请求边界临时组装,不进入基础配置、日志或 Sense API。
- 接入成功后幂等建立路由;冷启动恢复 desired=running,明确 stopped 路径保持停止;稳定路径只刷新状态,不重复下发或增加版本。
- 复用 GoAdmin JWT/Casbin/操作审计、迁移、动态菜单和 go-admin-ui BasicLayout、Element Plus Descriptions/Table/Tag/Button/MessageBox。
## 验收结果
| 标准 | 结果 |
|---|---|
| 已验证 Profile 幂等建立/更新路径 | 单元测试及接入内部端口通过 |
| 进程未启动、启动中、失败和恢复可定位 | 状态机、真实二进制和页面反馈通过 |
| 冷启动恢复 running 路径 | PostgreSQL 17 重建 Service 测试通过 |
| reader、上游、退避、下次重试、孤儿闸可观察 | DTO、页面、退避和外部进程测试通过 |
| 配置、日志和 Sense API 不泄漏凭据 | loopback 边界、无秘密模型/响应及扫描通过 |
| MediaMTX 失败不破坏 Profile | 失败持久化测试确认 Profile 保持 ready |
## 测试
- `go test ./...`、`go vet ./...`、`go build ./...`:通过。
- `go test -race ./app/sense/media ./app/sense/reconcile`:通过。
- 前端 lint:0 error,32 条上游/目录命名 warning。
- 前端单测:10 suites、35 tests 通过。
- 前端生产构建:通过,6 条上游继承 warning。
- 真实 MediaMTX v1.19.3:Control API 就绪、path add/patch/delete 通过。
- PostgreSQL 17:desired route 保存后重建 Service 并恢复为 ready 通过。
- PostgreSQL 17 定向迁移:3 个菜单、12 条角色策略、1 条迁移记录通过。
- Wiki 镜像检查:通过。
- 已处理测试安全问题:MediaMTX 自动 TLS 文件的工作目录已固定到外部配置目录;临时证书/私钥未进入最终提交或远端。
- 未验证部分:未连接客户真实摄像机和现场网络;真实上游持续拉流、reader 变化、端口 ACL 和目标浏览器留待授权现场验收。
- 已知相邻问题:空白 PostgreSQL 执行完整上游迁移链时,在到达 #67 前被旧 `sys_config` 初始化字段长度问题中止;#67 定向迁移已通过,空库安装链应由 #70/#71 单独复核,不在本工单混改上游初始化。
## 回退
回退 PR #85 可移除运行时、路由和页面入口;已验证 Device/Profile 不受影响。媒体路由表保留可审计状态,停止/回退不删除摄像机数据。
## 验收
- 用户于 2026-08-14 明确验收通过。
## 相关提交
- 19f9bfa feat: 重建 MediaMTX 生命周期与状态对账 (#67)
- 4149cc4 docs: 记录 Sense 视频服务生命周期 (#67)
- 263a68c docs: 归档工单 #67 待验收证据
+8
View File
@@ -111,6 +111,14 @@
{
"page": "Task-65-Sense设备台账与凭据边界",
"path": "docs/task/65-Sense设备台账与凭据边界.md"
},
{
"page": "Task-66-Sense视频接入与Profile",
"path": "docs/task/66-Sense视频接入与Profile.md"
},
{
"page": "Task-67-Sense视频服务生命周期与状态对账",
"path": "docs/task/67-Sense视频服务生命周期与状态对账.md"
}
]
}