Files
go-mxl-player/internal/playback/sync_worker.go
T
Dmitry Sergeev aac8c664d8 sync worker
2026-08-31 09:14:38 +03:00

180 lines
3.6 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(status Status) {
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(Status{
Unit: UnitSync,
State: state,
Attempt: attemptNumber,
})
attemptAudioSink := &stabilityAudioSink{
sink: w.audioSink,
onStable: func() {
w.emit(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(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(Status{
Unit: UnitSync,
State: StateStopping,
})
w.emit(Status{
Unit: UnitSync,
State: StateIdle,
})
return ctx.Err()
}
if err != nil {
w.emit(Status{
Unit: UnitSync,
State: StateFailed,
Attempt: attemptNumber,
FailedAttempts: latestRetry.FailedAttempts,
Err: err,
})
return err
}
w.emit(Status{
Unit: UnitSync,
State: StateIdle,
})
return nil
}