//go:build mxl_integration package mxlfabrics_test import ( "errors" "os" "sync" "testing" "time" "github.com/qvest-digital/go-mxl/fabrics" "github.com/qvest-digital/go-mxl/mxl" ) const testVideoFlowID = "5fbec3b1-1b0f-417d-9059-8b94a47197ed" const testVideoFlow = `{ "description": "MXL Player Fabrics SHM integration test", "id": "5fbec3b1-1b0f-417d-9059-8b94a47197ed", "format": "urn:x-nmos:format:video", "label": "Fabrics SHM test video", "tags": { "urn:x-nmos:tag:grouphint/v1.0": [ "mxl-player-fabrics-test:Video" ] }, "parents": [], "media_type": "video/v210", "grain_rate": { "numerator": 25, "denominator": 1 }, "frame_width": 1920, "frame_height": 1080, "interlace_mode": "progressive", "colorspace": "BT709", "components": [ { "name": "Y", "width": 1920, "height": 1080, "bit_depth": 10 }, { "name": "Cb", "width": 960, "height": 1080, "bit_depth": 10 }, { "name": "Cr", "width": 960, "height": 1080, "bit_depth": 10 } ] }` func newTestDomain(t *testing.T) *mxl.Instance { t.Helper() domain, err := os.MkdirTemp("/dev/shm", "mxl-player-fabrics-*") if err != nil { t.Fatalf("create temporary MXL domain: %v", err) } t.Cleanup(func() { if err := os.RemoveAll(domain); err != nil { t.Errorf("remove temporary MXL domain: %v", err) } }) instance, err := mxl.NewInstance(domain, "") if err != nil { t.Fatalf("create MXL instance: %v", err) } t.Cleanup(func() { if err := instance.Close(); err != nil { t.Errorf("close MXL instance: %v", err) } }) return instance } func TestGrainTransferSHM(t *testing.T) { sourceInstance := newTestDomain(t) targetInstance := newTestDomain(t) sourceWriter, _, err := sourceInstance.NewWriter(testVideoFlow) if err != nil { t.Fatalf("create source writer: %v", err) } t.Cleanup(func() { if err := sourceWriter.Close(); err != nil { t.Errorf("close source writer: %v", err) } }) sourceReader, err := sourceInstance.NewReader(testVideoFlowID) if err != nil { t.Fatalf("create source reader: %v", err) } t.Cleanup(func() { if err := sourceReader.Close(); err != nil { t.Errorf("close source reader: %v", err) } }) targetWriter, _, err := targetInstance.NewWriter(testVideoFlow) if err != nil { t.Fatalf("create target writer: %v", err) } t.Cleanup(func() { if err := targetWriter.Close(); err != nil { t.Errorf("close target writer: %v", err) } }) sourceFabrics, err := fabrics.NewInstance(sourceInstance) if err != nil { t.Fatalf("create source Fabrics instance: %v", err) } t.Cleanup(func() { if err := sourceFabrics.Close(); err != nil { t.Errorf("close source Fabrics instance: %v", err) } }) targetFabrics, err := fabrics.NewInstance(targetInstance) if err != nil { t.Fatalf("create target Fabrics instance: %v", err) } t.Cleanup(func() { if err := targetFabrics.Close(); err != nil { t.Errorf("close target Fabrics instance: %v", err) } }) target, err := targetFabrics.NewTarget() if err != nil { t.Fatalf("create target: %v", err) } t.Cleanup(func() { if err := target.Close(); err != nil { t.Errorf("close target: %v", err) } }) targetInfo, err := target.Setup(fabrics.TargetConfig{ Interface: requireSHMInterface(t, targetFabrics), Writer: targetWriter, }) if err != nil { t.Fatalf("set up target: %v", err) } t.Cleanup(func() { if err := targetInfo.Close(); err != nil { t.Errorf("close target info: %v", err) } }) initiator, err := sourceFabrics.NewInitiator() if err != nil { t.Fatalf("create initiator: %v", err) } t.Cleanup(func() { if err := initiator.Close(); err != nil { t.Errorf("close initiator: %v", err) } }) if err := initiator.Setup(fabrics.InitiatorConfig{ Interface: requireSHMInterface(t, sourceFabrics), Reader: sourceReader, }); err != nil { t.Fatalf("set up initiator: %v", err) } if err := initiator.AddTarget(targetInfo); err != nil { t.Fatalf("add target: %v", err) } connectDeadline := time.Now().Add(5 * time.Second) for { if time.Now().After(connectDeadline) { t.Fatal("SHM initiator did not become ready") } // The target may also need to advance its endpoint state. _, err := target.ReadGrainNonBlocking() if err != nil && !errors.Is(err, fabrics.ErrNotReady) { t.Fatalf("progress target during setup: %v", err) } err = initiator.MakeProgressNonBlocking() if err == nil { break } if !errors.Is(err, fabrics.ErrNotReady) { t.Fatalf("progress initiator during setup: %v", err) } time.Sleep(time.Millisecond) } index := mxl.CurrentIndex(sourceWriter.Config().Common.GrainRate) sourceGrain, err := sourceWriter.OpenGrain(index) if err != nil { t.Fatalf("open source grain: %v", err) } for offset := range sourceGrain.Payload { sourceGrain.Payload[offset] = byte((uint64(offset) + index) & 0xff) } totalSlices := sourceGrain.TotalSlices if err := sourceGrain.Commit(totalSlices, 0); err != nil { t.Fatalf("commit source grain: %v", err) } deadline := time.Now().Add(5 * time.Second) received := make(chan uint64, 1) targetErrors := make(chan error, 1) stopTarget := make(chan struct{}) var targetWG sync.WaitGroup targetWG.Add(1) go func() { defer targetWG.Done() for { select { case <-stopTarget: return default: } receivedIndex, err := target.ReadGrainNonBlocking() switch { case err == nil: received <- receivedIndex return case errors.Is(err, fabrics.ErrNotReady): time.Sleep(time.Millisecond) default: targetErrors <- err return } } }() defer func() { close(stopTarget) targetWG.Wait() }() // The SHM provider can apply backpressure while the peer progresses its // endpoint. ErrNotReady means the write was not queued, so it is safe to // make progress and retry the transfer. for { err := initiator.TransferGrain(index, 0, totalSlices) if err == nil { break } if !errors.Is(err, fabrics.ErrNotReady) { t.Fatalf("transfer grain: %v", err) } if time.Now().After(deadline) { t.Fatal("timed out enqueueing SHM grain transfer") } err = initiator.MakeProgressNonBlocking() if err != nil && !errors.Is(err, fabrics.ErrNotReady) { t.Fatalf("make initiator progress while enqueueing: %v", err) } time.Sleep(time.Millisecond) } var receivedIndex uint64 for { if time.Now().After(deadline) { t.Fatal("SHM grain transfer timed out") } err := initiator.MakeProgressNonBlocking() if err != nil && !errors.Is(err, fabrics.ErrNotReady) { t.Fatalf("make initiator progress: %v", err) } select { case receivedIndex = <-received: goto receivedGrain case err := <-targetErrors: t.Fatalf("read target completion: %v", err) default: time.Sleep(time.Millisecond) } } receivedGrain: if receivedIndex != index { t.Fatalf( "received grain index = %d, want %d", receivedIndex, index, ) } // Fabrics writes directly into the target writer's mapped grain memory. // Open the completed slot to inspect those bytes without relying on the // local reader head, which the transfer itself does not advance. receivedGrain, err := targetWriter.OpenGrain(index) if err != nil { t.Fatalf("open transferred grain: %v", err) } defer func() { if err := receivedGrain.Cancel(); err != nil { t.Errorf("cancel transferred grain: %v", err) } }() if len(receivedGrain.Payload) == 0 { t.Fatal("transferred grain has an empty payload") } for offset := 0; offset < len(receivedGrain.Payload); offset += 4096 { want := byte((uint64(offset) + index) & 0xff) if got := receivedGrain.Payload[offset]; got != want { t.Fatalf( "payload[%d] = %d, want %d", offset, got, want, ) } } }