flow restart processing
This commit is contained in:
+87
-18
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
]
|
||||
}
|
||||
Reference in New Issue
Block a user