From 1083f3b83ca2ca70a9eecffac5921f5e9ba403ff Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Fri, 28 Aug 2026 21:41:23 +0300 Subject: [PATCH] add local MXL audio adapter --- imgui.ini | 4 +- internal/adapter/mxl/audio.go | 211 +++++++++++++++++++++++ internal/adapter/mxl/audio_test.go | 264 +++++++++++++++++++++++++++++ 3 files changed, 477 insertions(+), 2 deletions(-) create mode 100644 internal/adapter/mxl/audio.go create mode 100644 internal/adapter/mxl/audio_test.go diff --git a/imgui.ini b/imgui.ini index e9bc32f..8609d51 100644 --- a/imgui.ini +++ b/imgui.ini @@ -14,7 +14,7 @@ Size=200,200 Collapsed=0 [Window][Connection] -Pos=74,351 -Size=799,226 +Pos=979,700 +Size=695,233 Collapsed=0 diff --git a/internal/adapter/mxl/audio.go b/internal/adapter/mxl/audio.go new file mode 100644 index 0000000..4373242 --- /dev/null +++ b/internal/adapter/mxl/audio.go @@ -0,0 +1,211 @@ +package mxladapter + +import ( + "context" + "errors" + "fmt" + "math/bits" + "time" + + "mxl-player/internal/playback" + "mxl-player/internal/source" + + mxl "github.com/qvest-digital/go-mxl/mxl" +) + +const ( + DefaultAudioReadTimeout = 20 * time.Millisecond + DefaultAudioBatchDuration = 10 * time.Millisecond + DefaultAudioUnavailableAfter = 2 * time.Second + DefaultAudioTemporaryDelay = 10 * time.Millisecond +) + +var ErrInvalidAudioBatch = errors.New("invalid audio batch configuration") + +type AudioFactory struct { + ReadTimeout time.Duration + BatchDuration time.Duration + UnavailableAfter time.Duration +} + +type localAudioSource interface { + ReadAudioOnceCtx(context.Context, uint64, time.Duration) (source.AudioFrame, error) + Rate() mxl.Rational + Channels() uint64 + Close() error +} + +type audioReader struct { + source localAudioSource + readTimeout time.Duration + batch uint64 + channels uint64 + rateNumerator int64 + rateDenominator int64 + unavailableAfter time.Duration + retryDelay time.Duration + now func() time.Time + wait temporaryWaitFunc +} + +var _ playback.AudioReaderFactory = AudioFactory{} +var _ playback.AudioReader = (*audioReader)(nil) + +func audioBatchSize( + rateNumerator int64, + rateDenominator int64, + duration time.Duration, +) (uint64, error) { + if rateNumerator <= 0 || rateDenominator <= 0 || duration <= 0 { + return 0, fmt.Errorf( + "%w: rate=%d/%d duration=%s", + ErrInvalidAudioBatch, + rateNumerator, + rateDenominator, + duration, + ) + } + + denominator := uint64(rateDenominator) + seconds := uint64(time.Second) + if denominator > ^uint64(0)/seconds { + return 0, fmt.Errorf("%w: denominator overflow", ErrInvalidAudioBatch) + } + denominator *= seconds + + high, low := bits.Mul64(uint64(rateNumerator), uint64(duration)) + if high >= denominator { + return 0, fmt.Errorf("%w: sample count overflow", ErrInvalidAudioBatch) + } + batch, _ := bits.Div64(high, low, denominator) + if batch == 0 { + batch = 1 + } + return batch, nil +} + +func (f AudioFactory) OpenAudio( + ctx context.Context, + config playback.FeedConfig, +) (playback.AudioReader, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + if err := config.Validate(); err != nil { + return nil, &source.SourceError{ + Op: "validate audio feed", + Kind: source.ErrorKindInvalidConfig, + Err: err, + } + } + if !config.IsConfigured() { + return nil, &source.SourceError{ + Op: "validate audio feed", + Kind: source.ErrorKindInvalidConfig, + Err: errors.New("audio feed is not configured"), + } + } + + src, err := source.OpenAudio(config.Domain, config.UUID) + if err != nil { + return nil, fmt.Errorf("open local MXL audio: %w", err) + } + closeOnError := true + defer func() { + if closeOnError { + _ = src.Close() + } + }() + + if err := ctx.Err(); err != nil { + return nil, err + } + + readTimeout := f.ReadTimeout + if readTimeout <= 0 { + readTimeout = DefaultAudioReadTimeout + } + batchDuration := f.BatchDuration + if batchDuration <= 0 { + batchDuration = DefaultAudioBatchDuration + } + unavailableAfter := f.UnavailableAfter + if unavailableAfter <= 0 { + unavailableAfter = DefaultAudioUnavailableAfter + } + + rate := src.Rate() + batch, err := audioBatchSize(rate.Num, rate.Den, batchDuration) + if err != nil { + return nil, &source.SourceError{ + Op: "calculate audio batch", + Kind: source.ErrorKindInvalidConfig, + Err: err, + } + } + + reader := &audioReader{ + source: src, + readTimeout: readTimeout, + batch: batch, + channels: src.Channels(), + rateNumerator: rate.Num, + rateDenominator: rate.Den, + unavailableAfter: unavailableAfter, + retryDelay: DefaultAudioTemporaryDelay, + now: time.Now, + wait: waitForTemporaryRetry, + } + closeOnError = false + return reader, nil +} + +func (r *audioReader) ReadAudio(ctx context.Context) (playback.AudioFrame, error) { + var unavailableSince time.Time + + for { + frame, err := r.source.ReadAudioOnceCtx(ctx, r.batch, r.readTimeout) + if err == nil { + return playback.AudioFrame{ + Index: frame.Index, + SampleCount: frame.SampleCount, + Channels: frame.Channels, + SampleRateNumerator: r.rateNumerator, + SampleRateDenominator: r.rateDenominator, + Samples: frame.Samples, + }, nil + } + if ctx.Err() != nil { + return playback.AudioFrame{}, ctx.Err() + } + if source.KindOf(err) != source.ErrorKindTemporary { + return playback.AudioFrame{}, err + } + + now := r.now() + if unavailableSince.IsZero() { + unavailableSince = now + } else if now.Sub(unavailableSince) >= r.unavailableAfter { + return playback.AudioFrame{}, &source.SourceError{ + Op: "read local MXL audio", + Kind: source.ErrorKindUnavailable, + Err: fmt.Errorf( + "no audio data for %s: %w", + r.unavailableAfter, + err, + ), + } + } + + if err := r.wait(ctx, r.retryDelay); err != nil { + if ctx.Err() != nil { + return playback.AudioFrame{}, ctx.Err() + } + return playback.AudioFrame{}, err + } + } +} + +func (r *audioReader) Close() error { + return r.source.Close() +} diff --git a/internal/adapter/mxl/audio_test.go b/internal/adapter/mxl/audio_test.go new file mode 100644 index 0000000..412cbdc --- /dev/null +++ b/internal/adapter/mxl/audio_test.go @@ -0,0 +1,264 @@ +package mxladapter + +import ( + "context" + "errors" + "math" + "testing" + "time" + + "mxl-player/internal/playback" + "mxl-player/internal/source" + + mxl "github.com/qvest-digital/go-mxl/mxl" +) + +type localAudioReadResult struct { + frame source.AudioFrame + err error +} + +type fakeLocalAudioSource struct { + results []localAudioReadResult + calls int + batches []uint64 + timeouts []time.Duration + closed bool + closeErr error + rate mxl.Rational + channels uint64 +} + +func (s *fakeLocalAudioSource) ReadAudioOnceCtx( + _ context.Context, + batch uint64, + timeout time.Duration, +) (source.AudioFrame, error) { + s.batches = append(s.batches, batch) + s.timeouts = append(s.timeouts, timeout) + if s.calls >= len(s.results) { + return source.AudioFrame{}, errors.New("unexpected local audio read") + } + result := s.results[s.calls] + s.calls++ + return result.frame, result.err +} + +func (s *fakeLocalAudioSource) Rate() mxl.Rational { return s.rate } +func (s *fakeLocalAudioSource) Channels() uint64 { return s.channels } + +func (s *fakeLocalAudioSource) Close() error { + s.closed = true + return s.closeErr +} + +func temporaryAudioError(cause error) error { + return &source.SourceError{ + Op: "read audio", + Kind: source.ErrorKindTemporary, + Err: cause, + } +} + +func TestAudioBatchSize(t *testing.T) { + tests := []struct { + name string + num int64 + den int64 + duration time.Duration + want uint64 + wantErr bool + }{ + {name: "ten milliseconds at 48kHz", num: 48000, den: 1, duration: 10 * time.Millisecond, want: 480}, + {name: "fraction rounds down", num: 30000, den: 1001, duration: time.Second, want: 29}, + {name: "minimum one sample", num: 1, den: 1, duration: time.Nanosecond, want: 1}, + {name: "zero numerator", den: 1, duration: time.Second, wantErr: true}, + {name: "zero denominator", num: 48000, duration: time.Second, wantErr: true}, + {name: "zero duration", num: 48000, den: 1, wantErr: true}, + {name: "result overflow", num: math.MaxInt64, den: 1, duration: time.Duration(math.MaxInt64), wantErr: true}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := audioBatchSize(tt.num, tt.den, tt.duration) + if tt.wantErr { + if !errors.Is(err, ErrInvalidAudioBatch) { + t.Fatalf("audioBatchSize() error = %v, want %v", err, ErrInvalidAudioBatch) + } + return + } + if err != nil { + t.Fatalf("audioBatchSize() error = %v", err) + } + if got != tt.want { + t.Fatalf("audioBatchSize() = %d, want %d", got, tt.want) + } + }) + } +} + +func TestAudioFactoryOpenAudioCanceled(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + reader, err := (AudioFactory{}).OpenAudio(ctx, playback.FeedConfig{}) + if reader != nil { + t.Fatal("OpenAudio() reader is not nil after cancellation") + } + if !errors.Is(err, context.Canceled) { + t.Fatalf("OpenAudio() error = %v, want context.Canceled", err) + } +} + +func TestAudioFactoryOpenAudioRejectsInvalidConfig(t *testing.T) { + configs := []playback.FeedConfig{ + {}, + {Domain: "/audio", Active: true}, + {UUID: "audio", Active: true}, + } + for _, config := range configs { + reader, err := (AudioFactory{}).OpenAudio(context.Background(), config) + if reader != nil { + t.Fatalf("OpenAudio(%#v) reader is not nil", config) + } + if source.KindOf(err) != source.ErrorKindInvalidConfig { + t.Fatalf("OpenAudio(%#v) error kind = %v, want invalid config", config, source.KindOf(err)) + } + if ShouldRetry(err) { + t.Fatalf("ShouldRetry(OpenAudio(%#v)) = true", config) + } + } +} + +func TestAudioReaderTemporaryFailureThenFrameWithoutCopy(t *testing.T) { + samples := [][]byte{{1, 2, 3, 4}, {5, 6, 7, 8}} + want := source.AudioFrame{ + Index: 42, + SampleCount: 1, + Channels: 2, + Samples: samples, + } + localSource := &fakeLocalAudioSource{ + results: []localAudioReadResult{ + {err: temporaryAudioError(errors.New("early"))}, + {frame: want}, + }, + } + waits := 0 + reader := &audioReader{ + source: localSource, + readTimeout: 25 * time.Millisecond, + batch: 480, + channels: 2, + rateNumerator: 48000, + rateDenominator: 1, + unavailableAfter: 2 * time.Second, + retryDelay: 10 * time.Millisecond, + now: func() time.Time { return time.Unix(100, 0) }, + wait: func(context.Context, time.Duration) error { + waits++ + return nil + }, + } + + got, err := reader.ReadAudio(context.Background()) + if err != nil { + t.Fatalf("ReadAudio() error = %v", err) + } + if localSource.calls != 2 || waits != 1 { + t.Fatalf("reads = %d, waits = %d; want 2, 1", localSource.calls, waits) + } + if got.Index != want.Index || got.SampleCount != want.SampleCount || got.Channels != want.Channels { + t.Fatalf("frame metadata = %#v, want %#v", got, want) + } + if got.SampleRateNumerator != 48000 || got.SampleRateDenominator != 1 { + t.Fatalf("sample rate = %d/%d, want 48000/1", got.SampleRateNumerator, got.SampleRateDenominator) + } + for channel := range samples { + if &got.Samples[channel][0] != &samples[channel][0] { + t.Fatalf("channel %d samples were copied", channel) + } + } + for _, batch := range localSource.batches { + if batch != 480 { + t.Fatalf("read batch = %d, want 480", batch) + } + } +} + +func TestAudioReaderProlongedTemporaryFailureBecomesUnavailable(t *testing.T) { + lastCause := errors.New("last timeout") + localSource := &fakeLocalAudioSource{ + results: []localAudioReadResult{ + {err: temporaryAudioError(errors.New("first timeout"))}, + {err: temporaryAudioError(errors.New("second timeout"))}, + {err: temporaryAudioError(lastCause)}, + }, + } + times := []time.Time{time.Unix(100, 0), time.Unix(101, 0), time.Unix(102, 0)} + nowCall := 0 + waits := 0 + reader := &audioReader{ + source: localSource, + readTimeout: 20 * time.Millisecond, + batch: 480, + unavailableAfter: 2 * time.Second, + retryDelay: 10 * time.Millisecond, + now: func() time.Time { + result := times[nowCall] + nowCall++ + return result + }, + wait: func(context.Context, time.Duration) error { + waits++ + return nil + }, + } + + _, err := reader.ReadAudio(context.Background()) + if source.KindOf(err) != source.ErrorKindUnavailable { + t.Fatalf("ReadAudio() error kind = %v, want unavailable", source.KindOf(err)) + } + if !errors.Is(err, lastCause) { + t.Fatalf("ReadAudio() error = %v, want cause %v", err, lastCause) + } + if !ShouldRetry(err) { + t.Fatal("ShouldRetry(ReadAudio()) = false, want true") + } + if localSource.calls != 3 || waits != 2 { + t.Fatalf("reads = %d, waits = %d; want 3, 2", localSource.calls, waits) + } +} + +func TestAudioReaderCancellationDuringTemporaryWait(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + localSource := &fakeLocalAudioSource{ + results: []localAudioReadResult{{err: temporaryAudioError(errors.New("early"))}}, + } + reader := &audioReader{ + source: localSource, + batch: 480, + unavailableAfter: 2 * time.Second, + now: func() time.Time { return time.Unix(100, 0) }, + wait: func(ctx context.Context, _ time.Duration) error { + cancel() + return ctx.Err() + }, + } + + _, err := reader.ReadAudio(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("ReadAudio() error = %v, want context.Canceled", err) + } +} + +func TestAudioReaderCloseDelegates(t *testing.T) { + closeErr := errors.New("close failed") + source := &fakeLocalAudioSource{closeErr: closeErr} + reader := &audioReader{source: source} + + err := reader.Close() + if !errors.Is(err, closeErr) || !source.closed { + t.Fatalf("Close() = %v, closed=%t; want %v, true", err, source.closed, closeErr) + } +}