From 84572fb88c569268841b985d2cbbb16ec06ac868 Mon Sep 17 00:00:00 2001 From: Dmitry Sergeev Date: Sat, 22 Aug 2026 14:59:54 +0300 Subject: [PATCH] and the video comes into chat --- cmd/{ => mxl-player}/main.go | 192 ++++++++++++++++++++++- cmd/mxl-player/shaders.go | 9 ++ cmd/mxl-player/shaders/decode.frag | 76 +++++++++ cmd/mxl-player/shaders/decode.frag.spv | Bin 0 -> 7436 bytes cmd/mxl-player/shaders/triangle.vert | 12 ++ cmd/mxl-player/shaders/triangle.vert.spv | Bin 0 -> 1164 bytes cmd/mxl-reader/main.go | 61 +++++++ go.mod | 6 +- go.sum | 4 +- internal/source/source.go | 166 ++++++++++++++++++++ nmos/video.json | 41 +++++ 11 files changed, 554 insertions(+), 13 deletions(-) rename cmd/{ => mxl-player}/main.go (65%) create mode 100644 cmd/mxl-player/shaders.go create mode 100644 cmd/mxl-player/shaders/decode.frag create mode 100644 cmd/mxl-player/shaders/decode.frag.spv create mode 100644 cmd/mxl-player/shaders/triangle.vert create mode 100644 cmd/mxl-player/shaders/triangle.vert.spv create mode 100644 cmd/mxl-reader/main.go create mode 100644 internal/source/source.go create mode 100644 nmos/video.json diff --git a/cmd/main.go b/cmd/mxl-player/main.go similarity index 65% rename from cmd/main.go rename to cmd/mxl-player/main.go index 4123dfd..04dbf53 100644 --- a/cmd/main.go +++ b/cmd/mxl-player/main.go @@ -1,9 +1,14 @@ package main import ( + "context" + "errors" + "flag" "fmt" "log" + "mxl-player/internal/source" "runtime" + "time" "unsafe" vk "github.com/christerso/vulkan-go/vk" @@ -13,8 +18,8 @@ import ( const ( APP_NAME = "MXL Player" APP_VER = "0.0.1" - WIN_WIDTH int32 = 1280 - WIN_HEIGHT int32 = 720 + WIN_WIDTH int32 = 1920 + WIN_HEIGHT int32 = 1080 ) const ( @@ -41,10 +46,10 @@ var ( func sdlError() string { return cstr(sdlGetError()) } -var loaded = false +var sdlLoaded = false func loadSDLMissing() error { - if loaded { + if sdlLoaded { return nil } h, err := purego.Dlopen("libSDL3.so.0", purego.RTLD_NOW|purego.RTLD_GLOBAL) @@ -60,7 +65,7 @@ func loadSDLMissing() error { purego.RegisterLibFunc(&sdlVulkanCreateSurface, h, "SDL_Vulkan_CreateSurface") purego.RegisterLibFunc(&sdlPollEvent, h, "SDL_PollEvent") purego.RegisterLibFunc(&sdlGetWindowSizeInPixels, h, "SDL_GetWindowSizeInPixels") - loaded = true + sdlLoaded = true return nil } @@ -85,6 +90,11 @@ func cbytes(s string) *byte { } func main() { + // TODO: remove flags defaults + mxlDomain := flag.String("d", "/dev/shm/mxl", "MXL domain") + mxlVideoFlowID := flag.String("v", "5fbec3b1-1b0f-417d-9059-8b94a47197ed", "MXL video flow UUID") + flag.Parse() + runtime.LockOSThread() if err := loadSDLMissing(); err != nil { panic(err) @@ -94,7 +104,7 @@ func main() { return } - windowHandler := sdlCreateWindow(cbytes(APP_NAME), WIN_WIDTH, WIN_HEIGHT, windowVulkan|windowResizable) + windowHandler := sdlCreateWindow(cbytes(fmt.Sprintf("%s %s", APP_NAME, APP_VER)), WIN_WIDTH, WIN_HEIGHT, windowVulkan|windowResizable) if windowHandler == 0 { sdlQuit() log.Fatalf("SDL_CreateWindow: %s", sdlError()) @@ -173,7 +183,8 @@ func main() { var vkColorSpace uint32 formats, _ := vkPhysDevice.SurfaceFormats(vkSurf) for _, f := range formats { - if f.Format == vk.FormatB8G8R8A8Srgb && f.ColorSpace == vk.ColorSpaceSRGBNonlinear { + // Prefer 10-bit RGB (A2B10G10R10_UNORM = 64) to preserve V210's 10 bits + if f.Format == vk.Format(64) && f.ColorSpace == vk.ColorSpaceSRGBNonlinear { vkFormat = f.Format vkColorSpace = f.ColorSpace } @@ -208,6 +219,97 @@ func main() { panic(err.Error()) } + // MXL Source + mxlSrc, err := source.Open(*mxlDomain, *mxlVideoFlowID) + if err != nil { + log.Fatalf("source: %v\n", err) + } + defer 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) + // Staging buffer: host-visible, persistently mapped. The reader goroutine + // writes V210 bytes here; the GPU copies from it. + staging, err := vkDevice.CreateBuffer(vkPhysDevice, vk.BufferConfig{ + Size: frameSize, + Usage: vk.BufferUsageTransferSrc, + Properties: vk.MemoryHostVisible | vk.MemoryHostCoherent, + Map: true, + }) + if err != nil { + panic(err) + } + defer vkDevice.DestroyBuffer(staging) + // Device-local V210 buffer: fast for the GPU to read (M3 compute), CPU can't + // write it. Filled each frame by a CopyBuffer from staging. + v210Buf, err := vkDevice.CreateBuffer(vkPhysDevice, vk.BufferConfig{ + Size: frameSize, + Usage: vk.BufferUsageTransferDst | vk.BufferUsageStorageBuffer, + Properties: vk.MemoryDeviceLocal, + Map: false, + }) + if err != nil { + panic(err) + } + defer vkDevice.DestroyBuffer(v210Buf) + + // Decode pipeline: fullscreen triangle, fragment reads V210 from the + // staging buffer and writes 10-bit RGB to the swapchain color attachment. + vertModule, err := vkDevice.CreateShaderModule(vertSPV) + if err != nil { + panic(err) + } + defer vkDevice.DestroyShaderModule(vertModule) + + fragModule, err := vkDevice.CreateShaderModule(fragSPV) + if err != nil { + panic(err) + } + defer vkDevice.DestroyShaderModule(fragModule) + // Descriptor set layout: binding 0 = storage buffer (V210), fragment stage. + decodeDSL, err := vkDevice.CreateDescriptorSetLayout([]vk.DescriptorBinding{ + {Binding: 0, Type: vk.DescriptorStorageBuffer, Count: 1, Stages: vk.ShaderStageFragment}, + }) + if err != nil { + panic(err) + } + defer vkDevice.DestroyDescriptorSetLayout(decodeDSL) + // Pipeline layout: the set layout + push constants {width,height,strideBytes}. + decodeLayout, err := vkDevice.CreatePipelineLayout([]vk.DescriptorSetLayout{decodeDSL}, vk.ShaderStageFragment, 12) + if err != nil { + panic(err) + } + defer vkDevice.DestroyPipelineLayout(decodeLayout) + decodePipeline, err := vkDevice.CreateGraphicsPipeline(vk.GraphicsPipelineConfig{ + Layout: decodeLayout, + RenderPass: vkRenderPass, + VertexShader: vertModule, + FragShader: fragModule, + Topology: vk.TopologyTriangleList, + PolygonMode: vk.PolygonFill, + CullMode: vk.CullNone, + FrontFace: vk.FrontFaceCounterClockwise, + }) + if err != nil { + panic(err) + } + defer vkDevice.DestroyPipeline(decodePipeline) + // Descriptor pool + set, bound once to the staging buffer + descPool, err := vkDevice.CreateDescriptorPool(1, map[vk.DescriptorType]uint32{ + vk.DescriptorStorageBuffer: 1, + }) + if err != nil { + panic(err) + } + defer vkDevice.DestroyDescriptorPool(descPool) + + decodeSet, err := vkDevice.AllocateDescriptorSet(descPool, decodeDSL) + if err != nil { + panic(err) + } + vkDevice.UpdateBufferDescriptor(decodeSet, 0, vk.DescriptorStorageBuffer, staging.Buffer, 0, vk.WholeSize) + + // Swapchain var ( vkExtent vk.Extent2D vkSwapchain vk.SwapchainKHR @@ -350,6 +452,46 @@ func main() { defer vkDevice.DestroyFence(inFlight) defer vkDevice.WaitIdle() + // Reader goroutine: stages V210 bytes into the staging buffer. + // Single-flight handshake: it only writes staging after the render thread + // grants permission (post-WaitFence), so the GPU is never reading it. + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + grant := make(chan struct{}) + staged := make(chan struct{}) + + go func() { + for { + select { + case <-grant: + case <-ctx.Done(): + return + } + f, err := mxlSrc.NextCtx(ctx, 200*time.Millisecond) + if err != nil { + if errors.Is(err, context.Canceled) { + return + } + log.Printf("source: %v", err) + cancel() + return + } + // Copy the borrowed payload into staging BEFORE the next read + // invalidates it. Within the grain's valid lifetime + vk.CopyToMapped(staging.Mapped, f.Payload) + // fmt.Println("staged", f.Index) + select { + case staged <- struct{}{}: + case <-ctx.Done(): + return + } + } + }() + + type decodePushConstants struct { + Width, Height, StrideBytes, WinW, WinH uint32 + } + // Core Loop running := true for running { @@ -382,7 +524,22 @@ func main() { panic(err) } - // Record: clear color & depth attachments + // 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 + } + select { + case <-staged: + case <-ctx.Done(): + running = false + continue + } + + // Record: copy staged V210 into device-local buffer, then clear cmd := vkCommands[0] if err := cmd.Reset(); err != nil { panic(err) @@ -390,6 +547,7 @@ func main() { if err := cmd.Begin(vk.CommandBufferOneTimeSubmit); err != nil { panic(err) } + cmd.CopyBuffer(staging.Buffer, v210Buf.Buffer, frameSize) cmd.BeginRenderPass( vkRenderPass, fbs[imageIndex], @@ -399,6 +557,24 @@ func main() { vk.ClearDepthStencil(1.0, 0), }, ) + cmd.SetViewport(vk.Viewport{ + X: 0, Y: 0, + Width: float32(vkExtent.Width), + Height: float32(vkExtent.Height), + MinDepth: 0, MaxDepth: 1, + }) + cmd.SetScissor(vk.Rect2D{Offset: vk.Offset2D{X: 0, Y: 0}, Extent: vkExtent}) + cmd.BindPipeline(decodePipeline) + cmd.BindDescriptorSet(decodeLayout, 0, decodeSet) + decodePC := decodePushConstants{ + Width: mxlSrc.Width(), + Height: mxlSrc.Height(), + StrideBytes: mxlSrc.Stride(), + WinW: vkExtent.Width, + WinH: vkExtent.Height, + } + cmd.PushConstants(decodeLayout, vk.ShaderStageFragment, 0, unsafe.Pointer(&decodePC), 20) + cmd.Draw(3, 1, 0, 0) cmd.EndRenderPass() if err := cmd.End(); err != nil { panic(err) diff --git a/cmd/mxl-player/shaders.go b/cmd/mxl-player/shaders.go new file mode 100644 index 0000000..00d5cfa --- /dev/null +++ b/cmd/mxl-player/shaders.go @@ -0,0 +1,9 @@ +package main + +import _ "embed" + +//go:embed shaders/triangle.vert.spv +var vertSPV []byte + +//go:embed shaders/decode.frag.spv +var fragSPV []byte diff --git a/cmd/mxl-player/shaders/decode.frag b/cmd/mxl-player/shaders/decode.frag new file mode 100644 index 0000000..ef6e03e --- /dev/null +++ b/cmd/mxl-player/shaders/decode.frag @@ -0,0 +1,76 @@ +#version 450 + +layout(set = 0, binding = 0, std430) readonly buffer V210 { + uint words[]; +}; + +layout(push_constant) uniform PC { + uint width; + uint height; + uint strideBytes; + uint winW; + uint winH; +} pc; + +layout(location = 0) out vec4 fragColor; + +void main() { + // Letterbox: fit video into the window, preserving aspect ratio. + float sx = float(pc.winW) / float(pc.width); + float sy = float(pc.winH) / float(pc.height); + float scale = min(sx, sy); + float dispW = float(pc.width) * scale; + float dispH = float(pc.height) * scale; + float offX = (float(pc.winW) - dispW) * 0.5; + float offY = (float(pc.winH) - dispH) * 0.5; + float fbx = gl_FragCoord.x - offX; + float fby = gl_FragCoord.y - offY; + if (fbx < 0.0 || fbx >= dispW || fby < 0.0 || fby >= dispH) { + fragColor = vec4(0.0, 0.0, 0.0, 1.0); + return; + } + // Framebuffer y is bottom-origin; flip to video top-origin. + uint x = uint(fbx / scale); + uint y = uint(fby / scale); + + // V210: 6 pixels per group of 4 words; 3 ten-bit components per word + // (bits 0-9 / 10-19 / 20-29). Stream order: Cb Y Cr Y Cb Y Cr Y ... + uint group = x / 6u; + uint sub = x % 6u; + uint wordsPerLine = pc.strideBytes / 4u; + uint base = y * wordsPerLine + group * 4u; + + uint w0 = words[base + 0u]; + uint w1 = words[base + 1u]; + uint w2 = words[base + 2u]; + uint w3 = words[base + 3u]; + + uint cb0 = (w0 ) & 0x3FFu; + uint y0 = (w0 >> 10u) & 0x3FFu; + uint cr0 = (w0 >> 20u) & 0x3FFu; + uint y1 = (w1 ) & 0x3FFu; + uint cb1 = (w1 >> 10u) & 0x3FFu; + uint y2 = (w1 >> 20u) & 0x3FFu; + uint cr1 = (w2 ) & 0x3FFu; + uint y3 = (w2 >> 10u) & 0x3FFu; + uint cb2 = (w2 >> 20u) & 0x3FFu; + uint y4 = (w3 ) & 0x3FFu; + uint cr2 = (w3 >> 10u) & 0x3FFu; + uint y5 = (w3 >> 20u) & 0x3FFu; + + float Y, Cb, Cr; + if (sub == 0u) { Y = float(y0); Cb = float(cb0); Cr = float(cr0); } + else if (sub == 1u) { Y = float(y1); Cb = float(cb0); Cr = float(cr0); } + else if (sub == 2u) { Y = float(y2); Cb = float(cb1); Cr = float(cr1); } + else if (sub == 3u) { Y = float(y3); Cb = float(cb1); Cr = float(cr1); } + else if (sub == 4u) { Y = float(y4); Cb = float(cb2); Cr = float(cr2); } + else { Y = float(y5); Cb = float(cb2); Cr = float(cr2); } + + float yf = (Y - 64.0) / 876.0; + float uf = (Cb - 512.0) / 896.0; + float vf = (Cr - 512.0) / 896.0; + float r = yf + 1.5748 * vf; + float g = yf - 0.1873 * uf - 0.4681 * vf; + float b = yf + 1.8556 * uf; + fragColor = vec4(clamp(r, 0.0, 1.0), clamp(g, 0.0, 1.0), clamp(b, 0.0, 1.0), 1.0); +} diff --git a/cmd/mxl-player/shaders/decode.frag.spv b/cmd/mxl-player/shaders/decode.frag.spv new file mode 100644 index 0000000000000000000000000000000000000000..d31e964dfc4fd9d54ddcbc42d8509effcf7d407c GIT binary patch literal 7436 zcmZ9Q37l4C6~-?NGi<{qn;STQ;sPQdAc9D&Lp26Lxnwmm3&YgFOf$nW7Kw?ZmSvh* znH3pYWLZ>XmKGF=Nw!%Ot+sDwyGxn=|M$CRcz@shaX8QOob#Udyzh7KeZToC`cCSf zWqq>2*^q2OpDe#dWPMRGw5r_K)~%_V*45oKZPx7RCJfCg@;q~fW&N{$cx}Vl){O?1 zVvX2v^3<6CR6*4&JHViR*`TbpuD*8V%G$bR^^F}J^Vs7VZ}7T?f8+=8XK_ao^=$c^$mH+>k3Je3P*mC48%x z%`oS7@wp|u3*6W+J=b@Gd-Ay|vOVgZCH-FYqR!qte@33aPraylf1Y2`A5bsq@6YpR z=J^NJXO-qXl;@Z9!|Jn3`eR@}9Q!TW2E3`jCxbT^xCTsm zoQ5%lc^7$>r z55|TSyb(Vh^D_Qcd==(6`}6a)?uz2G2jLRqf=?{bY$U68 z4HFn+eeoEObJyh7FMrqgzEK-Zy@9#Te_Z2>iDy3-2w*jt6`1$gO)l-Z5(Kh1`2b&G?N4PUpR?(4vn! z;I1?Ho#5vYg5M4IUI=~@j1_W`H4BdqL)vESKspR z!pJ`Ze)A&>FAP2v{^*`3d-+WGMLV8YIEUxRb-G7huHk8vO1yjFeh$RK`7fBHjXCU&cLe(s*YX_smJ-o=qwDyh-`aGVU3Y#yv+;?pc!Z zU1fZC8Q)XJJzLW8d&{_IOd9u`Nx5fD%J-M?gJt|s89$tG>v;Y|J-ug8$~}it?pc&_ z&!d!kCZ)WljC(evanC3DBJM)pgPtALnD560+><_sBQW)j;(0l)&|-|wmScQp#+cDy z^%ye-tacsq_{@#PtfPJ*PlEToTGTudY)#+0QS&6Qdel4_tQKQD+Z^M&ImVmy${8X?x)$E$_(=fI$_vRY6j?byF=&knnB#14`TD1!smC1WfYoA->ELt@zIJ8is5{40yjsk` z*Q|Vw^!}KIrXF+f)hbhqIp%=v^_pTG=Yh>pkF}i-R`(or-kEr{n0G$dwWasN1#tD4 zcL7)}=DiSX|1s~vJZAl4^S%hI9`nuvt2ys-d@Z&VTZTp56@?!B#c=zLd0qlGN8Pn7 z!mCBii@?^b!=lE^!0N$Qg6$=0t^%8*ZcX3QYLWkPu=CtU-+mu>1*Wzdi=0=2%?W-n z*w5v`V$3CAbJXuDv`fL(G8W@513Q1r^(wHMbH%vJ!HzW+<6aGRT=3U`{d{8n*MiMa zk2>{WYZ-frb+5&HJU^>kY(rVa~CW{g-blv~IllorM

aufnTEjkkf-BLD4RHJ{Bzyp3QQ4qSS@ON z1MIyMeLn{$N8aPy<@Z-Uj*z8`_BM~y>ZwV3l;;Izi0aJ97WZ^PB2#$m8p z)c6kA_u(Gy*ERU>VrsrWJX^kx_wsBp_C4%pm}iJM>ihugeH`CAeh60c8F>sZkK7-D zeQqN6$6z&ckKpBz`xCJDcI5sPtY+>}y!;t@u}^!lw;20#aJsf%z}4z7*Y-Hx%e5K% zCH4o*I^wAFE3nT?^zdu2TC7bTxxWFYYx^x+E!HNF+~0xIwf!Eh7HgBcHv6?mxh4u{ODDvrl`nw;1~`aJsgC!_|D3ySAtBUarm9f3W{!))7aYr@=n!(Ze%& z&at)0BllTwy0-tq)naY($o(HUUE6bTwOE_nwb`dV*;|Z#9-OYtKa|wAVXmz&+{?8Y z>w^u%tRs#(72t#9#oGG8)$GCNO&+1)+Tpt_GwS{ z7Gnp2eZR!m!Em+cNglb?;QdMN5V%_OB#+#o;IyYHhmU$m3qq%OA(P z2E9DqA$sSo!jGVq`x$?8jD)NCds%!x9t~FaH#2+k_XxG9F&2Czsc`~a&EM3b#z|mx ze^;}HzhS6Fjd9?kNsaMvwY0`5aP_F+?;dJVV*=RTf=>dQ6Mt`>3Z96mU(I>o`cDHp z&X|4G;MJnW>EN{HnQ(iLnrDIiEl@r7a0=LQ#-h&I;Ix-%aDOvYk8$UK9c#>9Cgatj z#td-U%enBhmsxOscT|t{%?3Npm^G*4)uQGcuxqa77Z3O2JTU+A-x!Ro1~*}mKNrou yBL93a|MGle<@xi_oHO$0gZY=|8;kr^{I;?f{Q@+5u&%kzZQqSW%)jsEVlMz_O`!+? literal 0 HcmV?d00001 diff --git a/cmd/mxl-player/shaders/triangle.vert b/cmd/mxl-player/shaders/triangle.vert new file mode 100644 index 0000000..4900f5a --- /dev/null +++ b/cmd/mxl-player/shaders/triangle.vert @@ -0,0 +1,12 @@ +#version 450 + +// Fullscreen triangle, no vertex buffer. One triangle covers the viewport. +vec2 positions[3] = vec2[]( + vec2(-1.0, -1.0), + vec2( 3.0, -1.0), + vec2(-1.0, 3.0) +); + +void main() { + gl_Position = vec4(positions[gl_VertexIndex], 0.0, 1.0); +} diff --git a/cmd/mxl-player/shaders/triangle.vert.spv b/cmd/mxl-player/shaders/triangle.vert.spv new file mode 100644 index 0000000000000000000000000000000000000000..8a2628eff60e250283dc1528ad299c324d1c1ba8 GIT binary patch literal 1164 zcmYk3U279j5Qa~aZd+U1T5DTBYTa1B&{Dim5k#%1tU{obg11Xa7Fk%gA=!$0)qf`V ztGp3>o+KM{!kd|S=bfE7bDE9a`4DEpQdkZT!sx7p226mP8``9O+}xd&wBM(lUN0R~KZ-0Z@-j=ic|Yq^ z`L|5n!jvvAJH=UdS`eX_?iGb7T%?W|Z?Ddac_#s%Qi#07O)ah^H&0%A( zh2EDmHn%~%fQ^CY(Kx<-p$zem2+3uurctv_A$N=Jl5a9nv=KQ*!Ph$zk_uScT(@6h3~23+r#MmwfcMb z_IB21?9w;$jvo0xFcJ0@k#FE`UZSx7O 0 { + lines = f.Size / src.Stride() + } + fmt.Printf("idx=%d size=%d (%dx%d, %d lines) invalid=%v\n", + f.Index, f.Size, f.Width, f.Height, lines, f.Invalid) + grains++ + if *count > 0 && grains >= *count { + fmt.Printf("done: %d grains read\n", grains) + return + } + } +} diff --git a/go.mod b/go.mod index 9ac7495..a1a1008 100644 --- a/go.mod +++ b/go.mod @@ -1,10 +1,10 @@ module mxl-player -go 1.25.0 +go 1.26.0 require ( github.com/christerso/vulkan-go v0.0.0-20260618152204-bff25e5b7646 - github.com/jupiterrider/purego-sdl3 v0.0.0-20260514083405-8523da70a041 + github.com/ebitengine/purego v0.10.2 ) -require github.com/ebitengine/purego v0.10.2 // indirect +require github.com/qvest-digital/go-mxl v1.1.0-rc.1 diff --git a/go.sum b/go.sum index bc0ed0e..971176c 100644 --- a/go.sum +++ b/go.sum @@ -2,5 +2,5 @@ github.com/christerso/vulkan-go v0.0.0-20260618152204-bff25e5b7646 h1:xfc6rFxs+U github.com/christerso/vulkan-go v0.0.0-20260618152204-bff25e5b7646/go.mod h1:DIo6W2yFav1DflNKc0W0Oc7mecbdBg2btsaHhIYmwUA= github.com/ebitengine/purego v0.10.2 h1:W809HbnvzAxgdm+aOvlSekrM16wGCdT/e76+9tS7gzE= github.com/ebitengine/purego v0.10.2/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= -github.com/jupiterrider/purego-sdl3 v0.0.0-20260514083405-8523da70a041 h1:sDTZNtan3t8NsIM15J4vvuFuqxN6qLtWYhmHnLy5TjU= -github.com/jupiterrider/purego-sdl3 v0.0.0-20260514083405-8523da70a041/go.mod h1:pRSNvzaSfMxcVPHP5VKo5VF5SAFp78pPIGlBIoI8KBw= +github.com/qvest-digital/go-mxl v1.1.0-rc.1 h1:CUg29Q5c/JA6pqyA5gP3vu7JBr8lRXbOOxC+r0W572s= +github.com/qvest-digital/go-mxl v1.1.0-rc.1/go.mod h1:cYzyT+S/AONytsKr/DozzFQP2OrNrg/JsQOSjvwn+Rk= diff --git a/internal/source/source.go b/internal/source/source.go new file mode 100644 index 0000000..e03e149 --- /dev/null +++ b/internal/source/source.go @@ -0,0 +1,166 @@ +package source + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + mxl "github.com/qvest-digital/go-mxl/mxl" +) + +type flowDef struct { + 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"` +} + +type Frame struct { + Index uint64 + Width uint32 + Height uint32 + Stride uint32 + Size uint32 + Invalid bool + 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 + idx uint64 +} + +func Open(domain, flowID string) (*Source, error) { + inst, err := mxl.NewInstance(domain, "") + if err != nil { + return nil, fmt.Errorf("NewInstance: %w", err) + } + r, err := inst.NewReader(flowID) + if err != nil { + inst.Close() + return nil, fmt.Errorf("NewReader: %w", err) + } + info, err := r.Info() + if err != nil { + r.Close() + inst.Close() + return nil, fmt.Errorf("Info: %w", err) + } + def, err := inst.FlowDef(flowID) + if err != nil { + r.Close() + inst.Close() + return nil, fmt.Errorf("FlowDef: %w", err) + } + var fd flowDef + if err := json.Unmarshal([]byte(def), &fd); err != nil { + r.Close() + inst.Close() + return nil, fmt.Errorf("parse flow def: %w", err) + } + rate := info.Config.Common.GrainRate + idx := mxl.CurrentIndex(rate) + if idx == mxl.UndefinedIndex { + r.Close() + inst.Close() + return nil, fmt.Errorf("invalid grain rate: %d/%d", rate.Num, rate.Den) + } + 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), + idx: idx, + }, nil +} + +func (s *Source) Close() error { + _ = s.reader.Close() + return s.inst.Close() +} + +func (s *Source) Next(timeout time.Duration) (Frame, error) { + for { + g, err := s.reader.GetGrain(s.idx, timeout) + switch { + case err == nil: + f := Frame{ + Index: g.Index, + Width: s.width, + Height: s.height, + Stride: s.stride, + Size: g.GrainSize, + Invalid: g.Invalid(), + Payload: g.Payload, + } + s.idx++ + return f, nil + case errors.Is(err, mxl.ErrTimeout): + s.idx = mxl.CurrentIndex(s.rate) + case errors.Is(err, mxl.ErrOutOfRangeEarly): + time.Sleep(10 * time.Millisecond) + case errors.Is(err, mxl.ErrOutOfRangeLate): + s.idx = mxl.CurrentIndex(s.rate) + default: + return Frame{}, fmt.Errorf("GetGrain: %w", err) + } + } +} + +func (s *Source) NextCtx(ctx context.Context, timeout time.Duration) (Frame, error) { + for { + select { + case <-ctx.Done(): + return Frame{}, ctx.Err() + default: + } + g, err := s.reader.GetGrain(s.idx, timeout) + switch { + case err == nil: + f := Frame{ + Index: g.Index, + Width: s.width, + Height: s.height, + Stride: s.stride, + Size: g.GrainSize, + Invalid: g.Invalid(), + Payload: g.Payload, + } + s.idx++ + return f, nil + case errors.Is(err, mxl.ErrTimeout): + s.idx = mxl.CurrentIndex(s.rate) + case errors.Is(err, mxl.ErrOutOfRangeEarly): + time.Sleep(10 * time.Millisecond) + case errors.Is(err, mxl.ErrOutOfRangeLate): + s.idx = mxl.CurrentIndex(s.rate) + default: + return Frame{}, fmt.Errorf("GetGrain: %w", 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 } diff --git a/nmos/video.json b/nmos/video.json new file mode 100644 index 0000000..e850f18 --- /dev/null +++ b/nmos/video.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": 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 + } + ] +}