[SEN] 修复实时监看按需拉流循环等待 (#54) #55

Open
ila wants to merge 34 commits from agent/codex/54-sense-liveview-on-demand into agent/codex/51-sense-auto-liveview
20 changed files with 514 additions and 24 deletions
@@ -20,6 +20,7 @@ type Source struct {
}
type PathStatus struct {
Name string `json:"name"`
Exists bool `json:"exists"`
Ready bool `json:"ready"`
Readers int `json:"readers"`
BytesReceived int64 `json:"bytes_received"`
@@ -46,7 +47,15 @@ func (c *HTTPController) Apply(ctx context.Context, source Source) error {
}
payload := map[string]any{"source": parsed.String(), "sourceOnDemand": true, "rtspTransport": "tcp"}
data, _ := json.Marshal(payload)
endpoint := c.base + "/v3/config/paths/replace/" + url.PathEscape(source.Path)
configured, err := c.configured(ctx, source.Path)
if err != nil {
return err
}
action := "add"
if configured {
action = "replace"
}
endpoint := c.base + "/v3/config/paths/" + action + "/" + url.PathEscape(source.Path)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(data))
if err != nil {
return err
@@ -63,6 +72,39 @@ func (c *HTTPController) Apply(ctx context.Context, source Source) error {
}
return nil
}
func (c *HTTPController) configured(ctx context.Context, path string) (bool, error) {
endpoint := c.base + "/v3/config/paths/get/" + url.PathEscape(path)
deadline := time.Now().Add(5 * time.Second)
var res *http.Response
for {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil)
if err != nil {
return false, err
}
res, err = c.client.Do(req)
if err == nil {
break
}
if !time.Now().Before(deadline) {
return false, fmt.Errorf("mediamtx control unavailable: %w", err)
}
select {
case <-ctx.Done():
return false, ctx.Err()
case <-time.After(100 * time.Millisecond):
}
}
defer res.Body.Close()
io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
if res.StatusCode == http.StatusNotFound {
return false, nil
}
if res.StatusCode < 200 || res.StatusCode >= 300 {
return false, fmt.Errorf("mediamtx config lookup returned %d", res.StatusCode)
}
return true, nil
}
func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.base+"/v3/paths/get/"+url.PathEscape(path), nil)
if err != nil {
@@ -74,7 +116,7 @@ func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, e
}
defer res.Body.Close()
if res.StatusCode == http.StatusNotFound {
return PathStatus{Name: path}, nil
return PathStatus{Name: path, Exists: false}, nil
}
if res.StatusCode < 200 || res.StatusCode >= 300 {
return PathStatus{}, fmt.Errorf("mediamtx status returned %d", res.StatusCode)
@@ -88,5 +130,6 @@ func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, e
if err := json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(&raw); err != nil {
return PathStatus{}, err
}
return PathStatus{Name: raw.Name, Ready: raw.Ready, Readers: len(raw.Readers), BytesReceived: raw.BytesReceived}, nil
return PathStatus{Name: raw.Name, Exists: true, Ready: raw.Ready, Readers: len(raw.Readers), BytesReceived: raw.BytesReceived}, nil
}
@@ -3,18 +3,30 @@ package mediamtx
import (
"context"
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
type roundTripFunc func(*http.Request) (*http.Response, error)
func (f roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { return f(request) }
func TestApplyBuildsCredentialSourceOnlyInTransientBody(t *testing.T) {
var body map[string]any
var appliedPath string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.Contains(r.URL.String(), "password") {
t.Fatal("credential leaked in control URL")
}
if r.Method == http.MethodGet {
http.NotFound(w, r)
return
}
appliedPath = r.URL.Path
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
t.Fatal(err)
}
@@ -28,6 +40,50 @@ func TestApplyBuildsCredentialSourceOnlyInTransientBody(t *testing.T) {
if body["source"] != "rtsp://fixture-user:fixture-password@camera.invalid/main" {
t.Fatalf("body=%#v", body)
}
if appliedPath != "/v3/config/paths/add/sense_test" {
t.Fatalf("applied path = %q", appliedPath)
}
}
func TestApplyReplacesExistingPath(t *testing.T) {
var appliedPath string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodGet {
w.WriteHeader(http.StatusOK)
return
}
appliedPath = r.URL.Path
w.WriteHeader(http.StatusOK)
}))
defer server.Close()
if err := NewHTTPController(server.URL).Apply(context.Background(), Source{Path: "sense_test", URI: "rtsp://camera.invalid/main"}); err != nil {
t.Fatal(err)
}
if appliedPath != "/v3/config/paths/replace/sense_test" {
t.Fatalf("applied path = %q", appliedPath)
}
}
func TestApplyWaitsForNewlyStartedControlAPI(t *testing.T) {
attempts := 0
client := NewHTTPController("http://127.0.0.1:9997")
client.client.Transport = roundTripFunc(func(request *http.Request) (*http.Response, error) {
attempts++
if attempts < 3 {
return nil, errors.New("connection refused")
}
status := http.StatusNotFound
if request.Method == http.MethodPost {
status = http.StatusOK
}
return &http.Response{StatusCode: status, Body: io.NopCloser(strings.NewReader("")), Header: make(http.Header)}, nil
})
if err := client.Apply(context.Background(), Source{Path: "sense_test", URI: "rtsp://camera.invalid/main"}); err != nil {
t.Fatal(err)
}
if attempts != 4 {
t.Fatalf("attempts = %d", attempts)
}
}
func TestApplyRejectsCredentialURI(t *testing.T) {
client := NewHTTPController("http://127.0.0.1")
@@ -35,3 +91,4 @@ func TestApplyRejectsCredentialURI(t *testing.T) {
t.Fatal("credential URI accepted")
}
}
+38
View File
@@ -48,6 +48,21 @@ func Require(permission string, next http.Handler) http.Handler {
})
}
// RequireSameOriginFrame authenticates browser iframe navigation that cannot
// attach the X-Product header used by API clients. Fetch Metadata keeps this
// narrow: only a same-origin iframe navigation with a valid Sense session is
// accepted.
func RequireSameOriginFrame(permission string, next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
module := activeModule.Load()
if module == nil {
platform.WriteError(w, &platform.APIError{Status: http.StatusServiceUnavailable, Code: "identity_not_ready", Message: "身份服务尚未就绪"})
return
}
module.RequireSameOriginFrame(permission, next).ServeHTTP(w, r)
})
}
func PrincipalFromContext(ctx context.Context) (Principal, bool) {
principal, ok := ctx.Value(principalContextKey{}).(Principal)
return principal, ok
@@ -84,6 +99,28 @@ func (m *Module) Require(permission string, next http.Handler) http.Handler {
}))
}
func (m *Module) RequireSameOriginFrame(permission string, next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Header.Get("Sec-Fetch-Site") != "same-origin" ||
r.Header.Get("Sec-Fetch-Mode") != "navigate" ||
r.Header.Get("Sec-Fetch-Dest") != "iframe" {
platform.WriteError(w, &platform.APIError{Status: http.StatusUnauthorized, Code: "unauthorized", Message: "登录状态无效"})
return
}
cookie, err := r.Cookie(sessionCookieName)
if err != nil {
platform.WriteError(w, &platform.APIError{Status: http.StatusUnauthorized, Code: "unauthorized", Message: "登录状态无效"})
return
}
principal, err := m.service.Authenticate(r.Context(), cookie.Value)
if err != nil || !HasPermission(principal.Role, permission) {
platform.WriteError(w, &platform.APIError{Status: http.StatusUnauthorized, Code: "unauthorized", Message: "登录状态无效"})
return
}
next.ServeHTTP(w, r.WithContext(context.WithValue(r.Context(), principalContextKey{}, principal)))
})
}
func (m *Module) bootstrap(w http.ResponseWriter, r *http.Request) {
var request struct {
Username string `json:"username"`
@@ -220,3 +257,4 @@ func SecureCookieFromEnvironment(value string, memoryMode bool) bool {
}
return !strings.EqualFold(value, "false")
}
@@ -21,6 +21,9 @@ func TestHTTPLoginCookieAndProductBoundary(t *testing.T) {
}
app := platform.NewApp(platform.Config{DatabaseMode: platform.DatabaseModeMemory}, nil, slog.New(slog.NewTextHandler(io.Discard, nil)))
NewModule(service, cfg).Register(app)
app.Handle("GET /same-origin-frame", RequireSameOriginFrame(PermissionMediaRead, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
})))
bootstrap := httptest.NewRequest(http.MethodPost, "/api/v1/identity/bootstrap", bytes.NewBufferString(`{"username":"admin","display_name":"管理员","password":"StrongPass2026"}`))
bootstrap.Header.Set("X-Sense-Bootstrap-Token", "bootstrap-test")
@@ -56,4 +59,41 @@ func TestHTTPLoginCookieAndProductBoundary(t *testing.T) {
if senseResult.Code != http.StatusOK {
t.Fatalf("Sense product header status = %d", senseResult.Code)
}
meWithoutProduct := httptest.NewRequest(http.MethodGet, "/api/v1/identity/me", nil)
meWithoutProduct.AddCookie(cookies[0])
missingProductResult := httptest.NewRecorder()
app.Handler().ServeHTTP(missingProductResult, meWithoutProduct)
if missingProductResult.Code != http.StatusUnauthorized {
t.Fatalf("missing product header status = %d", missingProductResult.Code)
}
frame := httptest.NewRequest(http.MethodGet, "/same-origin-frame", nil)
frame.AddCookie(cookies[0])
frame.Header.Set("Sec-Fetch-Site", "same-origin")
frame.Header.Set("Sec-Fetch-Mode", "navigate")
frame.Header.Set("Sec-Fetch-Dest", "iframe")
frameResult := httptest.NewRecorder()
app.Handler().ServeHTTP(frameResult, frame)
if frameResult.Code != http.StatusNoContent {
t.Fatalf("same-origin frame status = %d, body = %s", frameResult.Code, frameResult.Body.String())
}
frame.Header.Set("Sec-Fetch-Site", "cross-site")
crossSiteResult := httptest.NewRecorder()
app.Handler().ServeHTTP(crossSiteResult, frame)
if crossSiteResult.Code != http.StatusUnauthorized {
t.Fatalf("cross-site frame status = %d", crossSiteResult.Code)
}
missingCookie := httptest.NewRequest(http.MethodGet, "/same-origin-frame", nil)
missingCookie.Header.Set("Sec-Fetch-Site", "same-origin")
missingCookie.Header.Set("Sec-Fetch-Mode", "navigate")
missingCookie.Header.Set("Sec-Fetch-Dest", "iframe")
missingCookieResult := httptest.NewRecorder()
app.Handler().ServeHTTP(missingCookieResult, missingCookie)
if missingCookieResult.Code != http.StatusUnauthorized {
t.Fatalf("missing-cookie frame status = %d", missingCookieResult.Code)
}
}
+7 -2
View File
@@ -14,7 +14,7 @@ func (m *Module) Register(app *platform.App) {
app.Handle("GET /api/v1/liveview/routes", identity.Require(identity.PermissionMediaRead, http.HandlerFunc(m.routes)))
app.Handle("POST /api/v1/liveview/sessions", identity.Require(identity.PermissionMediaRead, http.HandlerFunc(m.create)))
app.Handle("GET /api/v1/liveview/sessions/{id}", identity.Require(identity.PermissionMediaRead, http.HandlerFunc(m.get)))
app.Handle("GET /api/v1/liveview/sessions/{id}/player", identity.Require(identity.PermissionMediaRead, http.HandlerFunc(m.player)))
app.Handle("GET /api/v1/liveview/sessions/{id}/player", identity.RequireSameOriginFrame(identity.PermissionMediaRead, http.HandlerFunc(m.player)))
}
func (m *Module) routes(w http.ResponseWriter, r *http.Request) {
items, err := m.service.Routes(r.Context())
@@ -61,6 +61,11 @@ func (m *Module) player(w http.ResponseWriter, r *http.Request) {
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
w.Header().Set("Cache-Control", "no-store")
w.Header().Set("Content-Security-Policy", "default-src 'none'; frame-src http: https:; style-src 'unsafe-inline'")
// This endpoint is the authenticated, same-origin wrapper loaded by the
// live-view page. Keep the global DENY policy everywhere else, and allow
// only Sense itself to embed this wrapper.
w.Header().Set("X-Frame-Options", "SAMEORIGIN")
w.Header().Set("Content-Security-Policy", "default-src 'none'; frame-ancestors 'self'; frame-src http: https:; style-src 'unsafe-inline'")
_ = playerTemplate.Execute(w, target)
}
@@ -0,0 +1,40 @@
package liveview
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"yovision.local/sense/app/sense/media"
)
func TestPlayerAllowsOnlySameOriginEmbedding(t *testing.T) {
service, err := NewService("http://127.0.0.1:8889", time.Minute)
if err != nil {
t.Fatal(err)
}
route := media.Route{ID: "device:main", Path: "sense_device_main", Desired: "running", Actual: "ready"}
service.route = func(context.Context, string) (media.Route, error) { return route, nil }
service.refresh = func(context.Context, string) (media.Route, error) { return route, nil }
service.sessions["view_test"] = Session{ID: "view_test", RouteID: route.ID, ExpiresAt: time.Now().Add(time.Minute)}
req := httptest.NewRequest(http.MethodGet, "/api/v1/liveview/sessions/view_test/player", nil)
req.SetPathValue("id", "view_test")
res := httptest.NewRecorder()
NewModule(service).player(res, req)
if res.Code != http.StatusOK {
t.Fatalf("status = %d", res.Code)
}
if got := res.Header().Get("X-Frame-Options"); got != "SAMEORIGIN" {
t.Fatalf("X-Frame-Options = %q", got)
}
csp := res.Header().Get("Content-Security-Policy")
if !strings.Contains(csp, "frame-ancestors 'self'") {
t.Fatalf("Content-Security-Policy = %q", csp)
}
}
+3 -2
View File
@@ -47,6 +47,7 @@ type Service struct {
base *url.URL
ttl time.Duration
route func(context.Context, string) (media.Route, error)
refresh func(context.Context, string) (media.Route, error)
routes func(context.Context) ([]media.Route, error)
mu sync.RWMutex
sessions map[string]Session
@@ -64,7 +65,7 @@ func NewService(rawBase string, ttl time.Duration) (*Service, error) {
if ttl <= 0 || ttl > 10*time.Minute {
ttl = 2 * time.Minute
}
return &Service{base: parsed, ttl: ttl, route: media.PlaybackRoute, routes: media.PlaybackRoutes, sessions: map[string]Session{}, now: time.Now}, nil
return &Service{base: parsed, ttl: ttl, route: media.PlaybackRoute, refresh: media.RefreshPlaybackRoute, routes: media.PlaybackRoutes, sessions: map[string]Session{}, now: time.Now}, nil
}
func (s *Service) Routes(ctx context.Context) ([]Route, error) {
items, err := s.routes(ctx)
@@ -117,7 +118,7 @@ func (s *Service) Get(ctx context.Context, owner, id string) (Session, error) {
if !ok || session.OwnerID != owner || !s.now().Before(session.ExpiresAt) {
return Session{}, fmt.Errorf("playback session expired")
}
route, err := s.route(ctx, session.RouteID)
route, err := s.refresh(ctx, session.RouteID)
if err != nil {
return Session{}, err
}
@@ -44,3 +44,28 @@ func TestTTLIsBounded(t *testing.T) {
t.Fatalf("ttl=%v", service.ttl)
}
}
func TestSessionPollRefreshesOnDemandMediaState(t *testing.T) {
service, err := NewService("http://127.0.0.1:8889", time.Minute)
if err != nil {
t.Fatal(err)
}
route := media.Route{ID: "device:main", DeviceID: "device", ProfileToken: "main", Path: "sense_device_main", Desired: "running", Actual: "waiting"}
service.route = func(context.Context, string) (media.Route, error) { return route, nil }
service.refresh = func(context.Context, string) (media.Route, error) {
updated := route
updated.Actual = "ready"
updated.Detail = "上游拉流正常"
updated.Readers = 1
return updated, nil
}
session, err := service.Create(context.Background(), "owner", route.ID)
if err != nil || session.Status != "waiting" || session.PlayerURL == "" {
t.Fatalf("session=%#v err=%v", session, err)
}
updated, err := service.Get(context.Background(), "owner", session.ID)
if err != nil || updated.Status != "ready" || updated.Detail != "上游拉流正常" {
t.Fatalf("updated=%#v err=%v", updated, err)
}
}
@@ -19,6 +19,13 @@ func PlaybackRoute(ctx context.Context, id string) (Route, error) {
}
return service.store.Get(ctx, id)
}
func RefreshPlaybackRoute(ctx context.Context, id string) (Route, error) {
service := activeService.Load()
if service == nil {
return Route{}, fmt.Errorf("media service is not ready")
}
return service.Refresh(ctx, id)
}
func PlaybackRoutes(ctx context.Context) ([]Route, error) {
service := activeService.Load()
if service == nil {
+67
View File
@@ -24,6 +24,32 @@ func NewService(store Store, process mediamtx.Process, controller mediamtx.Contr
return &Service{store: store, process: process, controller: controller, profile: admission.VerifiedProfile, credential: device.ReadRTSPCredential, now: time.Now}
}
// Restore recreates desired routes after MediaMTX starts with its base config.
// Individual route failures are persisted as safe business states and do not
// prevent the Sense management plane from starting.
func (s *Service) Restore(ctx context.Context) error {
items, err := s.store.List(ctx)
if err != nil {
return err
}
for _, route := range items {
if route.Desired != "running" {
continue
}
if _, reconcileErr := s.Reconcile(ctx, identity.Principal{}, route.ID); reconcileErr != nil {
route.Actual = "apply_failed"
route.Detail = "媒体路径恢复失败,请检查视频服务"
route.Readers = 0
route.Version++
route.UpdatedAt = s.now().UTC()
if saveErr := s.store.Save(ctx, route); saveErr != nil {
return saveErr
}
}
}
return nil
}
func (s *Service) ConfigureProfiles(ctx context.Context, actor identity.Principal, result admission.Result) admission.MediaOutcome {
configured := 0
ready := 0
@@ -108,6 +134,10 @@ func (s *Service) Reconcile(ctx context.Context, actor identity.Principal, id st
if err != nil {
route.Actual = "unconverged"
route.Detail = "尚未取得媒体状态"
} else if !status.Exists {
route.Actual = "apply_failed"
route.Detail = "媒体路径不存在,请重新对账"
route.Readers = 0
} else if status.Ready {
route.Actual = "ready"
route.Detail = "上游拉流正常"
@@ -125,6 +155,43 @@ func (s *Service) Reconcile(ctx context.Context, actor identity.Principal, id st
identity.RecordAudit(ctx, actor.UserID, "media.reconcile", id, "success", map[string]any{"actual": route.Actual})
return route, nil
}
// Refresh reads the MediaMTX runtime state without reapplying configuration or
// exposing the source URI. It is safe to call while a playback session polls.
func (s *Service) Refresh(ctx context.Context, id string) (Route, error) {
route, err := s.store.Get(ctx, id)
if err != nil {
return Route{}, err
}
if route.Desired != "running" {
return route, nil
}
status, err := s.controller.Status(ctx, route.Path)
if err != nil {
return route, nil
}
actual := "waiting"
detail := "等待播放器连接并按需拉流"
if !status.Exists {
actual = "apply_failed"
detail = "媒体路径不存在,请重新对账"
} else if status.Ready {
actual = "ready"
detail = "上游拉流正常"
}
if route.Actual == actual && route.Detail == detail && route.Readers == status.Readers {
return route, nil
}
route.Actual = actual
route.Detail = detail
route.Readers = status.Readers
route.Version++
route.UpdatedAt = s.now().UTC()
if err := s.store.Save(ctx, route); err != nil {
return Route{}, err
}
return route, nil
}
func (s *Service) Stop(ctx context.Context, actor identity.Principal, id string) (Route, error) {
route, err := s.store.Get(ctx, id)
if err != nil {
+51 -1
View File
@@ -40,14 +40,22 @@ type fakeController struct {
status mediamtx.PathStatus
}
type refreshController struct{ status mediamtx.PathStatus }
func (f refreshController) Apply(context.Context, mediamtx.Source) error { return nil }
func (f refreshController) Status(context.Context, string) (mediamtx.PathStatus, error) {
return f.status, nil
}
func (f fakeController) Apply(context.Context, mediamtx.Source) error { return f.applyErr }
func (f fakeController) Status(context.Context, string) (mediamtx.PathStatus, error) {
if f.status.Name == "error" {
return mediamtx.PathStatus{}, errors.New("timeout")
}
if !f.status.Ready {
return mediamtx.PathStatus{Ready: true, Readers: 2}, nil
return mediamtx.PathStatus{Exists: true, Ready: true, Readers: 2}, nil
}
f.status.Exists = true
return f.status, nil
}
@@ -142,5 +150,47 @@ func TestConfigureProfilesIsIdempotentAndReportsMediaState(t *testing.T) {
}
}
func TestRefreshTracksOnDemandReaderWithoutReapplyingRoute(t *testing.T) {
store := NewMemoryStore()
route := Route{ID: "device:main", DeviceID: "device", ProfileToken: "main", Path: "sense_device_main", Desired: "running", Actual: "waiting", Detail: "等待上游拉流", Version: 2}
if err := store.Save(context.Background(), route); err != nil {
t.Fatal(err)
}
service := NewService(store, &fakeProcess{}, refreshController{status: mediamtx.PathStatus{Name: route.Path, Exists: true, Ready: true, Readers: 1}})
result, err := service.Refresh(context.Background(), route.ID)
if err != nil || result.Actual != "ready" || result.Readers != 1 || result.Detail != "上游拉流正常" {
t.Fatalf("result=%#v err=%v", result, err)
}
unchanged, err := service.Refresh(context.Background(), route.ID)
if err != nil || unchanged.Version != result.Version {
t.Fatalf("unchanged=%#v err=%v", unchanged, err)
}
}
func TestRefreshReportsMissingPathInsteadOfWaiting(t *testing.T) {
store := NewMemoryStore()
route := Route{ID: "device:main", Path: "sense_device_main", Desired: "running", Actual: "waiting", Version: 1}
if err := store.Save(context.Background(), route); err != nil {
t.Fatal(err)
}
service := NewService(store, &fakeProcess{}, refreshController{status: mediamtx.PathStatus{Name: route.Path, Exists: false}})
result, err := service.Refresh(context.Background(), route.ID)
if err != nil || result.Actual != "apply_failed" || !strings.Contains(result.Detail, "不存在") {
t.Fatalf("result=%#v err=%v", result, err)
}
}
func TestRestoreReconcilesDesiredRoutes(t *testing.T) {
process := &fakeProcess{}
service, route := preparedService(t, process, fakeController{})
if err := service.Restore(context.Background()); err != nil {
t.Fatal(err)
}
restored, err := service.store.Get(context.Background(), route.ID)
if err != nil || !process.State().Running || restored.Actual != "ready" {
t.Fatalf("restored=%#v process=%#v err=%v", restored, process.State(), err)
}
}
func rtspReady() rtsp.Result { return rtsp.Result{Status: "ready"} }
+4 -1
View File
@@ -18,7 +18,9 @@ func init() {
process := mediamtx.NewSupervisor(os.Getenv("SENSE_MEDIAMTX_BINARY"), os.Getenv("SENSE_MEDIAMTX_CONFIG"), 3)
controller := mediamtx.NewHTTPController(valueOr("SENSE_MEDIAMTX_API", "http://127.0.0.1:9997"))
app.RegisterMigration(platform.Migration{Version: 2026081203, Name: "sense_media", SQL: media.MigrationSQL})
media.NewModule(media.NewService(store, process, controller)).Register(app)
service := media.NewService(store, process, controller)
media.NewModule(service).Register(app)
app.RegisterStartup(service.Restore)
return nil
})
}
@@ -28,3 +30,4 @@ func valueOr(key, fallback string) string {
}
return fallback
}
+4
View File
@@ -47,6 +47,9 @@ func Run() error {
if err := app.ApplyMigrations(context.Background()); err != nil {
return err
}
if err := app.Start(context.Background()); err != nil {
return fmt.Errorf("start Sense modules: %w", err)
}
server := &http.Server{
Addr: cfg.HTTPAddress,
@@ -77,3 +80,4 @@ func Run() error {
return serveErr
}
}
+16
View File
@@ -1,6 +1,7 @@
package platform
import (
"context"
"database/sql"
"io/fs"
"log/slog"
@@ -17,6 +18,20 @@ type App struct {
logger *slog.Logger
mux *http.ServeMux
migrations []Migration
startups []func(context.Context) error
}
func (a *App) RegisterStartup(startup func(context.Context) error) {
a.startups = append(a.startups, startup)
}
func (a *App) Start(ctx context.Context) error {
for _, startup := range a.startups {
if err := startup(ctx); err != nil {
return err
}
}
return nil
}
func NewApp(cfg Config, database *sql.DB, logger *slog.Logger) *App {
@@ -101,3 +116,4 @@ func requestSecurityHeaders(next http.Handler) http.Handler {
next.ServeHTTP(w, r)
})
}
@@ -1,14 +1,15 @@
<template>
<div class="stream-player">
<div v-if="state === 'loading'" class="player-state"><el-icon class="is-loading" size="42"><Loading /></el-icon><strong>正在打开视频</strong><span>通常需要几秒钟</span></div>
<div v-else-if="state !== 'ready'" class="player-state"><el-icon size="46"><WarningFilled /></el-icon><strong>{{ title }}</strong><span>{{ detail }}</span><el-button type="primary" @click="$emit('retry')">重新连接</el-button></div>
<iframe v-else :key="playerUrl" :src="playerUrl" title="Sense 单路实时视频" allow="autoplay; fullscreen" @load="$emit('loaded')" />
<iframe v-if="playable" :key="playerUrl" :src="playerUrl" title="Sense 单路实时视频" allow="autoplay; fullscreen" @load="$emit('loaded')" />
<div v-else-if="state === 'loading'" class="player-state"><el-icon class="is-loading" size="42"><Loading /></el-icon><strong>正在打开视频</strong><span>通常需要几秒钟</span></div>
<div v-else class="player-state"><el-icon size="46"><WarningFilled /></el-icon><strong>{{ title }}</strong><span>{{ detail }}</span><el-button type="primary" @click="$emit('retry')">重新连接</el-button></div>
</div>
</template>
<script setup>
import{computed}from'vue'
const props=defineProps({state:{type:String,default:'loading'},detail:{type:String,default:''},playerUrl:{type:String,default:''}});defineEmits(['retry','loaded'])
const playable=computed(()=>Boolean(props.playerUrl)&&['waiting','ready'].includes(props.state))
const title=computed(()=>({stopped:'视频已停止',process_failed:'视频服务未启动',apply_failed:'视频配置失败',unconverged:'视频状态未同步',waiting:'正在等待视频',expired:'播放会话已过期',offline:'视频已断开'})[props.state]||'暂时无法播放')
</script>
<style scoped>.stream-player{position:relative;width:100%;aspect-ratio:16/9;overflow:hidden;border-radius:4px;background:#101419}.stream-player iframe{width:100%;height:100%;border:0}.player-state{position:absolute;inset:0;display:flex;flex-direction:column;align-items:center;justify-content:center;gap:12px;color:#c9cdd4}.player-state strong{color:#fff;font-size:18px}.player-state span{max-width:70%;text-align:center;font-size:13px}</style>
+4 -4
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: 9aa7be78bf7e988594d38e9e6b264db4edf71fc9
synchronized_at: 2026-08-13T04:16:00Z
wiki_revision: e6bb3132a752cd8036cab29429cbf3f49e85bfbf
synchronized_at: 2026-08-13T14:06:00Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -102,8 +102,8 @@ Sense 后端功能以 `Sense/server/app/sense/` 为根,并通过 `Sense/server
- `identity/`:Sense 独立账户、bcrypt 密码、会话、四角色 RBAC 与统一审计;签发者和受众只属于 Sense。
- `device/`:Device 台账、状态、分页和 AES-256-GCM 凭据保险箱;ONVIF 与 RTSP 凭据可分离或显式复用,读取模型只返回两组凭据是否已配置。
- `adapters/onvif/`、`adapters/rtsp/`、`admission/`:获准网卡上的受控发现、手工 ONVIF 接入、Media 服务发现、Basic/Digest 认证、Profile/StreamUri 读取和 RTSP 验证。跨主机 Media 地址固定回用户已授权的 Device Service origin;跨主机 RTSP URI 只替换为授权主机并保留报告端口与路径。接入结果和脱敏 Profile 持久化到 PostgreSQL,重启后可恢复。
- `adapters/mediamtx/`、`media/`:外部 MediaMTX 进程所有权、localhost Control API、媒体期望态与实际态对账;接入验证出可用 Profile 后自动按 Device/Profile 幂等建立并对账媒体路径。
- `liveview/`:绑定当前用户、最长两分钟的单路播放会话;列表投影设备名称、位置、码流名称与用途,只在内部保留 ID,不暴露源 URI 或摄像机秘密。
- `adapters/mediamtx/`、`media/`:外部 MediaMTX 进程所有权、localhost Control API、媒体期望态与实际态对账;接入验证出可用 Profile 后自动按 Device/Profile 幂等建立并对账媒体路径。首次路径使用 add、已有路径使用 replace,控制 API 启动竞态在限定时间内重试;Sense 完成数据库迁移后自动恢复 `desired=running` 路径。路径不存在必须报告配置失败,不能伪装为 waiting。
- `liveview/`:绑定当前用户、最长两分钟的单路播放会话;列表投影设备名称、位置、码流名称与用途,只在内部保留 ID,不暴露源 URI 或摄像机秘密。`waiting` 会话立即加载播放器以触发 MediaMTX 按需拉流,会话轮询只读刷新媒体实际状态。播放器包装页仅允许被 Sense 同源页面嵌入(`SAMEORIGIN` 与 `frame-ancestors 'self'`),其他页面继续使用全局 `DENY`。iframe 导航使用有效 Sense 会话 Cookie、媒体读取权限及同源 iframe Fetch Metadata 认证,不依赖浏览器导航无法附加的 `X-Product` 请求头;普通 API 仍要求产品头。
- `area/`:归一化多边形/方向警戒线、不可变版本、并发版本校验和分辨率变化后的重新校准。
前端在 `Sense/ui/src/{api,views,router/modules,components}/sense/` 使用对应模块;通用页面复用 Element Plus 表单、表格、分页、Dialog、Tag 和应用容器,只为播放器与区域画布新增局部业务组件。项目内 `identity.RecordAudit`、`device.ReadRTSPCredential`、`device.Describe`、`admission.VerifiedProfile`、`media.PlaybackRoute` 和 `area.ExportCurrent` 是窄适配端口,不是跨项目契约。
+4 -4
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: 00c48e986cb9d75daae0b6a7eafe0529e3aba36a
synchronized_at: 2026-08-13T04:16:00Z
wiki_revision: 06b68d16bf907decdaf85a813d417896f7920f4c
synchronized_at: 2026-08-13T14:06:00Z
<!-- gitea-wiki-mirror:end -->
# 业务规则与术语
@@ -76,9 +76,9 @@ synchronized_at: 2026-08-13T04:16:00Z
- **摄像机凭据**:ONVIF 与 RTSP 可使用不同账号,也可显式复用;两组均只写不读,使用仓库外 32 字节密钥分别加密。密码不得进入 URL、日志、审计、工单或响应。
- **受控发现**:ONVIF Discovery 默认关闭,只有显式设置获准本机 IP 才能发送发现;不得扫描未授权网段。
- **Profile**:主辅码流按分辨率分类并使用 RTSP 凭据分别验证;脱敏 Stream URI、验证状态和时间持久化,重启后保留。至少一个 Profile 验证成功后 Device 进入 `active`;认证失败、不可达、超时与校时问题使用可定位状态。
- **自动媒体路径**:接入检查中每个验证成功的 Profile 都按 Device/Profile 幂等建立并立即对账媒体路径;MediaMTX 不可用不回滚摄像机接入,接入结果记录“需要处理”并引导到视频服务排错。
- **自动媒体路径**:接入检查中每个验证成功的 Profile 都按 Device/Profile 幂等建立并立即对账媒体路径;首次配置使用 MediaMTX add、已有配置使用 replace,Sense 启动时从数据库恢复所有 `desired=running` 路径;路径不存在属于配置失败,不属于等待拉流。MediaMTX 不可用不回滚摄像机接入,接入结果记录“需要处理”并引导到视频服务排错。
- **MediaMTX**:保持外部进程。Sense 只停止自己启动并持有句柄的进程,最多自动重启三次;摄像机凭据只在 localhost 控制请求中瞬时组装,不持久化、不返回。
- **播放会话**:由当前 Sense 用户创建,最长两分钟;实时监看以设备名称、位置和主/子码流等业务标签供用户选择,内部 ID 只用于系统关联;设备分页和媒体路径不以 16 路作为硬上限,页面一次只打开一路流。
- **播放会话**:由当前 Sense 用户创建,最长两分钟;实时监看以设备名称、位置和主/子码流等业务标签供用户选择,内部 ID 只用于系统关联;`waiting` 表示等待第一个播放器触发按需拉流,不是播放失败,播放器连接后会话轮询刷新为实际状态;播放器导航必须携带有效 Sense 会话 Cookie,且 Fetch Metadata 必须表明是同源 iframe,普通 API 的 `X-Product` 边界不变;设备分页和媒体路径不以 16 路作为硬上限,页面一次只打开一路流。
- **区域版本**:坐标为 0..1 归一化值,并绑定 Device、Profile、宽高。每次发布或停用形成新版本;范围、点数、自交、退化、方向和期望版本由后端校验。Profile 分辨率变化后旧版本必须标记为需要重新校准。
<!-- sense-mvp:end -->
+6 -3
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: caa6704ed8aecb77e829504c2a739e30f5404aa9
synchronized_at: 2026-08-13T04:16:00Z
wiki_revision: b70e92180546609cb9b3b05b2c13c1a49676c03b
synchronized_at: 2026-08-13T14:05:00Z
<!-- gitea-wiki-mirror:end -->
# 故障排查
@@ -55,7 +55,10 @@ synchronized_at: 2026-08-13T04:16:00Z
| 接入成功但提示“媒体服务未就绪” / `process_failed` | 接入资料和 Profile 已保存,实时监看路径也已建立;检查仓库外 `SENSE_MEDIAMTX_BINARY`、基础配置和进程退出原因,修复后到视频服务重新对账。达到三次重启上限后需人工处理。 |
| `apply_failed` / `unconverged` | 检查 localhost Control API 是否启用并为 v3;确认媒体路径和外部进程状态。 |
| 实时监看没有设备 | 先在视频接入完成一次检查;系统会自动建立媒体路径。若已接入仍为空,检查接入结果的媒体状态和视频服务对账。 |
| 画面等待、断开或会话过期 | 查看实时监看中的业务设备名和码流状态,必要时到视频服务执行对账后重新连接;播放会话最长两分钟。 |
| 实时监看提示“127.0.0.1 拒绝了我们的连接请求” | 先确认播放器接口响应头;该接口必须是 `X-Frame-Options: SAMEORIGIN` 且 CSP 包含 `frame-ancestors 'self'`,普通页面仍应为 `DENY`。播放器接口还必须允许不带 `X-Product`、但带有效 Sense Cookie 和同源 iframe Fetch Metadata 的浏览器导航;普通 API 仍应拒绝缺少产品头的请求。旧运行包会在响应头或认证中间件处阻止 iframe,需重新打包并重启 Sense。 |
| 播放器提示 `stream not found` | MediaMTX 中没有对应配置路径。检查 Control API 配置列表是否包含 Sense 路径;新版本会在 Sense 启动时自动恢复 `desired=running` 路径,并在控制端口尚未就绪时限时重试。路径缺失必须显示配置失败,不能仅显示等待拉流。 |
| 实时监看持续“等待拉流” | `waiting` 会话应立即加载播放器以触发按需拉流。若连接后仍等待,检查 MediaMTX 路径 readers 和源状态;`401 Unauthorized` 表示摄像机拒绝当前 RTSP 凭据,应在设备管理修正凭据后重新执行视频接入,不能把凭据写入日志或地址。 |
| 画面等待、断开或会话过期 | 查看实时监看中的业务设备名和码流状态;会话轮询会只读刷新 MediaMTX 实际状态,无需手工对账。播放会话最长两分钟。 |
| 区域提示需要重新校准 | Profile 分辨率已变化,按新画面重新绘制并发布新版本,不能静默复用旧坐标。 |
自动测试不访问真实摄像头或未授权网络。PostgreSQL、MediaMTX、目标浏览器与实验室摄像机的联合验证必须在获准部署环境完成。
@@ -0,0 +1,86 @@
<!-- gitea-wiki-mirror:start -->
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Task-54-修复实时监看按需拉流循环等待
wiki_url: https://git.ilapage.cn/ila/yovision/wiki/Task-54-%E4%BF%AE%E5%A4%8D%E5%AE%9E%E6%97%B6%E7%9B%91%E7%9C%8B%E6%8C%89%E9%9C%80%E6%8B%89%E6%B5%81%E5%BE%AA%E7%8E%AF%E7%AD%89%E5%BE%85.-
wiki_revision: ff79f9dbc718ae570b01dace25304379124a99b4
synchronized_at: 2026-08-13T14:05:00Z
<!-- gitea-wiki-mirror:end -->
# 54 修复实时监看按需拉流循环等待
- 类型:缺陷
- 所属 Epic:#7
- 所属 MVP / 版本:#8
- 状态:待验收
- 日期:2026-08-13
- Gitea 工单:https://git.ilapage.cn/ila/yovision/issues/54
- Wiki 页面:Task-54-修复实时监看按需拉流循环等待
- Wiki revision:见本地镜像头
## 背景与目标
MediaMTX 使用按需拉流时,必须先有播放器读取路径才会连接摄像机。原实时监看页面只在状态变为“可播放”后才加载播放器,导致“等待拉流 → 没有读取者 → 始终等待”的循环。目标是在保持按需拉流和安全边界的前提下,让监看页主动成为读取者,并让会话状态收敛到真实媒体状态。
## 最终方案
- 实时监看收到有效播放地址且状态为“等待拉流”或“可播放”时都加载播放器;“等待拉流”保留明确状态提示,同时触发 MediaMTX 按需连接。
- 会话查询通过只读媒体端口刷新控制器状态,把等待状态安全收敛为可播放或需要处理,不重新应用路径配置。
- 播放器包装接口覆盖全局嵌入策略为 `SAMEORIGIN`,并用 CSP `frame-ancestors 'self'` 限制为 Sense 同源页面;其他接口继续保持 `X-Frame-Options: DENY`。
- iframe 导航无法附加 API 客户端的 `X-Product` 头,因此播放器路由改用窄导航认证:有效 Sense 会话 Cookie、媒体读取权限、`Sec-Fetch-Site: same-origin`、`Sec-Fetch-Mode: navigate` 和 `Sec-Fetch-Dest: iframe` 必须同时满足。普通 API 的产品头要求保持不变。
- MediaMTX 路径首次配置使用 add,已存在时使用 replace;控制 API 刚启动尚未监听时限时等待。
- Sense 完成数据库迁移后自动恢复所有 `desired=running` 路径;路径不存在显式标记为配置失败,不再误报 waiting。
- 继续保持 `sourceOnDemand`,无人观看时不占用摄像机连接与转码资源。
- 真实环境验证过程中发现已保存的 RTSP 凭据被摄像机拒绝(401);已通过现有 Sense 接口使用本地已获授权配置修正运行数据,未把凭据写入仓库、工单、Wiki 或日志证据。
## 修改文件
- `Sense/server/app/sense/media/service.go`、`route_port.go`:增加只读状态刷新、缺失状态和启动恢复能力。
- `Sense/server/app/sense/adapters/mediamtx/client.go`、`client_test.go`:实现 add/replace 幂等语义并等待控制 API 就绪。
- `Sense/server/cmd/sense/{modules_media.go,root.go}`、`internal/platform/app.go`:在迁移完成后执行媒体路径恢复。
- `Sense/server/app/sense/media/service_test.go`:验证刷新状态且不重复配置路径。
- `Sense/server/app/sense/liveview/service.go`:查询会话时刷新媒体状态。
- `Sense/server/app/sense/liveview/http.go`、`http_test.go`:允许并验证播放器包装页仅同源嵌入。
- `Sense/server/app/sense/identity/http.go`、`http_test.go`:增加并验证仅限同源 iframe 的 Cookie 导航认证,保留普通 API 产品边界。
- `Sense/server/app/sense/liveview/service_test.go`:验证等待状态收敛到可播放。
- `Sense/ui/src/components/sense/liveview/StreamPlayer.vue`:等待拉流时即加载播放器。
- Wiki 架构、业务规则与排错页面。
## 验收结果
| 验收标准 | 结果 |
|---|---|
| 等待拉流时主动创建播放器读取者 | 通过,真实 MediaMTX 路径观察到 1 个读取者 |
| 会话状态从 waiting 收敛为 ready | 通过,真实 API 轮询由 waiting 变为 ready |
| 实际媒体源可解码 | 通过,探测到 H.264 1920×1080 视频和 AAC 音频 |
| 不重复应用路径配置 | 通过,单元测试验证刷新仅调用状态查询 |
| 不泄露摄像机凭据和源地址 | 通过 |
| 播放器可被 Sense 同源嵌入且其他页面仍禁止嵌入 | 通过,运行包接口响应头验证通过 |
| iframe 不带 `X-Product` 时仍能安全认证 | 通过,同源 iframe + 有效 Cookie 返回 200;跨站返回 401;普通 API 缺少产品头返回 401 |
| MediaMTX 冷启动后恢复 Sense 路径 | 通过,配置列表自动出现 2 条 `sourceOnDemand` 路径,无需手工对账 |
| 两条摄像机路径可实际解码 | 通过,主码流 H.264 1920×1080 + AAC,子码流 H.264,均收到媒体字节 |
## 测试
- `go test ./...`:通过。
- `corepack pnpm lint`:0 error,806 个既有格式 warning。
- `corepack pnpm build`:通过,存在 3 个既有 webpack 体积 warning。
- Go 1.26.5 Windows 打包:通过。
- 运行包响应头:播放器接口为 `SAMEORIGIN` 且包含 `frame-ancestors 'self'`;首页仍为 `DENY`。
- 运行包浏览器式认证:不带 `X-Product`、带有效 Cookie 和同源 iframe Fetch Metadata 的播放器请求返回 200;跨站播放器与缺少产品头的普通 API 均返回 401。
- 冷启动恢复:MediaMTX 配置列表自动包含两条 Sense `sourceOnDemand` 路径;运行态列表也存在两条路径。
- 真实媒体链路:主码流 H.264 1920×1080 + AAC,子码流 H.264;两条路径分别通过本机 RTSP 读取并收到媒体字节。
- 单元测试覆盖路径首次 add、已有 replace、控制 API 启动等待、启动恢复及缺失路径不误报 waiting。
- **未验证部分**:最终浏览器中的可视画面需要用户在当前浏览器验收;自动化已验证媒体源、读取者与会话状态收敛。
## 遗留问题
- 严格 Harness 仍受既有 #44 归档缺少“最终方案”章节影响。
## 相关提交
- `8d4e4c2` 修复实时监看按需拉流循环等待。
- `b3dfd23` 更新按需拉流长期 Wiki 镜像。
- `3696442` 允许播放器包装页仅被 Sense 同源嵌入。
- `8bf9a6d` 支持同源播放器导航认证并保持普通 API 产品边界。
- `d7cd3a5` 恢复 MediaMTX 按需路径并处理启动竞态。
+4
View File
@@ -143,6 +143,10 @@
{
"page": "Task-51-接入成功后自动建立媒体路由并进入实时监看",
"path": "docs/task/51-接入成功后自动建立媒体路由并进入实时监看.md"
},
{
"page": "Task-54-修复实时监看按需拉流循环等待",
"path": "docs/task/54-修复实时监看按需拉流循环等待.md"
}
]
}