fix: 修复实时监看按需拉流循环等待 (#54)
This commit is contained in:
@@ -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,27 @@ 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 {
|
||||
|
||||
@@ -125,6 +125,40 @@ 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.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,6 +40,13 @@ 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" {
|
||||
@@ -142,4 +149,21 @@ 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, 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 rtspReady() rtsp.Result { return rtsp.Result{Status: "ready"} }
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
<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>
|
||||
|
||||
|
||||
Reference in New Issue
Block a user