216 lines
5.2 KiB
Go
216 lines
5.2 KiB
Go
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 emittedRevision 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 ||
|
|
playlistSnapshot.Entry.Duration <= 0 ||
|
|
playlistSnapshot.Timing.Started ||
|
|
playlistSnapshot.Timing.Paused ||
|
|
playlistSnapshot.Revision == emittedRevision {
|
|
continue
|
|
}
|
|
|
|
sessionSnapshot, ok := c.session.Snapshot()
|
|
if !ok || !PlaylistEntryMatchesSession(
|
|
playlistSnapshot.Entry,
|
|
sessionSnapshot.Desired,
|
|
) {
|
|
continue
|
|
}
|
|
if !IsSessionPlaying(sessionSnapshot, c.statuses.SnapshotAll()) {
|
|
continue
|
|
}
|
|
|
|
ready := PlaylistReadiness{Revision: playlistSnapshot.Revision}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case c.output <- ready:
|
|
emittedRevision = 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 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)
|
|
}
|