From 38417584338d109eb6b2033cbebd784ae2554666 Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Sun, 23 Aug 2026 00:04:06 +0300 Subject: [PATCH] flow restart processing --- cmd/mxl-player/main.go | 105 ++++++++++++++++---- nmos/video.json => flow-def/video-1080.json | 0 flow-def/video-720.json | 41 ++++++++ 3 files changed, 128 insertions(+), 18 deletions(-) rename nmos/video.json => flow-def/video-1080.json (100%) create mode 100644 flow-def/video-720.json diff --git a/cmd/mxl-player/main.go b/cmd/mxl-player/main.go index 79785ea..405a865 100644 --- a/cmd/mxl-player/main.go +++ b/cmd/mxl-player/main.go @@ -13,6 +13,7 @@ import ( vk "github.com/christerso/vulkan-go/vk" "github.com/ebitengine/purego" + "github.com/qvest-digital/go-mxl/mxl" ) const ( @@ -229,7 +230,7 @@ func main() { if err != nil { log.Fatalf("source: %v\n", err) } - defer mxlSrc.Close() + defer func() { _ = mxlSrc.Close() }() frameSize := vk.DeviceSize(mxlSrc.Stride()) * vk.DeviceSize(mxlSrc.Height()) fmt.Printf("source: %dx%d stride=%d frameSize=%d\n", mxlSrc.Width(), mxlSrc.Height(), mxlSrc.Stride(), frameSize) @@ -465,8 +466,54 @@ func main() { // grants permission (post-WaitFence), so the GPU is never reading it ctx, cancel := context.WithCancel(context.Background()) defer cancel() - grant := make(chan struct{}) + grant := make(chan struct{}, 1) staged := make(chan uint64) + failed := make(chan struct{}) + + reopen := func() error { + if err := mxlSrc.Close(); err != nil { + return fmt.Errorf("close old source: %w", err) + } + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + s, err := source.Open(*mxlDomain, *mxlVideoFlowID) + if err == nil { + newSize := vk.DeviceSize(s.Stride()) * vk.DeviceSize(s.Height()) + if newSize != frameSize { + if err := vkDevice.WaitIdle(); err != nil { + return err + } + vkDevice.DestroyBuffer(v210Buf) + vkDevice.DestroyBuffer(staging) + staging, err = vkDevice.CreateBuffer(vkPhysDevice, vk.BufferConfig{ + Size: newSize, Usage: vk.BufferUsageTransferSrc, + Properties: vk.MemoryHostVisible | vk.MemoryHostCoherent, Map: true, + }) + if err != nil { + return err + } + v210Buf, err = vkDevice.CreateBuffer(vkPhysDevice, vk.BufferConfig{ + Size: newSize, Usage: vk.BufferUsageTransferDst | vk.BufferUsageStorageBuffer, + Properties: vk.MemoryDeviceLocal, Map: false, + }) + if err != nil { + return err + } + frameSize = newSize + vkDevice.UpdateBufferDescriptor(decodeSet, 0, vk.DescriptorStorageBuffer, staging.Buffer, 0, vk.WholeSize) + fmt.Printf("source: resolution changed, frameSize=%d\n", newSize) + } + mxlSrc = s + return nil + } + log.Printf("source: reopen retry: %v", err) + time.Sleep(500 * time.Millisecond) + } + } go func() { for { @@ -480,6 +527,20 @@ func main() { if errors.Is(err, context.Canceled) { return } + if errors.Is(err, mxl.ErrFlowInvalid) { + log.Printf("source: flow invalid, reopening") + if rerr := reopen(); rerr != nil { + log.Printf("source: reopen failed: %v", rerr) + cancel() + return + } + select { + case failed <- struct{}{}: + case <-ctx.Done(): + return + } + continue + } log.Printf("source: %v", err) cancel() return @@ -503,6 +564,7 @@ func main() { // Core Loop running := true resized := false + granted := false fullscreen := false var ( lastIndex uint64 @@ -546,7 +608,29 @@ func main() { } resized = false } - // Acquire the next swapchain image + if !granted { + select { + case grant <- struct{}{}: + granted = true + case <-ctx.Done(): + running = false + continue + } + } + var shownIndex uint64 + select { + case shownIndex = <-staged: + granted = false + case <-failed: + granted = false + case <-ctx.Done(): + running = false + continue + case <-time.After(100 * time.Millisecond): + continue + } + + // Acquire the next swapchain image. Right after we have a frame imageIndex, r := vkDevice.AcquireNextImage(vkSwapchain, imageAvailable, ^uint64(0)) if r == vk.ErrorOutOfDateKHR || r == vk.SuboptimalKHR { if err := recreateSwapChain(); err != nil { @@ -569,21 +653,6 @@ func main() { panic(err) } - // Grant the goroutine permission to write staging (GPU is idle now), - // then wait for it to stage a new payload. - select { - case grant <- struct{}{}: - case <-ctx.Done(): - running = false - continue - } - var shownIndex uint64 - select { - case <-staged: - case <-ctx.Done(): - running = false - continue - } if lastIndex != 0 && shownIndex > lastIndex { if g := shownIndex - lastIndex - 1; g > 0 { dropped += g diff --git a/nmos/video.json b/flow-def/video-1080.json similarity index 100% rename from nmos/video.json rename to flow-def/video-1080.json diff --git a/flow-def/video-720.json b/flow-def/video-720.json new file mode 100644 index 0000000..01d1e5d --- /dev/null +++ b/flow-def/video-720.json @@ -0,0 +1,41 @@ +{ + "description": "sample for mxl reader go player", + "id": "5fbec3b1-1b0f-417d-9059-8b94a47197ed", + "tags": { + "urn:x-nmos:tag:grouphint/v1.0": [ + "mxl-gst-testsrc pattern" + ] + }, + "format": "urn:x-nmos:format:video", + "label": "SMPTE bars test video", + "parents": [], + "media_type": "video/v210", + "grain_rate": { + "numerator": 25, + "denominator": 1 + }, + "frame_width": 1280, + "frame_height": 720, + "interlace_mode": "progressive", + "colorspace": "BT709", + "components": [ + { + "name": "Y", + "width": 1280, + "height": 720, + "bit_depth": 10 + }, + { + "name": "Cb", + "width": 640, + "height": 720, + "bit_depth": 10 + }, + { + "name": "Cr", + "width": 640, + "height": 720, + "bit_depth": 10 + } + ] +}