package source import ( "context" "encoding/json" "errors" "fmt" "time" mxl "github.com/qvest-digital/go-mxl/mxl" ) type flowDef struct { Label string `json:"label"` FrameWidth int `json:"frame_width"` FrameHeight int `json:"frame_height"` MediaType string `json:"media_type"` Colorspace string `json:"colorspace"` GrainRate struct { Numerator int64 `json:"numerator"` Denominator int64 `json:"denominator"` } `json:"grain_rate"` } func flowLabel(inst *mxl.Instance, flowID string) string { definition, err := inst.FlowDef(flowID) if err != nil { return "" } var metadata struct { Label string `json:"label"` } if json.Unmarshal([]byte(definition), &metadata) != nil { return "" } return metadata.Label } type Frame struct { Index uint64 Width uint32 Height uint32 Stride uint32 Size uint32 Invalid bool Label string FrameRateNumerator int64 FrameRateDenominator int64 Payload []byte } type Source struct { inst *mxl.Instance reader *mxl.Reader info mxl.FlowInfo def string rate mxl.Rational stride uint32 width uint32 height uint32 label string idx uint64 } func Open(domain, flowID string) (*Source, error) { inst, err := mxl.NewInstance(domain, "") if err != nil { return nil, wrapError("new MXL instance", ErrorKindUnavailable, err) } r, err := inst.NewReader(flowID) if err != nil { inst.Close() return nil, wrapError("open video reader", ErrorKindUnavailable, err) } info, err := r.Info() if err != nil { r.Close() inst.Close() return nil, wrapError("get video info", ErrorKindUnavailable, err) } def, err := inst.FlowDef(flowID) if err != nil { r.Close() inst.Close() return nil, wrapError("read video flow definition", ErrorKindUnavailable, err) } var fd flowDef if err := json.Unmarshal([]byte(def), &fd); err != nil { r.Close() inst.Close() return nil, wrapError("parse video flow definition JSON", ErrorKindInvalidConfig, err) } if fd.FrameWidth == 0 || fd.FrameHeight == 0 { r.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, errors.New("flow has no video dimensions"), ) } rate := info.Config.Common.GrainRate idx := mxl.CurrentIndex(rate) if idx == mxl.UndefinedIndex { r.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, fmt.Errorf("invalid grain rate: %d/%d", rate.Num, rate.Den), ) } if len(info.Config.Discrete.SliceSizes) == 0 { r.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, errors.New("video flow has no slice sizes"), ) } return &Source{ inst: inst, reader: r, info: info, def: def, rate: info.Config.Common.GrainRate, stride: info.Config.Discrete.SliceSizes[0], width: uint32(fd.FrameWidth), height: uint32(fd.FrameHeight), label: fd.Label, idx: idx, }, nil } func (s *Source) Close() error { _ = s.reader.Close() return s.inst.Close() } func (s *Source) NextCtx(ctx context.Context, timeout time.Duration) (Frame, error) { for { frame, err := s.ReadOnceCtx(ctx, timeout) if err == nil { return frame, nil } if ctx.Err() != nil { return Frame{}, ctx.Err() } 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() } } } } 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(), Label: s.label, FrameRateNumerator: s.rate.Num, FrameRateDenominator: s.rate.Den, 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 } func (s *Source) Width() uint32 { return s.width } func (s *Source) Height() uint32 { return s.height } func (s *Source) GrainCount() uint32 { return s.info.Config.Discrete.GrainCount } func (s *Source) Format() mxl.DataFormat { return s.info.Config.Common.Format } type AudioSource struct { inst *mxl.Instance r *mxl.Reader info mxl.FlowInfo rate mxl.Rational chans uint64 idx uint64 label string } type AudioFrame struct { Index uint64 SampleCount uint64 Channels uint64 Label string Samples [][]byte // per-channel byte slices (F32, deinterleaved) } func OpenAudio(domain, flowID string) (*AudioSource, error) { inst, err := mxl.NewInstance(domain, "") if err != nil { return nil, wrapError("new MXL instance", ErrorKindUnavailable, err) } r, err := inst.NewReader(flowID) if err != nil { inst.Close() return nil, wrapError("open audio reader", ErrorKindUnavailable, err) } info, err := r.Info() if err != nil { r.Close() inst.Close() return nil, wrapError("get audio info", ErrorKindUnavailable, err) } if info.Config.Common.Format.IsDiscrete() { r.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, errors.New("audio flow is discrete"), ) } channels := uint64(info.Config.Continuous.ChannelCount) if channels == 0 { r.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, errors.New("audio flow has no channels"), ) } rate := info.Config.Common.GrainRate if rate.Num <= 0 || rate.Den <= 0 { r.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, fmt.Errorf("invalid audio rate: %d/%d", rate.Num, rate.Den), ) } idx := info.Runtime.HeadIndex if idx == 0 { r.Close() inst.Close() return nil, wrapError( "open audio reader", ErrorKindUnavailable, errors.New("audio flow has no producer data"), ) } return &AudioSource{ inst: inst, r: r, info: info, rate: rate, chans: channels, idx: idx, label: flowLabel(inst, flowID), }, nil } func (s *AudioSource) ReadAudioOnceCtx( ctx context.Context, batch uint64, timeout time.Duration, ) (AudioFrame, error) { if err := ctx.Err(); err != nil { return AudioFrame{}, err } value, err := s.r.GetSamples(s.idx, int(batch), timeout) if ctxErr := ctx.Err(); ctxErr != nil { return AudioFrame{}, ctxErr } switch { case err == nil: samples := make([][]byte, s.chans) for channel := uint64(0); channel < s.chans; channel++ { first, second, _ := value.ChannelFragments(channel) if len(second) > 0 { samples[channel] = append(first, second...) } else { samples[channel] = first } } frame := AudioFrame{ Index: s.idx, SampleCount: batch, Channels: s.chans, Label: s.label, Samples: samples, } s.idx += batch return frame, nil case errors.Is(err, mxl.ErrOutOfRangeEarly): return AudioFrame{}, wrapError( "read audio", 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() } if KindOf(err) != ErrorKindTemporary { return AudioFrame{}, err } timer := time.NewTimer(10 * time.Millisecond) select { case <-timer.C: case <-ctx.Done(): if !timer.Stop() { select { case <-timer.C: default: } } return AudioFrame{}, ctx.Err() } } } func (s *AudioSource) Close() error { _ = s.r.Close() return s.inst.Close() } func (s *AudioSource) Rate() mxl.Rational { return s.rate } func (s *AudioSource) Channels() uint64 { return s.chans } type SyncSource struct { inst *mxl.Instance vr *mxl.Reader ar *mxl.Reader group *mxl.SyncGroup rate mxl.Rational // video rate aRate mxl.Rational chans uint64 idx uint64 width, height, stride uint32 videoLabel, audioLabel string } // OpenSameDomainSync opens a native MXL synchronization group. // Both feeds must belong to the supplied domain because go-mxl sync groups // cannot contain readers from different MXL instances. func OpenSameDomainSync(domain, videoFlow, audioFlow string) (*SyncSource, error) { inst, err := mxl.NewInstance(domain, "") if err != nil { return nil, wrapError("new MXL instance", ErrorKindUnavailable, err) } vr, err := inst.NewReader(videoFlow) if err != nil { inst.Close() return nil, wrapError("open video reader", ErrorKindUnavailable, err) } ar, err := inst.NewReader(audioFlow) if err != nil { vr.Close() inst.Close() return nil, wrapError("open audio reader", ErrorKindUnavailable, err) } vInfo, err := vr.Info() if err != nil { ar.Close() vr.Close() inst.Close() return nil, wrapError("get video info", ErrorKindUnavailable, err) } if !vInfo.Config.Common.Format.IsDiscrete() { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, errors.New("video flow is continuous"), ) } aInfo, err := ar.Info() if err != nil { ar.Close() vr.Close() inst.Close() return nil, wrapError("get audio info", ErrorKindUnavailable, err) } if aInfo.Config.Common.Format.IsDiscrete() { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, errors.New("audio flow is discrete"), ) } channels := uint64(aInfo.Config.Continuous.ChannelCount) if channels == 0 { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, errors.New("audio flow has no channels"), ) } aRate := aInfo.Config.Common.GrainRate if aRate.Num <= 0 || aRate.Den <= 0 { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate audio flow", ErrorKindInvalidConfig, fmt.Errorf("invalid audio rate: %d/%d", aRate.Num, aRate.Den), ) } def, err := inst.FlowDef(videoFlow) if err != nil { ar.Close() vr.Close() inst.Close() return nil, wrapError("read video flow definition", ErrorKindUnavailable, err) } var fd flowDef if err := json.Unmarshal([]byte(def), &fd); err != nil { ar.Close() vr.Close() inst.Close() return nil, wrapError("parse video flow definition JSON", ErrorKindInvalidConfig, err) } if fd.FrameWidth <= 0 || fd.FrameHeight <= 0 { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, fmt.Errorf( "invalid video dimensions: %dx%d", fd.FrameWidth, fd.FrameHeight, ), ) } if len(vInfo.Config.Discrete.SliceSizes) == 0 { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, errors.New("video flow has no slice sizes"), ) } vRate := vInfo.Config.Common.GrainRate idx := mxl.CurrentIndex(vRate) if idx == mxl.UndefinedIndex { ar.Close() vr.Close() inst.Close() return nil, wrapError( "validate video flow", ErrorKindInvalidConfig, fmt.Errorf("invalid grain rate: %d/%d", vRate.Num, vRate.Den), ) } group, err := inst.NewSyncGroup() if err != nil { ar.Close() vr.Close() inst.Close() return nil, wrapError( "create native sync group", ErrorKindUnavailable, err, ) } if err := group.AddReader(vr); err != nil { group.Close() ar.Close() vr.Close() inst.Close() return nil, wrapError( "add video reader to native sync group", ErrorKindUnavailable, err, ) } if err := group.AddReader(ar); err != nil { group.Close() ar.Close() vr.Close() inst.Close() return nil, wrapError( "add audio reader to native sync group", ErrorKindUnavailable, err, ) } return &SyncSource{ inst: inst, vr: vr, ar: ar, group: group, rate: vRate, aRate: aRate, chans: channels, idx: idx, width: uint32(fd.FrameWidth), height: uint32(fd.FrameHeight), stride: vInfo.Config.Discrete.SliceSizes[0], videoLabel: fd.Label, audioLabel: flowLabel(inst, audioFlow), }, nil } func (s *SyncSource) Close() error { _ = s.group.Close() _ = s.ar.Close() _ = s.vr.Close() return s.inst.Close() } // NextSync reads one video frame and the audio interval between this video // timestamp and the next. Deriving the interval for every frame preserves // exact long-term timing for fractional video rates. func (s *SyncSource) NextSync(ctx context.Context, timeout time.Duration) (Frame, AudioFrame, error) { var timeouts int for { select { case <-ctx.Done(): return Frame{}, AudioFrame{}, ctx.Err() default: } ts := mxl.IndexToTimestamp(s.rate, s.idx) err := s.group.WaitForDataAt(ts, timeout) switch { case err == nil: // read video g, gerr := s.vr.GetGrain(s.idx, 50*time.Millisecond) if gerr != nil { if errors.Is(gerr, mxl.ErrFlowInvalid) { s.idx = mxl.CurrentIndex(s.rate) continue } if errors.Is(gerr, mxl.ErrOutOfRangeLate) { s.idx = mxl.CurrentIndex(s.rate) continue } if errors.Is(gerr, mxl.ErrOutOfRangeEarly) { select { case <-time.After(5 * time.Millisecond): case <-ctx.Done(): return Frame{}, AudioFrame{}, ctx.Err() } continue } return Frame{}, AudioFrame{}, fmt.Errorf("GetGrain: %w", gerr) } // read audio at the same timestamp aIdx := mxl.TimestampToIndex(s.aRate, ts) nextTimestamp := mxl.IndexToTimestamp(s.rate, s.idx+1) nextAudioIndex := mxl.TimestampToIndex(s.aRate, nextTimestamp) if nextAudioIndex <= aIdx { return Frame{}, AudioFrame{}, wrapError( "calculate synchronized audio interval", ErrorKindInvalidConfig, fmt.Errorf("invalid audio interval: %d..%d", aIdx, nextAudioIndex), ) } audioBatch := nextAudioIndex - aIdx av, aerr := s.ar.GetSamples(aIdx, int(audioBatch), 50*time.Millisecond) if aerr != nil { kind := ErrorKindUnavailable if errors.Is(aerr, mxl.ErrTimeout) || errors.Is(aerr, mxl.ErrOutOfRangeEarly) || errors.Is(aerr, mxl.ErrOutOfRangeLate) { kind = ErrorKindTemporary } return Frame{}, AudioFrame{}, wrapError( "read synchronized audio", kind, aerr, ) } vFrame := Frame{ Index: g.Index, Width: s.width, Height: s.height, Stride: s.stride, Size: g.GrainSize, Invalid: g.Invalid(), Payload: g.Payload, Label: s.videoLabel, FrameRateNumerator: s.rate.Num, FrameRateDenominator: s.rate.Den, } s.idx++ samples := make([][]byte, s.chans) for ch := uint64(0); ch < s.chans; ch++ { f1, f2, _ := av.ChannelFragments(ch) if len(f2) > 0 { samples[ch] = append(f1, f2...) } else { samples[ch] = f1 } } aFrame := AudioFrame{ Index: aIdx, SampleCount: audioBatch, Channels: s.chans, Samples: samples, Label: s.audioLabel, } return vFrame, aFrame, nil case errors.Is(err, mxl.ErrTimeout), errors.Is(err, mxl.ErrOutOfRangeEarly), errors.Is(err, mxl.ErrOutOfRangeLate): timeouts++ if timeouts > 10 { timeouts = 0 s.idx = mxl.CurrentIndex(s.rate) return Frame{}, AudioFrame{}, fmt.Errorf("sync: feeds not responding") } s.idx = mxl.CurrentIndex(s.rate) select { case <-time.After(5 * time.Millisecond): case <-ctx.Done(): return Frame{}, AudioFrame{}, ctx.Err() } default: return Frame{}, AudioFrame{}, fmt.Errorf("WaitForDataAt: %w", err) } } } func (s *SyncSource) Width() uint32 { return s.width } func (s *SyncSource) Height() uint32 { return s.height } func (s *SyncSource) Stride() uint32 { return s.stride } func (s *SyncSource) Rate() mxl.Rational { return s.rate } func (s *SyncSource) AudioRate() mxl.Rational { return s.aRate } func (s *SyncSource) Channels() uint64 { return s.chans }