diff --git a/idd/LGIdd/capture/CFrameScheduler.cpp b/idd/LGIdd/capture/CFrameScheduler.cpp index 91ba9579..187865d7 100644 --- a/idd/LGIdd/capture/CFrameScheduler.cpp +++ b/idd/LGIdd/capture/CFrameScheduler.cpp @@ -701,13 +701,12 @@ bool CFrameScheduler::TryFrameSubmitted(const Schedule& schedule, publication.deadlineSerial = schedule.deadlineSerial; publication.frameSerial = frameSerial; publication.deadline = schedule.deadline; - publication.committed = true; registered = true; } return registered; } -void CFrameScheduler::FramePublished(const Schedule& schedule, +void CFrameScheduler::FrameCommitted(const Schedule& schedule, uint32_t frameSerial, uint64_t now, bool periodic) { CSRWExclusiveLock lock(m_lock); @@ -721,6 +720,25 @@ void CFrameScheduler::FramePublished(const Schedule& schedule, if (publication) publication->committed = true; + if (periodic) + { + AdvanceCurrentDeadline(); + AdvanceDeadline(now); + } + Client * client = FindClient(m_schedule.clientID); + if (client) + client->nextDelivery = m_nextDeadline; + } +} + +void CFrameScheduler::FramePublished(const Schedule& schedule, + uint32_t frameSerial) +{ + CSRWExclusiveLock lock(m_lock); + if (m_scheduling && schedule.clientID == m_schedule.clientID && + schedule.generation == m_schedule.generation && + schedule.epoch == m_schedule.epoch) + { if (schedule.forceTicket > m_forceAckTicket) m_forceAckTicket = schedule.forceTicket; if (schedule.republishTicket > m_republishAckTicket) @@ -735,26 +753,16 @@ void CFrameScheduler::FramePublished(const Schedule& schedule, client->deliveredFrameValid = true; } ++m_publishedFrames; - if (periodic) - { - AdvanceCurrentDeadline(); - AdvanceDeadline(now); - } - if (client) - client->nextDelivery = m_nextDeadline; } } -void CFrameScheduler::FrameRetained(const Schedule& schedule, - uint64_t now, bool periodic) +void CFrameScheduler::FrameRetained(const Schedule& schedule) { bool wake = false; CSRWExclusiveLock lock(m_lock); if (m_scheduling && schedule.clientID == m_schedule.clientID && schedule.generation == m_schedule.generation && - schedule.epoch == m_schedule.epoch && - schedule.deadlineSerial == m_deadlineSerial && - schedule.deadline == m_nextDeadline) + schedule.epoch == m_schedule.epoch) { if (schedule.forceTicket > m_forceAckTicket) m_forceAckTicket = schedule.forceTicket; @@ -766,16 +774,6 @@ void CFrameScheduler::FrameRetained(const Schedule& schedule, ++m_republishRequestTicket; wake = true; } - - if (periodic) - { - AdvanceCurrentDeadline(); - AdvanceDeadline(now); - } - - Client * client = FindClient(m_schedule.clientID); - if (client) - client->nextDelivery = m_nextDeadline; } lock.Unlock(); @@ -783,6 +781,11 @@ void CFrameScheduler::FrameRetained(const Schedule& schedule, WakePublisher(); } +void CFrameScheduler::FrameFailed() +{ + ForceFrame(); +} + bool CFrameScheduler::TryFrameCompleted(const Schedule& schedule, uint32_t frameSerial, uint64_t completedAt) { diff --git a/idd/LGIdd/capture/CFrameScheduler.h b/idd/LGIdd/capture/CFrameScheduler.h index eed835ea..029bd505 100644 --- a/idd/LGIdd/capture/CFrameScheduler.h +++ b/idd/LGIdd/capture/CFrameScheduler.h @@ -172,10 +172,11 @@ public: void FrameMissed(const Schedule& schedule, uint64_t now, bool periodic); void FrameSuperseded(); bool TryFrameSubmitted(const Schedule& schedule, uint32_t frameSerial); - void FramePublished(const Schedule& schedule, uint32_t frameSerial, + void FrameCommitted(const Schedule& schedule, uint32_t frameSerial, uint64_t now, bool periodic); - void FrameRetained(const Schedule& schedule, uint64_t now, - bool periodic); + void FramePublished(const Schedule& schedule, uint32_t frameSerial); + void FrameRetained(const Schedule& schedule); + void FrameFailed(); void FrameRepublished(const Schedule& schedule, uint32_t frameSerial); bool TryFrameCompleted(const Schedule& schedule, uint32_t frameSerial, uint64_t completedAt); diff --git a/idd/LGIdd/transport/CFrameHub.cpp b/idd/LGIdd/transport/CFrameHub.cpp index 042cf999..34231bda 100644 --- a/idd/LGIdd/transport/CFrameHub.cpp +++ b/idd/LGIdd/transport/CFrameHub.cpp @@ -27,6 +27,13 @@ static bool ContentAfter(uint64_t value, uint64_t current) static_cast(value - current) > 0); } +static bool TokenMatches(const FrameToken& left, const FrameToken& right) +{ + return left.backend == right.backend && + left.epoch == right.epoch && left.slot == right.slot && + left.serial == right.serial; +} + CFrameHub::CFrameHub() { m_wakeEvent = CreateEvent(nullptr, FALSE, FALSE, nullptr); @@ -89,8 +96,14 @@ bool CFrameHub::Bind(BackendId backend, uint32_t epoch, bool primary, selected->backend.store(backend, std::memory_order_release); selected->epoch.store(epoch, std::memory_order_release); selected->primary = primary; + selected->outstanding.store(0, std::memory_order_release); + SetEvent(selected->drained); + } + { + CSRWExclusiveLock lock(selected->laneLock); selected->needsFullCopy = true; selected->lastContent = 0; + selected->lastTerminalContent = 0; selected->blockedContent = 0; selected->blockedFrameSize = 0; selected->blockedAllocation = false; @@ -101,11 +114,10 @@ bool CFrameHub::Bind(BackendId backend, uint32_t epoch, bool primary, selected->frameType = FRAME_TYPE_INVALID; for (Sink::ResourceLane& lane : selected->lanes) lane = {}; - selected->outstanding.store(0, std::memory_order_release); - SetEvent(selected->drained); - target.SetFrameScheduleEvent(m_wakeEvent); } + target.SetFrameEvents(this); + target.SetFrameScheduleEvent(m_wakeEvent); { CSRWExclusiveLock lock(m_listLock); selected->active.store(true, std::memory_order_release); @@ -135,21 +147,52 @@ void CFrameHub::Unbind(BackendId backend, uint32_t epoch) if (!selected) return; + IFrameSink * target; { CSRWExclusiveLock lock(selected->callLock); - if (selected->target) - selected->target->SetFrameScheduleEvent(nullptr); + target = selected->target; } + if (target) + target->SetFrameScheduleEvent(nullptr); + + FrameToken cancel[FRAME_SINK_BUFFERS] = {}; + unsigned cancelCount = 0; + { + CSRWExclusiveLock lock(selected->laneLock); + for (Sink::ResourceLane& lane : selected->lanes) + { + if (lane.phase == Sink::ResourceLane::PENDING) + cancel[cancelCount++] = lane.token; + else if (lane.phase == Sink::ResourceLane::FILLING) + lane.cancelRequested = true; + } + } + for (unsigned i = 0; i < cancelCount; ++i) + for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) + { + CSRWSharedLock lock(selected->laneLock); + if (!TokenMatches(selected->lanes[lane].token, cancel[i])) + continue; + lock.Unlock(); + CancelLane(*selected, lane, cancel[i]); + break; + } WaitForSingleObject(selected->drained, INFINITE); + if (target) + target->SetFrameEvents(nullptr); { CSRWExclusiveLock lock(selected->callLock); selected->target = nullptr; selected->backend.store(0, std::memory_order_release); selected->epoch.store(0, std::memory_order_release); selected->primary = false; + } + { + CSRWExclusiveLock lock(selected->laneLock); selected->needsFullCopy = true; selected->lastContent = 0; + selected->lastTerminalContent = 0; selected->blockedContent = 0; selected->blockedFrameSize = 0; selected->blockedAllocation = false; @@ -186,7 +229,10 @@ void CFrameHub::ReleaseTarget(Batch& batch, BatchTarget& target) return; target.releasePending = true; if (target.resourceLane < FRAME_SINK_BUFFERS) + { + CSRWExclusiveLock lock(target.sink->laneLock); target.sink->lanes[target.resourceLane].busy = false; + } if (target.sink->outstanding.fetch_sub( 1, std::memory_order_acq_rel) == 1) SetEvent(target.sink->drained); @@ -199,37 +245,276 @@ void CFrameHub::ReleaseTarget(Batch& batch, BatchTarget& target) batch.active = false; } -void CFrameHub::CompleteTarget( - Batch& batch, BatchTarget& target, bool succeeded) +bool CFrameHub::DetachTarget(Batch& batch, BatchTarget& target, + FrameToken& token) { - CSRWExclusiveLock call(target.sink->callLock); - if (target.sink->target && - target.sink->backend.load(std::memory_order_acquire) == - target.backend && - target.sink->epoch.load(std::memory_order_acquire) == target.epoch) - { - target.sink->target->CompleteFrameBuffer(target.localSlot, succeeded); - if (succeeded) - { - if (ContentAfter(target.content, target.sink->lastContent)) - { - target.sink->lastContent = target.content; - target.sink->pitch = target.pitch; - target.sink->width = target.width; - target.sink->height = target.height; - target.sink->format = target.format; - target.sink->frameType = target.frameType; - target.sink->needsFullCopy = false; - target.sink->blockedContent = 0; - target.sink->blockedFrameSize = 0; - target.sink->blockedAllocation = false; - } + if (!target.active || target.releasePending || + target.resourceLane >= FRAME_SINK_BUFFERS) + return false; - } - else - target.sink->needsFullCopy = true; + token.backend = target.backend; + token.epoch = target.epoch; + token.slot = target.localSlot; + token.serial = batch.serial; + { + CSRWExclusiveLock lock(target.sink->laneLock); + Sink::ResourceLane& lane = + target.sink->lanes[target.resourceLane]; + if (!lane.busy || lane.phase != Sink::ResourceLane::IDLE) + return false; + lane.token = token; + lane.schedule = target.schedule; + lane.content = target.content; + lane.captureTime = target.captureTime; + lane.postProcessTime = target.postProcessTime; + lane.copyTime = target.copyTime; + lane.readyTime = target.readyTime; + lane.holdTime = target.holdTime; + lane.filledAt = target.filledAt; + lane.workStart = target.workStart; + lane.pitch = target.pitch; + lane.width = target.width; + lane.height = target.height; + lane.format = target.format; + lane.frameType = target.frameType; + lane.phase = Sink::ResourceLane::FILLING; + lane.timingValid = target.timingValid; + lane.resultPending = false; + lane.callActive = false; + lane.cancelRequested = false; + } + + target.releasePending = true; + target.active = false; + for (unsigned i = 0; i < batch.count; ++i) + if (batch.targets[i].active) + return true; + batch.active = false; + return true; +} + +void CFrameHub::FillLane(Sink& sink, unsigned laneIndex, + const FrameToken& token) +{ + if (laneIndex >= FRAME_SINK_BUFFERS) + return; + + { + CSRWExclusiveLock lock(sink.laneLock); + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (lane.phase != Sink::ResourceLane::FILLING || + !TokenMatches(lane.token, token)) + return; + lane.callActive = true; + } + + IFrameSink * target = sink.target; + const FrameFill fill = target ? target->FrameFilled(token) : + FrameFill::REJECTED; + + bool terminal = false; + bool cancel = false; + FrameDone result = FrameDone::FAILED; + uint64_t readyAt = 0; + { + CSRWExclusiveLock lock(sink.laneLock); + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (lane.phase != Sink::ResourceLane::FILLING || + !TokenMatches(lane.token, token)) + return; + lane.callActive = false; + if (lane.resultPending) + { + terminal = true; + result = lane.pendingResult; + readyAt = lane.pendingReadyAt; + lane.resultPending = false; + } + else if (fill == FrameFill::READY) + { + terminal = true; + result = FrameDone::READY; + } + else if (fill == FrameFill::REJECTED) + terminal = true; + else + { + lane.phase = Sink::ResourceLane::PENDING; + cancel = lane.cancelRequested || + !sink.active.load(std::memory_order_acquire); + } + } + + if (terminal) + CompleteLane(sink, laneIndex, token, result, readyAt); + else if (cancel) + CancelLane(sink, laneIndex, token); +} + +void CFrameHub::CancelLane(Sink& sink, unsigned laneIndex, + const FrameToken& token) +{ + if (laneIndex >= FRAME_SINK_BUFFERS) + return; + + { + CSRWExclusiveLock lock(sink.laneLock); + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (lane.phase != Sink::ResourceLane::PENDING || + !TokenMatches(lane.token, token)) + return; + lane.phase = Sink::ResourceLane::CANCELING; + lane.callActive = true; + } + + IFrameSink * target = sink.target; + if (target) + target->CancelFrame(token); + + FrameDone result = FrameDone::SUPERSEDED; + uint64_t readyAt = 0; + { + CSRWExclusiveLock lock(sink.laneLock); + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (lane.phase != Sink::ResourceLane::CANCELING || + !TokenMatches(lane.token, token)) + return; + lane.callActive = false; + if (lane.resultPending) + { + result = lane.pendingResult; + readyAt = lane.pendingReadyAt; + lane.resultPending = false; + } + } + CompleteLane(sink, laneIndex, token, result, readyAt); +} + +void CFrameHub::CompleteLane(Sink& sink, unsigned laneIndex, + const FrameToken& token, FrameDone result, uint64_t readyAt) +{ + if (laneIndex >= FRAME_SINK_BUFFERS) + return; + + Sink::ResourceLane completed; + { + CSRWExclusiveLock lock(sink.laneLock); + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (!lane.busy || lane.phase == Sink::ResourceLane::IDLE || + lane.phase == Sink::ResourceLane::COMPLETING || + !TokenMatches(lane.token, token)) + return; + lane.phase = Sink::ResourceLane::COMPLETING; + completed = lane; + } + + if (!readyAt) + readyAt = CFrameScheduler::Nanotime(); + IFrameSink * target = sink.target; + if (target) + { + if (result == FrameDone::READY && completed.timingValid) + { + uint64_t readyTime = completed.readyTime; + if (readyAt >= completed.filledAt) + readyTime += readyAt - completed.filledAt; + target->SetFrameTiming(completed.localSlot, completed.captureTime, + completed.postProcessTime, completed.copyTime, readyTime, + completed.holdTime, completed.schedule, readyAt); + if (completed.workStart && readyAt >= completed.workStart) + target->TryRecordFrameTiming(readyAt - completed.workStart); + } + target->CompleteFrameBuffer(completed.localSlot, result); + } + + { + CSRWExclusiveLock lock(sink.laneLock); + const bool newestTerminal = + ContentAfter(completed.content, sink.lastTerminalContent); + if (newestTerminal) + sink.lastTerminalContent = completed.content; + if (result == FrameDone::READY && + (newestTerminal || + completed.content == sink.lastTerminalContent)) + { + if (ContentAfter(completed.content, sink.lastContent)) + { + sink.lastContent = completed.content; + sink.pitch = completed.pitch; + sink.width = completed.width; + sink.height = completed.height; + sink.format = completed.format; + sink.frameType = completed.frameType; + sink.needsFullCopy = false; + sink.blockedContent = 0; + sink.blockedFrameSize = 0; + sink.blockedAllocation = false; + } + } + else if (newestTerminal) + sink.needsFullCopy = true; + Sink::ResourceLane& lane = sink.lanes[laneIndex]; + if (lane.phase == Sink::ResourceLane::COMPLETING && + TokenMatches(lane.token, token)) + { + lane.token = {}; + lane.schedule = {}; + lane.content = 0; + lane.phase = Sink::ResourceLane::IDLE; + lane.timingValid = false; + lane.resultPending = false; + lane.callActive = false; + lane.cancelRequested = false; + lane.busy = false; + } + } + if (sink.outstanding.fetch_sub(1, std::memory_order_acq_rel) == 1) + SetEvent(sink.drained); + SetEvent(m_wakeEvent); +} + +void CFrameHub::OnFrameDone(const FrameToken& token, FrameDone result, + uint64_t readyAt) +{ + if (!token.backend || !token.epoch || !token.serial) + return; + + for (Sink& sink : m_sinks) + { + if (sink.backend.load(std::memory_order_acquire) != token.backend || + sink.epoch.load(std::memory_order_acquire) != token.epoch) + continue; + + unsigned laneIndex = FRAME_SINK_BUFFERS; + { + CSRWExclusiveLock lock(sink.laneLock); + for (unsigned i = 0; i < FRAME_SINK_BUFFERS; ++i) + { + Sink::ResourceLane& lane = sink.lanes[i]; + if (!lane.busy || !TokenMatches(lane.token, token) || + lane.phase == Sink::ResourceLane::IDLE || + lane.phase == Sink::ResourceLane::COMPLETING) + continue; + if (lane.callActive) + { + if (!lane.resultPending) + { + lane.pendingResult = result; + lane.pendingReadyAt = readyAt; + lane.resultPending = true; + } + return; + } + if (lane.phase == Sink::ResourceLane::PENDING || + lane.phase == Sink::ResourceLane::CANCELING) + laneIndex = i; + break; + } + } + if (laneIndex < FRAME_SINK_BUFFERS) + CompleteLane(sink, laneIndex, token, result, readyAt); + return; } - ReleaseTarget(batch, target); } size_t CFrameHub::GetMaxFrameSize() const @@ -281,11 +566,12 @@ bool CFrameHub::NeedsFrame() const if (!refs[i].sink->active.load(std::memory_order_acquire) || !refs[i].sink->target) continue; + const size_t maxFrameSize = refs[i].sink->target->GetMaxFrameSize(); + CSRWSharedLock state(refs[i].sink->laneLock); if (refs[i].sink->blockedContent == newest && refs[i].sink->blockedFrameSize && (refs[i].sink->blockedAllocation || - refs[i].sink->target->GetMaxFrameSize() < - refs[i].sink->blockedFrameSize)) + maxFrameSize < refs[i].sink->blockedFrameSize)) continue; if (refs[i].sink->needsFullCopy || ContentAfter(newest, refs[i].sink->lastContent)) @@ -485,18 +771,24 @@ bool CFrameHub::PrepareFrameBatch(const FramePlan& plan, const size_t maxFrameSize = sink.target->GetMaxFrameSize(); if (!maxFrameSize || frameSize > maxFrameSize) { - sink.blockedContent = contentSerial; - sink.blockedFrameSize = frameSize; + CSRWExclusiveLock lock(sink.laneLock); + sink.blockedContent = contentSerial; + sink.blockedFrameSize = frameSize; sink.blockedAllocation = false; continue; } - const bool layoutChanged = sink.pitch != pitch || - sink.width != dstFormat.width || sink.height != dstFormat.height || - sink.format != dstFormat.desc.Format || - sink.frameType != dstFormat.format; - const bool forceFull = sink.needsFullCopy || layoutChanged || - sink.lastContent + 1 != contentSerial; + bool forceFull; + { + CSRWSharedLock lock(sink.laneLock); + const bool layoutChanged = sink.pitch != pitch || + sink.width != dstFormat.width || + sink.height != dstFormat.height || + sink.format != dstFormat.desc.Format || + sink.frameType != dstFormat.format; + forceFull = sink.needsFullCopy || layoutChanged || + sink.lastContent + 1 != contentSerial; + } const RECT * damage = forceFull ? nullptr : dirtyRects; const unsigned damageCount = forceFull ? 0 : nbDirtyRects; SinkTarget result = sink.target->PrepareFrameBuffer( @@ -506,99 +798,107 @@ bool CFrameHub::PrepareFrameBatch(const FramePlan& plan, { if (result.mem) { - sink.blockedContent = contentSerial; - sink.blockedFrameSize = frameSize; - sink.blockedAllocation = true; + { + CSRWExclusiveLock lock(sink.laneLock); + sink.blockedContent = contentSerial; + sink.blockedFrameSize = frameSize; + sink.blockedAllocation = true; + } sink.target->AbortFrameBuffer(result.slot); } continue; } - unsigned resourceLane = FRAME_SINK_BUFFERS; - for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) { - const Sink::ResourceLane& candidate = sink.lanes[lane]; - if (!candidate.busy && candidate.valid && - candidate.mem == result.mem && - candidate.heapOffset == result.heapOffset && - candidate.localSlot == result.slot && - candidate.direct == request.primary) - { - resourceLane = lane; - break; - } - } - if (resourceLane == FRAME_SINK_BUFFERS) + CSRWExclusiveLock lock(sink.laneLock); + unsigned resourceLane = FRAME_SINK_BUFFERS; for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) - if (!sink.lanes[lane].busy && !sink.lanes[lane].valid) + { + const Sink::ResourceLane& candidate = sink.lanes[lane]; + if (!candidate.busy && candidate.valid && + candidate.mem == result.mem && + candidate.heapOffset == result.heapOffset && + candidate.localSlot == result.slot && + candidate.direct == request.primary) { resourceLane = lane; break; } - if (resourceLane == FRAME_SINK_BUFFERS) - for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) - if (!sink.lanes[lane].busy) - { - resourceLane = lane; - break; - } - if (resourceLane == FRAME_SINK_BUFFERS) - { - sink.target->AbortFrameBuffer(result.slot); - continue; - } - - Sink::ResourceLane& lane = sink.lanes[resourceLane]; - lane.mem = result.mem; - lane.heapOffset = result.heapOffset; - lane.localSlot = result.slot; - lane.direct = request.primary; - lane.valid = true; - lane.busy = true; - - unsigned index = batch->count; - if (sink.primary && index) - { - for (unsigned move = index; move != 0; --move) - { - batch->targets[move] = batch->targets[move - 1]; - prepared.targets[move] = prepared.targets[move - 1]; } - index = 0; - } - ++batch->count; - BatchTarget& target = batch->targets[index]; - target.sink = &sink; - target.backend = request.backend; - target.epoch = request.epoch; - target.localSlot = result.slot; - target.resourceLane = resourceLane; - target.schedule = request.commitSchedule; - target.deliverySchedule = request.schedule; - target.content = contentSerial; - target.pitch = pitch; - target.width = dstFormat.width; - target.height = dstFormat.height; - target.format = dstFormat.desc.Format; - target.frameType = dstFormat.format; - target.periodic = request.periodic; - target.active = true; - if (sink.outstanding.fetch_add( - 1, std::memory_order_acq_rel) == 0) - ResetEvent(sink.drained); + if (resourceLane == FRAME_SINK_BUFFERS) + for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) + if (!sink.lanes[lane].busy && !sink.lanes[lane].valid) + { + resourceLane = lane; + break; + } + if (resourceLane == FRAME_SINK_BUFFERS) + for (unsigned lane = 0; lane < FRAME_SINK_BUFFERS; ++lane) + if (!sink.lanes[lane].busy) + { + resourceLane = lane; + break; + } + if (resourceLane == FRAME_SINK_BUFFERS) + { + lock.Unlock(); + sink.target->AbortFrameBuffer(result.slot); + continue; + } - PreparedFrameBuffer& output = prepared.targets[index]; - output.token.sink = request.backend; - output.token.epoch = request.epoch; - output.token.slot = result.slot; - output.token.serial = batch->serial; - output.resourceSlot = - request.sink * FRAME_SINK_BUFFERS + resourceLane; - output.mem = result.mem; - output.heapOffset = result.heapOffset; - output.capacity = result.capacity; - output.direct = request.primary; - output.fullCopy = forceFull || result.fullCopy; + Sink::ResourceLane& lane = sink.lanes[resourceLane]; + lane.mem = result.mem; + lane.heapOffset = result.heapOffset; + lane.localSlot = result.slot; + lane.direct = request.primary; + lane.valid = true; + lane.busy = true; + lane.phase = Sink::ResourceLane::IDLE; + + unsigned index = batch->count; + if (sink.primary && index) + { + for (unsigned move = index; move != 0; --move) + { + batch->targets[move] = batch->targets[move - 1]; + prepared.targets[move] = prepared.targets[move - 1]; + } + index = 0; + } + ++batch->count; + BatchTarget& target = batch->targets[index]; + target.sink = &sink; + target.backend = request.backend; + target.epoch = request.epoch; + target.localSlot = result.slot; + target.resourceLane = resourceLane; + target.schedule = request.commitSchedule; + target.deliverySchedule = request.schedule; + target.content = contentSerial; + target.pitch = pitch; + target.width = dstFormat.width; + target.height = dstFormat.height; + target.format = dstFormat.desc.Format; + target.frameType = dstFormat.format; + target.periodic = request.periodic; + target.active = true; + if (sink.outstanding.fetch_add( + 1, std::memory_order_acq_rel) == 0) + ResetEvent(sink.drained); + + PreparedFrameBuffer& output = prepared.targets[index]; + output.token.backend = request.backend; + output.token.epoch = request.epoch; + output.token.slot = result.slot; + output.token.serial = batch->serial; + output.resourceSlot = + request.sink * FRAME_SINK_BUFFERS + resourceLane; + output.mem = result.mem; + output.heapOffset = result.heapOffset; + output.capacity = result.capacity; + output.direct = request.primary; + output.fullCopy = forceFull || result.fullCopy; + } } prepared.count = batch->count; if (!batch->count) @@ -636,7 +936,10 @@ uint32_t CFrameHub::PublishFrameBatch(const FrameBatchToken& token) { if (valid) target.sink->target->AbortFrameBuffer(target.localSlot); - target.sink->needsFullCopy = true; + { + CSRWExclusiveLock state(target.sink->laneLock); + target.sink->needsFullCopy = true; + } ReleaseTarget(batch, target); continue; } @@ -658,6 +961,14 @@ void CFrameHub::CommitFrameBatch(const FrameBatchToken& token) if (token.slot >= FRAME_BATCHES) return; Batch& batch = m_batches[token.slot]; + struct FillRequest + { + Sink * sink; + unsigned lane; + FrameToken token; + bool succeeded; + } fill[FRAME_MAX_SINKS] = {}; + unsigned fillCount = 0; CSRWExclusiveLock lock(batch.lock); if (!BatchValid(batch, token)) return; @@ -678,8 +989,22 @@ void CFrameHub::CommitFrameBatch(const FrameBatchToken& token) } target.committed = true; if (target.completionPending) - CompleteTarget(batch, target, target.completionSucceeded); + { + FrameToken frameToken; + const unsigned lane = target.resourceLane; + Sink * sink = target.sink; + if (DetachTarget(batch, target, frameToken)) + fill[fillCount++] = { + sink, lane, frameToken, target.completionSucceeded }; + } } + lock.Unlock(); + for (unsigned i = 0; i < fillCount; ++i) + if (fill[i].succeeded) + FillLane(*fill[i].sink, fill[i].lane, fill[i].token); + else + CompleteLane(*fill[i].sink, fill[i].lane, fill[i].token, + FrameDone::FAILED, CFrameScheduler::Nanotime()); } void CFrameHub::AbortFrameBatch(const FrameBatchToken& token) @@ -701,7 +1026,10 @@ void CFrameHub::AbortFrameBatch(const FrameBatchToken& token) target.backend && target.sink->epoch.load(std::memory_order_acquire) == target.epoch) target.sink->target->AbortFrameBuffer(target.localSlot); - target.sink->needsFullCopy = true; + { + CSRWExclusiveLock state(target.sink->laneLock); + target.sink->needsFullCopy = true; + } ReleaseTarget(batch, target); } } @@ -730,7 +1058,10 @@ void CFrameHub::FailFrameBatch(const FrameBatchToken& token) else target.sink->target->AbortFrameBuffer(target.localSlot); } - target.sink->needsFullCopy = true; + { + CSRWExclusiveLock state(target.sink->laneLock); + target.sink->needsFullCopy = true; + } ReleaseTarget(batch, target); } } @@ -804,19 +1135,18 @@ void CFrameHub::SetFrameTargetTiming(const FrameBatchToken& token, if (token.slot >= FRAME_BATCHES) return; Batch& batch = m_batches[token.slot]; - CSRWSharedLock lock(batch.lock); + CSRWExclusiveLock lock(batch.lock); if (!BatchValid(batch, token) || index >= batch.count || !batch.targets[index].active) return; BatchTarget& target = batch.targets[index]; - CSRWSharedLock call(target.sink->callLock); - if (target.sink->target && - target.sink->backend.load(std::memory_order_acquire) == - target.backend && - target.sink->epoch.load(std::memory_order_acquire) == target.epoch) - target.sink->target->SetFrameTiming(target.localSlot, captureTime, - postProcessTime, copyTime, readyTime, holdTime, - target.schedule, completedAt); + target.captureTime = captureTime; + target.postProcessTime = postProcessTime; + target.copyTime = copyTime; + target.readyTime = readyTime; + target.holdTime = holdTime; + target.filledAt = completedAt; + target.timingValid = true; } void CFrameHub::TryRecordFrameTiming(const FrameBatchToken& token, @@ -825,17 +1155,13 @@ void CFrameHub::TryRecordFrameTiming(const FrameBatchToken& token, if (token.slot >= FRAME_BATCHES) return; Batch& batch = m_batches[token.slot]; - CSRWSharedLock lock(batch.lock); + CSRWExclusiveLock lock(batch.lock); if (!BatchValid(batch, token) || index >= batch.count || !batch.targets[index].active) return; BatchTarget& target = batch.targets[index]; - CSRWSharedLock call(target.sink->callLock); - if (target.sink->target && - target.sink->backend.load(std::memory_order_acquire) == - target.backend && - target.sink->epoch.load(std::memory_order_acquire) == target.epoch) - target.sink->target->TryRecordFrameTiming(duration); + if (target.filledAt >= duration) + target.workStart = target.filledAt - duration; } void CFrameHub::CompleteFrameTarget(const FrameBatchToken& token, @@ -855,7 +1181,17 @@ void CFrameHub::CompleteFrameTarget(const FrameBatchToken& token, target.completionSucceeded = succeeded; return; } - CompleteTarget(batch, target, succeeded); + FrameToken frameToken; + const unsigned lane = target.resourceLane; + Sink * sink = target.sink; + if (!DetachTarget(batch, target, frameToken)) + return; + lock.Unlock(); + if (succeeded) + FillLane(*sink, lane, frameToken); + else + CompleteLane(*sink, lane, frameToken, FrameDone::FAILED, + CFrameScheduler::Nanotime()); } void CFrameHub::ObserveFrame(uint64_t now) diff --git a/idd/LGIdd/transport/CFrameHub.h b/idd/LGIdd/transport/CFrameHub.h index 08e686c7..33209dbe 100644 --- a/idd/LGIdd/transport/CFrameHub.h +++ b/idd/LGIdd/transport/CFrameHub.h @@ -27,22 +27,55 @@ #include -class CFrameHub final : public IFrameTransport +class CFrameHub final : public IFrameTransport, public IFrameEvents { private: struct Sink { struct ResourceLane { + enum Phase + { + IDLE, + FILLING, + PENDING, + CANCELING, + COMPLETING, + }; + uint8_t * mem = nullptr; uint64_t heapOffset = 0; unsigned localSlot = 0; bool direct = false; bool valid = false; bool busy = false; + + FrameToken token = {}; + CFrameScheduler::Schedule schedule = {}; + uint64_t content = 0; + uint64_t captureTime = 0; + uint64_t postProcessTime = 0; + uint64_t copyTime = 0; + uint64_t readyTime = 0; + uint64_t holdTime = 0; + uint64_t filledAt = 0; + uint64_t workStart = 0; + unsigned pitch = 0; + unsigned width = 0; + unsigned height = 0; + DXGI_FORMAT format = DXGI_FORMAT_UNKNOWN; + FrameType frameType = FRAME_TYPE_INVALID; + Phase phase = IDLE; + FrameDone pendingResult = FrameDone::FAILED; + uint64_t pendingReadyAt = 0; + bool timingValid = false; + bool resultPending = false; + bool callActive = false; + bool cancelRequested = false; }; CSRWLock callLock; + CSRWLock laneLock; IFrameSink * target = nullptr; HANDLE drained = nullptr; std::atomic outstanding = 0; @@ -53,6 +86,7 @@ private: bool primary = false; bool needsFullCopy = true; uint64_t lastContent = 0; + uint64_t lastTerminalContent = 0; uint64_t blockedContent = 0; size_t blockedFrameSize = 0; bool blockedAllocation = false; @@ -87,6 +121,14 @@ private: bool completionPending = false; bool completionSucceeded = false; bool releasePending = false; + uint64_t captureTime = 0; + uint64_t postProcessTime = 0; + uint64_t copyTime = 0; + uint64_t readyTime = 0; + uint64_t holdTime = 0; + uint64_t filledAt = 0; + uint64_t workStart = 0; + bool timingValid = false; bool active = false; }; @@ -117,7 +159,14 @@ private: unsigned Snapshot(SinkRef refs[FRAME_MAX_SINKS]) const; bool BatchValid(const Batch& batch, const FrameBatchToken& token) const; void ReleaseTarget(Batch& batch, BatchTarget& target); - void CompleteTarget(Batch& batch, BatchTarget& target, bool succeeded); + bool DetachTarget(Batch& batch, BatchTarget& target, + FrameToken& token); + void FillLane(Sink& sink, unsigned laneIndex, + const FrameToken& token); + void CancelLane(Sink& sink, unsigned laneIndex, + const FrameToken& token); + void CompleteLane(Sink& sink, unsigned laneIndex, + const FrameToken& token, FrameDone result, uint64_t readyAt); public: CFrameHub(); @@ -130,6 +179,9 @@ public: IFrameSink& sink); void Unbind(BackendId backend, uint32_t epoch); + void OnFrameDone(const FrameToken& token, FrameDone result, + uint64_t readyAt) override; + size_t GetMaxFrameSize() const override; uint64_t NextContentSerial() override; void FrameProductReady(uint64_t contentSerial) override; diff --git a/idd/LGIdd/transport/IFrameSink.h b/idd/LGIdd/transport/IFrameSink.h index d57a7c8d..1712f9b7 100644 --- a/idd/LGIdd/transport/IFrameSink.h +++ b/idd/LGIdd/transport/IFrameSink.h @@ -28,11 +28,22 @@ #include #include +class IFrameEvents +{ +public: + virtual ~IFrameEvents() = default; + + virtual void OnFrameDone(const FrameToken& token, FrameDone result, + uint64_t readyAt) = 0; +}; + class IFrameSink { public: virtual ~IFrameSink() = default; + // Clearing events waits for callbacks through the previous pointer to end. + virtual void SetFrameEvents(IFrameEvents * events) = 0; virtual size_t GetMaxFrameSize() const = 0; virtual bool FrameBufferAvailable( @@ -60,7 +71,12 @@ public: bool deliveredToOwner) = 0; virtual void AbortFrameBuffer(unsigned slot) = 0; virtual void FailFrameBuffer(unsigned slot) = 0; - virtual void CompleteFrameBuffer(unsigned slot, bool succeeded) = 0; + // A terminal callback raised before this returns takes precedence over the + // return value. + virtual FrameFill FrameFilled(const FrameToken& token) = 0; + // On return, the sink no longer accesses the frame identified by token. + virtual void CancelFrame(const FrameToken& token) = 0; + virtual void CompleteFrameBuffer(unsigned slot, FrameDone result) = 0; virtual void SetFrameTiming(unsigned slot, uint64_t captureTime, uint64_t postProcessTime, uint64_t copyTime, uint64_t readyTime, uint64_t holdTime, const CFrameScheduler::Schedule& schedule, diff --git a/idd/LGIdd/transport/PreparedFrameBuffer.h b/idd/LGIdd/transport/PreparedFrameBuffer.h index 29705fe1..8b03681e 100644 --- a/idd/LGIdd/transport/PreparedFrameBuffer.h +++ b/idd/LGIdd/transport/PreparedFrameBuffer.h @@ -33,12 +33,26 @@ enum : unsigned struct FrameToken { - uint32_t sink = 0; + uint32_t backend = 0; uint32_t epoch = 0; uint32_t slot = 0; uint64_t serial = 0; }; +enum class FrameFill +{ + READY, + PENDING, + REJECTED, +}; + +enum class FrameDone +{ + READY, + FAILED, + SUPERSEDED, +}; + struct PreparedFrameBuffer { FrameToken token; diff --git a/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.cpp b/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.cpp index a81d148d..8ca83568 100644 --- a/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.cpp +++ b/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.cpp @@ -202,6 +202,7 @@ bool CLGMPFrameTransport::Setup(size_t alignSize) m_readyFrameIndex.store(-1, std::memory_order_release); m_deferredOwnerFrameIndex = -1; m_framePublishSequence = 0; + m_frameReadySequence = 0; memset(m_frameLastPublishSequence, 0, sizeof(m_frameLastPublishSequence)); for (FrameDelivery& delivery : m_frameDelivery) @@ -222,6 +223,7 @@ void CLGMPFrameTransport::DeInit() m_readyFrameIndex.store(-1, std::memory_order_release); m_deferredOwnerFrameIndex = -1; m_framePublishSequence = 0; + m_frameReadySequence = 0; memset(m_frameLastPublishSequence, 0, sizeof(m_frameLastPublishSequence)); memset(m_frameCompleted, 0, sizeof(m_frameCompleted)); @@ -845,6 +847,8 @@ SinkTarget CLGMPFrameTransport::PrepareFrameBuffer( flags |= FRAME_FLAG_TRUNCATED; fi->formatVer = m_formatVer; + if (!m_frameSerial) + ++m_frameSerial; fi->frameSerial = m_frameSerial++; fi->screenWidth = srcFormat.width; fi->screenHeight = srcFormat.height; @@ -1008,6 +1012,8 @@ bool CLGMPFrameTransport::PublishFrameBuffer(unsigned frameIndex, if (published) { m_frameLastPublishSequence[frameIndex] = ++m_framePublishSequence; + m_frameSchedule[frameIndex] = schedule; + m_frameDelivered[frameIndex] = deliveredToOwner; m_deferredOwnerFrameIndex = schedule.clientID && !deliveredToOwner ? static_cast(frameIndex) : -1; m_submittedFrameIndex.store( @@ -1122,15 +1128,14 @@ void CLGMPFrameTransport::CommitFrameBuffer(unsigned frameIndex, const CFrameScheduler::Schedule& schedule, bool periodic, bool deliveredToOwner) { + UNREFERENCED_PARAMETER(deliveredToOwner); + if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN) return; const uint64_t now = CFrameScheduler::Nanotime(); - if (deliveredToOwner) - m_frameScheduler.FramePublished( - schedule, m_frame[frameIndex]->frameSerial, now, periodic); - else - m_frameScheduler.FrameRetained(schedule, now, periodic); + m_frameScheduler.FrameCommitted( + schedule, m_frame[frameIndex]->frameSerial, now, periodic); } bool CLGMPFrameTransport::TryFrameSubmitted(unsigned frameIndex, @@ -1200,16 +1205,21 @@ void CLGMPFrameTransport::FailFrameBuffer(unsigned frameIndex) InterlockedExchange( (volatile LONG *)&m_frame[frameIndex]->timingValid, 0); FinalizeFrameBuffer(frameIndex); - CompleteFrameBuffer(frameIndex, false); + AbortFrameBuffer(frameIndex); } void CLGMPFrameTransport::CompleteFrameBuffer( - unsigned frameIndex, bool succeeded) + unsigned frameIndex, FrameDone result) { if (frameIndex >= LGMP_Q_FRAME_BUFFER_LEN) return; + const bool succeeded = result == FrameDone::READY; CSRWExclusiveLock lock(m_framePublishLock); + const CFrameScheduler::Schedule schedule = m_frameSchedule[frameIndex]; + const uint32_t frameSerial = m_frame[frameIndex]->frameSerial; + const bool delivered = m_frameDelivered[frameIndex]; + const uint64_t sequence = m_frameLastPublishSequence[frameIndex]; m_frameCompleted[frameIndex] = succeeded; if (!succeeded && m_deferredOwnerFrameIndex == static_cast(frameIndex)) @@ -1218,7 +1228,6 @@ void CLGMPFrameTransport::CompleteFrameBuffer( { // Completion callbacks may run out of order. Never replace a newer ready // frame with an older submission. - const uint64_t sequence = m_frameLastPublishSequence[frameIndex]; const LONG readyFrameIndex = m_readyFrameIndex.load(std::memory_order_acquire); if (sequence && @@ -1228,6 +1237,22 @@ void CLGMPFrameTransport::CompleteFrameBuffer( static_cast(frameIndex), std::memory_order_release); } m_frameInFlight[frameIndex].store(false, std::memory_order_release); + const bool newerThanReady = + sequence && sequence > m_frameReadySequence; + if (result == FrameDone::READY && newerThanReady) + m_frameReadySequence = sequence; + + if (result == FrameDone::READY && newerThanReady) + { + if (delivered) + m_frameScheduler.FramePublished(schedule, frameSerial); + else + m_frameScheduler.FrameRetained(schedule); + } + else if (result == FrameDone::FAILED && newerThanReady) + m_frameScheduler.FrameFailed(); + else if (result == FrameDone::SUPERSEDED && newerThanReady) + m_frameScheduler.FrameSuperseded(); } void CLGMPFrameTransport::SetFrameTiming(unsigned frameIndex, diff --git a/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.h b/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.h index acbad46e..c7cfa732 100644 --- a/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.h +++ b/idd/LGIdd/transport/lgmp/CLGMPFrameTransport.h @@ -101,7 +101,11 @@ private: bool m_frameCompleted[LGMP_Q_FRAME_BUFFER_LEN] = {}; CSRWLock m_framePublishLock; uint64_t m_framePublishSequence = 0; + uint64_t m_frameReadySequence = 0; uint64_t m_frameLastPublishSequence[LGMP_Q_FRAME_BUFFER_LEN] = {}; + CFrameScheduler::Schedule + m_frameSchedule[LGMP_Q_FRAME_BUFFER_LEN] = {}; + bool m_frameDelivered[LGMP_Q_FRAME_BUFFER_LEN] = {}; FrameDelivery m_frameDelivery[LGMP_Q_FRAME_BUFFER_LEN] = {}; OwnerDelivery m_ownerDelivery[LGMP_Q_FRAME_LEN] = {}; @@ -160,6 +164,7 @@ public: CLGMPFrameTransport(const CLGMPFrameTransport&) = delete; CLGMPFrameTransport& operator=(const CLGMPFrameTransport&) = delete; + void SetFrameEvents(IFrameEvents *) override {} size_t GetMaxFrameSize() const override { return m_maxFrameSize; } bool FrameBufferAvailable(const CFrameScheduler::Schedule& schedule, @@ -189,8 +194,12 @@ public: bool deliveredToOwner) override; void AbortFrameBuffer(unsigned frameIndex) override; void FailFrameBuffer(unsigned frameIndex) override; - void CompleteFrameBuffer( - unsigned frameIndex, bool succeeded) override; + FrameFill FrameFilled(const FrameToken&) override + { + return FrameFill::READY; + } + void CancelFrame(const FrameToken&) override {} + void CompleteFrameBuffer(unsigned frameIndex, FrameDone result) override; void SetFrameTiming(unsigned frameIndex, uint64_t captureTime, uint64_t postProcessTime, uint64_t copyTime, uint64_t readyTime, uint64_t holdTime, const CFrameScheduler::Schedule& schedule,