add bounded audio source read

This commit is contained in:
Dmitry Sergeev
2026-08-28 00:14:34 +03:00
parent d4041af119
commit 5d1a9a3218
+74 -23
View File
@@ -296,47 +296,98 @@ func OpenAudio(domain, flowID string) (*AudioSource, error) {
}, nil }, nil
} }
func (s *AudioSource) NextAudio(ctx context.Context, batch uint64, timeout time.Duration) (AudioFrame, error) { func (s *AudioSource) ReadAudioOnceCtx(
for { ctx context.Context,
select { batch uint64,
case <-ctx.Done(): timeout time.Duration,
return AudioFrame{}, ctx.Err() ) (AudioFrame, error) {
default: if err := ctx.Err(); err != nil {
return AudioFrame{}, err
} }
v, err := s.r.GetSamples(s.idx, int(batch), timeout)
value, err := s.r.GetSamples(s.idx, int(batch), timeout)
if ctxErr := ctx.Err(); ctxErr != nil {
return AudioFrame{}, ctxErr
}
switch { switch {
case err == nil: case err == nil:
samples := make([][]byte, s.chans) samples := make([][]byte, s.chans)
for ch := uint64(0); ch < s.chans; ch++ { for channel := uint64(0); channel < s.chans; channel++ {
f1, f2, _ := v.ChannelFragments(ch) first, second, _ := value.ChannelFragments(channel)
if len(f2) > 0 { if len(second) > 0 {
samples[ch] = append(f1, f2...) samples[channel] = append(first, second...)
} else { } else {
samples[ch] = f1 samples[channel] = first
} }
} }
f := AudioFrame{
frame := AudioFrame{
Index: s.idx, Index: s.idx,
SampleCount: batch, SampleCount: batch,
Channels: s.chans, Channels: s.chans,
Samples: samples, Samples: samples,
} }
s.idx += batch s.idx += batch
return f, nil return frame, nil
case errors.Is(err, mxl.ErrOutOfRangeEarly): case errors.Is(err, mxl.ErrOutOfRangeEarly):
select { return AudioFrame{}, wrapError(
case <-time.After(10 * time.Millisecond): "read audio",
case <-ctx.Done(): ErrorKindTemporary,
err,
)
case errors.Is(err, mxl.ErrOutOfRangeLate):
runtimeInfo, runtimeErr := s.r.Runtime()
if runtimeErr != nil {
return AudioFrame{}, wrapError(
"read audio runtime",
ErrorKindUnavailable,
runtimeErr,
)
}
s.idx = runtimeInfo.HeadIndex
return AudioFrame{}, wrapError(
"read audio",
ErrorKindTemporary,
err,
)
default:
return AudioFrame{}, wrapError(
"read audio",
ErrorKindUnavailable,
err,
)
}
}
func (s *AudioSource) NextAudio(
ctx context.Context,
batch uint64,
timeout time.Duration,
) (AudioFrame, error) {
for {
frame, err := s.ReadAudioOnceCtx(ctx, batch, timeout)
if err == nil {
return frame, nil
}
if ctx.Err() != nil {
return AudioFrame{}, ctx.Err() return AudioFrame{}, ctx.Err()
} }
case errors.Is(err, mxl.ErrOutOfRangeLate): if KindOf(err) != ErrorKindTemporary {
rt, rerr := s.r.Runtime() return AudioFrame{}, err
if rerr != nil {
return AudioFrame{}, fmt.Errorf("Runtime: %w", rerr)
} }
s.idx = rt.HeadIndex
timer := time.NewTimer(10 * time.Millisecond)
select {
case <-timer.C:
case <-ctx.Done():
if !timer.Stop() {
select {
case <-timer.C:
default: default:
return AudioFrame{}, fmt.Errorf("GetSamples: %w", err) }
}
return AudioFrame{}, ctx.Err()
} }
} }
} }