package playback import ( "context" "errors" "fmt" ) var ( ErrSyncFactoryRequired = errors.New("sync reader factory is required") ErrSyncVideoSinkRequired = errors.New("sync video sink is required") ErrSyncAudioSinkRequired = errors.New("sync audio sink is required") ErrSyncRetryDeciderRequired = errors.New("sync retry decider is required") ErrSyncFeedsInactive = errors.New("both sync feeds must be active") ) type SyncWorker struct { factory SyncReaderFactory videoSink VideoSink audioSink AudioSink retry RetryPolicy shouldRetry retryDecider observer StatusObserver wait waitFunc } func NewSyncWorker( factory SyncReaderFactory, videoSink VideoSink, audioSink AudioSink, retry RetryPolicy, shouldRetry func(error) bool, observer StatusObserver, ) (*SyncWorker, error) { if factory == nil { return nil, ErrSyncFactoryRequired } if videoSink == nil { return nil, ErrSyncVideoSinkRequired } if audioSink == nil { return nil, ErrSyncAudioSinkRequired } if err := retry.Validate(); err != nil { return nil, fmt.Errorf("validate sync retry policy: %w", err) } if shouldRetry == nil { return nil, ErrSyncRetryDeciderRequired } return &SyncWorker{ factory: factory, videoSink: videoSink, audioSink: audioSink, retry: retry, shouldRetry: shouldRetry, observer: observer, wait: waitForRetry, }, nil } func (w *SyncWorker) emit(ctx context.Context, status Status) { status.Generation = generationFromContext(ctx) if w.observer != nil { w.observer(status) } } func (w *SyncWorker) Run( ctx context.Context, videoConfig FeedConfig, audioConfig FeedConfig, ) error { if err := videoConfig.Validate(); err != nil { return fmt.Errorf("validate sync video config: %w", err) } if err := audioConfig.Validate(); err != nil { return fmt.Errorf("validate sync audio config: %w", err) } if !videoConfig.Active || !audioConfig.Active { return ErrSyncFeedsInactive } attemptNumber := 0 var latestRetry retryEvent attempt := func(ctx context.Context) (bool, error) { attemptNumber++ state := StateConnecting if attemptNumber > 1 { state = StateReconnecting } w.emit(ctx, Status{ Unit: UnitSync, State: state, Attempt: attemptNumber, }) attemptAudioSink := &stabilityAudioSink{ sink: w.audioSink, onStable: func() { w.emit(ctx, Status{ Unit: UnitSync, State: StatePlaying, Attempt: attemptNumber, }) }, } err := runSyncAttempt( ctx, w.factory, w.videoSink, attemptAudioSink, videoConfig, audioConfig, ) return attemptAudioSink.stable, err } decide := func(err error) bool { var videoErr *videoSinkError if errors.As(err, &videoErr) { return false } var audioErr *audioSinkError if errors.As(err, &audioErr) { return false } return w.shouldRetry(err) } observeRetry := func(event retryEvent) { latestRetry = event if !event.WillRetry { return } w.emit(ctx, Status{ Unit: UnitSync, State: StateReconnecting, Attempt: attemptNumber + 1, FailedAttempts: event.FailedAttempts, RetryIn: event.RetryIn, Err: event.Err, }) } err := runWithRetry( ctx, w.retry, attempt, decide, w.wait, observeRetry, ) if ctx.Err() != nil { w.emit(ctx, Status{ Unit: UnitSync, State: StateStopping, }) w.emit(ctx, Status{ Unit: UnitSync, State: StateIdle, }) return ctx.Err() } if err != nil { w.emit(ctx, Status{ Unit: UnitSync, State: StateFailed, Attempt: attemptNumber, FailedAttempts: latestRetry.FailedAttempts, Err: err, }) return err } w.emit(ctx, Status{ Unit: UnitSync, State: StateIdle, }) return nil }