From 179768ca4a768611afc8dedcbdcd4ed42318e684 Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Tue, 1 Sep 2026 23:35:41 +0300 Subject: [PATCH] playlist retry-exhaustion handling --- cmd/mxl-player/playlist_file.go | 15 +- cmd/mxl-player/playlist_file_test.go | 7 +- imgui.ini | 36 +--- internal/playback/playlist.go | 34 +++- internal/playback/playlist_controller.go | 48 +++++- internal/playback/playlist_failure_test.go | 181 +++++++++++++++++++++ internal/playback/playlist_readiness.go | 68 +++++++- playlists/sample-list-all.json | 1 + playlists/sample-list.json | 1 + 9 files changed, 346 insertions(+), 45 deletions(-) create mode 100644 internal/playback/playlist_failure_test.go diff --git a/cmd/mxl-player/playlist_file.go b/cmd/mxl-player/playlist_file.go index 2a6efb7..108280c 100644 --- a/cmd/mxl-player/playlist_file.go +++ b/cmd/mxl-player/playlist_file.go @@ -11,8 +11,9 @@ import ( ) type playlistFile struct { - Entries []playlistFileEntry `json:"entries"` - Loop bool `json:"loop"` + Entries []playlistFileEntry `json:"entries"` + Loop bool `json:"loop"` + OnFailure string `json:"on_failure"` } type playlistFileEntry struct { @@ -62,6 +63,16 @@ func decodePlaylistFile(reader io.Reader) (playback.Playlist, error) { Entries: make([]playback.PlaylistEntry, len(file.Entries)), Loop: file.Loop, } + switch file.OnFailure { + case "", "wait": + playlist.OnFailure = playback.PlaylistFailureWait + case "next": + playlist.OnFailure = playback.PlaylistFailureNext + default: + return playback.Playlist{}, fmt.Errorf( + "on_failure %q: %w", file.OnFailure, playback.ErrPlaylistFailurePolicy, + ) + } for index, entry := range file.Entries { duration := time.Duration(0) if entry.Duration != "" { diff --git a/cmd/mxl-player/playlist_file_test.go b/cmd/mxl-player/playlist_file_test.go index 9505baa..1cd15c5 100644 --- a/cmd/mxl-player/playlist_file_test.go +++ b/cmd/mxl-player/playlist_file_test.go @@ -14,6 +14,7 @@ import ( func TestDecodePlaylistFile(t *testing.T) { input := `{ "loop": true, + "on_failure": "next", "entries": [ { "name": "sync", @@ -44,7 +45,8 @@ func TestDecodePlaylistFile(t *testing.T) { t.Fatalf("decodePlaylistFile() error = %v", err) } want := playback.Playlist{ - Loop: true, + Loop: true, + OnFailure: playback.PlaylistFailureNext, Entries: []playback.PlaylistEntry{ { Name: "sync", @@ -69,7 +71,7 @@ func TestDecodePlaylistFile(t *testing.T) { }, }, } - if len(got.Entries) != len(want.Entries) || got.Loop != want.Loop { + if len(got.Entries) != len(want.Entries) || got.Loop != want.Loop || got.OnFailure != want.OnFailure { t.Fatalf("decodePlaylistFile() = %#v, want %#v", got, want) } for index := range want.Entries { @@ -100,6 +102,7 @@ func TestDecodePlaylistFileRejectsInvalidInput(t *testing.T) { {name: "malformed JSON", input: `{"entries": [`, wantText: "decode JSON"}, {name: "unknown field", input: `{"unknown": true}`, wantText: "unknown field"}, {name: "multiple roots", input: `{"entries": []} {"entries": []}`, wantText: "multiple root values"}, + {name: "invalid failure policy", input: `{"on_failure":"skip","entries":[]}`, wantErr: playback.ErrPlaylistFailurePolicy}, { name: "invalid duration", input: `{"entries":[{"video":{"domain":"/video","uuid":"video"},"duration":"later"}]}`, diff --git a/imgui.ini b/imgui.ini index 7bdb62d..bc38a87 100644 --- a/imgui.ini +++ b/imgui.ini @@ -1,39 +1,15 @@ [Window][Debug##Default] -Pos=519,181 -Size=400,398 -Collapsed=0 - -[Window][Test] Pos=60,60 -Size=251,92 -Collapsed=0 - -[Window][Stats] -Size=460,510 -Collapsed=0 - -[Window][Connection] -Pos=42,264 -Size=605,416 -Collapsed=0 - -[Window][Test slider] -Pos=1370,0 -Size=550,1080 -Collapsed=0 - -[Window][Settings] -Pos=730,0 -Size=550,720 +Size=400,400 Collapsed=0 [Window][Settings & Info] -Pos=580,0 -Size=700,720 -Collapsed=0 - -[Window][Seetings & Info] Pos=1220,0 Size=700,1080 Collapsed=0 +[Window][Stats] +Pos=0,0 +Size=460,510 +Collapsed=0 + diff --git a/internal/playback/playlist.go b/internal/playback/playlist.go index e29afc0..5825dda 100644 --- a/internal/playback/playlist.go +++ b/internal/playback/playlist.go @@ -11,8 +11,34 @@ var ( ErrPlaylistEntryEmpty = errors.New("playlist entry must contain at least one feed") ErrPlaylistSyncFeedsRequired = errors.New("synchronized playlist entry requires both video and audio feeds") ErrPlaylistDurationNegative = errors.New("playlist entry duration cannot be negative") + ErrPlaylistFailurePolicy = errors.New("invalid playlist failure policy") ) +type PlaylistFailurePolicy uint8 + +const ( + PlaylistFailureWait PlaylistFailurePolicy = iota + PlaylistFailureNext +) + +func (p PlaylistFailurePolicy) String() string { + switch p { + case PlaylistFailureWait: + return "wait" + case PlaylistFailureNext: + return "next" + default: + return fmt.Sprintf("PlaylistFailurePolicy(%d)", uint8(p)) + } +} + +func (p PlaylistFailurePolicy) Validate() error { + if p != PlaylistFailureWait && p != PlaylistFailureNext { + return ErrPlaylistFailurePolicy + } + return nil +} + type PlaylistFeed struct { Domain string UUID string @@ -27,8 +53,9 @@ type PlaylistEntry struct { } type Playlist struct { - Entries []PlaylistEntry - Loop bool + Entries []PlaylistEntry + Loop bool + OnFailure PlaylistFailurePolicy } func (f PlaylistFeed) IsConfigured() bool { @@ -65,6 +92,9 @@ func (e PlaylistEntry) Validate() error { } func (p Playlist) Validate() error { + if err := p.OnFailure.Validate(); err != nil { + return err + } for index, entry := range p.Entries { if err := entry.Validate(); err != nil { return fmt.Errorf("playlist entry %d: %w", index, err) diff --git a/internal/playback/playlist_controller.go b/internal/playback/playlist_controller.go index faafd2b..8446bba 100644 --- a/internal/playback/playlist_controller.go +++ b/internal/playback/playlist_controller.go @@ -7,10 +7,22 @@ import ( "time" ) -type PlaylistReadiness struct { +type PlaylistEventKind uint8 + +const ( + PlaylistEventReady PlaylistEventKind = iota + PlaylistEventFailed +) + +type PlaylistEvent struct { Revision uint64 + Kind PlaylistEventKind } +// PlaylistReadiness is retained as an alias for callers that only publish +// ready events. Its zero Kind is PlaylistEventReady. +type PlaylistReadiness = PlaylistEvent + type playlistTimer interface { C() <-chan time.Time Stop() bool @@ -162,6 +174,40 @@ func (c *PlaylistController) Run( readiness = nil continue } + if ready.Kind == PlaylistEventFailed { + if ready.Revision != revision { + continue + } + stopTimer() + // A failed entry must not retain a live or apparently active + // duration clock, even when the policy is to wait. + timing = NewPlaylistTiming(revision, timing.Duration) + if c.playlist.OnFailure != PlaylistFailureNext { + c.publish(state, revision, timing) + continue + } + next, sessionCommand, apply, err := ApplyPlaylistSelection( + c.playlist, + state, + PlaylistCommand{Kind: PlaylistNext}, + c.retry, + ) + if err != nil || !apply { + c.publish(state, revision, timing) + continue + } + select { + case <-ctx.Done(): + return ctx.Err() + case c.sessions <- sessionCommand: + } + revision++ + state = next + entry, _ := next.Entry(c.playlist) + timing = NewPlaylistTiming(revision, entry.Duration) + c.publish(state, revision, timing) + continue + } if timing.Paused && ready.Revision == timing.Revision && timing.Duration > 0 && diff --git a/internal/playback/playlist_failure_test.go b/internal/playback/playlist_failure_test.go new file mode 100644 index 0000000..20269c5 --- /dev/null +++ b/internal/playback/playlist_failure_test.go @@ -0,0 +1,181 @@ +package playback + +import ( + "context" + "testing" + "time" +) + +func TestIsSessionFailed(t *testing.T) { + const generation = 7 + video := FeedConfig{Domain: "/video", UUID: "video", Active: true} + audio := FeedConfig{Domain: "/audio", UUID: "audio", Active: true} + failedFeed := func(unit Unit, feed FeedConfig) Status { + return Status{ + Unit: unit, State: StateFailed, Generation: generation, Feed: feed, + } + } + independent := SessionSnapshot{ + Generation: generation, + Plan: SessionPlan{ + Topology: TopologyIndependent, + Video: video, + Audio: audio, + }, + } + synchronized := SessionSnapshot{ + Generation: generation, + Plan: SessionPlan{ + Topology: TopologySynchronized, + Sync: SyncPairConfig{Video: video, Audio: audio}, + }, + } + + tests := []struct { + name string + session SessionSnapshot + statuses PlaybackStatusSnapshot + want bool + }{ + { + name: "video failure in independent pair", + session: independent, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Video: failedFeed(UnitVideo, video), HasVideo: true, + }, + want: true, + }, + { + name: "audio failure in independent pair", + session: independent, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Audio: failedFeed(UnitAudio, audio), HasAudio: true, + }, + want: true, + }, + { + name: "sync failure", + session: synchronized, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Sync: Status{ + Unit: UnitSync, State: StateFailed, Generation: generation, + Pair: SyncPairConfig{Video: video, Audio: audio}, + }, + HasSync: true, + }, + want: true, + }, + { + name: "stale snapshot generation", + session: independent, + statuses: PlaybackStatusSnapshot{ + Generation: generation - 1, + Video: failedFeed(UnitVideo, video), HasVideo: true, + }, + }, + { + name: "wrong source", + session: independent, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Video: failedFeed(UnitVideo, FeedConfig{ + Domain: "/video", UUID: "other", Active: true, + }), + HasVideo: true, + }, + }, + { + name: "inactive failed unit is ignored", + session: SessionSnapshot{ + Generation: generation, + Plan: SessionPlan{ + Topology: TopologyIndependent, + Video: video, + Audio: FeedConfig{Domain: audio.Domain, UUID: audio.UUID}, + }, + }, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Audio: failedFeed(UnitAudio, audio), HasAudio: true, + }, + }, + { + name: "reconnecting has not exhausted retries", + session: independent, + statuses: PlaybackStatusSnapshot{ + Generation: generation, + Video: Status{ + Unit: UnitVideo, State: StateReconnecting, + Generation: generation, Feed: video, + }, + HasVideo: true, + }, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + if got := IsSessionFailed(test.session, test.statuses); got != test.want { + t.Fatalf("IsSessionFailed() = %v, want %v", got, test.want) + } + }) + } +} + +func TestPlaylistControllerFailurePolicy(t *testing.T) { + for _, test := range []struct { + name string + policy PlaylistFailurePolicy + wantAdvance bool + }{ + {name: "wait", policy: PlaylistFailureWait}, + {name: "next", policy: PlaylistFailureNext, wantAdvance: true}, + } { + t.Run(test.name, func(t *testing.T) { + playlist := navigationPlaylist(false) + playlist.OnFailure = test.policy + sessions := make(chan SessionCommand, 4) + controller, err := NewPlaylistController( + playlist, validPlaylistRetryPolicy(), sessions, + ) + if err != nil { + t.Fatalf("NewPlaylistController() error = %v", err) + } + commands := make(chan PlaylistCommand, 2) + events := make(chan PlaylistReadiness, 2) + ctx, cancel := context.WithCancel(context.Background()) + result := make(chan error, 1) + go func() { result <- controller.Run(ctx, commands, events) }() + + commands <- PlaylistCommand{Kind: PlaylistSelect, Index: 0} + <-sessions + events <- PlaylistEvent{Revision: 1, Kind: PlaylistEventFailed} + + if test.wantAdvance { + select { + case <-sessions: + case <-time.After(time.Second): + t.Fatal("failure did not advance playlist") + } + snapshot, _ := controller.Snapshot() + if snapshot.State.CurrentIndex != 1 || snapshot.Revision != 2 { + t.Fatalf("snapshot = %#v, want index 1 revision 2", snapshot) + } + } else { + select { + case command := <-sessions: + t.Fatalf("unexpected session command: %#v", command) + case <-time.After(20 * time.Millisecond): + } + } + + cancel() + if err := <-result; err != context.Canceled { + t.Fatalf("Run() error = %v, want context canceled", err) + } + }) + } +} diff --git a/internal/playback/playlist_readiness.go b/internal/playback/playlist_readiness.go index bd8d336..6670ea5 100644 --- a/internal/playback/playlist_readiness.go +++ b/internal/playback/playlist_readiness.go @@ -89,7 +89,8 @@ func (c *PlaylistReadinessCoordinator) Run(ctx context.Context) error { ticker := c.newTicker(c.interval) defer ticker.Stop() - var emittedRevision uint64 + var emittedReadyRevision uint64 + var emittedFailedRevision uint64 for { select { case <-ctx.Done(): @@ -99,11 +100,7 @@ func (c *PlaylistReadinessCoordinator) Run(ctx context.Context) error { playlistSnapshot, ok := c.playlist.Snapshot() if !ok || !playlistSnapshot.State.HasSelection || - playlistSnapshot.Revision == 0 || - playlistSnapshot.Entry.Duration <= 0 || - playlistSnapshot.Timing.Started || - playlistSnapshot.Timing.Paused || - playlistSnapshot.Revision == emittedRevision { + playlistSnapshot.Revision == 0 { continue } @@ -114,7 +111,30 @@ func (c *PlaylistReadinessCoordinator) Run(ctx context.Context) error { ) { continue } - if !IsSessionPlaying(sessionSnapshot, c.statuses.SnapshotAll()) { + statuses := c.statuses.SnapshotAll() + if IsSessionFailed(sessionSnapshot, statuses) { + if playlistSnapshot.Revision == emittedFailedRevision { + continue + } + failed := PlaylistEvent{ + Revision: playlistSnapshot.Revision, + Kind: PlaylistEventFailed, + } + select { + case <-ctx.Done(): + return ctx.Err() + case c.output <- failed: + emittedFailedRevision = playlistSnapshot.Revision + } + continue + } + if playlistSnapshot.Entry.Duration <= 0 || + playlistSnapshot.Timing.Started || + playlistSnapshot.Timing.Paused || + playlistSnapshot.Revision == emittedReadyRevision { + continue + } + if !IsSessionPlaying(sessionSnapshot, statuses) { continue } @@ -123,7 +143,7 @@ func (c *PlaylistReadinessCoordinator) Run(ctx context.Context) error { case <-ctx.Done(): return ctx.Err() case c.output <- ready: - emittedRevision = playlistSnapshot.Revision + emittedReadyRevision = playlistSnapshot.Revision } } } @@ -193,6 +213,38 @@ func IsSessionPlaying( } } +func IsSessionFailed( + session SessionSnapshot, + statuses PlaybackStatusSnapshot, +) bool { + if statuses.Generation != session.Generation { + return false + } + + switch session.Plan.Topology { + case TopologyIndependent: + return (session.Plan.Video.Active && statusIsFailed( + statuses.Video, statuses.HasVideo, session.Generation, session.Plan.Video, + )) || (session.Plan.Audio.Active && statusIsFailed( + statuses.Audio, statuses.HasAudio, session.Generation, session.Plan.Audio, + )) + case TopologySynchronized: + return statuses.HasSync && + statuses.Sync.Generation == session.Generation && + statuses.Sync.State == StateFailed && + sameSyncSource(statuses.Sync.Pair, session.Plan.Sync) + default: + return false + } +} + +func statusIsFailed(status Status, present bool, generation uint64, feed FeedConfig) bool { + return present && + status.Generation == generation && + status.State == StateFailed && + sameFeedSource(status.Feed, feed) +} + func statusIsPlaying( status Status, present bool, diff --git a/playlists/sample-list-all.json b/playlists/sample-list-all.json index de250aa..6214a86 100644 --- a/playlists/sample-list-all.json +++ b/playlists/sample-list-all.json @@ -1,5 +1,6 @@ { "loop": true, + "on_failure": "next", "entries": [ { "name": "timelapse", diff --git a/playlists/sample-list.json b/playlists/sample-list.json index f6fd9e1..51c60e5 100644 --- a/playlists/sample-list.json +++ b/playlists/sample-list.json @@ -1,5 +1,6 @@ { "loop": true, + "on_failure": "next", "entries": [ { "name": "timelapse",