[SEN] 修复实时监看按需拉流循环等待 (#54) #55
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"} }
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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>
|
||||
|
||||
|
||||
|
||||
@@ -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` 是窄适配端口,不是跨项目契约。
|
||||
|
||||
@@ -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 -->
|
||||
|
||||
|
||||
@@ -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 按需路径并处理启动竞态。
|
||||
|
||||
@@ -143,6 +143,10 @@
|
||||
{
|
||||
"page": "Task-51-接入成功后自动建立媒体路由并进入实时监看",
|
||||
"path": "docs/task/51-接入成功后自动建立媒体路由并进入实时监看.md"
|
||||
},
|
||||
{
|
||||
"page": "Task-54-修复实时监看按需拉流循环等待",
|
||||
"path": "docs/task/54-修复实时监看按需拉流循环等待.md"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user