fix: 恢复MediaMTX按需媒体路径 (#54)
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 = "上游拉流正常"
|
||||
}
|
||||
|
||||
@@ -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"} }
|
||||
|
||||
@@ -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
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user