[idd] transport: isolate interactive services

This commit is contained in:
Geoffrey McRae
2026-08-12 21:49:04 +10:00
parent 51f432db00
commit c92acb0a34
21 changed files with 2026 additions and 401 deletions

View File

@@ -27,6 +27,7 @@
#include <utility>
static const uint64_t RETRY_DELAY_MS = 500;
static const uint64_t SERVICE_RETRY_DELAY_MS = 250;
class CSourceEvents final : public ITransportEvents
{
@@ -353,7 +354,14 @@ bool CTransportManager::AddServices(Entry& entry)
uint32_t epoch = 0;
bool primary = false;
bool controlAdded = false;
bool controlFailed = false;
bool controlAbsent = false;
bool inputAdded = false;
bool inputFailed = false;
bool inputAbsent = false;
bool frameAdded = false;
bool frameAbsent = false;
uint64_t retryAt = 0;
{
CSRWSharedLock entryLock(entry.lock);
transport = entry.transport;
@@ -361,38 +369,149 @@ bool CTransportManager::AddServices(Entry& entry)
epoch = entry.epoch;
primary = entry.primary;
controlAdded = entry.controlAdded;
controlFailed = entry.controlFailed;
controlAbsent = entry.controlAbsent;
inputAdded = entry.inputAdded;
inputFailed = entry.inputFailed;
inputAbsent = entry.inputAbsent;
frameAdded = entry.frameAdded;
frameAbsent = entry.frameAbsent;
retryAt = entry.serviceRetryAt;
}
if (!transport)
return false;
return !primary;
if (!controlAdded)
const uint64_t now = GetTickCount64();
const bool attach = now >= retryAt;
bool controlRetry = false;
bool frameRetry = false;
bool inputRetry = false;
if (attach && !controlAdded && !controlFailed && !controlAbsent)
{
if (!m_control.Add(id, epoch, transport->Control()))
return false;
controlAdded = true;
CSRWExclusiveLock entryLock(entry.lock);
entry.controlAdded = true;
IControlSink * control = transport->Control();
if (control && m_control.Add(id, epoch, *control))
{
controlAdded = true;
CSRWExclusiveLock entryLock(entry.lock);
entry.controlAdded = true;
}
else if (control)
controlRetry = true;
else
{
controlAbsent = true;
CSRWExclusiveLock entryLock(entry.lock);
entry.controlAbsent = true;
}
}
if (!frameAdded &&
m_frames.Bind(id, epoch, primary, transport->FrameSink()))
if (attach && !frameAdded && !frameAbsent)
{
CSRWExclusiveLock entryLock(entry.lock);
entry.frameAdded = true;
return true;
IFrameSink * frame = transport->FrameSink();
if (frame && m_frames.Bind(id, epoch, primary, *frame))
{
CSRWExclusiveLock entryLock(entry.lock);
entry.frameAdded = true;
frameAdded = true;
}
else if (frame || primary)
frameRetry = true;
else
{
frameAbsent = true;
CSRWExclusiveLock entryLock(entry.lock);
entry.frameAbsent = true;
}
}
if (frameAdded)
return true;
if (attach && !inputAdded && !inputFailed && !inputAbsent)
{
IInputSource * input = transport->Input();
if (input && m_input.Bind(id, epoch, *input))
{
CSRWExclusiveLock entryLock(entry.lock);
entry.inputAdded = true;
inputAdded = true;
}
else if (input)
inputRetry = true;
else
{
inputAbsent = true;
CSRWExclusiveLock entryLock(entry.lock);
entry.inputAbsent = true;
}
}
m_control.Remove(id, epoch);
if (attach && (controlRetry || inputRetry || frameRetry))
{
CSRWExclusiveLock entryLock(entry.lock);
entry.controlAdded = false;
entry.serviceRetryAt = now + SERVICE_RETRY_DELAY_MS;
}
return !primary || frameAdded;
}
void CTransportManager::HandleServiceFailures()
{
ControlToken token;
while (m_control.TakeFailure(token))
{
Entry * entries[FRAME_MAX_SINKS] = {};
const unsigned count = Entries(entries);
for (unsigned i = 0; i < count; ++i)
{
Entry& entry = *entries[i];
bool restart = false;
{
CSRWExclusiveLock entryLock(entry.lock);
if (entry.id != token.backend || entry.epoch != token.epoch ||
!entry.controlAdded)
continue;
entry.controlAdded = false;
entry.controlFailed = true;
restart = !entry.exposed;
}
m_control.Remove(token.backend, token.epoch);
if (restart)
{
RemoveServices(entry);
ScheduleRetry(entry);
}
break;
}
}
SourceKey source;
while (m_input.TakeFailure(source))
{
Entry * entries[FRAME_MAX_SINKS] = {};
const unsigned count = Entries(entries);
for (unsigned i = 0; i < count; ++i)
{
Entry& entry = *entries[i];
bool restart = false;
{
CSRWExclusiveLock entryLock(entry.lock);
if (entry.id != source.backend || entry.epoch != source.epoch ||
!entry.inputAdded)
continue;
entry.inputAdded = false;
entry.inputFailed = true;
restart = !entry.exposed;
}
m_input.Unbind(source.backend, source.epoch);
if (restart)
{
RemoveServices(entry);
ScheduleRetry(entry);
}
break;
}
}
return false;
}
bool CTransportManager::SetupEntry(Entry& entry, size_t alignment)
@@ -400,9 +519,7 @@ bool CTransportManager::SetupEntry(Entry& entry, size_t alignment)
std::shared_ptr<ITransport> transport;
{
CSRWSharedLock entryLock(entry.lock);
if (entry.state == State::READY)
return true;
if (entry.state != State::INITIALIZED)
if (entry.state != State::INITIALIZED && entry.state != State::READY)
return false;
transport = entry.transport;
}
@@ -449,20 +566,25 @@ void CTransportManager::RemoveServices(Entry& entry)
uint32_t epoch = 0;
bool frameAdded = false;
bool controlAdded = false;
bool inputAdded = false;
{
CSRWExclusiveLock entryLock(entry.lock);
id = entry.id;
epoch = entry.epoch;
frameAdded = entry.frameAdded;
controlAdded = entry.controlAdded;
inputAdded = entry.inputAdded;
entry.frameAdded = false;
entry.controlAdded = false;
entry.inputAdded = false;
}
if (frameAdded)
m_frames.Unbind(id, epoch);
if (inputAdded)
m_input.Unbind(id, epoch);
if (controlAdded)
m_control.Remove(id, epoch);
if (frameAdded)
m_frames.Unbind(id, epoch);
}
void CTransportManager::RetryEntry(Entry& entry, uint64_t now,
@@ -489,6 +611,14 @@ void CTransportManager::RetryEntry(Entry& entry, uint64_t now,
entry.directMemory = DirectFrameBufferMemory {};
entry.directMemoryValid = false;
entry.setupDone = false;
entry.controlFailed = false;
entry.controlAbsent = false;
entry.inputFailed = false;
entry.inputAbsent = false;
entry.frameAbsent = false;
entry.serviceRetryAt = 0;
entry.recoveryPending = false;
entry.recovery = RecoveryUpdate {};
++entry.epoch;
if (!entry.epoch)
++entry.epoch;
@@ -503,8 +633,12 @@ void CTransportManager::RetryEntry(Entry& entry, uint64_t now,
}
if (setup && !SetupEntry(entry, alignment))
{
RemoveServices(entry);
ScheduleRetry(entry);
CSRWSharedLock entryLock(entry.lock);
if (entry.state == State::FAILED)
{
entryLock.Unlock();
ScheduleRetry(entry);
}
}
}
@@ -566,7 +700,7 @@ ITransport::OpenResult CTransportManager::Open()
const unsigned count = Entries(entries);
OpenResult aggregate = Primary() ?
OpenResult::SUCCESS : OpenResult::FAILURE;
for (unsigned i = 0; i < count && aggregate != OpenResult::FAILURE; ++i)
for (unsigned i = 0; i < count; ++i)
{
Entry& entry = *entries[i];
if (!BeginCall(entry, Call::LIFECYCLE, true))
@@ -593,8 +727,11 @@ ITransport::OpenResult CTransportManager::Open()
}
EndCall(entry, transport);
if (entry.required && result != OpenResult::SUCCESS)
aggregate = result;
if (entry.required && result == OpenResult::FAILURE)
aggregate = OpenResult::FAILURE;
else if (entry.required && result == OpenResult::RETRY &&
aggregate == OpenResult::SUCCESS)
aggregate = OpenResult::RETRY;
}
EndPhase();
@@ -609,7 +746,7 @@ bool CTransportManager::Initialize()
Entry * entries[FRAME_MAX_SINKS] = {};
const unsigned count = Entries(entries);
bool success = true;
for (unsigned i = 0; i < count && success; ++i)
for (unsigned i = 0; i < count; ++i)
{
Entry& entry = *entries[i];
if (!BeginCall(entry, Call::LIFECYCLE, true))
@@ -676,12 +813,12 @@ bool CTransportManager::Setup(size_t alignment)
Entry * entries[FRAME_MAX_SINKS] = {};
const unsigned count = Entries(entries);
bool success = initialized;
for (unsigned i = 0; i < count && success; ++i)
for (unsigned i = 0; i < count; ++i)
{
Entry& entry = *entries[i];
if (!BeginCall(entry, Call::LIFECYCLE, true))
{
if (entry.required)
if (entry.primary)
success = false;
continue;
}
@@ -694,12 +831,16 @@ bool CTransportManager::Setup(size_t alignment)
transport = entry.transport;
}
if (state == State::INITIALIZED && !SetupEntry(entry, alignment))
if ((state == State::INITIALIZED || state == State::READY) &&
!SetupEntry(entry, alignment))
{
RemoveServices(entry);
if (entry.required)
{
CSRWSharedLock entryLock(entry.lock);
state = entry.state;
}
if (entry.primary)
success = false;
else
else if (state == State::FAILED)
ScheduleRetry(entry);
}
{
@@ -709,14 +850,17 @@ bool CTransportManager::Setup(size_t alignment)
EndCall(entry, transport);
}
for (unsigned i = 0; i < count && success; ++i)
Entry * primary = Primary();
if (primary)
{
CSRWSharedLock entryLock(entries[i]->lock);
if (entries[i]->required && entries[i]->state != State::READY)
CSRWSharedLock entryLock(primary->lock);
if (primary->state != State::READY || !primary->frameAdded)
success = false;
}
else
success = false;
if (success)
if (initialized)
{
CSRWExclusiveLock managerLock(m_lock);
m_setup = true;
@@ -746,6 +890,7 @@ ITransport::ProcessResult CTransportManager::Process(
}
const uint64_t now = GetTickCount64();
HandleServiceFailures();
Entry * entries[FRAME_MAX_SINKS] = {};
const unsigned count = Entries(entries);
for (unsigned i = 0; i < count; ++i)
@@ -793,6 +938,9 @@ ITransport::ProcessResult CTransportManager::Process(
entry.setupDone = true;
}
if (setup)
SetupEntry(entry, alignment);
DrainRecovery(entry, transport);
CSourceEvents sourceEvents(id, epoch, events);
const ProcessResult result = transport->Process(sourceEvents);
@@ -880,6 +1028,11 @@ void CTransportManager::Stop()
return;
}
for (unsigned i = count; i > 0; --i)
RemoveServices(*entries[i - 1]);
m_input.Stop();
for (unsigned i = count; i > 0; --i)
{
Entry& entry = *entries[i - 1];
@@ -891,7 +1044,6 @@ void CTransportManager::Stop()
state = entry.state;
}
RemoveServices(entry);
if (transport && state != State::STOPPED)
transport->Stop();
{
@@ -952,8 +1104,9 @@ void CTransportManager::SyncRecovery()
}
}
void CTransportManager::RecoveryStatus(uint64_t session, uint32_t serial,
bool active, Recovery state, uint32_t error)
void CTransportManager::RecoveryStatus(const SourceKey& source,
uint64_t session, uint32_t serial, bool active,
Recovery state, uint32_t error)
{
{
CSRWSharedLock managerLock(m_lock);
@@ -970,7 +1123,8 @@ void CTransportManager::RecoveryStatus(uint64_t session, uint32_t serial,
bool call = false;
{
CSRWExclusiveLock entryLock(entry.lock);
if (entry.stopRequested || !entry.transport ||
if (entry.id != source.backend || entry.epoch != source.epoch ||
entry.stopRequested || !entry.transport ||
(entry.state != State::INITIALIZED && entry.state != State::READY))
continue;
@@ -995,6 +1149,7 @@ void CTransportManager::RecoveryStatus(uint64_t session, uint32_t serial,
if (call)
transport->RecoveryStatus(session, serial, active, state, error);
EndCall(entry, transport);
return;
}
}
@@ -1095,25 +1250,7 @@ IControlTransport& CTransportManager::Control()
return m_control;
}
IInputTransport * CTransportManager::Input()
IInputTransport& CTransportManager::Input()
{
Entry * primary = Primary();
if (!primary || !BeginPhase(Phase::ACCESS, true))
return nullptr;
if (!BeginCall(*primary, Call::ACCESS, true))
{
EndPhase();
return nullptr;
}
std::shared_ptr<ITransport> transport;
{
CSRWSharedLock entryLock(primary->lock);
transport = primary->transport;
}
Expose(*primary);
IInputTransport * input = transport ? transport->Input() : nullptr;
EndCall(*primary, transport);
EndPhase();
return input;
return m_input;
}