diff --git a/Sense/server/app/sense/adapters/mediamtx/client.go b/Sense/server/app/sense/adapters/mediamtx/client.go index c139c97..a9e54b9 100644 --- a/Sense/server/app/sense/adapters/mediamtx/client.go +++ b/Sense/server/app/sense/adapters/mediamtx/client.go @@ -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 } + diff --git a/Sense/server/app/sense/adapters/mediamtx/client_test.go b/Sense/server/app/sense/adapters/mediamtx/client_test.go index ce9ab6e..697b540 100644 --- a/Sense/server/app/sense/adapters/mediamtx/client_test.go +++ b/Sense/server/app/sense/adapters/mediamtx/client_test.go @@ -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") } } + diff --git a/Sense/server/app/sense/identity/http.go b/Sense/server/app/sense/identity/http.go index 977eee0..6a3f00d 100644 --- a/Sense/server/app/sense/identity/http.go +++ b/Sense/server/app/sense/identity/http.go @@ -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") } + diff --git a/Sense/server/app/sense/identity/http_test.go b/Sense/server/app/sense/identity/http_test.go index d96302c..b60396b 100644 --- a/Sense/server/app/sense/identity/http_test.go +++ b/Sense/server/app/sense/identity/http_test.go @@ -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) + } } + diff --git a/Sense/server/app/sense/liveview/http.go b/Sense/server/app/sense/liveview/http.go index 1e0c444..1dca78e 100644 --- a/Sense/server/app/sense/liveview/http.go +++ b/Sense/server/app/sense/liveview/http.go @@ -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) } + diff --git a/Sense/server/app/sense/liveview/http_test.go b/Sense/server/app/sense/liveview/http_test.go new file mode 100644 index 0000000..228ef14 --- /dev/null +++ b/Sense/server/app/sense/liveview/http_test.go @@ -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) + } +} + diff --git a/Sense/server/app/sense/liveview/service.go b/Sense/server/app/sense/liveview/service.go index dc412a4..fd9f243 100644 --- a/Sense/server/app/sense/liveview/service.go +++ b/Sense/server/app/sense/liveview/service.go @@ -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 } diff --git a/Sense/server/app/sense/liveview/service_test.go b/Sense/server/app/sense/liveview/service_test.go index 7bd13f9..db95b81 100644 --- a/Sense/server/app/sense/liveview/service_test.go +++ b/Sense/server/app/sense/liveview/service_test.go @@ -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) + } +} + diff --git a/Sense/server/app/sense/media/route_port.go b/Sense/server/app/sense/media/route_port.go index b266013..d48eaf0 100644 --- a/Sense/server/app/sense/media/route_port.go +++ b/Sense/server/app/sense/media/route_port.go @@ -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 { diff --git a/Sense/server/app/sense/media/service.go b/Sense/server/app/sense/media/service.go index b2afb49..89bd57e 100644 --- a/Sense/server/app/sense/media/service.go +++ b/Sense/server/app/sense/media/service.go @@ -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 { diff --git a/Sense/server/app/sense/media/service_test.go b/Sense/server/app/sense/media/service_test.go index dc55b3b..c4c6201 100644 --- a/Sense/server/app/sense/media/service_test.go +++ b/Sense/server/app/sense/media/service_test.go @@ -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"} } diff --git a/Sense/server/cmd/sense/modules_media.go b/Sense/server/cmd/sense/modules_media.go index da4e7fa..2bfdbc9 100644 --- a/Sense/server/cmd/sense/modules_media.go +++ b/Sense/server/cmd/sense/modules_media.go @@ -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 } + diff --git a/Sense/server/cmd/sense/root.go b/Sense/server/cmd/sense/root.go index 235d093..cd99c56 100644 --- a/Sense/server/cmd/sense/root.go +++ b/Sense/server/cmd/sense/root.go @@ -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 } } + diff --git a/Sense/server/internal/platform/app.go b/Sense/server/internal/platform/app.go index 397c1ed..d31068a 100644 --- a/Sense/server/internal/platform/app.go +++ b/Sense/server/internal/platform/app.go @@ -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) }) } + diff --git a/Sense/ui/src/components/sense/liveview/StreamPlayer.vue b/Sense/ui/src/components/sense/liveview/StreamPlayer.vue index 3dd5e38..594ba01 100644 --- a/Sense/ui/src/components/sense/liveview/StreamPlayer.vue +++ b/Sense/ui/src/components/sense/liveview/StreamPlayer.vue @@ -1,14 +1,15 @@ - 正在打开视频通常需要几秒钟 - {{ title }}{{ detail }}重新连接 - + + 正在打开视频通常需要几秒钟 + {{ title }}{{ detail }}重新连接 - + diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index ec67f74..449be2b 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -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 # 架构与代码地图 @@ -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` 是窄适配端口,不是跨项目契约。 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 170e1d7..db55f39 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -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 # 业务规则与术语 @@ -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 分辨率变化后旧版本必须标记为需要重新校准。 diff --git a/docs/06-troubleshooting.md b/docs/06-troubleshooting.md index e9983ac..70f620f 100644 --- a/docs/06-troubleshooting.md +++ b/docs/06-troubleshooting.md @@ -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 # 故障排查 @@ -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、目标浏览器与实验室摄像机的联合验证必须在获准部署环境完成。 diff --git a/docs/task/54-修复实时监看按需拉流循环等待.md b/docs/task/54-修复实时监看按需拉流循环等待.md new file mode 100644 index 0000000..6805ec7 --- /dev/null +++ b/docs/task/54-修复实时监看按需拉流循环等待.md @@ -0,0 +1,86 @@ + +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 + + +# 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 按需路径并处理启动竞态。 + diff --git a/wiki-docs.json b/wiki-docs.json index 5b91bd2..f8ccd6a 100644 --- a/wiki-docs.json +++ b/wiki-docs.json @@ -143,6 +143,10 @@ { "page": "Task-51-接入成功后自动建立媒体路由并进入实时监看", "path": "docs/task/51-接入成功后自动建立媒体路由并进入实时监看.md" + }, + { + "page": "Task-54-修复实时监看按需拉流循环等待", + "path": "docs/task/54-修复实时监看按需拉流循环等待.md" } ] }