playlist retry-exhaustion handling
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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 &&
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user