diff --git a/imgui.ini b/imgui.ini index dd9f37e..a71e8e8 100644 --- a/imgui.ini +++ b/imgui.ini @@ -14,7 +14,7 @@ Size=200,200 Collapsed=0 [Window][Connection] -Pos=177,324 -Size=640,354 +Pos=911,549 +Size=640,352 Collapsed=0 diff --git a/internal/playback/sync_slot.go b/internal/playback/sync_slot.go new file mode 100644 index 0000000..3f4bc06 --- /dev/null +++ b/internal/playback/sync_slot.go @@ -0,0 +1,122 @@ +package playback + +import ( + "context" + "errors" + "fmt" +) + +type SyncPairConfig struct { + Video FeedConfig + Audio FeedConfig +} + +var ( + ErrSyncWorkerRequired = errors.New("sync worker is required") + ErrSyncActivityMismatch = errors.New( + "synchronized video and audio must have matching active states", + ) +) + +func (c SyncPairConfig) Validate() error { + if err := c.Video.Validate(); err != nil { + return fmt.Errorf("video: %w", err) + } + if err := c.Audio.Validate(); err != nil { + return fmt.Errorf("audio: %w", err) + } + if c.Video.Active != c.Audio.Active { + return ErrSyncActivityMismatch + } + return nil +} + +func (c SyncPairConfig) Active() bool { + return c.Video.Active && c.Audio.Active +} + +type SyncSlot struct { + worker *SyncWorker +} + +func NewSyncSlot(worker *SyncWorker) (*SyncSlot, error) { + if worker == nil { + return nil, ErrSyncWorkerRequired + } + return &SyncSlot{ + worker: worker, + }, nil +} + +func (s *SyncSlot) Run( + ctx context.Context, + initial SyncPairConfig, + commands <-chan SyncPairConfig, +) error { + if err := initial.Validate(); err != nil { + return fmt.Errorf("validate initial sync config: %w", err) + } + + var ( + workerCancel context.CancelFunc + workerDone chan error + ) + + start := func(config SyncPairConfig) { + workerCtx, cancel := context.WithCancel(ctx) + done := make(chan error, 1) + + workerCancel = cancel + workerDone = done + + go func() { + done <- s.worker.Run(workerCtx, config.Video, config.Audio) + }() + } + + stop := func() { + if workerCancel == nil { + return + } + + workerCancel() + <-workerDone + + workerCancel = nil + workerDone = nil + } + + if initial.Active() { + start(initial) + } + + for { + select { + case <-ctx.Done(): + stop() + return ctx.Err() + + case config, ok := <-commands: + if !ok { + stop() + return nil + } + + if err := config.Validate(); err != nil { + // Ignore invalid commands without disturbing the current worker. + continue + } + + stop() + if config.Active() { + start(config) + } + + case <-workerDone: + // The worker stopped naturally or exhausted its retries. + workerCancel() + workerCancel = nil + workerDone = nil + } + } +} diff --git a/internal/playback/sync_slot_test.go b/internal/playback/sync_slot_test.go new file mode 100644 index 0000000..b8e9129 --- /dev/null +++ b/internal/playback/sync_slot_test.go @@ -0,0 +1,231 @@ +package playback + +import ( + "context" + "errors" + "sync" + "testing" + "time" +) + +type slotSyncFactory struct { + opened chan SyncPairConfig + + mu sync.Mutex + active int + maxActive int + closeCount int +} + +func newSlotSyncFactory() *slotSyncFactory { + return &slotSyncFactory{opened: make(chan SyncPairConfig, 8)} +} + +func (f *slotSyncFactory) OpenSync( + _ context.Context, + video FeedConfig, + audio FeedConfig, +) (SyncReader, error) { + f.mu.Lock() + f.active++ + if f.active > f.maxActive { + f.maxActive = f.active + } + f.mu.Unlock() + f.opened <- SyncPairConfig{Video: video, Audio: audio} + return &slotSyncReader{factory: f}, nil +} + +func (f *slotSyncFactory) counts() (active, maxActive, closeCount int) { + f.mu.Lock() + defer f.mu.Unlock() + return f.active, f.maxActive, f.closeCount +} + +type slotSyncReader struct { + factory *slotSyncFactory +} + +func (r *slotSyncReader) ReadSync(ctx context.Context) (SyncFrame, error) { + <-ctx.Done() + return SyncFrame{}, ctx.Err() +} + +func (r *slotSyncReader) Close() error { + r.factory.mu.Lock() + defer r.factory.mu.Unlock() + r.factory.active-- + r.factory.closeCount++ + return nil +} + +func newSlotTestSyncWorker(t *testing.T, factory SyncReaderFactory) *SyncWorker { + t.Helper() + worker, err := NewSyncWorker( + factory, + &fakeVideoSink{}, + &fakeAudioSink{}, + testRetryPolicy(1), + func(error) bool { return false }, + nil, + ) + if err != nil { + t.Fatalf("NewSyncWorker() error = %v", err) + } + return worker +} + +func receiveSyncSlotOpen(t *testing.T, opened <-chan SyncPairConfig) SyncPairConfig { + t.Helper() + select { + case config := <-opened: + return config + case <-time.After(time.Second): + t.Fatal("sync worker did not open") + return SyncPairConfig{} + } +} + +func testSyncPair(name string, active bool) SyncPairConfig { + return SyncPairConfig{ + Video: FeedConfig{Domain: "/video", UUID: name + "-video", Active: active}, + Audio: FeedConfig{Domain: "/audio", UUID: name + "-audio", Active: active}, + } +} + +func TestSyncPairConfigRejectsActivityMismatch(t *testing.T) { + config := testSyncPair("pair", true) + config.Audio.Active = false + if err := config.Validate(); !errors.Is(err, ErrSyncActivityMismatch) { + t.Fatalf("Validate() error = %v, want %v", err, ErrSyncActivityMismatch) + } +} + +func TestNewSyncSlotRequiresWorker(t *testing.T) { + slot, err := NewSyncSlot(nil) + if slot != nil { + t.Fatalf("NewSyncSlot(nil) slot = %#v, want nil", slot) + } + if !errors.Is(err, ErrSyncWorkerRequired) { + t.Fatalf("NewSyncSlot(nil) error = %v, want %v", err, ErrSyncWorkerRequired) + } +} + +func TestSyncSlotStartsInitialActivePair(t *testing.T) { + factory := newSlotSyncFactory() + slot, err := NewSyncSlot(newSlotTestSyncWorker(t, factory)) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + want := testSyncPair("first", true) + go func() { done <- slot.Run(ctx, want, make(chan SyncPairConfig)) }() + + if got := receiveSyncSlotOpen(t, factory.opened); got != want { + t.Fatalf("opened config = %#v, want %#v", got, want) + } + cancel() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) { + t.Fatalf("Run() error = %v, want context.Canceled", err) + } + case <-time.After(time.Second): + t.Fatal("Run() did not stop after cancellation") + } + active, _, closed := factory.counts() + if active != 0 || closed != 1 { + t.Fatalf("reader counts = active %d, closed %d; want 0, 1", active, closed) + } +} + +func TestSyncSlotReplacesWithoutOverlappingWorkers(t *testing.T) { + factory := newSlotSyncFactory() + slot, err := NewSyncSlot(newSlotTestSyncWorker(t, factory)) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + commands := make(chan SyncPairConfig) + done := make(chan error, 1) + first := testSyncPair("first", true) + second := testSyncPair("second", true) + go func() { done <- slot.Run(ctx, first, commands) }() + receiveSyncSlotOpen(t, factory.opened) + commands <- second + if got := receiveSyncSlotOpen(t, factory.opened); got != second { + t.Fatalf("replacement = %#v, want %#v", got, second) + } + cancel() + <-done + active, maxActive, closed := factory.counts() + if active != 0 || maxActive != 1 || closed != 2 { + t.Fatalf("counts = active %d, maximum %d, closed %d; want 0, 1, 2", active, maxActive, closed) + } +} + +func TestSyncSlotIgnoresInvalidCommand(t *testing.T) { + factory := newSlotSyncFactory() + slot, err := NewSyncSlot(newSlotTestSyncWorker(t, factory)) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + commands := make(chan SyncPairConfig) + done := make(chan error, 1) + initial := testSyncPair("first", true) + go func() { done <- slot.Run(ctx, initial, commands) }() + receiveSyncSlotOpen(t, factory.opened) + + invalid := testSyncPair("invalid", true) + invalid.Audio.Active = false + commands <- invalid + select { + case config := <-factory.opened: + t.Fatalf("invalid command opened config %#v", config) + case <-time.After(20 * time.Millisecond): + } + active, _, closed := factory.counts() + if active != 1 || closed != 0 { + t.Fatalf("invalid command disturbed reader: active %d, closed %d", active, closed) + } + cancel() + <-done +} + +func TestSyncSlotInactivePairStopsWithoutRestart(t *testing.T) { + factory := newSlotSyncFactory() + slot, err := NewSyncSlot(newSlotTestSyncWorker(t, factory)) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + commands := make(chan SyncPairConfig) + done := make(chan error, 1) + go func() { done <- slot.Run(ctx, testSyncPair("first", true), commands) }() + receiveSyncSlotOpen(t, factory.opened) + commands <- testSyncPair("first", false) + + deadline := time.Now().Add(time.Second) + for { + active, _, closed := factory.counts() + if active == 0 && closed == 1 { + break + } + if time.Now().After(deadline) { + t.Fatal("inactive pair did not stop reader") + } + time.Sleep(time.Millisecond) + } + close(commands) + select { + case err := <-done: + if err != nil { + t.Fatalf("Run() error = %v, want nil", err) + } + case <-time.After(time.Second): + t.Fatal("Run() did not stop after commands closed") + } +}