From aac8c664d8aad99a5fb6c71f450dab6076715f4f Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Mon, 31 Aug 2026 09:14:38 +0300 Subject: [PATCH] sync worker --- imgui.ini | 2 +- internal/playback/sync_worker.go | 179 ++++++++++++++++++ internal/playback/sync_worker_test.go | 261 ++++++++++++++++++++++++++ 3 files changed, 441 insertions(+), 1 deletion(-) create mode 100644 internal/playback/sync_worker.go create mode 100644 internal/playback/sync_worker_test.go diff --git a/imgui.ini b/imgui.ini index a8c20f3..dd9f37e 100644 --- a/imgui.ini +++ b/imgui.ini @@ -14,7 +14,7 @@ Size=200,200 Collapsed=0 [Window][Connection] -Pos=427,302 +Pos=177,324 Size=640,354 Collapsed=0 diff --git a/internal/playback/sync_worker.go b/internal/playback/sync_worker.go new file mode 100644 index 0000000..f058e98 --- /dev/null +++ b/internal/playback/sync_worker.go @@ -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 +} diff --git a/internal/playback/sync_worker_test.go b/internal/playback/sync_worker_test.go new file mode 100644 index 0000000..b7ff47d --- /dev/null +++ b/internal/playback/sync_worker_test.go @@ -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) + } + } +}