From 6b37a34a12954e87b7f632b7a165ee54adaada9c Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Thu, 27 Aug 2026 10:39:48 +0300 Subject: [PATCH] add bounded video source read --- internal/source/source.go | 104 +++++++++++++++++++++++++++++--------- 1 file changed, 80 insertions(+), 24 deletions(-) diff --git a/internal/source/source.go b/internal/source/source.go index 7bed226..16a6db4 100644 --- a/internal/source/source.go +++ b/internal/source/source.go @@ -120,37 +120,93 @@ func (s *Source) Close() error { func (s *Source) NextCtx(ctx context.Context, timeout time.Duration) (Frame, error) { for { - select { - case <-ctx.Done(): + frame, err := s.ReadOnceCtx(ctx, timeout) + if err == nil { + return frame, nil + } + if ctx.Err() != nil { return Frame{}, ctx.Err() - default: } - g, err := s.reader.GetGrain(s.idx, timeout) - switch { - case err == nil: - f := Frame{ - Index: g.Index, - Width: s.width, - Height: s.height, - Stride: s.stride, - Size: g.GrainSize, - Invalid: g.Invalid(), - Payload: g.Payload, + if KindOf(err) != ErrorKindTemporary { + return Frame{}, err + } + + if errors.Is(err, mxl.ErrOutOfRangeEarly) { + select { + case <-time.After(10 * time.Millisecond): + case <-ctx.Done(): + return Frame{}, ctx.Err() } - s.idx++ - return f, nil - case errors.Is(err, mxl.ErrTimeout): - s.idx = mxl.CurrentIndex(s.rate) - case errors.Is(err, mxl.ErrOutOfRangeEarly): - time.Sleep(10 * time.Millisecond) - case errors.Is(err, mxl.ErrOutOfRangeLate): - s.idx = mxl.CurrentIndex(s.rate) - default: - return Frame{}, fmt.Errorf("GetGrain: %w", err) } } } +func (s *Source) ReadOnceCtx( + ctx context.Context, + timeout time.Duration, +) (Frame, error) { + if err := ctx.Err(); err != nil { + return Frame{}, err + } + + g, err := s.reader.GetGrain(s.idx, timeout) + if ctx.Err() != nil { + return Frame{}, ctx.Err() + } + + switch { + case err == nil: + frame := Frame{ + Index: g.Index, + Width: s.width, + Height: s.height, + Stride: s.stride, + Size: g.GrainSize, + Invalid: g.Invalid(), + Payload: g.Payload, + } + s.idx++ + return frame, nil + + case errors.Is(err, mxl.ErrTimeout): + s.idx = mxl.CurrentIndex(s.rate) + return Frame{}, wrapError( + "read video", + ErrorKindTemporary, + err, + ) + + case errors.Is(err, mxl.ErrOutOfRangeEarly): + return Frame{}, wrapError( + "read video", + ErrorKindTemporary, + err, + ) + + case errors.Is(err, mxl.ErrOutOfRangeLate): + s.idx = mxl.CurrentIndex(s.rate) + return Frame{}, wrapError( + "read video", + ErrorKindTemporary, + err, + ) + + case errors.Is(err, mxl.ErrFlowInvalid): + return Frame{}, wrapError( + "read video", + ErrorKindUnavailable, + err, + ) + + default: + return Frame{}, wrapError( + "read video", + ErrorKindUnavailable, + err, + ) + } +} + func (s *Source) FlowDef() string { return s.def } func (s *Source) Rate() mxl.Rational { return s.rate } func (s *Source) Stride() uint32 { return s.stride }