From 05f356fa1080207ef85136f158182998dc7997bb Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Mon, 31 Aug 2026 08:27:37 +0300 Subject: [PATCH] add synchronized playback attempt --- internal/playback/sync_attempt.go | 56 ++++++ internal/playback/sync_attempt_test.go | 241 +++++++++++++++++++++++++ 2 files changed, 297 insertions(+) create mode 100644 internal/playback/sync_attempt.go create mode 100644 internal/playback/sync_attempt_test.go diff --git a/internal/playback/sync_attempt.go b/internal/playback/sync_attempt.go new file mode 100644 index 0000000..995fc67 --- /dev/null +++ b/internal/playback/sync_attempt.go @@ -0,0 +1,56 @@ +package playback + +import ( + "context" + "errors" + "fmt" +) + +func runSyncAttempt( + ctx context.Context, + factory SyncReaderFactory, + videoSink VideoSink, + audioSink AudioSink, + videoConfig FeedConfig, + audioConfig FeedConfig, +) (resultErr error) { + reader, err := factory.OpenSync( + ctx, + videoConfig, + audioConfig, + ) + if err != nil { + return fmt.Errorf("open sync group: %w", err) + } + + defer func() { + if closeErr := reader.Close(); closeErr != nil { + closeErr = fmt.Errorf("close sync group: %w", closeErr) + resultErr = errors.Join(resultErr, closeErr) + } + }() + + for { + frame, err := reader.ReadSync(ctx) + if err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + return fmt.Errorf("read sync group: %w", err) + } + + if err := videoSink.ConsumeVideo(ctx, frame.Video); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + return &videoSinkError{err: err} + } + + if err := audioSink.ConsumeAudio(ctx, frame.Audio); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + return &audioSinkError{err: err} + } + } +} diff --git a/internal/playback/sync_attempt_test.go b/internal/playback/sync_attempt_test.go new file mode 100644 index 0000000..4a9e1b7 --- /dev/null +++ b/internal/playback/sync_attempt_test.go @@ -0,0 +1,241 @@ +package playback + +import ( + "context" + "errors" + "testing" +) + +type fakeSyncFactory struct { + reader SyncReader + err error + calls int + videoConfig FeedConfig + audioConfig FeedConfig +} + +func (f *fakeSyncFactory) OpenSync( + _ context.Context, + videoConfig FeedConfig, + audioConfig FeedConfig, +) (SyncReader, error) { + f.calls++ + f.videoConfig = videoConfig + f.audioConfig = audioConfig + return f.reader, f.err +} + +type fakeSyncReader struct { + frames []SyncFrame + readErr error + closeErr error + readCalls int + closed bool + read func(context.Context) (SyncFrame, error) +} + +func (r *fakeSyncReader) ReadSync(ctx context.Context) (SyncFrame, error) { + r.readCalls++ + if r.read != nil { + return r.read(ctx) + } + if len(r.frames) == 0 { + return SyncFrame{}, r.readErr + } + frame := r.frames[0] + r.frames = r.frames[1:] + return frame, nil +} + +func (r *fakeSyncReader) Close() error { + r.closed = true + return r.closeErr +} + +type orderedVideoSink struct { + order *[]string + err error + frame VideoFrame +} + +func (s *orderedVideoSink) ConsumeVideo(_ context.Context, frame VideoFrame) error { + *s.order = append(*s.order, "video") + s.frame = frame + return s.err +} + +type orderedAudioSink struct { + order *[]string + err error + frame AudioFrame +} + +func (s *orderedAudioSink) ConsumeAudio(_ context.Context, frame AudioFrame) error { + *s.order = append(*s.order, "audio") + s.frame = frame + return s.err +} + +func TestRunSyncAttemptOpenFailure(t *testing.T) { + openErr := errors.New("open failed") + factory := &fakeSyncFactory{err: openErr} + videoConfig := FeedConfig{Domain: "/mxl", UUID: "video", Active: true} + audioConfig := FeedConfig{Domain: "/mxl", UUID: "audio", Active: true} + + err := runSyncAttempt( + context.Background(), + factory, + &fakeVideoSink{}, + &fakeAudioSink{}, + videoConfig, + audioConfig, + ) + + if !errors.Is(err, openErr) { + t.Fatalf("runSyncAttempt() error = %v, want %v", err, openErr) + } + if factory.calls != 1 || factory.videoConfig != videoConfig || factory.audioConfig != audioConfig { + t.Fatalf( + "factory call = %d, video %#v, audio %#v", + factory.calls, + factory.videoConfig, + factory.audioConfig, + ) + } +} + +func TestRunSyncAttemptConsumesBorrowedPairInOrderThenReturnsReadError(t *testing.T) { + readErr := errors.New("sync read failed") + videoPayload := []byte{1, 2, 3, 4} + audioSamples := [][]byte{{5, 6, 7, 8}} + want := SyncFrame{ + Video: VideoFrame{Index: 10, Payload: videoPayload}, + Audio: AudioFrame{Index: 20, Samples: audioSamples}, + } + reader := &fakeSyncReader{frames: []SyncFrame{want}, readErr: readErr} + var order []string + videoSink := &orderedVideoSink{order: &order} + audioSink := &orderedAudioSink{order: &order} + + err := runSyncAttempt( + context.Background(), + &fakeSyncFactory{reader: reader}, + videoSink, + audioSink, + FeedConfig{}, + FeedConfig{}, + ) + + if !errors.Is(err, readErr) { + t.Fatalf("runSyncAttempt() error = %v, want %v", err, readErr) + } + if len(order) != 2 || order[0] != "video" || order[1] != "audio" { + t.Fatalf("sink order = %v, want [video audio]", order) + } + if &videoSink.frame.Payload[0] != &videoPayload[0] { + t.Fatal("video payload was copied") + } + if &audioSink.frame.Samples[0][0] != &audioSamples[0][0] { + t.Fatal("audio samples were copied") + } + if reader.readCalls != 2 || !reader.closed { + t.Fatalf("reader calls = %d, closed = %t; want 2, true", reader.readCalls, reader.closed) + } +} + +func TestRunSyncAttemptVideoSinkFailureSkipsAudio(t *testing.T) { + sinkErr := errors.New("video output failed") + reader := &fakeSyncReader{frames: []SyncFrame{{}}} + var order []string + + err := runSyncAttempt( + context.Background(), + &fakeSyncFactory{reader: reader}, + &orderedVideoSink{order: &order, err: sinkErr}, + &orderedAudioSink{order: &order}, + FeedConfig{}, + FeedConfig{}, + ) + + if !errors.Is(err, sinkErr) { + t.Fatalf("runSyncAttempt() error = %v, want %v", err, sinkErr) + } + var typedErr *videoSinkError + if !errors.As(err, &typedErr) { + t.Fatalf("runSyncAttempt() error type = %T, want *videoSinkError", err) + } + if len(order) != 1 || order[0] != "video" { + t.Fatalf("sink order = %v, want [video]", order) + } +} + +func TestRunSyncAttemptAudioSinkFailureFollowsVideo(t *testing.T) { + sinkErr := errors.New("audio output failed") + reader := &fakeSyncReader{frames: []SyncFrame{{}}} + var order []string + + err := runSyncAttempt( + context.Background(), + &fakeSyncFactory{reader: reader}, + &orderedVideoSink{order: &order}, + &orderedAudioSink{order: &order, err: sinkErr}, + FeedConfig{}, + FeedConfig{}, + ) + + if !errors.Is(err, sinkErr) { + t.Fatalf("runSyncAttempt() error = %v, want %v", err, sinkErr) + } + var typedErr *audioSinkError + if !errors.As(err, &typedErr) { + t.Fatalf("runSyncAttempt() error type = %T, want *audioSinkError", err) + } + if len(order) != 2 || order[0] != "video" || order[1] != "audio" { + t.Fatalf("sink order = %v, want [video audio]", order) + } +} + +func TestRunSyncAttemptCancellationClosesReader(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + reader := &fakeSyncReader{ + read: func(ctx context.Context) (SyncFrame, error) { + cancel() + return SyncFrame{}, ctx.Err() + }, + } + + err := runSyncAttempt( + ctx, + &fakeSyncFactory{reader: reader}, + &fakeVideoSink{}, + &fakeAudioSink{}, + FeedConfig{}, + FeedConfig{}, + ) + + if !errors.Is(err, context.Canceled) { + t.Fatalf("runSyncAttempt() error = %v, want context.Canceled", err) + } + if !reader.closed { + t.Fatal("reader was not closed") + } +} + +func TestRunSyncAttemptJoinsReadAndCloseErrors(t *testing.T) { + readErr := errors.New("read failed") + closeErr := errors.New("close failed") + reader := &fakeSyncReader{readErr: readErr, closeErr: closeErr} + + err := runSyncAttempt( + context.Background(), + &fakeSyncFactory{reader: reader}, + &fakeVideoSink{}, + &fakeAudioSink{}, + FeedConfig{}, + FeedConfig{}, + ) + + if !errors.Is(err, readErr) || !errors.Is(err, closeErr) { + t.Fatalf("runSyncAttempt() error = %v, want read and close errors", err) + } +}