187 lines
3.9 KiB
Go
187 lines
3.9 KiB
Go
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,
|
|
pair SyncPairConfig,
|
|
status Status,
|
|
) {
|
|
status.Generation = generationFromContext(ctx)
|
|
status.Pair = pair
|
|
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
|
|
}
|
|
pair := SyncPairConfig{Video: videoConfig, Audio: audioConfig}
|
|
|
|
attemptNumber := 0
|
|
var latestRetry retryEvent
|
|
|
|
attempt := func(ctx context.Context) (bool, error) {
|
|
attemptNumber++
|
|
|
|
state := StateConnecting
|
|
if attemptNumber > 1 {
|
|
state = StateReconnecting
|
|
}
|
|
w.emit(ctx, pair, Status{
|
|
Unit: UnitSync,
|
|
State: state,
|
|
Attempt: attemptNumber,
|
|
})
|
|
|
|
attemptAudioSink := &stabilityAudioSink{
|
|
sink: w.audioSink,
|
|
onStable: func() {
|
|
w.emit(ctx, pair, 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, pair, 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, pair, Status{
|
|
Unit: UnitSync,
|
|
State: StateStopping,
|
|
})
|
|
w.emit(ctx, pair, Status{
|
|
Unit: UnitSync,
|
|
State: StateIdle,
|
|
})
|
|
return ctx.Err()
|
|
}
|
|
if err != nil {
|
|
w.emit(ctx, pair, Status{
|
|
Unit: UnitSync,
|
|
State: StateFailed,
|
|
Attempt: attemptNumber,
|
|
FailedAttempts: latestRetry.FailedAttempts,
|
|
Err: err,
|
|
})
|
|
return err
|
|
}
|
|
w.emit(ctx, pair, Status{
|
|
Unit: UnitSync,
|
|
State: StateIdle,
|
|
})
|
|
return nil
|
|
}
|