[client/idd/obs] scheduler: isolate timing owner delivery

This commit is contained in:
Geoffrey McRae
2026-08-05 14:58:21 +10:00
parent 0b9b805530
commit c7ed04412e
21 changed files with 1572 additions and 362 deletions

View File

@@ -135,7 +135,9 @@ typedef struct LG_TransportFrame
uint64_t serial; uint64_t serial;
uint64_t timestamp; uint64_t timestamp;
uint32_t scheduleGeneration; uint32_t scheduleGeneration;
uint32_t scheduleEpoch;
LG_TransportFrameFlags flags; LG_TransportFrameFlags flags;
bool scheduleOwner;
// Backend-owned immutable metadata, valid until releaseFrame. // Backend-owned immutable metadata, valid until releaseFrame.
const LG_TransportFrameFormat * format; const LG_TransportFrameFormat * format;
@@ -207,6 +209,7 @@ typedef struct LG_TransportControl
uint64_t targetSlack; uint64_t targetSlack;
int64_t phaseError; int64_t phaseError;
uint32_t feedbackFrameSerial; uint32_t feedbackFrameSerial;
uint32_t feedbackScheduleEpoch;
uint32_t lease; uint32_t lease;
} }
frameSchedule; frameSchedule;

View File

@@ -171,10 +171,13 @@ bool egl_texBufferStreamInit(EGL_Texture ** texture, EGL_TexType type,
{ {
case EGL_TEXTYPE_BUFFER_STREAM: case EGL_TEXTYPE_BUFFER_STREAM:
case EGL_TEXTYPE_FRAMEBUFFER: case EGL_TEXTYPE_FRAMEBUFFER:
case EGL_TEXTYPE_DMABUF:
this->texCount = 2; this->texCount = 2;
break; break;
case EGL_TEXTYPE_DMABUF:
this->texCount = LGMP_Q_FRAME_BUFFER_LEN;
break;
case EGL_TEXTYPE_BUFFER_MAP: case EGL_TEXTYPE_BUFFER_MAP:
this->texCount = 1; this->texCount = 1;
break; break;

View File

@@ -22,9 +22,10 @@
#include "texture.h" #include "texture.h"
#include "texture_util.h" #include "texture_util.h"
#include "common/LGMPConfig.h"
#include "common/locking.h" #include "common/locking.h"
#define EGL_TEX_BUFFER_MAX 2 #define EGL_TEX_BUFFER_MAX LGMP_Q_FRAME_BUFFER_LEN
typedef struct TextureBuffer typedef struct TextureBuffer
{ {

View File

@@ -42,7 +42,7 @@ typedef struct TexDMABUF
EGLDisplay display; EGLDisplay display;
struct FdImage images[2]; struct FdImage images[EGL_TEX_BUFFER_MAX];
int lastIndex; int lastIndex;
int renderIndex; int renderIndex;
@@ -220,11 +220,28 @@ static bool egl_texDMABUFUpdate(EGL_Texture * texture,
DEBUG_ASSERT(update->type == EGL_TEXTYPE_DMABUF); DEBUG_ASSERT(update->type == EGL_TEXTYPE_DMABUF);
struct FdImage *fdImage = struct FdImage *fdImage = NULL;
(this->images[0].fd == update->dmaFD) ? &this->images[0] : for (int i = 0; i < ARRAY_LENGTH(this->images); ++i)
(this->images[1].fd == update->dmaFD) ? &this->images[1] : if (this->images[i].fd == update->dmaFD)
(this->images[0].fd == -1) ? &this->images[0] : {
&this->images[1]; fdImage = &this->images[i];
break;
}
if (!fdImage)
for (int i = 0; i < ARRAY_LENGTH(this->images); ++i)
if (this->images[i].fd == -1)
{
fdImage = &this->images[i];
break;
}
if (!fdImage)
{
DEBUG_ERROR("No free DMABUF image slot");
return false;
}
EGLImage image = fdImage->image; EGLImage image = fdImage->image;
if (unlikely(image == EGL_NO_IMAGE)) if (unlikely(image == EGL_NO_IMAGE))
{ {
@@ -255,7 +272,7 @@ static bool egl_texDMABUFUpdate(EGL_Texture * texture,
fdImage->fd = update->dmaFD; fdImage->fd = update->dmaFD;
fdImage->image = image; fdImage->image = image;
int slot = (fdImage == &this->images[0]) ? 0 : 1; const int slot = (int)(fdImage - this->images);
fdImage->texIndex = slot; fdImage->texIndex = slot;
GLsync sync = 0; GLsync sync = 0;
INTERLOCKED_SECTION(parent->copyLock, INTERLOCKED_SECTION(parent->copyLock,
@@ -287,7 +304,7 @@ static bool egl_texDMABUFUpdate(EGL_Texture * texture,
INTERLOCKED_SECTION(parent->copyLock, INTERLOCKED_SECTION(parent->copyLock,
{ {
fdImage->frameToken = update->frameToken; fdImage->frameToken = update->frameToken;
this->lastIndex = (fdImage == &this->images[0]) ? 0 : 1; this->lastIndex = (int)(fdImage - this->images);
}); });
return true; return true;

View File

@@ -43,6 +43,7 @@ static struct
_Atomic(bool) supported; _Atomic(bool) supported;
_Atomic(bool) active; _Atomic(bool) active;
bool controlPending; bool controlPending;
bool immediatePending;
_Atomic(uint32_t) generation; _Atomic(uint32_t) generation;
_Atomic(uint64_t) period; _Atomic(uint64_t) period;
uint64_t lastSend; uint64_t lastSend;
@@ -50,6 +51,7 @@ static struct
int64_t phaseError; int64_t phaseError;
uint32_t feedbackFrameSerial; uint32_t feedbackFrameSerial;
uint32_t feedbackScheduleEpoch;
unsigned feedbackSamples; unsigned feedbackSamples;
bool feedbackDirty; bool feedbackDirty;
@@ -89,9 +91,11 @@ static bool sendSchedule(LG_TransportFrameScheduleFlags flags,
int64_t phaseError; int64_t phaseError;
uint32_t feedbackFrameSerial; uint32_t feedbackFrameSerial;
uint32_t feedbackScheduleEpoch;
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
phaseError = l_frameScheduler.phaseError; phaseError = l_frameScheduler.phaseError;
feedbackFrameSerial = l_frameScheduler.feedbackFrameSerial; feedbackFrameSerial = l_frameScheduler.feedbackFrameSerial;
feedbackScheduleEpoch = l_frameScheduler.feedbackScheduleEpoch;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
const LG_TransportControl control = { const LG_TransportControl control = {
@@ -103,6 +107,7 @@ static bool sendSchedule(LG_TransportFrameScheduleFlags flags,
.targetSlack = FRAME_SCHEDULER_TARGET_SLACK_NS, .targetSlack = FRAME_SCHEDULER_TARGET_SLACK_NS,
.phaseError = phaseError, .phaseError = phaseError,
.feedbackFrameSerial = feedbackFrameSerial, .feedbackFrameSerial = feedbackFrameSerial,
.feedbackScheduleEpoch = feedbackScheduleEpoch,
.lease = FRAME_SCHEDULER_LEASE_MS, .lease = FRAME_SCHEDULER_LEASE_MS,
}, },
}; };
@@ -118,8 +123,11 @@ static bool sendSchedule(LG_TransportFrameScheduleFlags flags,
} }
l_frameScheduler.controlPending = true; l_frameScheduler.controlPending = true;
if (flags & LG_TRANSPORT_FRAME_SCHEDULE_IMMEDIATE)
l_frameScheduler.immediatePending = false;
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
if (l_frameScheduler.feedbackFrameSerial == feedbackFrameSerial) if (l_frameScheduler.feedbackFrameSerial == feedbackFrameSerial &&
l_frameScheduler.feedbackScheduleEpoch == feedbackScheduleEpoch)
l_frameScheduler.feedbackDirty = false; l_frameScheduler.feedbackDirty = false;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
return true; return true;
@@ -142,6 +150,7 @@ void frameScheduler_start(LG_TransportFeatureFlags features)
features & LG_TRANSPORT_FEATURE_FRAME_SCHEDULE; features & LG_TRANSPORT_FEATURE_FRAME_SCHEDULE;
l_frameScheduler.active = false; l_frameScheduler.active = false;
l_frameScheduler.controlPending = false; l_frameScheduler.controlPending = false;
l_frameScheduler.immediatePending = true;
l_frameScheduler.lastSend = 0; l_frameScheduler.lastSend = 0;
l_frameScheduler.lastCadence = 0; l_frameScheduler.lastCadence = 0;
++l_frameScheduler.generation; ++l_frameScheduler.generation;
@@ -149,6 +158,7 @@ void frameScheduler_start(LG_TransportFeatureFlags features)
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
l_frameScheduler.phaseError = 0; l_frameScheduler.phaseError = 0;
l_frameScheduler.feedbackFrameSerial = 0; l_frameScheduler.feedbackFrameSerial = 0;
l_frameScheduler.feedbackScheduleEpoch = 0;
l_frameScheduler.feedbackSamples = 0; l_frameScheduler.feedbackSamples = 0;
l_frameScheduler.feedbackDirty = false; l_frameScheduler.feedbackDirty = false;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
@@ -162,6 +172,7 @@ void frameScheduler_stop(void)
l_frameScheduler.supported = false; l_frameScheduler.supported = false;
l_frameScheduler.active = false; l_frameScheduler.active = false;
l_frameScheduler.controlPending = false; l_frameScheduler.controlPending = false;
l_frameScheduler.immediatePending = false;
} }
void frameScheduler_update(void) void frameScheduler_update(void)
@@ -181,9 +192,11 @@ void frameScheduler_update(void)
{ {
l_frameScheduler.active = false; l_frameScheduler.active = false;
l_frameScheduler.period = 0; l_frameScheduler.period = 0;
l_frameScheduler.immediatePending = true;
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
l_frameScheduler.phaseError = 0; l_frameScheduler.phaseError = 0;
l_frameScheduler.feedbackFrameSerial = 0; l_frameScheduler.feedbackFrameSerial = 0;
l_frameScheduler.feedbackScheduleEpoch = 0;
l_frameScheduler.feedbackSamples = 0; l_frameScheduler.feedbackSamples = 0;
l_frameScheduler.feedbackDirty = false; l_frameScheduler.feedbackDirty = false;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
@@ -204,9 +217,11 @@ void frameScheduler_update(void)
{ {
l_frameScheduler.period = period; l_frameScheduler.period = period;
++l_frameScheduler.generation; ++l_frameScheduler.generation;
l_frameScheduler.immediatePending = true;
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
l_frameScheduler.phaseError = 0; l_frameScheduler.phaseError = 0;
l_frameScheduler.feedbackFrameSerial = 0; l_frameScheduler.feedbackFrameSerial = 0;
l_frameScheduler.feedbackScheduleEpoch = 0;
l_frameScheduler.feedbackSamples = 0; l_frameScheduler.feedbackSamples = 0;
l_frameScheduler.feedbackDirty = false; l_frameScheduler.feedbackDirty = false;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
@@ -215,7 +230,8 @@ void frameScheduler_update(void)
l_frameScheduler.period = l_frameScheduler.period =
(l_frameScheduler.period * 7 + period) / 8; (l_frameScheduler.period * 7 + period) / 8;
if (l_frameScheduler.active && !reset) if (l_frameScheduler.active && !reset &&
!l_frameScheduler.immediatePending)
{ {
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
const bool feedbackDirty = l_frameScheduler.feedbackDirty; const bool feedbackDirty = l_frameScheduler.feedbackDirty;
@@ -230,6 +246,8 @@ void frameScheduler_update(void)
LG_TRANSPORT_FRAME_SCHEDULE_ACTIVE; LG_TRANSPORT_FRAME_SCHEDULE_ACTIVE;
if (reset) if (reset)
flags |= LG_TRANSPORT_FRAME_SCHEDULE_RESET; flags |= LG_TRANSPORT_FRAME_SCHEDULE_RESET;
if (l_frameScheduler.immediatePending)
flags |= LG_TRANSPORT_FRAME_SCHEDULE_IMMEDIATE;
if (sendSchedule(flags, l_frameScheduler.period)) if (sendSchedule(flags, l_frameScheduler.period))
{ {
@@ -239,9 +257,9 @@ void frameScheduler_update(void)
} }
void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation, void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation,
uint64_t measuredPhase) uint32_t scheduleEpoch, uint64_t measuredPhase)
{ {
if (!generation) if (!generation || !scheduleEpoch)
return; return;
LG_LOCK(l_frameScheduler.lock); LG_LOCK(l_frameScheduler.lock);
@@ -259,6 +277,13 @@ void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation,
return; return;
} }
if (l_frameScheduler.feedbackScheduleEpoch &&
l_frameScheduler.feedbackScheduleEpoch != scheduleEpoch)
{
l_frameScheduler.phaseError = 0;
l_frameScheduler.feedbackSamples = 0;
}
int64_t error = measuredPhase > FRAME_SCHEDULER_TARGET_SLACK_NS ? int64_t error = measuredPhase > FRAME_SCHEDULER_TARGET_SLACK_NS ?
(int64_t)(measuredPhase - FRAME_SCHEDULER_TARGET_SLACK_NS) : (int64_t)(measuredPhase - FRAME_SCHEDULER_TARGET_SLACK_NS) :
-(int64_t)(FRAME_SCHEDULER_TARGET_SLACK_NS - measuredPhase); -(int64_t)(FRAME_SCHEDULER_TARGET_SLACK_NS - measuredPhase);
@@ -277,6 +302,7 @@ void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation,
++l_frameScheduler.feedbackSamples; ++l_frameScheduler.feedbackSamples;
l_frameScheduler.feedbackFrameSerial = (uint32_t)frameSerial; l_frameScheduler.feedbackFrameSerial = (uint32_t)frameSerial;
l_frameScheduler.feedbackScheduleEpoch = scheduleEpoch;
l_frameScheduler.feedbackDirty = true; l_frameScheduler.feedbackDirty = true;
LG_UNLOCK(l_frameScheduler.lock); LG_UNLOCK(l_frameScheduler.lock);
} }

View File

@@ -31,6 +31,6 @@ void frameScheduler_start(LG_TransportFeatureFlags features);
void frameScheduler_stop(void); void frameScheduler_stop(void);
void frameScheduler_update(void); void frameScheduler_update(void);
void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation, void frameScheduler_feedback(uint64_t frameSerial, uint32_t generation,
uint64_t measuredPhase); uint32_t scheduleEpoch, uint64_t measuredPhase);
#endif #endif

View File

@@ -214,8 +214,10 @@ struct FrameTimingRecord
LG_RendererFrameToken token; LG_RendererFrameToken token;
uint64_t frameSerial; uint64_t frameSerial;
uint32_t scheduleGeneration; uint32_t scheduleGeneration;
uint32_t scheduleEpoch;
unsigned readyMask; unsigned readyMask;
bool producerValid; bool producerValid;
bool scheduleOwner;
uint64_t captureTime; uint64_t captureTime;
uint64_t postProcessTime; uint64_t postProcessTime;
@@ -326,8 +328,9 @@ static void frameTimingCancel(LG_RendererFrameToken token)
} }
static void frameTimingQueue(LG_RendererFrameToken token, uint64_t frameSerial, static void frameTimingQueue(LG_RendererFrameToken token, uint64_t frameSerial,
uint32_t scheduleGeneration, uint64_t importTime, uint32_t scheduleGeneration, uint32_t scheduleEpoch, bool scheduleOwner,
uint64_t importWaitTime, uint64_t dispatchStart, uint64_t queueStart) uint64_t importTime, uint64_t importWaitTime, uint64_t dispatchStart,
uint64_t queueStart)
{ {
INTERLOCKED_SECTION(l_frameTiming.lock, { INTERLOCKED_SECTION(l_frameTiming.lock, {
struct FrameTimingRecord * record = frameTimingRecord(token); struct FrameTimingRecord * record = frameTimingRecord(token);
@@ -342,6 +345,8 @@ static void frameTimingQueue(LG_RendererFrameToken token, uint64_t frameSerial,
record->queueStart = queueStart; record->queueStart = queueStart;
record->frameSerial = frameSerial; record->frameSerial = frameSerial;
record->scheduleGeneration = scheduleGeneration; record->scheduleGeneration = scheduleGeneration;
record->scheduleEpoch = scheduleEpoch;
record->scheduleOwner = scheduleOwner;
if (record->timestamp < queueStart) if (record->timestamp < queueStart)
record->timestamp = queueStart; record->timestamp = queueStart;
} }
@@ -390,6 +395,8 @@ static void frameTimingFinishRender(const LG_RendererFrameTiming * timing,
uint64_t feedbackFrameSerial = 0; uint64_t feedbackFrameSerial = 0;
uint64_t feedbackQueueStart = 0; uint64_t feedbackQueueStart = 0;
uint32_t feedbackGeneration = 0; uint32_t feedbackGeneration = 0;
uint32_t feedbackEpoch = 0;
bool feedbackOwner = false;
LG_LOCK(l_frameTiming.lock); LG_LOCK(l_frameTiming.lock);
if (l_frameTiming.retireToken <= timing->frameToken) if (l_frameTiming.retireToken <= timing->frameToken)
@@ -411,6 +418,8 @@ static void frameTimingFinishRender(const LG_RendererFrameTiming * timing,
feedbackFrameSerial = record->frameSerial; feedbackFrameSerial = record->frameSerial;
feedbackGeneration = record->scheduleGeneration; feedbackGeneration = record->scheduleGeneration;
feedbackEpoch = record->scheduleEpoch;
feedbackOwner = record->scheduleOwner;
feedbackQueueStart = record->queueStart; feedbackQueueStart = record->queueStart;
if (unlikely( if (unlikely(
@@ -434,11 +443,12 @@ static void frameTimingFinishRender(const LG_RendererFrameTiming * timing,
} }
LG_UNLOCK(l_frameTiming.lock); LG_UNLOCK(l_frameTiming.lock);
if (g_state.jitRender && feedbackFrameSerial && feedbackGeneration && if (g_state.jitRender && feedbackOwner && feedbackFrameSerial &&
feedbackQueueStart && prepareStart >= feedbackQueueStart) feedbackGeneration && feedbackEpoch && feedbackQueueStart &&
prepareStart >= feedbackQueueStart)
{ {
frameScheduler_feedback( frameScheduler_feedback(
feedbackFrameSerial, feedbackGeneration, feedbackFrameSerial, feedbackGeneration, feedbackEpoch,
prepareStart - feedbackQueueStart); prepareStart - feedbackQueueStart);
} }
} }
@@ -936,7 +946,8 @@ int main_frameThread(void * unused)
break; break;
} }
if (frame.serial == frameSerial && g_state.formatValid) if (frame.serial == frameSerial && g_state.formatValid &&
!frame.scheduleOwner)
{ {
g_state.transportOps->releaseFrame(g_state.transport, &frame); g_state.transportOps->releaseFrame(g_state.transport, &frame);
continue; continue;
@@ -1113,8 +1124,8 @@ int main_frameThread(void * unused)
memory_order_release); memory_order_release);
#endif #endif
frameTimingQueue(frameToken, frame.serial, frame.scheduleGeneration, frameTimingQueue(frameToken, frame.serial, frame.scheduleGeneration,
g_state.frameImportTime, g_state.frameImportWaitTime, frame.scheduleEpoch, frame.scheduleOwner, g_state.frameImportTime,
dispatchStart, queueStart); g_state.frameImportWaitTime, dispatchStart, queueStart);
if (g_state.jitRender) if (g_state.jitRender)
{ {

View File

@@ -50,6 +50,8 @@ struct LG_Transport
struct IVSHMEM shm; struct IVSHMEM shm;
PLGMPClient client; PLGMPClient client;
PLGMPClientQueue frameQueue; PLGMPClientQueue frameQueue;
PLGMPClientQueue ownerFrameQueue[LGMP_Q_FRAME_LEN];
PLGMPClientQueue pendingFrameQueue;
PLGMPClientQueue pointerQueue; PLGMPClientQueue pointerQueue;
LG_Lock pointerLock; LG_Lock pointerLock;
@@ -58,13 +60,15 @@ struct LG_Transport
bool allowDMA; bool allowDMA;
bool connected; bool connected;
bool framePending; bool framePending;
bool frameScheduleSupported;
const KVMFRFrame * pendingFrame; const KVMFRFrame * pendingFrame;
uint32_t clientID; uint32_t clientID;
uint32_t frameSerial; uint32_t frameSerial;
bool frameSerialValid;
bool formatValid; bool formatValid;
LG_TransportFrameFormat format; LG_TransportFrameFormat format;
struct DMAFrameInfo dma[LGMP_Q_FRAME_LEN]; struct DMAFrameInfo dma[LGMP_Q_FRAME_BUFFER_LEN];
uint8_t * pointerData; uint8_t * pointerData;
size_t pointerDataSize; size_t pointerDataSize;
}; };
@@ -180,7 +184,7 @@ static bool lgmp_create(LG_Transport ** result)
this->framePollInterval = framePoll; this->framePollInterval = framePoll;
this->cursorPollInterval = cursorPoll; this->cursorPollInterval = cursorPoll;
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
this->dma[i].fd = -1; this->dma[i].fd = -1;
LG_LOCK_INIT(this->pointerLock); LG_LOCK_INIT(this->pointerLock);
@@ -210,24 +214,38 @@ static bool lgmp_create(LG_Transport ** result)
static void lgmp_stopFrame(struct LG_Transport * this) static void lgmp_stopFrame(struct LG_Transport * this)
{ {
if (this->framePending && this->frameQueue) if (this->framePending && this->pendingFrameQueue)
{ {
const LGMP_STATUS status = lgmpClientMessageDone(this->frameQueue); const LGMP_STATUS status =
lgmpClientMessageDone(this->pendingFrameQueue);
if (status != LGMP_OK) if (status != LGMP_OK)
DEBUG_WARN("Failed to release pending LGMP frame: %s", DEBUG_WARN("Failed to release pending LGMP frame: %s",
lgmpStatusString(status)); lgmpStatusString(status));
} }
this->framePending = false; this->framePending = false;
this->pendingFrame = NULL; this->pendingFrame = NULL;
this->pendingFrameQueue = NULL;
LGMP_STATUS status = lgmpClientUnsubscribe(&this->frameQueue); LGMP_STATUS status = lgmpClientUnsubscribe(&this->frameQueue);
if (status != LGMP_OK) if (status != LGMP_OK)
{ {
DEBUG_WARN("Failed to unsubscribe from the LGMP frame queue: %s", DEBUG_WARN("Failed to unsubscribe from the shared LGMP frame queue: %s",
lgmpStatusString(status)); lgmpStatusString(status));
this->frameQueue = NULL; this->frameQueue = NULL;
} }
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
status = lgmpClientUnsubscribe(&this->ownerFrameQueue[i]);
if (status != LGMP_OK)
{
DEBUG_WARN("Failed to unsubscribe from owner LGMP frame queue %u: %s",
i, lgmpStatusString(status));
this->ownerFrameQueue[i] = NULL;
}
}
this->frameSerial = 0; this->frameSerial = 0;
this->frameSerialValid = false;
this->formatValid = false; this->formatValid = false;
} }
@@ -252,7 +270,7 @@ static void lgmp_closeQueues(struct LG_Transport * this)
static void lgmp_closeDMA(struct LG_Transport * this) static void lgmp_closeDMA(struct LG_Transport * this)
{ {
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
{ {
if (this->dma[i].fd >= 0) if (this->dma[i].fd >= 0)
close(this->dma[i].fd); close(this->dma[i].fd);
@@ -371,7 +389,10 @@ static LG_TransportStatus lgmp_connect(LG_Transport * this,
if (!lgmp_parseSession(data, size, session)) if (!lgmp_parseSession(data, size, session))
return LG_TRANSPORT_INVALID_VERSION; return LG_TRANSPORT_INVALID_VERSION;
this->connected = true; this->connected = true;
this->frameScheduleSupported =
session->features & LG_TRANSPORT_FEATURE_FRAME_SCHEDULE;
this->frameSerial = 0; this->frameSerial = 0;
this->frameSerialValid = false;
this->formatValid = false; this->formatValid = false;
return LG_TRANSPORT_OK; return LG_TRANSPORT_OK;
@@ -393,6 +414,7 @@ static void lgmp_disconnect(LG_Transport * this)
lgmp_closeQueues(this); lgmp_closeQueues(this);
lgmp_closeDMA(this); lgmp_closeDMA(this);
this->connected = false; this->connected = false;
this->frameScheduleSupported = false;
this->clientID = 0; this->clientID = 0;
} }
@@ -457,11 +479,119 @@ static LG_TransportStatus lgmp_process(PLGMPClientQueue queue,
} }
} }
struct LGMPFrameMessage
{
PLGMPClientQueue queue;
LGMPMessage message;
const KVMFRFrame * frame;
bool owner;
};
static LG_TransportStatus lgmp_pollFrameQueue(PLGMPClientQueue queue,
bool owner, struct LGMPFrameMessage * result)
{
const LGMP_STATUS advance = lgmpClientAdvanceToLast(queue);
switch (advance)
{
case LGMP_OK:
break;
case LGMP_ERR_QUEUE_EMPTY:
break;
case LGMP_ERR_INVALID_SESSION:
return LG_TRANSPORT_DISCONNECTED;
default:
DEBUG_ERROR("lgmpClientAdvanceToLast failed: %s",
lgmpStatusString(advance));
return LG_TRANSPORT_ERROR;
}
const LG_TransportStatus status = lgmp_process(queue, 0, &result->message);
if (status == LG_TRANSPORT_OK)
{
result->queue = queue;
result->owner = owner || result->message.udata != 0;
}
return status;
}
static LG_TransportStatus lgmp_doneFrameMessage(
struct LGMPFrameMessage * message)
{
if (!message || !message->queue)
return LG_TRANSPORT_OK;
const LGMP_STATUS status = lgmpClientMessageDone(message->queue);
message->queue = NULL;
message->frame = NULL;
if (status == LGMP_OK)
return LG_TRANSPORT_OK;
if (status == LGMP_ERR_INVALID_SESSION)
return LG_TRANSPORT_DISCONNECTED;
DEBUG_WARN("Failed to release discarded LGMP frame: %s",
lgmpStatusString(status));
return LG_TRANSPORT_ERROR;
}
static void lgmp_mergeFrameStatus(LG_TransportStatus status,
LG_TransportStatus * result)
{
if (status != LG_TRANSPORT_OK &&
(*result == LG_TRANSPORT_OK ||
status == LG_TRANSPORT_DISCONNECTED))
*result = status;
}
static bool lgmp_validateFrameMessage(struct LGMPFrameMessage * message)
{
if (message->message.size < sizeof(KVMFRFrame))
{
DEBUG_ERROR("LGMP frame payload is too small");
return false;
}
const KVMFRFrame * frame = (const KVMFRFrame *)message->message.mem;
const size_t frameDataSize = (size_t)frame->dataHeight * frame->pitch;
if (frame->type <= FRAME_TYPE_INVALID || frame->type >= FRAME_TYPE_MAX ||
frame->offset > message->message.size - sizeof(FrameBuffer) ||
frameDataSize >
message->message.size - frame->offset - sizeof(FrameBuffer))
{
DEBUG_ERROR("LGMP frame payload contains invalid dimensions or offsets");
return false;
}
message->frame = frame;
return true;
}
static bool lgmp_frameSerialNewer(uint32_t lhs, uint32_t rhs)
{
return lhs != rhs &&
(uint32_t)(lhs - rhs) < UINT32_C(0x80000000);
}
static void lgmp_selectNewestFrameMessage(
struct LGMPFrameMessage * candidate,
struct LGMPFrameMessage ** selected)
{
if (!candidate->frame)
return;
if (!*selected ||
lgmp_frameSerialNewer(candidate->frame->frameSerial,
(*selected)->frame->frameSerial) ||
(candidate->frame->frameSerial == (*selected)->frame->frameSerial &&
candidate->owner && !(*selected)->owner))
*selected = candidate;
}
static int lgmp_getDMA(struct LG_Transport * this, const KVMFRFrame * frame, static int lgmp_getDMA(struct LG_Transport * this, const KVMFRFrame * frame,
size_t dataSize) size_t dataSize)
{ {
struct DMAFrameInfo * dma = NULL; struct DMAFrameInfo * dma = NULL;
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
if (this->dma[i].frame == frame) if (this->dma[i].frame == frame)
{ {
dma = &this->dma[i]; dma = &this->dma[i];
@@ -474,7 +604,7 @@ static int lgmp_getDMA(struct LG_Transport * this, const KVMFRFrame * frame,
} }
if (!dma) if (!dma)
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
if (!this->dma[i].frame) if (!this->dma[i].frame)
{ {
dma = &this->dma[i]; dma = &this->dma[i];
@@ -502,48 +632,144 @@ static LG_TransportStatus lgmp_nextFrame(LG_Transport * this, bool useDMA,
if (this->framePending) if (this->framePending)
return LG_TRANSPORT_ERROR; return LG_TRANSPORT_ERROR;
LG_TransportStatus status = lgmp_subscribe(this->client, LGMP_Q_FRAME, if (this->frameScheduleSupported)
&this->frameQueue); for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (status != LG_TRANSPORT_OK)
{ {
if (status == LG_TRANSPORT_TIMEOUT) const LG_TransportStatus status = lgmp_subscribe(this->client,
usleep(1000); LGMP_Q_FRAME_OWNER + i, &this->ownerFrameQueue[i]);
if (status != LG_TRANSPORT_OK && status != LG_TRANSPORT_TIMEOUT)
return status; return status;
} }
LGMPMessage message; const LG_TransportStatus sharedSubscribe = lgmp_subscribe(this->client,
status = lgmp_process(this->frameQueue, this->framePollInterval, &message); LGMP_Q_FRAME, &this->frameQueue);
if (status != LG_TRANSPORT_OK) if (sharedSubscribe != LG_TRANSPORT_OK &&
return status; sharedSubscribe != LG_TRANSPORT_TIMEOUT)
return sharedSubscribe;
if (message.size < sizeof(KVMFRFrame)) struct LGMPFrameMessage owner[LGMP_Q_FRAME_LEN] = {};
struct LGMPFrameMessage shared = {};
LG_TransportStatus sharedStatus = LG_TRANSPORT_TIMEOUT;
LG_TransportStatus ownerStatus[LGMP_Q_FRAME_LEN];
/* Select the newest frame across the shared and owner delivery lanes. */
if (this->frameQueue)
sharedStatus = lgmp_pollFrameQueue(this->frameQueue, false, &shared);
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{ {
lgmpClientMessageDone(this->frameQueue); ownerStatus[i] = LG_TRANSPORT_TIMEOUT;
DEBUG_ERROR("LGMP frame payload is too small"); if (this->ownerFrameQueue[i])
return LG_TRANSPORT_ERROR; ownerStatus[i] = lgmp_pollFrameQueue(
this->ownerFrameQueue[i], true, &owner[i]);
} }
const KVMFRFrame * frame = (const KVMFRFrame *)message.mem; LG_TransportStatus failure = LG_TRANSPORT_OK;
const size_t frameDataSize = (size_t)frame->dataHeight * frame->pitch; if (sharedStatus != LG_TRANSPORT_OK &&
if (frame->type <= FRAME_TYPE_INVALID || frame->type >= FRAME_TYPE_MAX || sharedStatus != LG_TRANSPORT_TIMEOUT)
frame->offset > message.size - sizeof(FrameBuffer) || failure = sharedStatus;
frameDataSize > message.size - frame->offset - sizeof(FrameBuffer)) for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (ownerStatus[i] != LG_TRANSPORT_OK &&
ownerStatus[i] != LG_TRANSPORT_TIMEOUT &&
(failure == LG_TRANSPORT_OK ||
ownerStatus[i] == LG_TRANSPORT_DISCONNECTED))
failure = ownerStatus[i];
if (failure != LG_TRANSPORT_OK)
{ {
lgmpClientMessageDone(this->frameQueue); lgmp_mergeFrameStatus(
DEBUG_ERROR("LGMP frame payload contains invalid dimensions or offsets"); lgmp_doneFrameMessage(&shared), &failure);
return LG_TRANSPORT_ERROR; for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
lgmp_mergeFrameStatus(
lgmp_doneFrameMessage(&owner[i]), &failure);
return failure;
} }
if (frame->frameSerial == this->frameSerial && this->frameSerial) bool malformed = false;
LG_TransportStatus releaseFailure = LG_TRANSPORT_OK;
if (sharedStatus == LG_TRANSPORT_OK &&
!lgmp_validateFrameMessage(&shared))
{ {
lgmpClientMessageDone(this->frameQueue); malformed = true;
return LG_TRANSPORT_TIMEOUT; releaseFailure = lgmp_doneFrameMessage(&shared);
} }
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (ownerStatus[i] == LG_TRANSPORT_OK &&
!lgmp_validateFrameMessage(&owner[i]))
{
malformed = true;
const LG_TransportStatus done =
lgmp_doneFrameMessage(&owner[i]);
lgmp_mergeFrameStatus(done, &releaseFailure);
}
if (releaseFailure != LG_TRANSPORT_OK)
{
lgmp_mergeFrameStatus(
lgmp_doneFrameMessage(&shared), &releaseFailure);
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
lgmp_mergeFrameStatus(
lgmp_doneFrameMessage(&owner[i]), &releaseFailure);
return releaseFailure;
}
struct LGMPFrameMessage * selected = NULL;
lgmp_selectNewestFrameMessage(&shared, &selected);
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
lgmp_selectNewestFrameMessage(&owner[i], &selected);
if (!selected)
{
if (this->framePollInterval)
usleep(this->framePollInterval);
return malformed ? LG_TRANSPORT_ERROR : LG_TRANSPORT_TIMEOUT;
}
if (selected != &shared)
releaseFailure = lgmp_doneFrameMessage(&shared);
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (selected != &owner[i])
{
const LG_TransportStatus done =
lgmp_doneFrameMessage(&owner[i]);
lgmp_mergeFrameStatus(done, &releaseFailure);
}
if (releaseFailure != LG_TRANSPORT_OK)
{
lgmp_mergeFrameStatus(
lgmp_doneFrameMessage(selected), &releaseFailure);
return releaseFailure;
}
const KVMFRFrame * frame = selected->frame;
const bool equalSerial = this->frameSerialValid &&
frame->frameSerial == this->frameSerial;
if (this->frameSerialValid && !equalSerial &&
!lgmp_frameSerialNewer(frame->frameSerial, this->frameSerial))
{
const LG_TransportStatus done = lgmp_doneFrameMessage(selected);
return done == LG_TRANSPORT_OK ? LG_TRANSPORT_TIMEOUT : done;
}
if (equalSerial && !selected->owner)
{
const LG_TransportStatus done = lgmp_doneFrameMessage(selected);
return done == LG_TRANSPORT_OK ? LG_TRANSPORT_TIMEOUT : done;
}
const bool fullDamage = !this->frameSerialValid || equalSerial ||
frame->frameSerial != this->frameSerial + 1;
this->frameSerial = frame->frameSerial; this->frameSerial = frame->frameSerial;
this->frameSerialValid = true;
memset(result, 0, sizeof(*result)); memset(result, 0, sizeof(*result));
result->serial = frame->frameSerial; result->serial = frame->frameSerial;
result->scheduleGeneration = message.udata; if (selected->owner)
{
const uint64_t scheduleToken = selected->message.udata;
result->scheduleGeneration = (uint32_t)(scheduleToken >> 32);
result->scheduleEpoch = (uint32_t)scheduleToken;
result->scheduleOwner = true;
}
if (frame->flags & FRAME_FLAG_BLOCK_SCREENSAVER) if (frame->flags & FRAME_FLAG_BLOCK_SCREENSAVER)
result->flags |= LG_TRANSPORT_FRAME_BLOCK_SCREENSAVER; result->flags |= LG_TRANSPORT_FRAME_BLOCK_SCREENSAVER;
if (frame->flags & FRAME_FLAG_REQUEST_ACTIVATION) if (frame->flags & FRAME_FLAG_REQUEST_ACTIVATION)
@@ -596,22 +822,24 @@ static LG_TransportStatus lgmp_nextFrame(LG_Transport * this, bool useDMA,
result->dmaFD = lgmp_getDMA(this, frame, dataSize); result->dmaFD = lgmp_getDMA(this, frame, dataSize);
if (result->dmaFD < 0) if (result->dmaFD < 0)
{ {
lgmpClientMessageDone(this->frameQueue); const LG_TransportStatus done = lgmp_doneFrameMessage(selected);
DEBUG_ERROR("Failed to obtain a DMA buffer for the frame"); DEBUG_ERROR("Failed to obtain a DMA buffer for the frame");
return LG_TRANSPORT_ERROR; return done == LG_TRANSPORT_DISCONNECTED ?
done : LG_TRANSPORT_ERROR;
} }
} }
if (frame->damageRectsCount <= KVMFR_MAX_DAMAGE_RECTS) if (!fullDamage && frame->damageRectsCount <= KVMFR_MAX_DAMAGE_RECTS)
{ {
result->damageRects = frame->damageRects; result->damageRects = frame->damageRects;
result->damageRectsCount = frame->damageRectsCount; result->damageRectsCount = frame->damageRectsCount;
} }
else else if (!fullDamage)
DEBUG_WARN("Invalid damage rectangles, forcing a full update"); DEBUG_WARN("Invalid damage rectangles, forcing a full update");
this->framePending = true; this->framePending = true;
this->pendingFrame = frame; this->pendingFrame = frame;
this->pendingFrameQueue = selected->queue;
return LG_TRANSPORT_OK; return LG_TRANSPORT_OK;
} }
@@ -652,10 +880,11 @@ static void lgmp_getFrameTiming(LG_Transport * this,
static void lgmp_releaseFrame(LG_Transport * this, LG_TransportFrame * frame) static void lgmp_releaseFrame(LG_Transport * this, LG_TransportFrame * frame)
{ {
if (this->framePending && this->frameQueue) if (this->framePending && this->pendingFrameQueue)
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(this->pendingFrameQueue);
this->framePending = false; this->framePending = false;
this->pendingFrame = NULL; this->pendingFrame = NULL;
this->pendingFrameQueue = NULL;
memset(frame, 0, sizeof(*frame)); memset(frame, 0, sizeof(*frame));
} }
@@ -721,7 +950,7 @@ static LG_TransportStatus lgmp_nextPointer(LG_Transport * this,
if (status == LG_TRANSPORT_OK) if (status == LG_TRANSPORT_OK)
{ {
memcpy(this->pointerData, message.mem, needed); memcpy(this->pointerData, message.mem, needed);
pointerFlags = message.udata; pointerFlags = (uint32_t)message.udata;
} }
lgmpClientMessageDone(this->pointerQueue); lgmpClientMessageDone(this->pointerQueue);
} }
@@ -805,6 +1034,8 @@ static LG_TransportStatus lgmp_sendControl(LG_Transport * this,
.targetSlack = control->frameSchedule.targetSlack, .targetSlack = control->frameSchedule.targetSlack,
.phaseError = control->frameSchedule.phaseError, .phaseError = control->frameSchedule.phaseError,
.feedbackFrameSerial = control->frameSchedule.feedbackFrameSerial, .feedbackFrameSerial = control->frameSchedule.feedbackFrameSerial,
.feedbackScheduleEpoch =
control->frameSchedule.feedbackScheduleEpoch,
.lease = control->frameSchedule.lease, .lease = control->frameSchedule.lease,
}; };
memcpy(buffer, &message, sizeof(message)); memcpy(buffer, &message, sizeof(message));

View File

@@ -30,7 +30,7 @@
#include "LGMPConfig.h" #include "LGMPConfig.h"
#define KVMFR_MAGIC "KVMFR---" #define KVMFR_MAGIC "KVMFR---"
#define KVMFR_VERSION 26 #define KVMFR_VERSION 27
// Fallback used by producers that cannot report the source display's SDR // Fallback used by producers that cannot report the source display's SDR
// white level. IDD frames override this with IDDCX_METADATA2::SdrWhiteLevel. // white level. IDD frames override this with IDDCX_METADATA2::SdrWhiteLevel.
@@ -296,8 +296,9 @@ typedef struct KVMFRFrameSchedule
// Positive when the frame arrived early, negative when it arrived late. // Positive when the frame arrived early, negative when it arrived late.
int64_t phaseError; // ready-to-render phase error (ns) int64_t phaseError; // ready-to-render phase error (ns)
uint32_t feedbackFrameSerial; uint32_t feedbackFrameSerial;
uint32_t feedbackScheduleEpoch;
uint32_t lease; // lease duration (ms) uint32_t lease; // lease duration (ms)
uint8_t reserved[16]; uint8_t reserved[12];
} }
KVMFRFrameSchedule; KVMFRFrameSchedule;

View File

@@ -23,8 +23,13 @@
#define LGMP_Q_POINTER 1 #define LGMP_Q_POINTER 1
#define LGMP_Q_FRAME 2 #define LGMP_Q_FRAME 2
// Base ID for LGMP_Q_FRAME_LEN independent owner-delivery queues.
#define LGMP_Q_FRAME_OWNER 3
// Two delivery lanes plus a spare buffer let the timing owner continue to
// alternate buffers while a secondary client holds the shared delivery.
#define LGMP_Q_FRAME_LEN 2 #define LGMP_Q_FRAME_LEN 2
#define LGMP_Q_FRAME_BUFFER_LEN 3
#define LGMP_Q_POINTER_LEN 32 #define LGMP_Q_POINTER_LEN 32
#endif #endif

View File

@@ -30,7 +30,7 @@ class CFrameBufferPool
{ {
CSwapChainProcessor * m_swapChain; CSwapChainProcessor * m_swapChain;
CFrameBufferResource m_buffers[LGMP_Q_FRAME_LEN]; CFrameBufferResource m_buffers[LGMP_Q_FRAME_BUFFER_LEN];
public: public:
void Init(CSwapChainProcessor * swapChain); void Init(CSwapChainProcessor * swapChain);

View File

@@ -90,8 +90,7 @@ bool CFrameScheduler::ElectOwner(uint64_t now)
if (!client.active || client.expiry <= now) if (!client.active || client.expiry <= now)
{ {
client.active = false; client.active = false;
fastest = nullptr; continue;
break;
} }
if (!fastest || client.period < fastest->period) if (!fastest || client.period < fastest->period)
@@ -108,6 +107,7 @@ bool CFrameScheduler::ElectOwner(uint64_t now)
const uint32_t oldClientID = m_schedule.clientID; const uint32_t oldClientID = m_schedule.clientID;
const uint32_t oldGeneration = m_schedule.generation; const uint32_t oldGeneration = m_schedule.generation;
const uint32_t oldEpoch = m_schedule.epoch;
const uint64_t oldPeriod = m_schedule.period; const uint64_t oldPeriod = m_schedule.period;
const uint64_t oldSlack = m_schedule.targetSlack; const uint64_t oldSlack = m_schedule.targetSlack;
if (!fastest) if (!fastest)
@@ -119,6 +119,7 @@ bool CFrameScheduler::ElectOwner(uint64_t now)
{ {
m_schedule.clientID = fastest->clientID; m_schedule.clientID = fastest->clientID;
m_schedule.generation = fastest->generation; m_schedule.generation = fastest->generation;
m_schedule.epoch = oldEpoch;
m_schedule.period = fastest->period; m_schedule.period = fastest->period;
m_schedule.targetSlack = fastest->targetSlack; m_schedule.targetSlack = fastest->targetSlack;
m_scheduling = true; m_scheduling = true;
@@ -127,8 +128,18 @@ bool CFrameScheduler::ElectOwner(uint64_t now)
const bool ownerChanged = oldClientID != m_schedule.clientID; const bool ownerChanged = oldClientID != m_schedule.clientID;
if (ownerChanged || oldGeneration != m_schedule.generation) if (ownerChanged || oldGeneration != m_schedule.generation)
{ {
if (m_scheduling)
{
if (!++m_epoch)
++m_epoch;
m_schedule.epoch = m_epoch;
}
m_nextDeadline = m_scheduling ? now + m_schedule.period : 0; m_nextDeadline = m_scheduling ? now + m_schedule.period : 0;
m_forceNext = m_scheduling; m_forceNext = m_scheduling;
m_republishNext = m_scheduling;
if (fastest)
fastest->lastFeedbackFrameSerial = 0;
m_lastPublishedFrameSerial = 0; m_lastPublishedFrameSerial = 0;
m_lastPhaseError = 0; m_lastPhaseError = 0;
@@ -138,8 +149,8 @@ bool CFrameScheduler::ElectOwner(uint64_t now)
m_lastLogPublished = m_publishedFrames; m_lastLogPublished = m_publishedFrames;
if (ownerChanged && m_scheduling) if (ownerChanged && m_scheduling)
DEBUG_INFO("Frame timing owner %u generation %u at %.3f Hz", DEBUG_INFO("Frame timing owner %u generation %u epoch %u at %.3f Hz",
m_schedule.clientID, m_schedule.generation, m_schedule.clientID, m_schedule.generation, m_schedule.epoch,
1000000000.0 / m_schedule.period); 1000000000.0 / m_schedule.period);
else if (ownerChanged && oldClientID) else if (ownerChanged && oldClientID)
DEBUG_INFO("Frame timing owner released; using push delivery"); DEBUG_INFO("Frame timing owner released; using push delivery");
@@ -157,6 +168,8 @@ void CFrameScheduler::Reset()
m_schedule = {}; m_schedule = {};
m_scheduling = false; m_scheduling = false;
m_forceNext = false; m_forceNext = false;
m_republishNext = false;
m_epoch = 0;
m_lastArrival = 0; m_lastArrival = 0;
m_guestPeriod = 0; m_guestPeriod = 0;
@@ -198,11 +211,15 @@ void CFrameScheduler::UpdateSubscribers(const uint32_t * clientIDs,
} }
if (client) if (client)
{
client->subscribed = true; client->subscribed = true;
client->subscriptionSeen = true;
}
} }
for (Client& client : m_clients) for (Client& client : m_clients)
if (client.clientID && !client.subscribed) if (client.clientID && !client.subscribed &&
(client.subscriptionSeen || !client.active || client.expiry <= now))
client = {}; client = {};
const bool changed = ElectOwner(now); const bool changed = ElectOwner(now);
@@ -211,8 +228,8 @@ void CFrameScheduler::UpdateSubscribers(const uint32_t * clientIDs,
WakePublisher(); WakePublisher();
} }
bool CFrameScheduler::UpdateSchedule(const KVMFRFrameSchedule& schedule, bool CFrameScheduler::UpdateSchedule(uint32_t sourceClientID,
uint64_t now) const KVMFRFrameSchedule& schedule, uint64_t now)
{ {
static const KVMFRFrameScheduleFlags validFlags = static const KVMFRFrameScheduleFlags validFlags =
KVMFR_FRAME_SCHEDULE_ACTIVE | KVMFR_FRAME_SCHEDULE_ACTIVE |
@@ -220,19 +237,20 @@ bool CFrameScheduler::UpdateSchedule(const KVMFRFrameSchedule& schedule,
KVMFR_FRAME_SCHEDULE_RESET | KVMFR_FRAME_SCHEDULE_RESET |
KVMFR_FRAME_SCHEDULE_IMMEDIATE; KVMFR_FRAME_SCHEDULE_IMMEDIATE;
if (!schedule.clientID || schedule.flags & ~validFlags) if (!sourceClientID || schedule.clientID != sourceClientID ||
schedule.flags & ~validFlags)
return false; return false;
AcquireSRWLockExclusive(&m_lock);
Client * client = FindClient(schedule.clientID);
bool wake = false;
if (schedule.flags & KVMFR_FRAME_SCHEDULE_RELEASE) if (schedule.flags & KVMFR_FRAME_SCHEDULE_RELEASE)
{ {
if (client && client->subscribed) AcquireSRWLockExclusive(&m_lock);
Client * client = FindClient(schedule.clientID);
bool wake = false;
if (client)
{ {
client->active = false; client->active = false;
client->expiry = 0; client->expiry = 0;
client->immediate = false;
wake = ElectOwner(now); wake = ElectOwner(now);
} }
ReleaseSRWLockExclusive(&m_lock); ReleaseSRWLockExclusive(&m_lock);
@@ -241,12 +259,6 @@ bool CFrameScheduler::UpdateSchedule(const KVMFRFrameSchedule& schedule,
return true; return true;
} }
if (!client || !client->subscribed)
{
ReleaseSRWLockExclusive(&m_lock);
return false;
}
if (!(schedule.flags & KVMFR_FRAME_SCHEDULE_ACTIVE) || if (!(schedule.flags & KVMFR_FRAME_SCHEDULE_ACTIVE) ||
schedule.period < MIN_PERIOD_NS || schedule.period < MIN_PERIOD_NS ||
schedule.period > MAX_PERIOD_NS || schedule.period > MAX_PERIOD_NS ||
@@ -254,24 +266,47 @@ bool CFrameScheduler::UpdateSchedule(const KVMFRFrameSchedule& schedule,
schedule.phaseError > static_cast<int64_t>(schedule.period) || schedule.phaseError > static_cast<int64_t>(schedule.period) ||
schedule.phaseError < -static_cast<int64_t>(schedule.period) || schedule.phaseError < -static_cast<int64_t>(schedule.period) ||
schedule.lease < MIN_LEASE_MS || schedule.lease > MAX_LEASE_MS) schedule.lease < MIN_LEASE_MS || schedule.lease > MAX_LEASE_MS)
return false;
AcquireSRWLockExclusive(&m_lock);
Client * client = FindClient(schedule.clientID);
if (!client)
for (Client& candidate : m_clients)
if (!candidate.clientID)
{
candidate.clientID = schedule.clientID;
client = &candidate;
break;
}
if (!client)
{ {
ReleaseSRWLockExclusive(&m_lock); ReleaseSRWLockExclusive(&m_lock);
return false; return false;
} }
bool wake = false;
if (client->generation != schedule.generation) if (client->generation != schedule.generation)
{
client->lastFeedbackFrameSerial = 0; client->lastFeedbackFrameSerial = 0;
client->immediate = false;
}
client->generation = schedule.generation; client->generation = schedule.generation;
client->period = schedule.period; client->period = schedule.period;
client->targetSlack = schedule.targetSlack; client->targetSlack = schedule.targetSlack;
client->expiry = now + static_cast<uint64_t>(schedule.lease) * 1000000; client->expiry = now + static_cast<uint64_t>(schedule.lease) * 1000000;
client->active = true; client->active = true;
if (schedule.flags & KVMFR_FRAME_SCHEDULE_IMMEDIATE) if (schedule.flags & KVMFR_FRAME_SCHEDULE_IMMEDIATE)
client->immediate = true;
wake |= ElectOwner(now);
if (m_scheduling && client->clientID == m_schedule.clientID &&
client->generation == m_schedule.generation && client->immediate)
{ {
m_forceNext = true; m_forceNext = true;
m_republishNext = true;
wake = true; wake = true;
} }
wake |= ElectOwner(now);
wake |= ApplyFeedback(*client, schedule); wake |= ApplyFeedback(*client, schedule);
ReleaseSRWLockExclusive(&m_lock); ReleaseSRWLockExclusive(&m_lock);
if (wake) if (wake)
@@ -284,6 +319,7 @@ bool CFrameScheduler::ApplyFeedback(Client& client,
{ {
if (!m_scheduling || client.clientID != m_schedule.clientID || if (!m_scheduling || client.clientID != m_schedule.clientID ||
schedule.generation != m_schedule.generation || schedule.generation != m_schedule.generation ||
schedule.feedbackScheduleEpoch != m_schedule.epoch ||
!schedule.feedbackFrameSerial || !schedule.feedbackFrameSerial ||
(client.lastFeedbackFrameSerial && (client.lastFeedbackFrameSerial &&
static_cast<int32_t>(schedule.feedbackFrameSerial - static_cast<int32_t>(schedule.feedbackFrameSerial -
@@ -366,11 +402,12 @@ void CFrameScheduler::ForceFrame()
} }
bool CFrameScheduler::GetPublishTarget(uint64_t now, uint64_t& target, bool CFrameScheduler::GetPublishTarget(uint64_t now, uint64_t& target,
uint32_t& generation, bool& periodic) Schedule& schedule, bool& periodic, bool& republish)
{ {
target = now; target = now;
generation = 0; schedule = {};
periodic = false; periodic = false;
republish = false;
AcquireSRWLockExclusive(&m_lock); AcquireSRWLockExclusive(&m_lock);
if (!m_scheduling) if (!m_scheduling)
{ {
@@ -378,7 +415,8 @@ bool CFrameScheduler::GetPublishTarget(uint64_t now, uint64_t& target,
return true; return true;
} }
generation = m_schedule.generation; schedule = m_schedule;
republish = m_republishNext;
if (m_forceNext) if (m_forceNext)
{ {
@@ -412,13 +450,20 @@ void CFrameScheduler::FrameSuperseded()
ReleaseSRWLockExclusive(&m_lock); ReleaseSRWLockExclusive(&m_lock);
} }
void CFrameScheduler::FramePublished(uint32_t generation, void CFrameScheduler::FramePublished(const Schedule& schedule,
uint32_t frameSerial, uint64_t now, bool periodic) uint32_t frameSerial, uint64_t now, bool periodic)
{ {
AcquireSRWLockExclusive(&m_lock); AcquireSRWLockExclusive(&m_lock);
if (m_scheduling && generation == m_schedule.generation) if (m_scheduling && schedule.clientID == m_schedule.clientID &&
schedule.generation == m_schedule.generation &&
schedule.epoch == m_schedule.epoch)
{ {
m_forceNext = false; m_forceNext = false;
m_republishNext = false;
Client * client = FindClient(m_schedule.clientID);
if (client)
client->immediate = false;
m_lastPublishedFrameSerial = frameSerial; m_lastPublishedFrameSerial = frameSerial;
++m_publishedFrames; ++m_publishedFrames;
if (periodic) if (periodic)
@@ -430,6 +475,38 @@ void CFrameScheduler::FramePublished(uint32_t generation,
ReleaseSRWLockExclusive(&m_lock); ReleaseSRWLockExclusive(&m_lock);
} }
void CFrameScheduler::FrameRepublished(const Schedule& schedule,
uint32_t frameSerial)
{
AcquireSRWLockExclusive(&m_lock);
if (m_scheduling && schedule.clientID == m_schedule.clientID &&
schedule.generation == m_schedule.generation &&
schedule.epoch == m_schedule.epoch)
{
m_republishNext = false;
Client * client = FindClient(m_schedule.clientID);
if (client)
client->immediate = false;
m_lastPublishedFrameSerial = frameSerial;
++m_publishedFrames;
}
ReleaseSRWLockExclusive(&m_lock);
}
void CFrameScheduler::FrameDelivered(const uint32_t * clientIDs,
unsigned count)
{
AcquireSRWLockExclusive(&m_lock);
for (unsigned i = 0; i < count; ++i)
{
Client * client = FindClient(clientIDs[i]);
if (client)
client->immediate = false;
}
ReleaseSRWLockExclusive(&m_lock);
}
void CFrameScheduler::LogStatistics(uint64_t now) void CFrameScheduler::LogStatistics(uint64_t now)
{ {
AcquireSRWLockExclusive(&m_lock); AcquireSRWLockExclusive(&m_lock);

View File

@@ -36,6 +36,7 @@ public:
{ {
uint32_t clientID; uint32_t clientID;
uint32_t generation; uint32_t generation;
uint32_t epoch;
uint64_t period; uint64_t period;
uint64_t targetSlack; uint64_t targetSlack;
}; };
@@ -50,7 +51,9 @@ private:
uint64_t expiry; uint64_t expiry;
uint32_t lastFeedbackFrameSerial; uint32_t lastFeedbackFrameSerial;
bool subscribed; bool subscribed;
bool subscriptionSeen;
bool active; bool active;
bool immediate;
}; };
mutable SRWLOCK m_lock = SRWLOCK_INIT; mutable SRWLOCK m_lock = SRWLOCK_INIT;
@@ -59,6 +62,8 @@ private:
Schedule m_schedule = {}; Schedule m_schedule = {};
bool m_scheduling = false; bool m_scheduling = false;
bool m_forceNext = false; bool m_forceNext = false;
bool m_republishNext = false;
uint32_t m_epoch = 0;
uint64_t m_lastArrival = 0; uint64_t m_lastArrival = 0;
uint64_t m_guestPeriod = 0; uint64_t m_guestPeriod = 0;
@@ -91,16 +96,19 @@ public:
void Reset(); void Reset();
void UpdateSubscribers(const uint32_t * clientIDs, unsigned count, void UpdateSubscribers(const uint32_t * clientIDs, unsigned count,
uint64_t now); uint64_t now);
bool UpdateSchedule(const KVMFRFrameSchedule& schedule, uint64_t now); bool UpdateSchedule(uint32_t sourceClientID,
const KVMFRFrameSchedule& schedule, uint64_t now);
bool GetSchedule(Schedule& schedule) const; bool GetSchedule(Schedule& schedule) const;
HANDLE GetWakeEvent() const { return m_wakeEvent; } HANDLE GetWakeEvent() const { return m_wakeEvent; }
void ObserveFrame(uint64_t now); void ObserveFrame(uint64_t now);
void ForceFrame(); void ForceFrame();
bool GetPublishTarget(uint64_t now, uint64_t& target, bool GetPublishTarget(uint64_t now, uint64_t& target,
uint32_t& generation, bool& periodic); Schedule& schedule, bool& periodic, bool& republish);
void FrameSuperseded(); void FrameSuperseded();
void FramePublished(uint32_t generation, uint32_t frameSerial, void FramePublished(const Schedule& schedule, uint32_t frameSerial,
uint64_t now, bool periodic); uint64_t now, bool periodic);
void FrameRepublished(const Schedule& schedule, uint32_t frameSerial);
void FrameDelivered(const uint32_t * clientIDs, unsigned count);
void RecordFrameTiming(uint64_t duration); void RecordFrameTiming(uint64_t duration);
void LogStatistics(uint64_t now); void LogStatistics(uint64_t now);
}; };

View File

@@ -44,6 +44,22 @@ static const struct LGMPQueueConfig POINTER_QUEUE_CONFIG =
1000 //subTimeout 1000 //subTimeout
}; };
static uint64_t FrameScheduleToken(
const CFrameScheduler::Schedule& schedule)
{
return static_cast<uint64_t>(schedule.generation) << 32 |
schedule.epoch;
}
static bool FrameScheduleMatches(
const CFrameScheduler::Schedule& a,
const CFrameScheduler::Schedule& b)
{
return a.clientID == b.clientID &&
a.generation == b.generation &&
a.epoch == b.epoch;
}
static const UINT IDDCX_VERSION_1_10 = 0x1A00; static const UINT IDDCX_VERSION_1_10 = 0x1A00;
static const UINT64 FRAME_BYTES_PER_PIXEL = 4; static const UINT64 FRAME_BYTES_PER_PIXEL = 4;
@@ -879,7 +895,7 @@ bool CIndirectDeviceContext::GetResolutionMemoryRequirements(
return false; return false;
ivshmemSize = frameMemoryStart + ivshmemSize = frameMemoryStart +
frameAllocationSize * LGMP_Q_FRAME_LEN; frameAllocationSize * LGMP_Q_FRAME_BUFFER_LEN;
return true; return true;
} }
@@ -1026,6 +1042,23 @@ bool CIndirectDeviceContext::InitializeLGMP()
return false; return false;
} }
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
const struct LGMPQueueConfig config =
{
LGMP_Q_FRAME_OWNER + i, //queueID
LGMP_Q_FRAME_LEN, //numMessages
1000 //subTimeout
};
if ((status = lgmpHostQueueNew(
m_lgmp, config, &m_frameOwnerQueue[i])) != LGMP_OK)
{
DEBUG_ERROR("lgmpHostQueueCreate Failed (Frame Owner %u): %s",
i, lgmpStatusString(status));
return false;
}
}
if ((status = lgmpHostQueueNew(m_lgmp, POINTER_QUEUE_CONFIG, &m_pointerQueue)) != LGMP_OK) if ((status = lgmpHostQueueNew(m_lgmp, POINTER_QUEUE_CONFIG, &m_pointerQueue)) != LGMP_OK)
{ {
DEBUG_ERROR("lgmpHostQueueCreate Failed (Pointer): %s", lgmpStatusString(status)); DEBUG_ERROR("lgmpHostQueueCreate Failed (Pointer): %s", lgmpStatusString(status));
@@ -1110,7 +1143,7 @@ bool CIndirectDeviceContext::SetupLGMP(size_t alignSize)
const size_t alignmentMask = m_alignSize - 1; const size_t alignmentMask = m_alignSize - 1;
size_t frameAllocationSize = size_t frameAllocationSize =
(available - alignmentPadding) / LGMP_Q_FRAME_LEN; (available - alignmentPadding) / LGMP_Q_FRAME_BUFFER_LEN;
frameAllocationSize &= ~alignmentMask; frameAllocationSize &= ~alignmentMask;
if (frameAllocationSize <= m_alignSize || if (frameAllocationSize <= m_alignSize ||
frameAllocationSize > UINT32_MAX) frameAllocationSize > UINT32_MAX)
@@ -1127,7 +1160,7 @@ bool CIndirectDeviceContext::SetupLGMP(size_t alignSize)
(unsigned int)(maxFrameSize / 1048576)); (unsigned int)(maxFrameSize / 1048576));
LGMP_STATUS status; LGMP_STATUS status;
for (int i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (int i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
{ {
if ((status = lgmpHostMemAllocAligned(m_lgmp, if ((status = lgmpHostMemAllocAligned(m_lgmp,
(uint32_t)frameAllocationSize, (uint32_t)frameAllocationSize,
@@ -1151,6 +1184,13 @@ bool CIndirectDeviceContext::SetupLGMP(size_t alignSize)
m_maxFrameSize = maxFrameSize; m_maxFrameSize = maxFrameSize;
m_publishedFrameIndex.store(-1, std::memory_order_release); m_publishedFrameIndex.store(-1, std::memory_order_release);
m_frameResendPending = false; m_frameResendPending = false;
m_framePublishSequence = 0;
memset(m_frameLastPublishSequence, 0,
sizeof(m_frameLastPublishSequence));
for (FrameDelivery& delivery : m_frameDelivery)
delivery = {};
for (OwnerDelivery& delivery : m_ownerDelivery)
delivery = {};
WDF_TIMER_CONFIG config; WDF_TIMER_CONFIG config;
WDF_TIMER_CONFIG_INIT_PERIODIC(&config, WDF_TIMER_CONFIG_INIT_PERIODIC(&config,
@@ -1199,6 +1239,13 @@ void CIndirectDeviceContext::DeInitLGMP()
m_frameScheduler.Reset(); m_frameScheduler.Reset();
m_publishedFrameIndex.store(-1, std::memory_order_release); m_publishedFrameIndex.store(-1, std::memory_order_release);
m_frameResendPending = false; m_frameResendPending = false;
m_framePublishSequence = 0;
memset(m_frameLastPublishSequence, 0,
sizeof(m_frameLastPublishSequence));
for (FrameDelivery& delivery : m_frameDelivery)
delivery = {};
for (OwnerDelivery& delivery : m_ownerDelivery)
delivery = {};
return; return;
} }
@@ -1213,12 +1260,19 @@ void CIndirectDeviceContext::DeInitLGMP()
AcquireSRWLockExclusive(&m_framePublishLock); AcquireSRWLockExclusive(&m_framePublishLock);
m_publishedFrameIndex.store(-1, std::memory_order_release); m_publishedFrameIndex.store(-1, std::memory_order_release);
m_frameResendPending = false; m_frameResendPending = false;
m_framePublishSequence = 0;
memset(m_frameLastPublishSequence, 0,
sizeof(m_frameLastPublishSequence));
for (FrameDelivery& delivery : m_frameDelivery)
delivery = {};
for (OwnerDelivery& delivery : m_ownerDelivery)
delivery = {};
ReleaseSRWLockExclusive(&m_framePublishLock); ReleaseSRWLockExclusive(&m_framePublishLock);
for (int i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (int i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
m_frameInFlight[i].store(false, std::memory_order_release); m_frameInFlight[i].store(false, std::memory_order_release);
for (int i = 0; i < LGMP_Q_FRAME_LEN; ++i) for (int i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
lgmpHostMemFree(&m_frameMemory[i]); lgmpHostMemFree(&m_frameMemory[i]);
for (int i = 0; i < LGMP_Q_POINTER_LEN; ++i) for (int i = 0; i < LGMP_Q_POINTER_LEN; ++i)
lgmpHostMemFree(&m_pointerMemory[i]); lgmpHostMemFree(&m_pointerMemory[i]);
@@ -1266,16 +1320,35 @@ void CIndirectDeviceContext::LGMPTimer()
const uint64_t now = CFrameScheduler::Nanotime(); const uint64_t now = CFrameScheduler::Nanotime();
uint32_t clientIDs[LGMP_MAX_CLIENTS] = {}; uint32_t clientIDs[LGMP_MAX_CLIENTS] = {};
unsigned clientCount = 0; unsigned clientCount = 0;
status = lgmpHostGetClientIDs(m_frameQueue, clientIDs, &clientCount); LGMP_STATUS subscriberStatus = lgmpHostGetClientIDs(
if (status == LGMP_OK) m_frameQueue, clientIDs, &clientCount);
m_frameScheduler.UpdateSubscribers(clientIDs, clientCount, now); for (unsigned queueIndex = 0;
else subscriberStatus == LGMP_OK && queueIndex < LGMP_Q_FRAME_LEN;
DEBUG_WARN("Failed to query LGMP frame subscribers: %s", ++queueIndex)
lgmpStatusString(status)); {
uint32_t queueClientIDs[LGMP_MAX_CLIENTS] = {};
unsigned queueClientCount = 0;
subscriberStatus = lgmpHostGetClientIDs(
m_frameOwnerQueue[queueIndex], queueClientIDs, &queueClientCount);
unsigned commonCount = 0;
for (unsigned i = 0;
subscriberStatus == LGMP_OK && i < clientCount;
++i)
for (unsigned candidate = 0; candidate < queueClientCount; ++candidate)
if (clientIDs[i] == queueClientIDs[candidate])
{
clientIDs[commonCount++] = clientIDs[i];
break;
}
clientCount = commonCount;
}
uint8_t data[LGMP_MSGS_SIZE]; uint8_t data[LGMP_MSGS_SIZE];
size_t size; size_t size;
while ((status = lgmpHostReadData(m_pointerQueue, &data, &size)) == LGMP_OK) uint32_t sourceClientID;
while ((status = lgmpHostReadDataWithSource(
m_pointerQueue, &data, &size, &sourceClientID)) == LGMP_OK)
{ {
KVMFRMessage * msg = (KVMFRMessage *)data; KVMFRMessage * msg = (KVMFRMessage *)data;
switch (msg->type) switch (msg->type)
@@ -1296,10 +1369,28 @@ void CIndirectDeviceContext::LGMPTimer()
case KVMFR_MESSAGE_FRAME_SCHEDULE: case KVMFR_MESSAGE_FRAME_SCHEDULE:
{ {
if (size != sizeof(KVMFRFrameSchedule) || const KVMFRFrameSchedule * frameSchedule =
!m_frameScheduler.UpdateSchedule( reinterpret_cast<KVMFRFrameSchedule *>(msg);
*reinterpret_cast<KVMFRFrameSchedule *>(msg), now)) const bool valid = size == sizeof(*frameSchedule) &&
m_frameScheduler.UpdateSchedule(
sourceClientID, *frameSchedule, now);
if (!valid)
DEBUG_WARN("Ignoring invalid KVMFR frame schedule"); DEBUG_WARN("Ignoring invalid KVMFR frame schedule");
else if ((frameSchedule->flags &
(KVMFR_FRAME_SCHEDULE_ACTIVE |
KVMFR_FRAME_SCHEDULE_IMMEDIATE)) ==
(KVMFR_FRAME_SCHEDULE_ACTIVE |
KVMFR_FRAME_SCHEDULE_IMMEDIATE))
{
CFrameScheduler::Schedule owner = {};
if (!m_frameScheduler.GetSchedule(owner) ||
owner.clientID != sourceClientID)
{
AcquireSRWLockExclusive(&m_framePublishLock);
m_frameResendPending = true;
ReleaseSRWLockExclusive(&m_framePublishLock);
}
}
break; break;
} }
} }
@@ -1307,28 +1398,22 @@ void CIndirectDeviceContext::LGMPTimer()
lgmpHostAckData(m_pointerQueue); lgmpHostAckData(m_pointerQueue);
} }
if (subscriberStatus == LGMP_OK)
m_frameScheduler.UpdateSubscribers(clientIDs, clientCount, now);
else
DEBUG_WARN("Failed to query LGMP frame subscribers: %s",
lgmpStatusString(subscriberStatus));
m_frameScheduler.LogStatistics(now); m_frameScheduler.LogStatistics(now);
AcquireSRWLockExclusive(&m_framePublishLock);
if (lgmpHostQueueNewSubs(m_frameQueue)) if (lgmpHostQueueNewSubs(m_frameQueue))
{
AcquireSRWLockExclusive(&m_framePublishLock);
m_frameResendPending = true; m_frameResendPending = true;
if (m_frameResendPending && m_monitor &&
lgmpHostQueuePending(m_frameQueue) == 0)
{
const LONG frameIndex =
m_publishedFrameIndex.load(std::memory_order_acquire);
if (frameIndex >= 0)
{
status = lgmpHostQueuePost(
m_frameQueue, 0, m_frameMemory[frameIndex]);
if (status == LGMP_OK)
m_frameResendPending = false;
else if (status != LGMP_ERR_QUEUE_FULL)
DEBUG_ERROR("Failed to resend frame: %s", lgmpStatusString(status));
}
}
ReleaseSRWLockExclusive(&m_framePublishLock); ReleaseSRWLockExclusive(&m_framePublishLock);
}
ProcessFrameDeliveries();
if (lgmpHostQueueNewSubs(m_pointerQueue)) if (lgmpHostQueueNewSubs(m_pointerQueue))
{ {
@@ -1337,18 +1422,347 @@ void CIndirectDeviceContext::LGMPTimer()
} }
} }
bool CIndirectDeviceContext::FrameBufferAvailable() const bool CIndirectDeviceContext::PostSharedFrame(unsigned frameIndex,
uint32_t excludeClientID)
{ {
if (!m_lgmp || !m_frameQueue || if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN ||
lgmpHostQueuePending(m_frameQueue) != 0)
return false;
uint32_t clientIDs[LGMP_MAX_CLIENTS] = {};
unsigned clientCount = 0;
LGMP_STATUS status =
lgmpHostGetClientIDs(m_frameQueue, clientIDs, &clientCount);
if (status != LGMP_OK)
{
DEBUG_ERROR("Failed to query shared frame subscribers: %s",
lgmpStatusString(status));
return false;
}
unsigned targetCount = 0;
for (unsigned i = 0; i < clientCount; ++i)
{
bool excluded = clientIDs[i] == excludeClientID;
for (const OwnerDelivery& delivery : m_ownerDelivery)
if (delivery.active && delivery.clientID == clientIDs[i])
{
excluded = true;
break;
}
if (!excluded)
clientIDs[targetCount++] = clientIDs[i];
}
if (!targetCount)
{
m_frameDelivery[frameIndex].sharedOwnerToken = 0;
m_frameDelivery[frameIndex].sharedOwnerClientID = 0;
m_frameDelivery[frameIndex].sharedOwnerPending = false;
m_frameDelivery[frameIndex].sharedPending = false;
m_frameDelivery[frameIndex].sharedDelivered = true;
m_frameResendPending = false;
return true;
}
unsigned recipientCount = 0;
status = lgmpHostQueuePostForClients(
m_frameQueue, 0, m_frameMemory[frameIndex],
clientIDs, targetCount, &recipientCount);
if (status != LGMP_OK)
{
if (status != LGMP_ERR_QUEUE_FULL)
DEBUG_ERROR("Failed to publish shared frame: %s",
lgmpStatusString(status));
return false;
}
if (!recipientCount)
{
m_frameDelivery[frameIndex].sharedOwnerToken = 0;
m_frameDelivery[frameIndex].sharedOwnerClientID = 0;
m_frameDelivery[frameIndex].sharedOwnerPending = false;
m_frameDelivery[frameIndex].sharedPending = false;
m_frameDelivery[frameIndex].sharedDelivered = true;
m_frameResendPending = false;
return true;
}
m_frameDelivery[frameIndex].sharedOwnerToken = 0;
m_frameDelivery[frameIndex].sharedOwnerClientID = 0;
m_frameDelivery[frameIndex].sharedOwnerPending = false;
m_frameDelivery[frameIndex].sharedPending = true;
m_frameDelivery[frameIndex].sharedDelivered = true;
m_frameResendPending = false;
m_frameScheduler.FrameDelivered(clientIDs, targetCount);
return true;
}
bool CIndirectDeviceContext::PostSharedOwnerFrame(unsigned frameIndex,
const CFrameScheduler::Schedule& schedule)
{
if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN || !schedule.clientID ||
lgmpHostQueuePending(m_frameQueue) >= LGMP_Q_FRAME_LEN) lgmpHostQueuePending(m_frameQueue) >= LGMP_Q_FRAME_LEN)
return false; return false;
unsigned recipientCount = 0;
const LGMP_STATUS status = lgmpHostQueuePostForClients(
m_frameQueue, FrameScheduleToken(schedule), m_frameMemory[frameIndex],
&schedule.clientID, 1, &recipientCount);
if (status != LGMP_OK || !recipientCount)
{
if (status != LGMP_OK && status != LGMP_ERR_QUEUE_FULL)
DEBUG_ERROR("Failed to publish shared owner frame: %s",
lgmpStatusString(status));
return false;
}
FrameDelivery& delivery = m_frameDelivery[frameIndex];
delivery.sharedOwnerToken = FrameScheduleToken(schedule);
delivery.sharedOwnerClientID = schedule.clientID;
delivery.sharedOwnerPending = true;
delivery.sharedPending = true;
if (!delivery.sharedDelivered)
m_frameResendPending = true;
return true;
}
void CIndirectDeviceContext::ProcessFrameDeliveries()
{
if (!m_frameQueue)
return;
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (!m_frameOwnerQueue[i])
return;
AcquireSRWLockExclusive(&m_framePublishLock);
for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
{
if (m_frameDelivery[i].sharedOwnerPending &&
!lgmpHostQueueMessagePending(
m_frameQueue, m_frameMemory[i],
m_frameDelivery[i].sharedOwnerToken))
m_frameDelivery[i].sharedOwnerPending = false;
if (m_frameDelivery[i].sharedPending &&
!lgmpHostQueuePayloadPending(m_frameQueue, m_frameMemory[i]))
m_frameDelivery[i].sharedPending = false;
}
CFrameScheduler::Schedule schedule = {};
const bool scheduled = m_frameScheduler.GetSchedule(schedule);
const uint64_t scheduleToken = scheduled ?
FrameScheduleToken(schedule) : 0;
const LONG publishedFrameIndex = const LONG publishedFrameIndex =
m_publishedFrameIndex.load(std::memory_order_acquire); m_publishedFrameIndex.load(std::memory_order_acquire);
const unsigned frameIndex =
static_cast<unsigned>(publishedFrameIndex + 1) % LGMP_Q_FRAME_LEN;
return !m_frameInFlight[frameIndex].load(std::memory_order_acquire); int newestAcked = -1;
uint32_t newestAckedClient = 0;
for (unsigned queueIndex = 0;
queueIndex < LGMP_Q_FRAME_LEN;
++queueIndex)
{
OwnerDelivery& owner = m_ownerDelivery[queueIndex];
if (!owner.active ||
lgmpHostQueuePayloadPending(
m_frameOwnerQueue[queueIndex],
m_frameMemory[owner.frameIndex]))
continue;
const unsigned frameIndex = owner.frameIndex;
const bool latestDelivery =
static_cast<LONG>(frameIndex) == publishedFrameIndex &&
m_frameLastPublishSequence[frameIndex] == m_framePublishSequence;
const bool authoritative = (!scheduled && latestDelivery) ||
(scheduled && owner.clientID == schedule.clientID &&
owner.token == scheduleToken);
if (authoritative &&
(newestAcked < 0 ||
m_frameLastPublishSequence[frameIndex] >
m_frameLastPublishSequence[newestAcked]))
{
newestAcked = static_cast<int>(frameIndex);
newestAckedClient = owner.clientID;
}
else if (!authoritative && !latestDelivery)
m_frameResendPending = true;
m_frameDelivery[frameIndex].ownerQueueMask &=
~(1U << queueIndex);
owner = {};
}
if (newestAcked >= 0)
{
FrameDelivery& delivery = m_frameDelivery[newestAcked];
const bool needsShared = !delivery.sharedDelivered;
const bool latest = newestAcked == publishedFrameIndex &&
m_frameLastPublishSequence[newestAcked] == m_framePublishSequence;
if (needsShared &&
(!latest ||
m_frameInFlight[newestAcked].load(std::memory_order_acquire) ||
lgmpHostQueuePending(m_frameQueue) != 0))
{
m_frameResendPending = true;
}
else if (needsShared &&
!PostSharedFrame(
static_cast<unsigned>(newestAcked), newestAckedClient))
{
m_frameResendPending = true;
}
}
if (m_frameResendPending &&
lgmpHostQueuePending(m_frameQueue) == 0)
{
const LONG frameIndex =
m_publishedFrameIndex.load(std::memory_order_acquire);
if (frameIndex >= 0 &&
!m_frameDelivery[frameIndex].sharedPending &&
!m_frameInFlight[frameIndex].load(std::memory_order_acquire))
{
FrameDelivery& delivery = m_frameDelivery[frameIndex];
bool ownerDelivered = !scheduled;
if (scheduled)
{
ownerDelivered =
delivery.sharedOwnerClientID == schedule.clientID &&
delivery.sharedOwnerToken == scheduleToken;
for (const OwnerDelivery& owner : m_ownerDelivery)
if (owner.active &&
owner.clientID == schedule.clientID &&
owner.token == scheduleToken &&
owner.frameIndex == static_cast<unsigned>(frameIndex))
{
ownerDelivered = true;
break;
}
}
const bool reserveSharedOwner = scheduled &&
delivery.sharedOwnerClientID == schedule.clientID &&
delivery.sharedOwnerToken == scheduleToken &&
FindAvailableOwnerQueue(static_cast<unsigned>(frameIndex)) < 0;
if (ownerDelivered && !reserveSharedOwner)
{
const uint32_t excludeClientID = scheduled ?
schedule.clientID : delivery.sharedOwnerClientID;
PostSharedFrame(
static_cast<unsigned>(frameIndex), excludeClientID);
}
}
}
ReleaseSRWLockExclusive(&m_framePublishLock);
}
int CIndirectDeviceContext::FindAvailableOwnerQueue(
unsigned preferredIndex) const
{
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
const unsigned queueIndex =
(preferredIndex + i) % LGMP_Q_FRAME_LEN;
if (!m_ownerDelivery[queueIndex].active &&
m_frameOwnerQueue[queueIndex] &&
lgmpHostQueuePending(m_frameOwnerQueue[queueIndex]) == 0)
return static_cast<int>(queueIndex);
}
return -1;
}
int CIndirectDeviceContext::FindOwnerDelivery(uint32_t clientID) const
{
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (m_ownerDelivery[i].active &&
m_ownerDelivery[i].clientID == clientID)
return static_cast<int>(i);
return -1;
}
int CIndirectDeviceContext::FindSharedOwnerDelivery(
uint32_t clientID) const
{
for (unsigned i = 0; i < LGMP_Q_FRAME_BUFFER_LEN; ++i)
if (m_frameDelivery[i].sharedOwnerPending &&
m_frameDelivery[i].sharedOwnerClientID == clientID)
return static_cast<int>(i);
return -1;
}
bool CIndirectDeviceContext::HasOwnerDelivery(uint32_t clientID) const
{
return FindOwnerDelivery(clientID) >= 0 ||
FindSharedOwnerDelivery(clientID) >= 0;
}
int CIndirectDeviceContext::FindAvailableFrameBuffer() const
{
const LONG publishedFrameIndex =
m_publishedFrameIndex.load(std::memory_order_acquire);
int available = -1;
uint64_t newestPublish = 0;
for (unsigned frameIndex = 0;
frameIndex < LGMP_Q_FRAME_BUFFER_LEN;
++frameIndex)
{
if (static_cast<LONG>(frameIndex) == publishedFrameIndex ||
m_frameInFlight[frameIndex].load(std::memory_order_acquire) ||
m_frameDelivery[frameIndex].ownerQueueMask ||
m_frameDelivery[frameIndex].sharedPending ||
lgmpHostQueuePayloadPending(
m_frameQueue, m_frameMemory[frameIndex]))
continue;
bool ownerPending = false;
for (unsigned queueIndex = 0;
!ownerPending && queueIndex < LGMP_Q_FRAME_LEN;
++queueIndex)
ownerPending = lgmpHostQueuePayloadPending(
m_frameOwnerQueue[queueIndex], m_frameMemory[frameIndex]);
if (ownerPending)
continue;
if (available < 0 ||
m_frameLastPublishSequence[frameIndex] > newestPublish)
{
available = static_cast<int>(frameIndex);
newestPublish = m_frameLastPublishSequence[frameIndex];
}
}
return available;
}
bool CIndirectDeviceContext::FrameBufferAvailable(
const CFrameScheduler::Schedule& schedule)
{
if (!m_lgmp || !m_frameQueue)
return false;
for (unsigned i = 0; i < LGMP_Q_FRAME_LEN; ++i)
if (!m_frameOwnerQueue[i])
return false;
AcquireSRWLockShared(&m_framePublishLock);
bool deliveryAvailable;
if (schedule.clientID)
deliveryAvailable = !HasOwnerDelivery(schedule.clientID) &&
(FindAvailableOwnerQueue(0) >= 0 ||
lgmpHostQueuePending(m_frameQueue) < LGMP_Q_FRAME_LEN);
else
deliveryAvailable = lgmpHostQueuePending(m_frameQueue) == 0;
const bool available = deliveryAvailable &&
FindAvailableFrameBuffer() >= 0;
ReleaseSRWLockShared(&m_framePublishLock);
return available;
} }
void CIndirectDeviceContext::ProcessFrameQueue() void CIndirectDeviceContext::ProcessFrameQueue()
@@ -1362,6 +1776,9 @@ void CIndirectDeviceContext::ProcessFrameQueue()
if (status != LGMP_OK && status != LGMP_ERR_CORRUPTED) if (status != LGMP_OK && status != LGMP_ERR_CORRUPTED)
DEBUG_ERROR("lgmpHostProcess Failed: %s", lgmpStatusString(status)); DEBUG_ERROR("lgmpHostProcess Failed: %s", lgmpStatusString(status));
if (status == LGMP_OK)
ProcessFrameDeliveries();
} }
CIndirectDeviceContext::PreparedFrameBuffer CIndirectDeviceContext::PrepareFrameBuffer( CIndirectDeviceContext::PreparedFrameBuffer CIndirectDeviceContext::PrepareFrameBuffer(
@@ -1375,24 +1792,28 @@ CIndirectDeviceContext::PreparedFrameBuffer CIndirectDeviceContext::PrepareFrame
const unsigned dataHeight = dstFormat.dataHeight ? const unsigned dataHeight = dstFormat.dataHeight ?
dstFormat.dataHeight : dstFormat.desc.Height; dstFormat.dataHeight : dstFormat.desc.Height;
if (!FrameBufferAvailable())
return result;
if (dstFormat.format == FRAME_TYPE_INVALID) if (dstFormat.format == FRAME_TYPE_INVALID)
{ {
DEBUG_ERROR("Unsupported frame format, skipping frame"); DEBUG_ERROR("Unsupported frame format, skipping frame");
return result; return result;
} }
const LONG publishedFrameIndex = AcquireSRWLockExclusive(&m_framePublishLock);
m_publishedFrameIndex.load(std::memory_order_acquire); const int availableFrameIndex = FindAvailableFrameBuffer();
const unsigned frameIndex =
static_cast<unsigned>(publishedFrameIndex + 1) % LGMP_Q_FRAME_LEN;
bool expected = false; bool expected = false;
if (!m_frameInFlight[frameIndex].compare_exchange_strong( const bool acquired = availableFrameIndex >= 0 &&
expected, true, std::memory_order_acq_rel)) m_frameInFlight[availableFrameIndex].compare_exchange_strong(
expected, true, std::memory_order_acq_rel);
if (acquired)
m_frameDelivery[availableFrameIndex] = {};
const bool fullCopy = acquired &&
(!m_frameLastPublishSequence[availableFrameIndex] ||
m_framePublishSequence >
m_frameLastPublishSequence[availableFrameIndex] + 1);
ReleaseSRWLockExclusive(&m_framePublishLock);
if (!acquired)
return result; return result;
const unsigned frameIndex = static_cast<unsigned>(availableFrameIndex);
if (m_width != dataWidth || if (m_width != dataWidth ||
m_height != dataHeight || m_height != dataHeight ||
@@ -1518,45 +1939,198 @@ CIndirectDeviceContext::PreparedFrameBuffer CIndirectDeviceContext::PrepareFrame
result.frameIndex = frameIndex; result.frameIndex = frameIndex;
result.mem = fb->data; result.mem = fb->data;
result.fullCopy = fullCopy;
return result; return result;
} }
bool CIndirectDeviceContext::PublishFrameBuffer(unsigned frameIndex, bool CIndirectDeviceContext::PublishFrameBuffer(unsigned frameIndex,
uint32_t scheduleGeneration) const CFrameScheduler::Schedule& schedule)
{ {
if (!m_frameQueue || frameIndex >= LGMP_Q_FRAME_LEN) if (!m_frameQueue || frameIndex >= LGMP_Q_FRAME_BUFFER_LEN)
return false; return false;
AcquireSRWLockExclusive(&m_framePublishLock); AcquireSRWLockExclusive(&m_framePublishLock);
const LGMP_STATUS status = CFrameScheduler::Schedule currentSchedule = {};
lgmpHostQueuePost( const bool scheduling =
m_frameQueue, scheduleGeneration, m_frameMemory[frameIndex]); m_frameScheduler.GetSchedule(currentSchedule);
if (status == LGMP_OK) if (scheduling != (schedule.clientID != 0) ||
(scheduling && !FrameScheduleMatches(schedule, currentSchedule)))
{ {
ReleaseSRWLockExclusive(&m_framePublishLock);
return false;
}
LGMP_STATUS status = LGMP_OK;
bool published = false;
if (schedule.clientID)
{
if (HasOwnerDelivery(schedule.clientID))
status = LGMP_ERR_QUEUE_FULL;
else
{
const int ownerQueueIndex = FindAvailableOwnerQueue(frameIndex);
if (ownerQueueIndex < 0)
published = PostSharedOwnerFrame(frameIndex, schedule);
else
{
unsigned recipientCount = 0;
status = lgmpHostQueuePostForClients(
m_frameOwnerQueue[ownerQueueIndex], FrameScheduleToken(schedule),
m_frameMemory[frameIndex],
&schedule.clientID, 1, &recipientCount);
if (status == LGMP_OK && recipientCount)
{
const unsigned queueIndex =
static_cast<unsigned>(ownerQueueIndex);
OwnerDelivery& owner = m_ownerDelivery[queueIndex];
owner.token = FrameScheduleToken(schedule);
owner.clientID = schedule.clientID;
owner.frameIndex = frameIndex;
owner.active = true;
FrameDelivery& delivery = m_frameDelivery[frameIndex];
delivery.ownerQueueMask |= 1U << queueIndex;
delivery.sharedDelivered = false;
published = true;
if (!PostSharedFrame(frameIndex, schedule.clientID))
m_frameResendPending = true;
}
}
}
}
else
published = PostSharedFrame(frameIndex, 0);
if (published)
{
m_frameLastPublishSequence[frameIndex] = ++m_framePublishSequence;
m_publishedFrameIndex.store( m_publishedFrameIndex.store(
static_cast<LONG>(frameIndex), std::memory_order_release); static_cast<LONG>(frameIndex), std::memory_order_release);
m_frameResendPending = false;
} }
ReleaseSRWLockExclusive(&m_framePublishLock); ReleaseSRWLockExclusive(&m_framePublishLock);
if (status != LGMP_OK) if (!published)
{ {
DEBUG_ERROR("Failed to publish frame: %s", lgmpStatusString(status)); if (status != LGMP_OK && status != LGMP_ERR_QUEUE_FULL)
DEBUG_ERROR("Failed to publish frame: %s",
lgmpStatusString(status));
return false; return false;
} }
return true; return true;
} }
void CIndirectDeviceContext::CommitFrameBuffer(unsigned frameIndex, bool CIndirectDeviceContext::RepublishFrameBuffer(
uint32_t scheduleGeneration, bool periodic) const CFrameScheduler::Schedule& schedule)
{ {
if (frameIndex >= LGMP_Q_FRAME_LEN) if (!schedule.clientID)
return false;
AcquireSRWLockExclusive(&m_framePublishLock);
CFrameScheduler::Schedule currentSchedule = {};
if (!m_frameScheduler.GetSchedule(currentSchedule) ||
!FrameScheduleMatches(schedule, currentSchedule))
{
ReleaseSRWLockExclusive(&m_framePublishLock);
return false;
}
const LONG frameIndex =
m_publishedFrameIndex.load(std::memory_order_acquire);
if (frameIndex < 0 ||
m_frameInFlight[frameIndex].load(std::memory_order_acquire))
{
ReleaseSRWLockExclusive(&m_framePublishLock);
return false;
}
const uint64_t scheduleToken = FrameScheduleToken(schedule);
const int existingSharedDelivery =
FindSharedOwnerDelivery(schedule.clientID);
if (existingSharedDelivery >= 0)
{
const FrameDelivery& delivery =
m_frameDelivery[existingSharedDelivery];
const bool delivered = existingSharedDelivery == frameIndex &&
delivery.sharedOwnerToken == scheduleToken;
const uint32_t frameSerial = m_frame[frameIndex]->frameSerial;
ReleaseSRWLockExclusive(&m_framePublishLock);
if (delivered)
m_frameScheduler.FrameRepublished(schedule, frameSerial);
return delivered;
}
const int existingDelivery =
FindOwnerDelivery(schedule.clientID);
if (existingDelivery >= 0)
{
const OwnerDelivery& owner = m_ownerDelivery[existingDelivery];
const bool delivered = owner.frameIndex ==
static_cast<unsigned>(frameIndex) &&
owner.token == scheduleToken;
const uint32_t frameSerial = m_frame[frameIndex]->frameSerial;
ReleaseSRWLockExclusive(&m_framePublishLock);
if (delivered)
m_frameScheduler.FrameRepublished(schedule, frameSerial);
return delivered;
}
const int ownerQueueIndex =
FindAvailableOwnerQueue(static_cast<unsigned>(frameIndex));
if (ownerQueueIndex < 0)
{
const bool published = PostSharedOwnerFrame(
static_cast<unsigned>(frameIndex), schedule);
const uint32_t frameSerial = m_frame[frameIndex]->frameSerial;
ReleaseSRWLockExclusive(&m_framePublishLock);
if (published)
m_frameScheduler.FrameRepublished(schedule, frameSerial);
return published;
}
const uint32_t frameSerial = m_frame[frameIndex]->frameSerial;
unsigned recipientCount = 0;
const LGMP_STATUS status = lgmpHostQueuePostForClients(
m_frameOwnerQueue[ownerQueueIndex], scheduleToken,
m_frameMemory[frameIndex],
&schedule.clientID, 1, &recipientCount);
if (status == LGMP_OK && recipientCount)
{
const unsigned queueIndex = static_cast<unsigned>(ownerQueueIndex);
OwnerDelivery& owner = m_ownerDelivery[queueIndex];
owner.token = scheduleToken;
owner.clientID = schedule.clientID;
owner.frameIndex = static_cast<unsigned>(frameIndex);
owner.active = true;
FrameDelivery& delivery = m_frameDelivery[frameIndex];
delivery.ownerQueueMask |= 1U << queueIndex;
}
ReleaseSRWLockExclusive(&m_framePublishLock);
if (status != LGMP_OK || !recipientCount)
{
if (status != LGMP_OK && status != LGMP_ERR_QUEUE_FULL)
DEBUG_ERROR("Failed to republish frame: %s",
lgmpStatusString(status));
return false;
}
m_frameScheduler.FrameRepublished(schedule, frameSerial);
return true;
}
void CIndirectDeviceContext::CommitFrameBuffer(unsigned frameIndex,
const CFrameScheduler::Schedule& schedule, bool periodic)
{
if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN)
return; return;
m_frameScheduler.FramePublished( m_frameScheduler.FramePublished(
scheduleGeneration, m_frame[frameIndex]->frameSerial, schedule, m_frame[frameIndex]->frameSerial,
CFrameScheduler::Nanotime(), periodic); CFrameScheduler::Nanotime(), periodic);
} }
@@ -1571,10 +2145,11 @@ void CIndirectDeviceContext::ForceFrame()
} }
bool CIndirectDeviceContext::GetPublishTarget(uint64_t now, bool CIndirectDeviceContext::GetPublishTarget(uint64_t now,
uint64_t& target, uint32_t& generation, bool& periodic) uint64_t& target, CFrameScheduler::Schedule& schedule, bool& periodic,
bool& republish)
{ {
return m_frameScheduler.GetPublishTarget( return m_frameScheduler.GetPublishTarget(
now, target, generation, periodic); now, target, schedule, periodic, republish);
} }
void CIndirectDeviceContext::FrameSuperseded() void CIndirectDeviceContext::FrameSuperseded()
@@ -1589,17 +2164,20 @@ void CIndirectDeviceContext::RecordFrameTiming(uint64_t duration)
void CIndirectDeviceContext::AbortFrameBuffer(unsigned frameIndex) void CIndirectDeviceContext::AbortFrameBuffer(unsigned frameIndex)
{ {
if (frameIndex >= LGMP_Q_FRAME_LEN) if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN)
return; return;
AcquireSRWLockExclusive(&m_framePublishLock);
m_frameBuffer[frameIndex]->wp = 0; m_frameBuffer[frameIndex]->wp = 0;
InterlockedExchange((volatile LONG *)&m_frame[frameIndex]->timingValid, 0); InterlockedExchange(
(volatile LONG *)&m_frame[frameIndex]->timingValid, 0);
m_frameInFlight[frameIndex].store(false, std::memory_order_release); m_frameInFlight[frameIndex].store(false, std::memory_order_release);
ReleaseSRWLockExclusive(&m_framePublishLock);
} }
void CIndirectDeviceContext::FailFrameBuffer(unsigned frameIndex) void CIndirectDeviceContext::FailFrameBuffer(unsigned frameIndex)
{ {
if (frameIndex >= LGMP_Q_FRAME_LEN) if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN)
return; return;
InterlockedExchange((volatile LONG *)&m_frame[frameIndex]->timingValid, 0); InterlockedExchange((volatile LONG *)&m_frame[frameIndex]->timingValid, 0);
@@ -1609,7 +2187,7 @@ void CIndirectDeviceContext::FailFrameBuffer(unsigned frameIndex)
void CIndirectDeviceContext::CompleteFrameBuffer(unsigned frameIndex) void CIndirectDeviceContext::CompleteFrameBuffer(unsigned frameIndex)
{ {
if (frameIndex < LGMP_Q_FRAME_LEN) if (frameIndex < LGMP_Q_FRAME_BUFFER_LEN)
m_frameInFlight[frameIndex].store(false, std::memory_order_release); m_frameInFlight[frameIndex].store(false, std::memory_order_release);
} }
@@ -1617,7 +2195,7 @@ void CIndirectDeviceContext::SetFrameTiming(unsigned frameIndex,
uint64_t captureTime, uint64_t postProcessTime, uint64_t copyTime, uint64_t captureTime, uint64_t postProcessTime, uint64_t copyTime,
uint64_t readyTime) uint64_t readyTime)
{ {
if (frameIndex >= LGMP_Q_FRAME_LEN) if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN)
return; return;
KVMFRFrame * frame = m_frame[frameIndex]; KVMFRFrame * frame = m_frame[frameIndex];

View File

@@ -89,6 +89,7 @@ private:
PLGMPHost m_lgmp = nullptr; PLGMPHost m_lgmp = nullptr;
WDFTIMER m_lgmpTimer = nullptr; WDFTIMER m_lgmpTimer = nullptr;
PLGMPHostQueue m_frameQueue = nullptr; PLGMPHostQueue m_frameQueue = nullptr;
PLGMPHostQueue m_frameOwnerQueue[LGMP_Q_FRAME_LEN] = {};
SRWLOCK m_lgmpProcessLock = SRWLOCK_INIT; SRWLOCK m_lgmpProcessLock = SRWLOCK_INIT;
CFrameScheduler m_frameScheduler; CFrameScheduler m_frameScheduler;
@@ -108,14 +109,37 @@ private:
size_t m_frameMemoryOffset = 0; size_t m_frameMemoryOffset = 0;
size_t m_maxFrameSize = 0; size_t m_maxFrameSize = 0;
std::atomic<LONG> m_publishedFrameIndex = -1; std::atomic<LONG> m_publishedFrameIndex = -1;
std::atomic<bool> m_frameInFlight[LGMP_Q_FRAME_LEN] = {}; std::atomic<bool> m_frameInFlight[LGMP_Q_FRAME_BUFFER_LEN] = {};
SRWLOCK m_framePublishLock = SRWLOCK_INIT; SRWLOCK m_framePublishLock = SRWLOCK_INIT;
bool m_frameResendPending = false; bool m_frameResendPending = false;
uint64_t m_framePublishSequence = 0;
uint64_t m_frameLastPublishSequence[LGMP_Q_FRAME_BUFFER_LEN] = {};
struct FrameDelivery
{
uint64_t sharedOwnerToken = 0;
unsigned ownerQueueMask = 0;
uint32_t sharedOwnerClientID = 0;
bool sharedOwnerPending = false;
bool sharedPending = false;
bool sharedDelivered = false;
};
struct OwnerDelivery
{
uint64_t token = 0;
uint32_t clientID = 0;
unsigned frameIndex = 0;
bool active = false;
};
FrameDelivery m_frameDelivery[LGMP_Q_FRAME_BUFFER_LEN] = {};
OwnerDelivery m_ownerDelivery[LGMP_Q_FRAME_LEN] = {};
uint32_t m_formatVer = 0; uint32_t m_formatVer = 0;
uint32_t m_frameSerial = 0; uint32_t m_frameSerial = 0;
PLGMPMemory m_frameMemory[LGMP_Q_FRAME_LEN] = {}; PLGMPMemory m_frameMemory[LGMP_Q_FRAME_BUFFER_LEN] = {};
KVMFRFrame * m_frame [LGMP_Q_FRAME_LEN] = {}; KVMFRFrame * m_frame [LGMP_Q_FRAME_BUFFER_LEN] = {};
FrameBuffer * m_frameBuffer[LGMP_Q_FRAME_LEN] = {}; FrameBuffer * m_frameBuffer[LGMP_Q_FRAME_BUFFER_LEN] = {};
unsigned m_width = 0; unsigned m_width = 0;
unsigned m_height = 0; unsigned m_height = 0;
@@ -155,6 +179,15 @@ private:
bool InitializeLGMP(); bool InitializeLGMP();
void DeInitLGMP(); void DeInitLGMP();
void LGMPTimer(); void LGMPTimer();
void ProcessFrameDeliveries();
int FindAvailableFrameBuffer() const;
int FindAvailableOwnerQueue(unsigned preferredIndex) const;
int FindOwnerDelivery(uint32_t clientID) const;
int FindSharedOwnerDelivery(uint32_t clientID) const;
bool HasOwnerDelivery(uint32_t clientID) const;
bool PostSharedFrame(unsigned frameIndex, uint32_t excludeClientID);
bool PostSharedOwnerFrame(unsigned frameIndex,
const CFrameScheduler::Schedule& schedule);
void ResendCursor(); void ResendCursor();
void SendColorTransform(); void SendColorTransform();
void InitializeEdid(); void InitializeEdid();
@@ -222,14 +255,21 @@ public:
{ {
unsigned frameIndex; unsigned frameIndex;
uint8_t* mem; uint8_t* mem;
bool fullCopy;
}; };
bool FrameBufferAvailable() const; bool FrameBufferAvailable(const CFrameScheduler::Schedule& schedule);
bool HasPublishedFrame() const
{
return m_publishedFrameIndex.load(std::memory_order_acquire) >= 0;
}
void ProcessFrameQueue(); void ProcessFrameQueue();
PreparedFrameBuffer PrepareFrameBuffer(unsigned pitch, const D12FrameFormat& srcFormat, const D12FrameFormat& dstFormat, const RECT * dirtyRects, unsigned nbDirtyRects); PreparedFrameBuffer PrepareFrameBuffer(unsigned pitch, const D12FrameFormat& srcFormat, const D12FrameFormat& dstFormat, const RECT * dirtyRects, unsigned nbDirtyRects);
bool PublishFrameBuffer(unsigned frameIndex, uint32_t scheduleGeneration); bool PublishFrameBuffer(unsigned frameIndex,
void CommitFrameBuffer(unsigned frameIndex, uint32_t scheduleGeneration, const CFrameScheduler::Schedule& schedule);
bool periodic); bool RepublishFrameBuffer(const CFrameScheduler::Schedule& schedule);
void CommitFrameBuffer(unsigned frameIndex,
const CFrameScheduler::Schedule& schedule, bool periodic);
void AbortFrameBuffer(unsigned frameIndex); void AbortFrameBuffer(unsigned frameIndex);
void FailFrameBuffer(unsigned frameIndex); void FailFrameBuffer(unsigned frameIndex);
void CompleteFrameBuffer(unsigned frameIndex); void CompleteFrameBuffer(unsigned frameIndex);
@@ -241,7 +281,7 @@ public:
void ObserveFrame(uint64_t now); void ObserveFrame(uint64_t now);
void ForceFrame(); void ForceFrame();
bool GetPublishTarget(uint64_t now, uint64_t& target, bool GetPublishTarget(uint64_t now, uint64_t& target,
uint32_t& generation, bool& periodic); CFrameScheduler::Schedule& schedule, bool& periodic, bool& republish);
void FrameSuperseded(); void FrameSuperseded();
HANDLE GetFrameScheduleEvent() const HANDLE GetFrameScheduleEvent() const
{ {

View File

@@ -35,7 +35,7 @@ static const uint64_t PUBLISH_RETRY_NS = 1000000ULL;
static const DWORD CANDIDATE_WAIT_MS = 2; static const DWORD CANDIDATE_WAIT_MS = 2;
static_assert(LGMP_Q_FRAME_LEN == 2, static_assert(LGMP_Q_FRAME_LEN == 2,
"IDD damage repair assumes two alternating frame buffers"); "IDD candidate pipeline assumes two slots");
class CSRWExclusiveLock class CSRWExclusiveLock
{ {
@@ -223,8 +223,30 @@ void CSwapChainProcessor::PublisherThread()
for (;;) for (;;)
{ {
const uint64_t now = CFrameScheduler::Nanotime();
uint64_t target;
CFrameScheduler::Schedule schedule;
bool periodic;
bool republish;
m_devContext->GetPublishTarget(
now, target, schedule, periodic, republish);
if (!HasReadyCandidate()) if (!HasReadyCandidate())
{ {
if (republish && m_devContext->HasPublishedFrame())
{
m_devContext->ProcessFrameQueue();
if (m_devContext->RepublishFrameBuffer(schedule))
continue;
ArmPublishTimer(m_publishTimer.Get(), PUBLISH_RETRY_NS);
if (WaitForMultipleObjects(
ARRAYSIZE(timerHandles), timerHandles, FALSE, INFINITE) ==
WAIT_OBJECT_0)
break;
continue;
}
if (m_publishTimer.Get()) if (m_publishTimer.Get())
CancelWaitableTimer(m_publishTimer.Get()); CancelWaitableTimer(m_publishTimer.Get());
if (WaitForMultipleObjects( if (WaitForMultipleObjects(
@@ -234,12 +256,6 @@ void CSwapChainProcessor::PublisherThread()
continue; continue;
} }
const uint64_t now = CFrameScheduler::Nanotime();
uint64_t target;
uint32_t generation;
bool periodic;
m_devContext->GetPublishTarget(now, target, generation, periodic);
if (target > now) if (target > now)
{ {
ArmPublishTimer(m_publishTimer.Get(), target - now); ArmPublishTimer(m_publishTimer.Get(), target - now);
@@ -252,9 +268,9 @@ void CSwapChainProcessor::PublisherThread()
const uint64_t publishStart = CFrameScheduler::Nanotime(); const uint64_t publishStart = CFrameScheduler::Nanotime();
m_devContext->ProcessFrameQueue(); m_devContext->ProcessFrameQueue();
if (!m_devContext->FrameBufferAvailable() || if (!m_devContext->FrameBufferAvailable(schedule) ||
!PublishNewestCandidate( !PublishNewestCandidate(
generation, periodic, publishStart)) schedule, periodic, publishStart))
{ {
ArmPublishTimer(m_publishTimer.Get(), PUBLISH_RETRY_NS); ArmPublishTimer(m_publishTimer.Get(), PUBLISH_RETRY_NS);
if (WaitForMultipleObjects( if (WaitForMultipleObjects(
@@ -907,7 +923,8 @@ void CSwapChainProcessor::SignalCandidateState()
} }
bool CSwapChainProcessor::PublishNewestCandidate( bool CSwapChainProcessor::PublishNewestCandidate(
uint32_t scheduleGeneration, bool periodic, uint64_t publishStart) const CFrameScheduler::Schedule& schedule, bool periodic,
uint64_t publishStart)
{ {
int selectedCandidate = -1; int selectedCandidate = -1;
uint64_t newestSequence = 0; uint64_t newestSequence = 0;
@@ -1007,7 +1024,7 @@ bool CSwapChainProcessor::PublishNewestCandidate(
RECT copyDirtyRects[LG_MAX_DIRTY_RECTS * 2] = {}; RECT copyDirtyRects[LG_MAX_DIRTY_RECTS * 2] = {};
unsigned nbCopyDirtyRects = 0; unsigned nbCopyDirtyRects = 0;
bool fullCopy = bool fullCopy = buffer.fullCopy ||
candidate.nbDirtyRects == 0 || nbPreviousDirtyRects == 0; candidate.nbDirtyRects == 0 || nbPreviousDirtyRects == 0;
if (!fullCopy) if (!fullCopy)
@@ -1060,7 +1077,7 @@ bool CSwapChainProcessor::PublishNewestCandidate(
// Reserve the LGMP message before submitting the copy. This makes post // Reserve the LGMP message before submitting the copy. This makes post
// failure recoverable without racing a very fast GPU completion callback. // failure recoverable without racing a very fast GPU completion callback.
if (!m_devContext->PublishFrameBuffer( if (!m_devContext->PublishFrameBuffer(
buffer.frameIndex, scheduleGeneration)) buffer.frameIndex, schedule))
{ {
copySlot->Cancel(); copySlot->Cancel();
m_devContext->AbortFrameBuffer(buffer.frameIndex); m_devContext->AbortFrameBuffer(buffer.frameIndex);
@@ -1104,7 +1121,7 @@ bool CSwapChainProcessor::PublishNewestCandidate(
ReleaseSRWLockExclusive(&m_damageLock); ReleaseSRWLockExclusive(&m_damageLock);
m_devContext->CommitFrameBuffer( m_devContext->CommitFrameBuffer(
buffer.frameIndex, scheduleGeneration, periodic); buffer.frameIndex, schedule, periodic);
unsigned superseded = 0; unsigned superseded = 0;
AcquireSRWLockExclusive(&m_candidateLock); AcquireSRWLockExclusive(&m_candidateLock);

View File

@@ -134,7 +134,8 @@ private:
static DWORD CALLBACK _PublisherThread(LPVOID arg); static DWORD CALLBACK _PublisherThread(LPVOID arg);
void PublisherThread(); void PublisherThread();
bool PublishNewestCandidate( bool PublishNewestCandidate(
uint32_t scheduleGeneration, bool periodic, uint64_t publishStart); const CFrameScheduler::Schedule& schedule, bool periodic,
uint64_t publishStart);
bool HasReadyCandidate(); bool HasReadyCandidate();
int AcquireCandidate(); int AcquireCandidate();
void ReleaseCandidate(unsigned candidateIndex); void ReleaseCandidate(unsigned candidateIndex);

View File

@@ -36,6 +36,7 @@ void lgFrameSchedulerInit(LGFrameScheduler * scheduler, bool supported,
{ {
memset(scheduler, 0, sizeof(*scheduler)); memset(scheduler, 0, sizeof(*scheduler));
scheduler->supported = supported; scheduler->supported = supported;
scheduler->immediatePending = supported;
scheduler->clientID = clientID; scheduler->clientID = clientID;
} }
@@ -61,42 +62,61 @@ void lgFrameSchedulerSetPeriod(LGFrameScheduler * scheduler,
++scheduler->generation; ++scheduler->generation;
scheduler->resetPending = true; scheduler->resetPending = true;
scheduler->immediatePending = true;
scheduler->phaseError = 0; scheduler->phaseError = 0;
scheduler->feedbackFrameSerial = 0; scheduler->feedbackFrameSerial = 0;
scheduler->feedbackScheduleEpoch = 0;
scheduler->feedbackSamples = 0; scheduler->feedbackSamples = 0;
scheduler->feedbackDirty = false; scheduler->feedbackDirty = false;
} }
void lgFrameSchedulerObserveFrame(LGFrameScheduler * scheduler, void lgFrameSchedulerRequestImmediate(LGFrameScheduler * scheduler)
uint32_t frameSerial, uint32_t generation, uint64_t readyTime)
{ {
if (!generation || if (scheduler->supported)
scheduler->immediatePending = true;
}
void lgFrameSchedulerObserveFrame(LGFrameScheduler * scheduler,
uint32_t frameSerial, uint32_t generation, uint32_t scheduleEpoch,
uint64_t readyTime)
{
if (!generation || !scheduleEpoch ||
(scheduler->readyFrameSerial == frameSerial && (scheduler->readyFrameSerial == frameSerial &&
scheduler->readyGeneration == generation)) scheduler->readyGeneration == generation &&
scheduler->readyScheduleEpoch == scheduleEpoch))
return; return;
scheduler->readyFrameSerial = frameSerial; scheduler->readyFrameSerial = frameSerial;
scheduler->readyGeneration = generation; scheduler->readyGeneration = generation;
scheduler->readyScheduleEpoch = scheduleEpoch;
scheduler->readyTime = readyTime; scheduler->readyTime = readyTime;
} }
void lgFrameSchedulerFeedback(LGFrameScheduler * scheduler, void lgFrameSchedulerFeedback(LGFrameScheduler * scheduler,
uint32_t frameSerial, uint32_t generation, uint64_t neededTime) uint32_t frameSerial, uint32_t generation, uint32_t scheduleEpoch,
uint64_t tickTime)
{ {
if (!scheduler->active || generation != scheduler->generation || if (!scheduler->active || generation != scheduler->generation ||
frameSerial == scheduler->feedbackFrameSerial) !scheduleEpoch ||
(frameSerial == scheduler->feedbackFrameSerial &&
scheduleEpoch == scheduler->feedbackScheduleEpoch) ||
scheduler->readyFrameSerial != frameSerial ||
scheduler->readyGeneration != generation ||
scheduler->readyScheduleEpoch != scheduleEpoch)
return; return;
uint64_t readyTime = neededTime; if (scheduler->feedbackScheduleEpoch &&
if (scheduler->readyFrameSerial == frameSerial && scheduler->feedbackScheduleEpoch != scheduleEpoch)
scheduler->readyGeneration == generation) {
readyTime = scheduler->readyTime; scheduler->phaseError = 0;
scheduler->feedbackSamples = 0;
}
const uint64_t measuredPhase = neededTime > readyTime ? const int64_t measuredPhase = tickTime >= scheduler->readyTime ?
neededTime - readyTime : 0; (int64_t)(tickTime - scheduler->readyTime) :
int64_t error = measuredPhase > FRAME_SCHEDULER_TARGET_SLACK_NS ? -(int64_t)(scheduler->readyTime - tickTime);
(int64_t)(measuredPhase - FRAME_SCHEDULER_TARGET_SLACK_NS) : int64_t error = measuredPhase -
-(int64_t)(FRAME_SCHEDULER_TARGET_SLACK_NS - measuredPhase); (int64_t)FRAME_SCHEDULER_TARGET_SLACK_NS;
const int64_t period = (int64_t)scheduler->period; const int64_t period = (int64_t)scheduler->period;
const int64_t limit = period / 2; const int64_t limit = period / 2;
if (error > limit) if (error > limit)
@@ -112,6 +132,7 @@ void lgFrameSchedulerFeedback(LGFrameScheduler * scheduler,
++scheduler->feedbackSamples; ++scheduler->feedbackSamples;
scheduler->feedbackFrameSerial = frameSerial; scheduler->feedbackFrameSerial = frameSerial;
scheduler->feedbackScheduleEpoch = scheduleEpoch;
scheduler->feedbackDirty = true; scheduler->feedbackDirty = true;
} }
@@ -124,15 +145,20 @@ void lgFrameSchedulerUpdate(LGFrameScheduler * scheduler,
const uint64_t interval = scheduler->feedbackDirty ? const uint64_t interval = scheduler->feedbackDirty ?
FRAME_SCHEDULER_FEEDBACK_NS : FRAME_SCHEDULER_RENEW_NS; FRAME_SCHEDULER_FEEDBACK_NS : FRAME_SCHEDULER_RENEW_NS;
if (scheduler->active && !scheduler->resetPending && if (scheduler->active && !scheduler->resetPending &&
!scheduler->immediatePending &&
now - scheduler->lastSend < interval) now - scheduler->lastSend < interval)
return; return;
KVMFRFrameScheduleFlags flags = KVMFR_FRAME_SCHEDULE_ACTIVE; KVMFRFrameScheduleFlags flags = KVMFR_FRAME_SCHEDULE_ACTIVE;
if (scheduler->resetPending) if (scheduler->resetPending)
flags |= KVMFR_FRAME_SCHEDULE_RESET; flags |= KVMFR_FRAME_SCHEDULE_RESET;
if (scheduler->immediatePending)
flags |= KVMFR_FRAME_SCHEDULE_IMMEDIATE;
const uint32_t feedbackFrameSerial = const uint32_t feedbackFrameSerial =
scheduler->feedbackFrameSerial; scheduler->feedbackFrameSerial;
const uint32_t feedbackScheduleEpoch =
scheduler->feedbackScheduleEpoch;
const KVMFRFrameSchedule message = { const KVMFRFrameSchedule message = {
.msg.type = KVMFR_MESSAGE_FRAME_SCHEDULE, .msg.type = KVMFR_MESSAGE_FRAME_SCHEDULE,
.clientID = scheduler->clientID, .clientID = scheduler->clientID,
@@ -142,6 +168,7 @@ void lgFrameSchedulerUpdate(LGFrameScheduler * scheduler,
.targetSlack = FRAME_SCHEDULER_TARGET_SLACK_NS, .targetSlack = FRAME_SCHEDULER_TARGET_SLACK_NS,
.phaseError = scheduler->phaseError, .phaseError = scheduler->phaseError,
.feedbackFrameSerial = feedbackFrameSerial, .feedbackFrameSerial = feedbackFrameSerial,
.feedbackScheduleEpoch = feedbackScheduleEpoch,
.lease = FRAME_SCHEDULER_LEASE_MS, .lease = FRAME_SCHEDULER_LEASE_MS,
}; };
@@ -150,7 +177,9 @@ void lgFrameSchedulerUpdate(LGFrameScheduler * scheduler,
scheduler->active = true; scheduler->active = true;
scheduler->resetPending = false; scheduler->resetPending = false;
scheduler->immediatePending = false;
scheduler->lastSend = now; scheduler->lastSend = now;
if (scheduler->feedbackFrameSerial == feedbackFrameSerial) if (scheduler->feedbackFrameSerial == feedbackFrameSerial &&
scheduler->feedbackScheduleEpoch == feedbackScheduleEpoch)
scheduler->feedbackDirty = false; scheduler->feedbackDirty = false;
} }

View File

@@ -31,6 +31,7 @@ typedef struct LGFrameScheduler
bool supported; bool supported;
bool active; bool active;
bool resetPending; bool resetPending;
bool immediatePending;
bool feedbackDirty; bool feedbackDirty;
uint32_t clientID; uint32_t clientID;
uint32_t generation; uint32_t generation;
@@ -39,10 +40,12 @@ typedef struct LGFrameScheduler
int64_t phaseError; int64_t phaseError;
uint32_t feedbackFrameSerial; uint32_t feedbackFrameSerial;
uint32_t feedbackScheduleEpoch;
unsigned feedbackSamples; unsigned feedbackSamples;
uint32_t readyFrameSerial; uint32_t readyFrameSerial;
uint32_t readyGeneration; uint32_t readyGeneration;
uint32_t readyScheduleEpoch;
uint64_t readyTime; uint64_t readyTime;
} }
LGFrameScheduler; LGFrameScheduler;
@@ -51,10 +54,13 @@ void lgFrameSchedulerInit(LGFrameScheduler * scheduler, bool supported,
uint32_t clientID); uint32_t clientID);
void lgFrameSchedulerSetPeriod(LGFrameScheduler * scheduler, void lgFrameSchedulerSetPeriod(LGFrameScheduler * scheduler,
uint64_t period); uint64_t period);
void lgFrameSchedulerRequestImmediate(LGFrameScheduler * scheduler);
void lgFrameSchedulerObserveFrame(LGFrameScheduler * scheduler, void lgFrameSchedulerObserveFrame(LGFrameScheduler * scheduler,
uint32_t frameSerial, uint32_t generation, uint64_t readyTime); uint32_t frameSerial, uint32_t generation, uint32_t scheduleEpoch,
uint64_t readyTime);
void lgFrameSchedulerFeedback(LGFrameScheduler * scheduler, void lgFrameSchedulerFeedback(LGFrameScheduler * scheduler,
uint32_t frameSerial, uint32_t generation, uint64_t neededTime); uint32_t frameSerial, uint32_t generation, uint32_t scheduleEpoch,
uint64_t tickTime);
void lgFrameSchedulerUpdate(LGFrameScheduler * scheduler, void lgFrameSchedulerUpdate(LGFrameScheduler * scheduler,
PLGMPClientQueue queue, uint64_t now); PLGMPClientQueue queue, uint64_t now);

240
obs/lg.c
View File

@@ -99,7 +99,9 @@ typedef struct
int bpp; int bpp;
struct IVSHMEM shmDev; struct IVSHMEM shmDev;
PLGMPClient lgmp; PLGMPClient lgmp;
PLGMPClientQueue frameQueue, pointerQueue; PLGMPClientQueue frameQueue;
PLGMPClientQueue frameOwnerQueue[LGMP_Q_FRAME_LEN];
PLGMPClientQueue pointerQueue;
gs_texture_t * texture; gs_texture_t * texture;
gs_texture_t * dstTexture; gs_texture_t * dstTexture;
uint8_t * texData; uint8_t * texData;
@@ -114,7 +116,7 @@ typedef struct
#if LIBOBS_API_MAJOR_VER >= 27 #if LIBOBS_API_MAJOR_VER >= 27
bool dmabuf; bool dmabuf;
bool dmabufTested; bool dmabufTested;
DMAFrameInfo dmaInfo[LGMP_Q_FRAME_LEN]; DMAFrameInfo dmaInfo[LGMP_Q_FRAME_BUFFER_LEN];
gs_texture_t * dmaTexture; gs_texture_t * dmaTexture;
#endif #endif
@@ -124,6 +126,8 @@ typedef struct
pthread_t frameThread, pointerThread; pthread_t frameThread, pointerThread;
os_sem_t * frameSem; os_sem_t * frameSem;
uint32_t frameSerial;
bool frameSerialValid;
pthread_mutex_t pointerLock; pthread_mutex_t pointerLock;
LGFrameScheduler frameScheduler; LGFrameScheduler frameScheduler;
@@ -169,6 +173,16 @@ typedef struct
} }
LGPlugin; LGPlugin;
typedef struct
{
LGMPMessage msg;
PLGMPClientQueue queue;
uint32_t generation;
uint32_t scheduleEpoch;
bool owner;
}
LGFrameMessage;
static void * frameThread(void * data); static void * frameThread(void * data);
static void * pointerThread(void * data); static void * pointerThread(void * data);
static void lgUpdate(void * data, obs_data_t * settings); static void lgUpdate(void * data, obs_data_t * settings);
@@ -188,6 +202,95 @@ static const char * lgGetName(void * unused)
return obs_module_text("Looking Glass Client"); return obs_module_text("Looking Glass Client");
} }
static bool lgFrameSerialNewer(uint32_t lhs, uint32_t rhs)
{
return (int32_t)(lhs - rhs) > 0;
}
static LGMP_STATUS lgFrameProcessNewest(LGPlugin * this,
LGFrameMessage * result)
{
struct
{
PLGMPClientQueue queue;
bool owner;
}
queues[LGMP_Q_FRAME_LEN + 1] =
{
{ this->frameQueue, false },
};
for (unsigned int i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
queues[i + 1].queue = this->frameOwnerQueue[i];
queues[i + 1].owner = true;
}
memset(result, 0, sizeof(*result));
for (unsigned int i = 0; i < ARRAY_LENGTH(queues); ++i)
{
if (!queues[i].queue)
continue;
LGMP_STATUS status = lgmpClientAdvanceToLast(queues[i].queue);
if (status != LGMP_OK && status != LGMP_ERR_QUEUE_EMPTY)
return status;
LGMPMessage msg;
status = lgmpClientProcess(queues[i].queue, &msg);
if (status == LGMP_ERR_QUEUE_EMPTY)
continue;
if (status != LGMP_OK)
return status;
const KVMFRFrame * frame = (const KVMFRFrame *)msg.mem;
const bool owner = queues[i].owner || msg.udata != 0;
if (!result->queue)
{
result->msg = msg;
result->queue = queues[i].queue;
result->owner = owner;
continue;
}
const KVMFRFrame * selected = (const KVMFRFrame *)result->msg.mem;
const bool replace = lgFrameSerialNewer(
frame->frameSerial, selected->frameSerial) ||
(frame->frameSerial == selected->frameSerial &&
owner && !result->owner);
PLGMPClientQueue discard = queues[i].queue;
if (replace)
{
discard = result->queue;
result->msg = msg;
result->queue = queues[i].queue;
result->owner = owner;
}
status = lgmpClientMessageDone(discard);
if (status != LGMP_OK)
return status;
}
if (result->owner)
{
result->generation = (uint32_t)(result->msg.udata >> 32);
result->scheduleEpoch = (uint32_t)result->msg.udata;
}
return result->queue ? LGMP_OK : LGMP_ERR_QUEUE_EMPTY;
}
static void lgFrameUnsubscribeOwnerQueues(LGPlugin * this)
{
for (unsigned int i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
if (this->frameOwnerQueue[i])
lgmpClientUnsubscribe(&this->frameOwnerQueue[i]);
this->frameOwnerQueue[i] = NULL;
}
}
static void * lgCreate(obs_data_t * settings, obs_source_t * context) static void * lgCreate(obs_data_t * settings, obs_source_t * context)
{ {
LGPlugin * this = bzalloc(sizeof(LGPlugin)); LGPlugin * this = bzalloc(sizeof(LGPlugin));
@@ -479,6 +582,12 @@ static obs_properties_t * lgGetProperties(void * data)
static void * frameThread(void * data) static void * frameThread(void * data)
{ {
LGPlugin * this = (LGPlugin *)data; LGPlugin * this = (LGPlugin *)data;
for (unsigned int i = 0; i < LGMP_Q_FRAME_LEN; ++i)
this->frameOwnerQueue[i] = NULL;
this->frameQueue = NULL;
this->frameSerial = 0;
this->frameSerialValid = false;
if (lgmpClientSubscribe( if (lgmpClientSubscribe(
this->lgmp, LGMP_Q_FRAME, &this->frameQueue) != LGMP_OK) this->lgmp, LGMP_Q_FRAME, &this->frameQueue) != LGMP_OK)
@@ -487,31 +596,47 @@ static void * frameThread(void * data)
return NULL; return NULL;
} }
if (this->frameScheduler.supported)
{
for (unsigned int i = 0; i < LGMP_Q_FRAME_LEN; ++i)
{
if (lgmpClientSubscribe(
this->lgmp, LGMP_Q_FRAME_OWNER + i,
&this->frameOwnerQueue[i]) == LGMP_OK)
continue;
lgFrameUnsubscribeOwnerQueues(this);
lgmpClientUnsubscribe(&this->frameQueue);
this->frameQueue = NULL;
this->state = STATE_STOPPING;
return NULL;
}
}
lgFrameSchedulerRequestImmediate(&this->frameScheduler);
this->state = STATE_RUNNING; this->state = STATE_RUNNING;
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
while(this->state == STATE_RUNNING) while(this->state == STATE_RUNNING)
{ {
LGMP_STATUS status;
os_sem_wait(this->frameSem); os_sem_wait(this->frameSem);
if ((status = lgmpClientAdvanceToLast(this->frameQueue)) != LGMP_OK) LGFrameMessage frameMessage;
{ const LGMP_STATUS status = lgFrameProcessNewest(this, &frameMessage);
if (status != LGMP_ERR_QUEUE_EMPTY) if (status != LGMP_OK && status != LGMP_ERR_QUEUE_EMPTY)
{ {
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
printf("lgmpClientAdvanceToLast: %s\n", lgmpStatusString(status)); printf("lgFrameProcessNewest: %s\n", lgmpStatusString(status));
break; break;
} }
}
uint64_t now = os_gettime_ns(); uint64_t now = os_gettime_ns();
if (status == LGMP_OK) if (status == LGMP_OK)
{ {
LGMPMessage msg; const KVMFRFrame * frame =
if (lgmpClientProcess(this->frameQueue, &msg) == LGMP_OK) (const KVMFRFrame *)frameMessage.msg.mem;
if (frameMessage.owner)
{ {
const KVMFRFrame * frame = (const KVMFRFrame *)msg.mem;
const FrameBuffer * fb = const FrameBuffer * fb =
(const FrameBuffer *)((const uint8_t *)frame + frame->offset); (const FrameBuffer *)((const uint8_t *)frame + frame->offset);
if (framebuffer_wait( if (framebuffer_wait(
@@ -519,7 +644,8 @@ static void * frameThread(void * data)
{ {
now = os_gettime_ns(); now = os_gettime_ns();
lgFrameSchedulerObserveFrame(&this->frameScheduler, lgFrameSchedulerObserveFrame(&this->frameScheduler,
frame->frameSerial, msg.udata, now); frame->frameSerial, frameMessage.generation,
frameMessage.scheduleEpoch, now);
} }
} }
} }
@@ -532,7 +658,9 @@ static void * frameThread(void * data)
usleep(1000); usleep(1000);
} }
lgFrameUnsubscribeOwnerQueues(this);
lgmpClientUnsubscribe(&this->frameQueue); lgmpClientUnsubscribe(&this->frameQueue);
this->frameQueue = NULL;
this->state = STATE_RESTARTING; this->state = STATE_RESTARTING;
return NULL; return NULL;
} }
@@ -584,7 +712,7 @@ static void * pointerThread(void * data)
if (msg.udata & CURSOR_FLAG_VISIBLE_VALID) if (msg.udata & CURSOR_FLAG_VISIBLE_VALID)
{ {
this->cursorVisible = this->hideMouse ? this->cursorVisible = this->hideMouse ?
0 : msg.udata & CURSOR_FLAG_VISIBLE; false : (msg.udata & CURSOR_FLAG_VISIBLE) != 0;
if (cursor->sdrWhiteLevel) if (cursor->sdrWhiteLevel)
atomic_store(&this->sdrWhiteLevel, cursor->sdrWhiteLevel); atomic_store(&this->sdrWhiteLevel, cursor->sdrWhiteLevel);
} }
@@ -966,7 +1094,7 @@ static void lgComputeColorMatrix(LGPlugin * this)
} }
static void lgFormatInit(LGPlugin * this, const KVMFRFrame * frame, static void lgFormatInit(LGPlugin * this, const KVMFRFrame * frame,
LGMPMessage * msg) LGMPMessage * msg, PLGMPClientQueue frameQueue)
{ {
this->formatVer = frame->formatVer; this->formatVer = frame->formatVer;
this->screenWidth = frame->screenWidth; this->screenWidth = frame->screenWidth;
@@ -1061,7 +1189,7 @@ static void lgFormatInit(LGPlugin * this, const KVMFRFrame * frame,
default: default:
printf("invalid type %d\n", this->type); printf("invalid type %d\n", this->type);
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(frameQueue);
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
obs_leave_graphics(); obs_leave_graphics();
return; return;
@@ -1101,7 +1229,7 @@ static void lgFormatInit(LGPlugin * this, const KVMFRFrame * frame,
if (!this->texture) if (!this->texture)
{ {
printf("create texture failed\n"); printf("create texture failed\n");
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(frameQueue);
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
obs_leave_graphics(); obs_leave_graphics();
return; return;
@@ -1143,7 +1271,7 @@ static void lgVideoTick(void * data, float seconds)
const uint64_t tickTime = os_gettime_ns(); const uint64_t tickTime = os_gettime_ns();
const uint64_t framePeriod = lgFramePeriod(); const uint64_t framePeriod = lgFramePeriod();
LGMP_STATUS status; LGMP_STATUS status;
LGMPMessage msg; LGFrameMessage frameMessage;
os_sem_wait(this->frameSem); os_sem_wait(this->frameSem);
if (this->state != STATE_RUNNING) if (this->state != STATE_RUNNING)
@@ -1269,17 +1397,7 @@ static void lgVideoTick(void * data, float seconds)
os_sem_post(this->cursorSem); os_sem_post(this->cursorSem);
} }
if ((status = lgmpClientAdvanceToLast(this->frameQueue)) != LGMP_OK) if ((status = lgFrameProcessNewest(this, &frameMessage)) != LGMP_OK)
{
if (status != LGMP_ERR_QUEUE_EMPTY)
{
os_sem_post(this->frameSem);
printf("lgmpClientAdvanceToLast: %s\n", lgmpStatusString(status));
return;
}
}
if ((status = lgmpClientProcess(this->frameQueue, &msg)) != LGMP_OK)
{ {
if (status == LGMP_ERR_QUEUE_EMPTY) if (status == LGMP_ERR_QUEUE_EMPTY)
{ {
@@ -1287,25 +1405,64 @@ static void lgVideoTick(void * data, float seconds)
return; return;
} }
printf("lgmpClientProcess: %s\n", lgmpStatusString(status)); printf("lgFrameProcessNewest: %s\n", lgmpStatusString(status));
this->state = STATE_STOPPING; this->state = STATE_STOPPING;
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
return; return;
} }
const KVMFRFrame * frame = (KVMFRFrame *)msg.mem; const KVMFRFrame * frame = (const KVMFRFrame *)frameMessage.msg.mem;
const bool equalSerial = this->frameSerialValid &&
frame->frameSerial == this->frameSerial;
if (this->frameSerialValid && !equalSerial &&
!lgFrameSerialNewer(frame->frameSerial, this->frameSerial))
{
lgmpClientMessageDone(frameMessage.queue);
os_sem_post(this->frameSem);
return;
}
if (equalSerial && !frameMessage.owner)
{
lgmpClientMessageDone(frameMessage.queue);
os_sem_post(this->frameSem);
return;
}
const FrameBuffer * fb =
(const FrameBuffer *)((const uint8_t *)frame + frame->offset);
bool frameComplete = false;
if (frameMessage.owner)
{
frameComplete = framebuffer_wait(
fb, (size_t)frame->dataHeight * frame->pitch);
if (!frameComplete)
{
lgmpClientMessageDone(frameMessage.queue);
os_sem_post(this->frameSem);
return;
}
const uint64_t readyTime = os_gettime_ns();
lgFrameSchedulerObserveFrame(&this->frameScheduler,
frame->frameSerial, frameMessage.generation,
frameMessage.scheduleEpoch, readyTime);
lgFrameSchedulerFeedback(&this->frameScheduler, lgFrameSchedulerFeedback(&this->frameScheduler,
frame->frameSerial, msg.udata, tickTime); frame->frameSerial, frameMessage.generation,
frameMessage.scheduleEpoch, tickTime);
}
this->frameSerial = frame->frameSerial;
this->frameSerialValid = true;
bool textureValid = (this->dmabufTested && this->dmabuf) || this->texture; bool textureValid = (this->dmabufTested && this->dmabuf) || this->texture;
if (!textureValid || this->formatVer != frame->formatVer) if (!textureValid || this->formatVer != frame->formatVer)
lgFormatInit(this, frame, &msg); lgFormatInit(this, frame, &frameMessage.msg, frameMessage.queue);
#if LIBOBS_API_MAJOR_VER >= 27 #if LIBOBS_API_MAJOR_VER >= 27
if (this->dmabuf) if (this->dmabuf)
{ {
obs_enter_graphics(); obs_enter_graphics();
DMAFrameInfo * fi = dmabufOpenDMAFrameInfo(this, &msg, frame, DMAFrameInfo * fi = dmabufOpenDMAFrameInfo(this, &frameMessage.msg, frame,
(size_t)frame->dataHeight * frame->pitch); (size_t)frame->dataHeight * frame->pitch);
bool importFailed = !fi; bool importFailed = !fi;
if (fi && !fi->texture) if (fi && !fi->texture)
@@ -1315,12 +1472,12 @@ static void lgVideoTick(void * data, float seconds)
if (!importFailed) if (!importFailed)
{ {
// wait for the frame to be complete before we try to use it // wait for the frame to be complete before we try to use it
FrameBuffer * fb = (FrameBuffer *)(((uint8_t*)frame) + frame->offset); if (!frameComplete)
const bool complete = framebuffer_wait( frameComplete = framebuffer_wait(
fb, (size_t)frame->dataHeight * frame->pitch); fb, (size_t)frame->dataHeight * frame->pitch);
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(frameMessage.queue);
if (!complete) if (!frameComplete)
{ {
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
return; return;
@@ -1333,18 +1490,17 @@ static void lgVideoTick(void * data, float seconds)
puts("Failed to create dmabuf texture, falling back to CPU upload"); puts("Failed to create dmabuf texture, falling back to CPU upload");
this->dmabuf = false; this->dmabuf = false;
lgFormatInit(this, frame, &msg); lgFormatInit(this, frame, &frameMessage.msg, frameMessage.queue);
} }
#endif #endif
if (!this->texture) if (!this->texture)
{ {
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(frameMessage.queue);
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
return; return;
} }
FrameBuffer * fb = (FrameBuffer *)(((uint8_t*)frame) + frame->offset);
framebuffer_read( framebuffer_read(
fb, fb,
this->texData , // dst this->texData , // dst
@@ -1355,7 +1511,7 @@ static void lgVideoTick(void * data, float seconds)
frame->pitch frame->pitch
); );
lgmpClientMessageDone(this->frameQueue); lgmpClientMessageDone(frameMessage.queue);
os_sem_post(this->frameSem); os_sem_post(this->frameSem);
obs_enter_graphics(); obs_enter_graphics();