Files
go-mxl-player/internal/adapter/mxlfabrics/grain_transfer_integration_test.go
T
Dmitry Sergeev 3ecca24fe9 MXL Fabrics tests
2026-09-06 16:35:54 +03:00

342 lines
7.7 KiB
Go

//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,
)
}
}
}