sync groups sample app
This commit is contained in:
@@ -0,0 +1,116 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"flag"
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"os"
|
||||||
|
"os/signal"
|
||||||
|
"syscall"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/qvest-digital/go-mxl/mxl"
|
||||||
|
)
|
||||||
|
|
||||||
|
func main() {
|
||||||
|
domain := flag.String("d", "/dev/shm/mxl", "MXL domain")
|
||||||
|
videoFlow := flag.String("v", "", "Video flow UUID")
|
||||||
|
audioFlow := flag.String("a", "", "Audio flow UUID")
|
||||||
|
flag.Parse()
|
||||||
|
if *videoFlow == "" || *audioFlow == "" {
|
||||||
|
log.Fatal("need both -v <video-flow> and -a <audio-flow>")
|
||||||
|
}
|
||||||
|
|
||||||
|
inst, err := mxl.NewInstance(*domain, "")
|
||||||
|
if err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
defer inst.Close()
|
||||||
|
|
||||||
|
vr, err := inst.NewReader(*videoFlow)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
defer vr.Close()
|
||||||
|
|
||||||
|
ar, err := inst.NewReader(*audioFlow)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
defer ar.Close()
|
||||||
|
|
||||||
|
vInfo, _ := vr.Info()
|
||||||
|
aInfo, _ := ar.Info()
|
||||||
|
vRate := vInfo.Config.Common.GrainRate
|
||||||
|
aRate := aInfo.Config.Common.GrainRate
|
||||||
|
aChans := aInfo.Config.Continuous.ChannelCount
|
||||||
|
fmt.Printf("video: %dx%d %d/%d | audio: %dch %d/%d\n",
|
||||||
|
vInfo.Config.Discrete.SliceSizes[0],
|
||||||
|
vInfo.Config.Discrete.GrainCount,
|
||||||
|
vRate.Num, vRate.Den,
|
||||||
|
aChans, aRate.Num, aRate.Den)
|
||||||
|
|
||||||
|
group, err := inst.NewSyncGroup()
|
||||||
|
if err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
defer group.Close()
|
||||||
|
if err := group.AddReader(vr); err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := group.AddReader(ar); err != nil {
|
||||||
|
log.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
stop := make(chan os.Signal, 1)
|
||||||
|
signal.Notify(stop, os.Interrupt, syscall.SIGTERM)
|
||||||
|
|
||||||
|
idx := mxl.CurrentIndex(vRate)
|
||||||
|
audioBatch := uint64(aRate.Num / (100 * aRate.Den)) // ~10ms
|
||||||
|
if audioBatch == 0 {
|
||||||
|
audioBatch = 1
|
||||||
|
}
|
||||||
|
|
||||||
|
var ticks int
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-stop:
|
||||||
|
fmt.Printf("\nstopped after %d ticks\n", ticks)
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
|
||||||
|
ts := mxl.IndexToTimestamp(vRate, idx)
|
||||||
|
err := group.WaitForDataAt(ts, 500*time.Millisecond)
|
||||||
|
switch {
|
||||||
|
case err == nil:
|
||||||
|
// Read video grain
|
||||||
|
g, gerr := vr.GetGrain(idx, 50*time.Millisecond)
|
||||||
|
// Read audio samples at the same timestamp
|
||||||
|
aIdx := mxl.TimestampToIndex(aRate, ts)
|
||||||
|
_, aerr := ar.GetSamples(aIdx, int(audioBatch), 50*time.Millisecond)
|
||||||
|
if gerr != nil {
|
||||||
|
log.Printf("grain: %v", gerr)
|
||||||
|
} else if aerr != nil {
|
||||||
|
log.Printf("samples: %v", aerr)
|
||||||
|
} else {
|
||||||
|
fmt.Printf("tick idx=%d ts=%d grainSize=%d audioIdx=%d\n",
|
||||||
|
idx, ts, g.GrainSize, aIdx)
|
||||||
|
ticks++
|
||||||
|
if ticks >= 10 {
|
||||||
|
fmt.Println("done")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
idx++
|
||||||
|
case errors.Is(err, mxl.ErrTimeout), errors.Is(err, mxl.ErrOutOfRangeEarly):
|
||||||
|
time.Sleep(5 * time.Millisecond)
|
||||||
|
case errors.Is(err, mxl.ErrOutOfRangeLate):
|
||||||
|
log.Printf("fell behind, resyncing")
|
||||||
|
idx = mxl.CurrentIndex(vRate)
|
||||||
|
default:
|
||||||
|
log.Fatalf("WaitForDataAt: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user