sync worker
This commit is contained in:
@@ -14,7 +14,7 @@ Size=200,200
|
|||||||
Collapsed=0
|
Collapsed=0
|
||||||
|
|
||||||
[Window][Connection]
|
[Window][Connection]
|
||||||
Pos=427,302
|
Pos=177,324
|
||||||
Size=640,354
|
Size=640,354
|
||||||
Collapsed=0
|
Collapsed=0
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,179 @@
|
|||||||
|
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
|
||||||
|
}
|
||||||
@@ -0,0 +1,261 @@
|
|||||||
|
package playback
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
type syncOpenResult struct {
|
||||||
|
reader SyncReader
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
type scriptedSyncFactory struct {
|
||||||
|
results []syncOpenResult
|
||||||
|
calls int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *scriptedSyncFactory) OpenSync(
|
||||||
|
context.Context,
|
||||||
|
FeedConfig,
|
||||||
|
FeedConfig,
|
||||||
|
) (SyncReader, error) {
|
||||||
|
if f.calls >= len(f.results) {
|
||||||
|
return nil, errors.New("unexpected sync open attempt")
|
||||||
|
}
|
||||||
|
result := f.results[f.calls]
|
||||||
|
f.calls++
|
||||||
|
return result.reader, result.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func activeSyncConfigs() (FeedConfig, FeedConfig) {
|
||||||
|
return FeedConfig{Domain: "/mxl", UUID: "video", Active: true},
|
||||||
|
FeedConfig{Domain: "/mxl", UUID: "audio", Active: true}
|
||||||
|
}
|
||||||
|
|
||||||
|
func newTestSyncWorker(
|
||||||
|
t *testing.T,
|
||||||
|
factory SyncReaderFactory,
|
||||||
|
videoSink VideoSink,
|
||||||
|
audioSink AudioSink,
|
||||||
|
maxAttempts int,
|
||||||
|
observer StatusObserver,
|
||||||
|
) *SyncWorker {
|
||||||
|
t.Helper()
|
||||||
|
worker, err := NewSyncWorker(
|
||||||
|
factory,
|
||||||
|
videoSink,
|
||||||
|
audioSink,
|
||||||
|
testRetryPolicy(maxAttempts),
|
||||||
|
func(error) bool { return true },
|
||||||
|
observer,
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewSyncWorker() error = %v", err)
|
||||||
|
}
|
||||||
|
worker.wait = func(context.Context, time.Duration) error { return nil }
|
||||||
|
return worker
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestNewSyncWorkerValidatesDependencies(t *testing.T) {
|
||||||
|
factory := &scriptedSyncFactory{}
|
||||||
|
videoSink := &fakeVideoSink{}
|
||||||
|
audioSink := &fakeAudioSink{}
|
||||||
|
retry := testRetryPolicy(3)
|
||||||
|
decide := func(error) bool { return true }
|
||||||
|
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
factory SyncReaderFactory
|
||||||
|
videoSink VideoSink
|
||||||
|
audioSink AudioSink
|
||||||
|
retry RetryPolicy
|
||||||
|
shouldRetry func(error) bool
|
||||||
|
wantErr error
|
||||||
|
}{
|
||||||
|
{"missing factory", nil, videoSink, audioSink, retry, decide, ErrSyncFactoryRequired},
|
||||||
|
{"missing video sink", factory, nil, audioSink, retry, decide, ErrSyncVideoSinkRequired},
|
||||||
|
{"missing audio sink", factory, videoSink, nil, retry, decide, ErrSyncAudioSinkRequired},
|
||||||
|
{"invalid retry", factory, videoSink, audioSink, RetryPolicy{}, decide, ErrInvalidRetryDelay},
|
||||||
|
{"missing decider", factory, videoSink, audioSink, retry, nil, ErrSyncRetryDeciderRequired},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
worker, err := NewSyncWorker(
|
||||||
|
tt.factory, tt.videoSink, tt.audioSink, tt.retry, tt.shouldRetry, nil,
|
||||||
|
)
|
||||||
|
if worker != nil {
|
||||||
|
t.Fatal("NewSyncWorker() worker is not nil")
|
||||||
|
}
|
||||||
|
if !errors.Is(err, tt.wantErr) {
|
||||||
|
t.Fatalf("NewSyncWorker() error = %v, want %v", err, tt.wantErr)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncWorkerRejectsInvalidOrInactiveFeeds(t *testing.T) {
|
||||||
|
video, audio := activeSyncConfigs()
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
video FeedConfig
|
||||||
|
audio FeedConfig
|
||||||
|
want error
|
||||||
|
}{
|
||||||
|
{"invalid video", FeedConfig{Active: true}, audio, ErrActiveFeedNotConfigured},
|
||||||
|
{"invalid audio", video, FeedConfig{Active: true}, ErrActiveFeedNotConfigured},
|
||||||
|
{"inactive video", FeedConfig{Domain: video.Domain, UUID: video.UUID}, audio, ErrSyncFeedsInactive},
|
||||||
|
{"inactive audio", video, FeedConfig{Domain: audio.Domain, UUID: audio.UUID}, ErrSyncFeedsInactive},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
factory := &scriptedSyncFactory{}
|
||||||
|
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 3, nil)
|
||||||
|
err := worker.Run(context.Background(), tt.video, tt.audio)
|
||||||
|
if !errors.Is(err, tt.want) {
|
||||||
|
t.Fatalf("Run() error = %v, want %v", err, tt.want)
|
||||||
|
}
|
||||||
|
if factory.calls != 0 {
|
||||||
|
t.Fatalf("factory calls = %d, want 0", factory.calls)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncWorkerExhaustsOpenRetries(t *testing.T) {
|
||||||
|
openErr := errors.New("sync producer unavailable")
|
||||||
|
factory := &scriptedSyncFactory{results: []syncOpenResult{{err: openErr}, {err: openErr}}}
|
||||||
|
var statuses []Status
|
||||||
|
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 2,
|
||||||
|
func(status Status) { statuses = append(statuses, status) })
|
||||||
|
video, audio := activeSyncConfigs()
|
||||||
|
|
||||||
|
err := worker.Run(context.Background(), video, audio)
|
||||||
|
if !errors.Is(err, openErr) {
|
||||||
|
t.Fatalf("Run() error = %v, want %v", err, openErr)
|
||||||
|
}
|
||||||
|
if factory.calls != 2 {
|
||||||
|
t.Fatalf("factory calls = %d, want 2", factory.calls)
|
||||||
|
}
|
||||||
|
want := []State{StateConnecting, StateReconnecting, StateReconnecting, StateFailed}
|
||||||
|
if len(statuses) != len(want) {
|
||||||
|
t.Fatalf("statuses = %+v, want %d entries", statuses, len(want))
|
||||||
|
}
|
||||||
|
for i, state := range want {
|
||||||
|
if statuses[i].Unit != UnitSync || statuses[i].State != state {
|
||||||
|
t.Errorf("status %d = %+v, want unit=%v state=%v", i, statuses[i], UnitSync, state)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncWorkerStablePairResetsRetryCounter(t *testing.T) {
|
||||||
|
readErr := errors.New("sync disconnected")
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
lastReader := &fakeSyncReader{read: func(ctx context.Context) (SyncFrame, error) {
|
||||||
|
cancel()
|
||||||
|
return SyncFrame{}, ctx.Err()
|
||||||
|
}}
|
||||||
|
factory := &scriptedSyncFactory{results: []syncOpenResult{
|
||||||
|
{reader: &fakeSyncReader{frames: []SyncFrame{{}}, readErr: readErr}},
|
||||||
|
{reader: &fakeSyncReader{frames: []SyncFrame{{}}, readErr: readErr}},
|
||||||
|
{reader: lastReader},
|
||||||
|
}}
|
||||||
|
var statuses []Status
|
||||||
|
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 2,
|
||||||
|
func(status Status) { statuses = append(statuses, status) })
|
||||||
|
video, audio := activeSyncConfigs()
|
||||||
|
|
||||||
|
err := worker.Run(ctx, video, audio)
|
||||||
|
if !errors.Is(err, context.Canceled) {
|
||||||
|
t.Fatalf("Run() error = %v, want context.Canceled", err)
|
||||||
|
}
|
||||||
|
if factory.calls != 3 {
|
||||||
|
t.Fatalf("factory calls = %d, want 3", factory.calls)
|
||||||
|
}
|
||||||
|
var retryFailures []int
|
||||||
|
for _, status := range statuses {
|
||||||
|
if status.State == StateReconnecting && status.RetryIn > 0 {
|
||||||
|
retryFailures = append(retryFailures, status.FailedAttempts)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(retryFailures) != 2 || retryFailures[0] != 1 || retryFailures[1] != 1 {
|
||||||
|
t.Fatalf("retry failure counts = %v, want [1 1]", retryFailures)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncWorkerDoesNotRetrySinkFailures(t *testing.T) {
|
||||||
|
sinkErr := errors.New("output failed")
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
videoSink VideoSink
|
||||||
|
audioSink AudioSink
|
||||||
|
}{
|
||||||
|
{"video", &fakeVideoSink{err: sinkErr}, &fakeAudioSink{}},
|
||||||
|
{"audio", &fakeVideoSink{}, &fakeAudioSink{err: sinkErr}},
|
||||||
|
}
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
reader := &fakeSyncReader{frames: []SyncFrame{{}}}
|
||||||
|
factory := &scriptedSyncFactory{results: []syncOpenResult{{reader: reader}}}
|
||||||
|
deciderCalls := 0
|
||||||
|
worker, err := NewSyncWorker(factory, tt.videoSink, tt.audioSink, testRetryPolicy(0),
|
||||||
|
func(error) bool { deciderCalls++; return true }, nil)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
worker.wait = func(context.Context, time.Duration) error { return nil }
|
||||||
|
video, audio := activeSyncConfigs()
|
||||||
|
err = worker.Run(context.Background(), video, audio)
|
||||||
|
if !errors.Is(err, sinkErr) {
|
||||||
|
t.Fatalf("Run() error = %v, want %v", err, sinkErr)
|
||||||
|
}
|
||||||
|
if factory.calls != 1 || deciderCalls != 0 || !reader.closed {
|
||||||
|
t.Fatalf("calls=%d deciderCalls=%d closed=%t", factory.calls, deciderCalls, reader.closed)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSyncWorkerEmitsPlayingThenStopsOnCancellation(t *testing.T) {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
reader := &fakeSyncReader{
|
||||||
|
frames: []SyncFrame{{}},
|
||||||
|
read: func(ctx context.Context) (SyncFrame, error) {
|
||||||
|
cancel()
|
||||||
|
return SyncFrame{}, ctx.Err()
|
||||||
|
},
|
||||||
|
}
|
||||||
|
// Preserve the first frame before switching to the cancellation callback.
|
||||||
|
readCalls := 0
|
||||||
|
reader.read = func(ctx context.Context) (SyncFrame, error) {
|
||||||
|
readCalls++
|
||||||
|
if readCalls == 1 {
|
||||||
|
return SyncFrame{}, nil
|
||||||
|
}
|
||||||
|
cancel()
|
||||||
|
return SyncFrame{}, ctx.Err()
|
||||||
|
}
|
||||||
|
factory := &scriptedSyncFactory{results: []syncOpenResult{{reader: reader}}}
|
||||||
|
var statuses []Status
|
||||||
|
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 3,
|
||||||
|
func(status Status) { statuses = append(statuses, status) })
|
||||||
|
video, audio := activeSyncConfigs()
|
||||||
|
|
||||||
|
err := worker.Run(ctx, video, audio)
|
||||||
|
if !errors.Is(err, context.Canceled) {
|
||||||
|
t.Fatalf("Run() error = %v, want context.Canceled", err)
|
||||||
|
}
|
||||||
|
want := []State{StateConnecting, StatePlaying, StateStopping, StateIdle}
|
||||||
|
if len(statuses) != len(want) {
|
||||||
|
t.Fatalf("statuses = %+v, want %v", statuses, want)
|
||||||
|
}
|
||||||
|
for i, state := range want {
|
||||||
|
if statuses[i].Unit != UnitSync || statuses[i].State != state {
|
||||||
|
t.Errorf("status %d = %+v, want unit=%v state=%v", i, statuses[i], UnitSync, state)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user