diff --git a/Sense/server/app/sense/adapters/mediamtx/client.go b/Sense/server/app/sense/adapters/mediamtx/client.go index c139c97..957ab33 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,5 @@ 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..122b192 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") diff --git a/Sense/server/app/sense/media/service.go b/Sense/server/app/sense/media/service.go index 513fb10..16222d9 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 = "上游拉流正常" @@ -142,7 +172,10 @@ func (s *Service) Refresh(ctx context.Context, id string) (Route, error) { } actual := "waiting" detail := "等待播放器连接并按需拉流" - if status.Ready { + if !status.Exists { + actual = "apply_failed" + detail = "媒体路径不存在,请重新对账" + } else if status.Ready { actual = "ready" detail = "上游拉流正常" } diff --git a/Sense/server/app/sense/media/service_test.go b/Sense/server/app/sense/media/service_test.go index 3f19420..cb5f87a 100644 --- a/Sense/server/app/sense/media/service_test.go +++ b/Sense/server/app/sense/media/service_test.go @@ -53,8 +53,9 @@ func (f fakeController) Status(context.Context, string) (mediamtx.PathStatus, er 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 } @@ -155,7 +156,7 @@ func TestRefreshTracksOnDemandReaderWithoutReapplyingRoute(t *testing.T) { if err := store.Save(context.Background(), route); err != nil { t.Fatal(err) } - service := NewService(store, &fakeProcess{}, refreshController{status: mediamtx.PathStatus{Name: route.Path, Ready: true, Readers: 1}}) + 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) @@ -166,4 +167,29 @@ func TestRefreshTracksOnDemandReaderWithoutReapplyingRoute(t *testing.T) { } } +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..0f7d917 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 }) } diff --git a/Sense/server/cmd/sense/root.go b/Sense/server/cmd/sense/root.go index 235d093..8803fc9 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, diff --git a/Sense/server/internal/platform/app.go b/Sense/server/internal/platform/app.go index 397c1ed..dea1acd 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 {