package playback import ( "context" "errors" "time" ) type PlaylistSnapshotSource interface { Snapshot() (PlaylistSnapshot, bool) } type SessionSnapshotSource interface { Snapshot() (SessionSnapshot, bool) } type PlaybackStatusSnapshotSource interface { SnapshotAll() PlaybackStatusSnapshot } type playlistReadinessTicker interface { C() <-chan time.Time Stop() } type playlistReadinessTickerFactory func(time.Duration) playlistReadinessTicker type realPlaylistReadinessTicker struct { ticker *time.Ticker } func (t realPlaylistReadinessTicker) C() <-chan time.Time { return t.ticker.C } func (t realPlaylistReadinessTicker) Stop() { t.ticker.Stop() } var ( ErrPlaylistSnapshotSourceRequired = errors.New("playlist snapshot source is required") ErrSessionSnapshotSourceRequired = errors.New("session snapshot source is required") ErrStatusSnapshotSourceRequired = errors.New("playback status snapshot source is required") ErrPlaylistReadinessOutputRequired = errors.New("playlist readiness output channel is required") ErrPlaylistReadinessInterval = errors.New("playlist readiness interval must be positive") ) type PlaylistReadinessCoordinator struct { playlist PlaylistSnapshotSource session SessionSnapshotSource statuses PlaybackStatusSnapshotSource output chan<- PlaylistReadiness interval time.Duration newTicker playlistReadinessTickerFactory } func NewPlaylistReadinessCoordinator( playlist PlaylistSnapshotSource, session SessionSnapshotSource, statuses PlaybackStatusSnapshotSource, output chan<- PlaylistReadiness, interval time.Duration, ) (*PlaylistReadinessCoordinator, error) { if playlist == nil { return nil, ErrPlaylistSnapshotSourceRequired } if session == nil { return nil, ErrSessionSnapshotSourceRequired } if statuses == nil { return nil, ErrStatusSnapshotSourceRequired } if output == nil { return nil, ErrPlaylistReadinessOutputRequired } if interval <= 0 { return nil, ErrPlaylistReadinessInterval } return &PlaylistReadinessCoordinator{ playlist: playlist, session: session, statuses: statuses, output: output, interval: interval, newTicker: func(interval time.Duration) playlistReadinessTicker { return realPlaylistReadinessTicker{ticker: time.NewTicker(interval)} }, }, nil } func (c *PlaylistReadinessCoordinator) Run(ctx context.Context) error { ticker := c.newTicker(c.interval) defer ticker.Stop() var emittedReadyRevision uint64 var emittedFailedRevision uint64 for { select { case <-ctx.Done(): return ctx.Err() case <-ticker.C(): playlistSnapshot, ok := c.playlist.Snapshot() if !ok || !playlistSnapshot.State.HasSelection || playlistSnapshot.Revision == 0 { continue } sessionSnapshot, ok := c.session.Snapshot() if !ok || !PlaylistEntryMatchesSession( playlistSnapshot.Entry, sessionSnapshot.Desired, ) { continue } statuses := c.statuses.SnapshotAll() if failure, failed := SessionFailureStatus(sessionSnapshot, statuses); failed { if playlistSnapshot.Revision == emittedFailedRevision { continue } failed := PlaylistEvent{ Revision: playlistSnapshot.Revision, Kind: PlaylistEventFailed, Failure: failure, } 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 } ready := PlaylistReadiness{Revision: playlistSnapshot.Revision} select { case <-ctx.Done(): return ctx.Err() case c.output <- ready: emittedReadyRevision = playlistSnapshot.Revision } } } } func PlaylistEntryMatchesSession(entry PlaylistEntry, session SessionConfig) bool { if err := entry.Validate(); err != nil { return false } return playlistFeedMatchesSession(entry.Video, session.Video) && playlistFeedMatchesSession(entry.Audio, session.Audio) && entry.SyncRequested == session.SyncRequested } func playlistFeedMatchesSession(playlist PlaylistFeed, session FeedConfig) bool { if !playlist.IsConfigured() { return !session.IsConfigured() && !session.Active } return session.Active && playlist.Domain == session.Domain && playlist.UUID == session.UUID } func IsSessionPlaying( session SessionSnapshot, statuses PlaybackStatusSnapshot, ) bool { if statuses.Generation != session.Generation { return false } switch session.Plan.Topology { case TopologyIndependent: hasActiveFeed := session.Plan.Video.Active || session.Plan.Audio.Active if !hasActiveFeed { return false } if session.Plan.Video.Active && !statusIsPlaying( statuses.Video, statuses.HasVideo, session.Generation, session.Plan.Video, ) { return false } if session.Plan.Audio.Active && !statusIsPlaying( statuses.Audio, statuses.HasAudio, session.Generation, session.Plan.Audio, ) { return false } return true case TopologySynchronized: return statuses.HasSync && statuses.Sync.Generation == session.Generation && statuses.Sync.State == StatePlaying && sameSyncSource(statuses.Sync.Pair, session.Plan.Sync) case TopologyIdle: return false default: return false } } func IsSessionFailed( session SessionSnapshot, statuses PlaybackStatusSnapshot, ) bool { _, failed := SessionFailureStatus(session, statuses) return failed } func SessionFailureStatus( session SessionSnapshot, statuses PlaybackStatusSnapshot, ) (Status, bool) { if statuses.Generation != session.Generation { return Status{}, false } switch session.Plan.Topology { case TopologyIndependent: if session.Plan.Video.Active && statusIsFailed( statuses.Video, statuses.HasVideo, session.Generation, session.Plan.Video, ) { return statuses.Video, true } if session.Plan.Audio.Active && statusIsFailed( statuses.Audio, statuses.HasAudio, session.Generation, session.Plan.Audio, ) { return statuses.Audio, true } return Status{}, false case TopologySynchronized: failed := statuses.HasSync && statuses.Sync.Generation == session.Generation && statuses.Sync.State == StateFailed && sameSyncSource(statuses.Sync.Pair, session.Plan.Sync) return statuses.Sync, failed default: return Status{}, 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, generation uint64, feed FeedConfig, ) bool { return present && status.Generation == generation && status.State == StatePlaying && sameFeedSource(status.Feed, feed) } func sameFeedSource(a, b FeedConfig) bool { return a.Domain == b.Domain && a.UUID == b.UUID } func sameSyncSource(a, b SyncPairConfig) bool { return sameFeedSource(a.Video, b.Video) && sameFeedSource(a.Audio, b.Audio) }