Fix render worker reuse and stabilize out-of-process video rendering
This commit is contained in:
+454
-154
@@ -21,6 +21,7 @@
|
||||
#include "renderworkerpool.h"
|
||||
|
||||
#include <QCoreApplication>
|
||||
#include <QDateTime>
|
||||
#include <QDir>
|
||||
#include <QFile>
|
||||
#include <QFileInfo>
|
||||
@@ -30,8 +31,10 @@
|
||||
#include <QTemporaryFile>
|
||||
#include <QXmlStreamWriter>
|
||||
#include <algorithm>
|
||||
#include <memory>
|
||||
#include <optional>
|
||||
#include <thread>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
#if defined(Q_OS_WIN)
|
||||
#include <windows.h>
|
||||
@@ -183,6 +186,17 @@ FramePtr DecodeInputFrame(DecoderCache *decoder_cache,
|
||||
FramePtr frame = decoder->RetrieveVideoFrame(retrieve);
|
||||
if (frame) {
|
||||
frame->set_timestamp(input.time);
|
||||
|
||||
// Ensure the frame carries the colorspace the color manager expects.
|
||||
// Decoders do not always set this on the returned frame, but the worker
|
||||
// needs it to build the correct OCIO transform.
|
||||
VideoParams frame_params = frame->video_params();
|
||||
if (frame_params.colorspace().isEmpty() &&
|
||||
!stream_data.colorspace().isEmpty()) {
|
||||
frame_params.set_colorspace(stream_data.colorspace());
|
||||
frame->set_video_params(frame_params);
|
||||
}
|
||||
|
||||
}
|
||||
return frame;
|
||||
}
|
||||
@@ -254,6 +268,26 @@ bool KillProcessById(qint64 process_id)
|
||||
#endif
|
||||
}
|
||||
|
||||
bool IsProcessAlive(qint64 process_id)
|
||||
{
|
||||
if (process_id <= 0) {
|
||||
return false;
|
||||
}
|
||||
|
||||
#if defined(Q_OS_WIN)
|
||||
HANDLE handle = OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, FALSE, DWORD(process_id));
|
||||
if (!handle) {
|
||||
return false;
|
||||
}
|
||||
DWORD exit_code = 0;
|
||||
const bool alive = GetExitCodeProcess(handle, &exit_code) && exit_code == STILL_ACTIVE;
|
||||
CloseHandle(handle);
|
||||
return alive;
|
||||
#else
|
||||
return ::kill(pid_t(process_id), 0) == 0;
|
||||
#endif
|
||||
}
|
||||
|
||||
QString WorkerProcessDetails(const QProcess *process)
|
||||
{
|
||||
if (!process) {
|
||||
@@ -323,9 +357,11 @@ bool ReadControlMessage(QProcess *process, QJsonObject *out, QString *error,
|
||||
} // namespace
|
||||
|
||||
RenderWorkerPool::RenderWorkerPool(DecoderCache *decoder_cache,
|
||||
const QString &gpu_backend,
|
||||
QObject *parent)
|
||||
: QThread(parent)
|
||||
, decoder_cache_(decoder_cache)
|
||||
, gpu_backend_(gpu_backend)
|
||||
{
|
||||
}
|
||||
|
||||
@@ -393,7 +429,7 @@ bool RenderWorkerPool::RemoveTicket(RenderTicketPtr ticket)
|
||||
|
||||
void RenderWorkerPool::Shutdown()
|
||||
{
|
||||
QVector<QString> queued_graph_paths;
|
||||
QVector<QString> graph_paths_to_clean;
|
||||
|
||||
{
|
||||
QMutexLocker locker(&mutex_);
|
||||
@@ -402,7 +438,6 @@ void RenderWorkerPool::Shutdown()
|
||||
if (job.ticket) {
|
||||
job.ticket->Cancel();
|
||||
}
|
||||
queued_graph_paths.append(job.graph_path);
|
||||
}
|
||||
queue_.clear();
|
||||
for (ActiveJob &active : active_jobs_) {
|
||||
@@ -411,10 +446,14 @@ void RenderWorkerPool::Shutdown()
|
||||
CancelActiveProcess(active.process_id);
|
||||
}
|
||||
}
|
||||
for (auto it = graph_cache_.begin(); it != graph_cache_.end(); ++it) {
|
||||
graph_paths_to_clean.append(it->path);
|
||||
}
|
||||
graph_cache_.clear();
|
||||
wait_.wakeAll();
|
||||
}
|
||||
|
||||
for (const QString &path : queued_graph_paths) {
|
||||
for (const QString &path : graph_paths_to_clean) {
|
||||
CleanupGraphFile(path);
|
||||
}
|
||||
|
||||
@@ -431,11 +470,12 @@ void RenderWorkerPool::run()
|
||||
active_jobs_.resize(worker_count);
|
||||
}
|
||||
|
||||
std::vector<std::vector<std::unique_ptr<PooledWorker>>> local_pools(worker_count);
|
||||
std::vector<std::thread> workers;
|
||||
workers.reserve(size_t(worker_count));
|
||||
for (int i = 0; i < worker_count; i++) {
|
||||
workers.emplace_back([this, i]() {
|
||||
WorkerLoop(i);
|
||||
workers.emplace_back([this, i, &local_pools]() {
|
||||
WorkerLoop(i, &local_pools[i]);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -443,11 +483,19 @@ void RenderWorkerPool::run()
|
||||
worker.join();
|
||||
}
|
||||
|
||||
for (auto &local_pool : local_pools) {
|
||||
ShutdownLocalPool(&local_pool);
|
||||
}
|
||||
|
||||
ClearGraphCache();
|
||||
|
||||
QMutexLocker locker(&mutex_);
|
||||
active_jobs_.clear();
|
||||
}
|
||||
|
||||
void RenderWorkerPool::WorkerLoop(int worker_index)
|
||||
void RenderWorkerPool::WorkerLoop(
|
||||
int worker_index,
|
||||
std::vector<std::unique_ptr<PooledWorker>> *local_pool)
|
||||
{
|
||||
while (true) {
|
||||
mutex_.lock();
|
||||
@@ -463,8 +511,7 @@ void RenderWorkerPool::WorkerLoop(int worker_index)
|
||||
queue_.pop_front();
|
||||
mutex_.unlock();
|
||||
|
||||
ProcessJob(job, worker_index);
|
||||
CleanupGraphFile(job.graph_path);
|
||||
ProcessJob(job, worker_index, local_pool);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -491,8 +538,26 @@ bool RenderWorkerPool::PrepareJob(RenderTicketPtr ticket,
|
||||
}
|
||||
|
||||
QString graph_path;
|
||||
if (!WriteGraphSnapshot(project, &graph_path)) {
|
||||
return false;
|
||||
bool wrote_new_snapshot = false;
|
||||
{
|
||||
const QUuid project_uuid = project->GetUuid();
|
||||
QMutexLocker locker(&mutex_);
|
||||
auto it = graph_cache_.find(project_uuid);
|
||||
if (it != graph_cache_.end() && !project->is_modified()) {
|
||||
graph_path = it->path;
|
||||
} else {
|
||||
if (it != graph_cache_.end()) {
|
||||
CleanupGraphFile(it->path);
|
||||
graph_cache_.erase(it);
|
||||
}
|
||||
locker.unlock();
|
||||
if (!WriteGraphSnapshot(project, &graph_path)) {
|
||||
return false;
|
||||
}
|
||||
wrote_new_snapshot = true;
|
||||
locker.relock();
|
||||
graph_cache_.insert(project_uuid, {graph_path});
|
||||
}
|
||||
}
|
||||
|
||||
job->ticket = ticket;
|
||||
@@ -500,6 +565,7 @@ bool RenderWorkerPool::PrepareJob(RenderTicketPtr ticket,
|
||||
job->graph_path = graph_path;
|
||||
job->node_token = QString::number(reinterpret_cast<quintptr>(params.node));
|
||||
job->input_frames = input_frames;
|
||||
Q_UNUSED(wrote_new_snapshot)
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -535,7 +601,9 @@ bool RenderWorkerPool::IsSupported(const RenderManager::RenderVideoParams ¶m
|
||||
params.video_params.is_valid();
|
||||
}
|
||||
|
||||
void RenderWorkerPool::ProcessJob(const Job &job, int worker_index)
|
||||
void RenderWorkerPool::ProcessJob(
|
||||
const Job &job, int worker_index,
|
||||
std::vector<std::unique_ptr<PooledWorker>> *local_pool)
|
||||
{
|
||||
const qint64 ticket_id = qint64(reinterpret_cast<quintptr>(job.ticket.get()));
|
||||
SetActiveWorker(worker_index, job.ticket, nullptr, ticket_id);
|
||||
@@ -547,8 +615,39 @@ void RenderWorkerPool::ProcessJob(const Job &job, int worker_index)
|
||||
return;
|
||||
}
|
||||
|
||||
std::unique_ptr<PooledWorker> worker = AcquireWorker(local_pool, job.graph_path);
|
||||
if (!worker) {
|
||||
qWarning() << "RenderWorkerPool failed to acquire worker for ticket"
|
||||
<< ticket_id;
|
||||
job.ticket->Finish();
|
||||
ClearActiveWorker(worker_index, 0);
|
||||
return;
|
||||
}
|
||||
|
||||
for (int attempt = 0; attempt < kMaxAttempts; attempt++) {
|
||||
const JobResult result = ProcessJobAttempt(job, worker_index, attempt);
|
||||
if (attempt > 0) {
|
||||
worker = AcquireWorker(local_pool, job.graph_path);
|
||||
if (!worker) {
|
||||
qWarning() << "RenderWorkerPool failed to acquire worker for retry"
|
||||
<< ticket_id;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
const JobResult result = ProcessJobAttempt(job, worker_index, attempt,
|
||||
worker.get());
|
||||
const qint64 worker_pid = worker && worker->process
|
||||
? worker->process->processId()
|
||||
: 0;
|
||||
const bool process_state_running = worker && worker->process &&
|
||||
worker->process->state() == QProcess::Running;
|
||||
const bool os_alive = worker_pid > 0 && IsProcessAlive(worker_pid);
|
||||
const bool worker_healthy = process_state_running || os_alive;
|
||||
const bool keep_alive = (result == JobResult::kFinished) && worker_healthy;
|
||||
|
||||
ReturnWorker(local_pool, std::move(worker), keep_alive);
|
||||
worker.reset();
|
||||
|
||||
if (result == JobResult::kFinished) {
|
||||
ClearActiveWorker(worker_index, 0);
|
||||
return;
|
||||
@@ -578,198 +677,241 @@ void RenderWorkerPool::ProcessJob(const Job &job, int worker_index)
|
||||
}
|
||||
|
||||
RenderWorkerPool::JobResult RenderWorkerPool::ProcessJobAttempt(
|
||||
const Job &job, int worker_index, int attempt_index)
|
||||
const Job &job, int worker_index, int attempt_index,
|
||||
PooledWorker *worker)
|
||||
{
|
||||
const qint64 ticket_id = qint64(reinterpret_cast<quintptr>(job.ticket.get()));
|
||||
if (job.ticket->IsCancelled()) {
|
||||
return JobResult::kCancelled;
|
||||
}
|
||||
|
||||
const int linesize = Frame::generate_linesize_bytes(
|
||||
kMaxWidth, PixelFormat::F32, VideoParams::kRGBAChannelCount);
|
||||
const size_t slot_bytes = size_t(linesize) * kMaxHeight;
|
||||
const size_t region_bytes = ipc::FrameSlotPool::BytesNeeded(kOutputSlots, slot_bytes);
|
||||
const QString shm_key =
|
||||
ipc::SharedMemoryRegion::MakeKey(QCoreApplication::applicationPid(),
|
||||
int((reinterpret_cast<quintptr>(job.ticket.get()) +
|
||||
attempt_index * 2) & 0xFFFF));
|
||||
|
||||
ipc::SharedMemoryRegion region;
|
||||
if (!region.Open(shm_key, region_bytes, ipc::SharedMemoryRegion::kCreate)) {
|
||||
qWarning() << "RenderWorkerPool failed to create shared memory"
|
||||
<< region.error();
|
||||
return JobResult::kFatalFailure;
|
||||
if (!worker || !worker->process) {
|
||||
return JobResult::kRetryableFailure;
|
||||
}
|
||||
ipc::FrameSlotPool output_pool =
|
||||
ipc::FrameSlotPool::Create(region.data(), kOutputSlots, slot_bytes);
|
||||
|
||||
const QString input_shm_key =
|
||||
job.input_frames.isEmpty()
|
||||
? QString()
|
||||
: ipc::SharedMemoryRegion::MakeKey(
|
||||
QCoreApplication::applicationPid(),
|
||||
int((reinterpret_cast<quintptr>(job.ticket.get()) +
|
||||
attempt_index * 2 + 1) & 0xFFFF));
|
||||
ipc::SharedMemoryRegion input_region;
|
||||
std::optional<ipc::FrameSlotPool> input_pool;
|
||||
QVector<int> input_slots;
|
||||
if (!job.input_frames.isEmpty()) {
|
||||
const uint32_t input_slot_count = uint32_t(job.input_frames.size());
|
||||
const size_t input_region_bytes =
|
||||
ipc::FrameSlotPool::BytesNeeded(input_slot_count, slot_bytes);
|
||||
if (!input_region.Open(input_shm_key, input_region_bytes,
|
||||
ipc::SharedMemoryRegion::kCreate)) {
|
||||
qWarning() << "RenderWorkerPool failed to create input shared memory"
|
||||
<< input_region.error();
|
||||
const qint64 worker_process_id = worker->process->processId();
|
||||
|
||||
const int output_width = job.params.force_size.width() > 0
|
||||
? job.params.force_size.width()
|
||||
: job.params.video_params.effective_width();
|
||||
const int output_height = job.params.force_size.height() > 0
|
||||
? job.params.force_size.height()
|
||||
: job.params.video_params.effective_height();
|
||||
const PixelFormat::Format output_format =
|
||||
job.params.force_format != PixelFormat::INVALID
|
||||
? PixelFormat::Format(job.params.force_format)
|
||||
: PixelFormat::F32;
|
||||
const int output_channels = job.params.force_channel_count > 0
|
||||
? job.params.force_channel_count
|
||||
: VideoParams::kRGBAChannelCount;
|
||||
const int output_linesize =
|
||||
Frame::generate_linesize_bytes(output_width, output_format,
|
||||
output_channels);
|
||||
const size_t estimated_output_slot_bytes =
|
||||
size_t(output_linesize) * size_t(output_height);
|
||||
const int f32_rgba_linesize =
|
||||
Frame::generate_linesize_bytes(output_width, PixelFormat::F32,
|
||||
VideoParams::kRGBAChannelCount);
|
||||
const size_t f32_rgba_slot_bytes =
|
||||
size_t(f32_rgba_linesize) * size_t(output_height);
|
||||
const size_t output_slot_bytes =
|
||||
std::max(estimated_output_slot_bytes, f32_rgba_slot_bytes);
|
||||
size_t input_slot_bytes = 0;
|
||||
for (const FramePtr &frame : job.input_frames) {
|
||||
if (frame && frame->is_allocated()) {
|
||||
input_slot_bytes =
|
||||
std::max(input_slot_bytes, size_t(frame->allocated_size()));
|
||||
}
|
||||
}
|
||||
const size_t output_region_bytes =
|
||||
ipc::FrameSlotPool::BytesNeeded(kOutputSlots, output_slot_bytes);
|
||||
|
||||
if (!worker->output_region.IsValid() ||
|
||||
worker->output_slot_bytes < output_slot_bytes) {
|
||||
if (worker->output_region.IsValid()) {
|
||||
worker->output_region.Close();
|
||||
worker->output_pool = ipc::FrameSlotPool();
|
||||
}
|
||||
if (worker->output_shm_key.isEmpty()) {
|
||||
worker->output_shm_key =
|
||||
ipc::SharedMemoryRegion::MakeKey(worker_process_id, 0) +
|
||||
QStringLiteral("-out");
|
||||
}
|
||||
if (!worker->output_region.Open(worker->output_shm_key,
|
||||
output_region_bytes,
|
||||
ipc::SharedMemoryRegion::kCreate)) {
|
||||
qWarning() << "RenderWorkerPool failed to create output shared memory"
|
||||
<< worker->output_region.error();
|
||||
return JobResult::kFatalFailure;
|
||||
} else {
|
||||
input_pool = ipc::FrameSlotPool::Create(input_region.data(),
|
||||
input_slot_count,
|
||||
slot_bytes);
|
||||
for (const FramePtr &frame : job.input_frames) {
|
||||
if (frame->allocated_size() > int(slot_bytes)) {
|
||||
qWarning() << "RenderWorkerPool decoded input frame exceeds slot size";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
}
|
||||
worker->output_pool = ipc::FrameSlotPool::Create(
|
||||
worker->output_region.data(), kOutputSlots, output_slot_bytes);
|
||||
worker->output_slot_bytes = output_slot_bytes;
|
||||
}
|
||||
const QString shm_key = worker->output_shm_key;
|
||||
ipc::FrameSlotPool &output_pool = worker->output_pool;
|
||||
|
||||
uint32_t slot = 0;
|
||||
if (!input_pool->Acquire(&slot)) {
|
||||
qWarning() << "RenderWorkerPool input pool had no free slot";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
|
||||
memcpy(input_pool->SlotData(slot), frame->const_data(),
|
||||
size_t(frame->allocated_size()));
|
||||
ipc::FrameSlotMeta *meta = input_pool->Meta(slot);
|
||||
meta->id = qint64(input_slots.size());
|
||||
meta->time_num = frame->timestamp().numerator();
|
||||
meta->time_den = frame->timestamp().denominator();
|
||||
meta->width = frame->width();
|
||||
meta->height = frame->height();
|
||||
meta->format = int32_t(frame->format());
|
||||
meta->channel_count = frame->channel_count();
|
||||
meta->linesize = frame->linesize_bytes();
|
||||
meta->data_size = frame->allocated_size();
|
||||
if (!input_pool->Publish(slot)) {
|
||||
qWarning() << "RenderWorkerPool failed to publish input slot";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
input_slots.append(int(slot));
|
||||
const uint32_t input_slot_count =
|
||||
job.input_frames.isEmpty() ? 0 : uint32_t(job.input_frames.size());
|
||||
if (input_slot_count > 0) {
|
||||
if (!worker->input_region.IsValid() ||
|
||||
worker->input_slot_bytes < input_slot_bytes ||
|
||||
worker->input_pool.slot_count() < input_slot_count) {
|
||||
if (worker->input_region.IsValid()) {
|
||||
worker->input_region.Close();
|
||||
worker->input_pool = ipc::FrameSlotPool();
|
||||
}
|
||||
|
||||
if (input_slots.size() != job.input_frames.size()) {
|
||||
qWarning() << "RenderWorkerPool failed to publish all input frames;"
|
||||
<< "aborting worker render";
|
||||
if (worker->input_shm_key.isEmpty()) {
|
||||
worker->input_shm_key =
|
||||
ipc::SharedMemoryRegion::MakeKey(worker_process_id, 1) +
|
||||
QStringLiteral("-in");
|
||||
}
|
||||
const size_t input_region_bytes =
|
||||
ipc::FrameSlotPool::BytesNeeded(input_slot_count, input_slot_bytes);
|
||||
if (!worker->input_region.Open(worker->input_shm_key,
|
||||
input_region_bytes,
|
||||
ipc::SharedMemoryRegion::kCreate)) {
|
||||
qWarning() << "RenderWorkerPool failed to create input shared memory"
|
||||
<< worker->input_region.error();
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
worker->input_pool = ipc::FrameSlotPool::Create(
|
||||
worker->input_region.data(), input_slot_count, input_slot_bytes);
|
||||
worker->input_slot_bytes = input_slot_bytes;
|
||||
}
|
||||
}
|
||||
const QString input_shm_key = worker->input_shm_key;
|
||||
ipc::FrameSlotPool &input_pool = worker->input_pool;
|
||||
QVector<int> input_slots;
|
||||
if (input_slot_count > 0) {
|
||||
for (const FramePtr &frame : job.input_frames) {
|
||||
if (frame->allocated_size() > int(worker->input_slot_bytes)) {
|
||||
qWarning() << "RenderWorkerPool decoded input frame exceeds slot size";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
|
||||
uint32_t slot = 0;
|
||||
if (!input_pool.Acquire(&slot)) {
|
||||
qWarning() << "RenderWorkerPool input pool had no free slot";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
|
||||
memcpy(input_pool.SlotData(slot), frame->const_data(),
|
||||
size_t(frame->allocated_size()));
|
||||
ipc::FrameSlotMeta *meta = input_pool.Meta(slot);
|
||||
meta->id = qint64(input_slots.size());
|
||||
meta->time_num = frame->timestamp().numerator();
|
||||
meta->time_den = frame->timestamp().denominator();
|
||||
meta->width = frame->width();
|
||||
meta->height = frame->height();
|
||||
meta->format = int32_t(frame->format());
|
||||
meta->channel_count = frame->channel_count();
|
||||
meta->linesize = frame->linesize_bytes();
|
||||
meta->data_size = frame->allocated_size();
|
||||
memset(meta->colorspace, 0, sizeof(meta->colorspace));
|
||||
const QString cs = frame->video_params().colorspace();
|
||||
if (!cs.isEmpty()) {
|
||||
const QByteArray cs_utf8 = cs.toUtf8();
|
||||
const size_t copy_len = qMin(
|
||||
static_cast<size_t>(cs_utf8.size()),
|
||||
sizeof(meta->colorspace) - 1);
|
||||
memcpy(meta->colorspace, cs_utf8.constData(), copy_len);
|
||||
meta->colorspace[copy_len] = '\0';
|
||||
}
|
||||
if (!input_pool.Publish(slot)) {
|
||||
qWarning() << "RenderWorkerPool failed to publish input slot";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
input_slots.append(int(slot));
|
||||
}
|
||||
|
||||
if (input_slots.size() != job.input_frames.size()) {
|
||||
qWarning() << "RenderWorkerPool failed to publish all input frames;"
|
||||
<< "aborting worker render";
|
||||
return JobResult::kFatalFailure;
|
||||
}
|
||||
}
|
||||
|
||||
QProcess worker;
|
||||
worker.setProgram(WorkerProgramPath());
|
||||
worker.start();
|
||||
if (!worker.waitForStarted(10000)) {
|
||||
qWarning() << "RenderWorkerPool failed to start worker"
|
||||
<< worker.errorString();
|
||||
return JobResult::kRetryableFailure;
|
||||
}
|
||||
const qint64 worker_process_id = worker.processId();
|
||||
|
||||
SetActiveWorker(worker_index, job.ticket, &worker, ticket_id);
|
||||
SetActiveWorker(worker_index, job.ticket, worker->process, ticket_id);
|
||||
if (job.ticket->IsCancelled()) {
|
||||
ipc::CancelMsg cancel;
|
||||
cancel.ticket_id = ticket_id;
|
||||
TryWriteControlMessage(&worker, cancel.ToJson());
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
TryWriteControlMessage(worker->process, cancel.ToJson());
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return JobResult::kCancelled;
|
||||
}
|
||||
|
||||
QString error;
|
||||
QJsonObject response;
|
||||
if (!ReadControlMessage(&worker, &response, &error)) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool did not receive startup handshake"
|
||||
<< error << worker.readAllStandardError();
|
||||
}
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
|
||||
ipc::HandshakeMsg handshake;
|
||||
handshake.protocol_version = kProtocolVersion;
|
||||
handshake.shm_key = shm_key;
|
||||
handshake.input_shm_key = input_slots.isEmpty() ? QString() : input_shm_key;
|
||||
handshake.input_slots = input_slots.size();
|
||||
handshake.output_slots = int(kOutputSlots);
|
||||
handshake.slot_data_bytes = qint64(slot_bytes);
|
||||
handshake.input_slot_data_bytes = input_slots.isEmpty() ? 0 : qint64(slot_bytes);
|
||||
if (!WriteControlMessage(&worker, handshake.ToJson())) {
|
||||
handshake.slot_data_bytes = qint64(output_slot_bytes);
|
||||
handshake.input_slot_data_bytes = input_slots.isEmpty()
|
||||
? 0
|
||||
: qint64(input_slot_bytes);
|
||||
if (!WriteControlMessage(worker->process, handshake.ToJson())) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool failed to send shared-memory handshake";
|
||||
}
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
|
||||
ipc::LoadGraphMsg load;
|
||||
load.path = job.graph_path;
|
||||
if (!WriteControlMessage(&worker, load.ToJson()) ||
|
||||
!ReadControlMessage(&worker, &response, &error)) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool failed to load graph in worker"
|
||||
<< error << worker.readAllStandardError();
|
||||
if (worker->loaded_graph_path != job.graph_path) {
|
||||
ipc::LoadGraphMsg load;
|
||||
load.path = job.graph_path;
|
||||
QString error;
|
||||
QJsonObject response;
|
||||
if (!WriteControlMessage(worker->process, load.ToJson()) ||
|
||||
!ReadControlMessage(worker->process, &response, &error)) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool failed to load graph in worker"
|
||||
<< error << worker->process->readAllStandardError();
|
||||
}
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
worker->loaded_graph_path = job.graph_path;
|
||||
|
||||
}
|
||||
ipc::RenderFrameMsg render;
|
||||
render.ticket_id = ticket_id;
|
||||
render.node_uuid = job.node_token;
|
||||
render.time_num = job.params.time.numerator();
|
||||
render.time_den = job.params.time.denominator();
|
||||
render.width = job.params.force_size.width();
|
||||
render.height = job.params.force_size.height();
|
||||
render.format = int(job.params.force_format);
|
||||
render.channel_count = job.params.force_channel_count;
|
||||
render.width = output_width;
|
||||
render.height = output_height;
|
||||
render.format = int(output_format);
|
||||
render.channel_count = output_channels;
|
||||
render.mode = int(job.params.mode);
|
||||
render.input_slot = input_slots.isEmpty() ? -1 : input_slots.front();
|
||||
render.input_slots = input_slots;
|
||||
|
||||
if (!WriteControlMessage(&worker, render.ToJson())) {
|
||||
if (!WriteControlMessage(worker->process, render.ToJson())) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool failed to send render_frame";
|
||||
}
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
|
||||
QString error;
|
||||
QJsonObject response;
|
||||
ipc::FrameReadyMsg ready;
|
||||
while (true) {
|
||||
if (!ReadControlMessage(&worker, &response, &error, 30000)) {
|
||||
if (!ReadControlMessage(worker->process, &response, &error, 30000)) {
|
||||
if (!job.ticket->IsCancelled()) {
|
||||
qWarning() << "RenderWorkerPool failed waiting for frame_ready"
|
||||
<< error << worker.readAllStandardError();
|
||||
<< error << worker->process->readAllStandardError();
|
||||
}
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return job.ticket->IsCancelled() ? JobResult::kCancelled
|
||||
: JobResult::kRetryableFailure;
|
||||
: JobResult::kRetryableFailure;
|
||||
}
|
||||
|
||||
if (ipc::FrameReadyMsg::FromJson(response, &ready)) {
|
||||
@@ -780,19 +922,22 @@ RenderWorkerPool::JobResult RenderWorkerPool::ProcessJobAttempt(
|
||||
if (job.ticket->IsCancelled()) {
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return JobResult::kCancelled;
|
||||
} else {
|
||||
FinishWithFrame(job.ticket, output_pool, uint32_t(ready.output_slot));
|
||||
}
|
||||
|
||||
uint32_t consumed_slot = 0;
|
||||
if (!output_pool.Consume(&consumed_slot)) {
|
||||
qWarning() << "RenderWorkerPool failed to consume output slot";
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
return JobResult::kRetryableFailure;
|
||||
}
|
||||
if (int(consumed_slot) != ready.output_slot) {
|
||||
qWarning() << "RenderWorkerPool output slot mismatch: consumed"
|
||||
<< consumed_slot << "expected" << ready.output_slot;
|
||||
}
|
||||
FinishWithFrame(job.ticket, output_pool, consumed_slot);
|
||||
output_pool.Release(consumed_slot);
|
||||
ClearActiveWorker(worker_index, worker_process_id);
|
||||
|
||||
QJsonObject shutdown;
|
||||
shutdown[QStringLiteral("type")] = ipc::msgtype::kShutdown;
|
||||
WriteControlMessage(&worker, shutdown);
|
||||
worker.closeWriteChannel();
|
||||
if (!worker.waitForFinished(5000)) {
|
||||
worker.kill();
|
||||
worker.waitForFinished();
|
||||
}
|
||||
return JobResult::kFinished;
|
||||
}
|
||||
|
||||
@@ -834,8 +979,163 @@ void RenderWorkerPool::ClearActiveWorker(int worker_index, qint64 process_id)
|
||||
|
||||
int RenderWorkerPool::WorkerCount() const
|
||||
{
|
||||
// GPU rendering is the bottleneck for video frames; too many workers just
|
||||
// multiply first-frame warmup (shader/OCIO cache creation) and compete for
|
||||
// the same GPU. Cap at a small number while still leaving cores free.
|
||||
const int ideal = QThread::idealThreadCount();
|
||||
return std::max(1, ideal - 2);
|
||||
return std::max(1, std::min(ideal - 2, 4));
|
||||
}
|
||||
|
||||
std::unique_ptr<RenderWorkerPool::PooledWorker> RenderWorkerPool::AcquireWorker(
|
||||
std::vector<std::unique_ptr<PooledWorker>> *local_pool,
|
||||
const QString &graph_path)
|
||||
{
|
||||
if (!local_pool) {
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
const qint64 now = QDateTime::currentMSecsSinceEpoch();
|
||||
|
||||
// Prefer an idle worker that already has the requested graph loaded.
|
||||
int best_index = -1;
|
||||
for (size_t i = 0; i < local_pool->size();) {
|
||||
PooledWorker *candidate = (*local_pool)[i].get();
|
||||
if (!candidate || !candidate->process) {
|
||||
local_pool->erase(local_pool->begin() + i);
|
||||
continue;
|
||||
}
|
||||
const bool candidate_state_running =
|
||||
candidate->process->state() == QProcess::Running;
|
||||
const bool candidate_os_alive =
|
||||
IsProcessAlive(candidate->process->processId());
|
||||
if (!candidate_state_running && !candidate_os_alive) {
|
||||
ShutdownWorker(candidate);
|
||||
local_pool->erase(local_pool->begin() + i);
|
||||
continue;
|
||||
}
|
||||
if (now - candidate->last_used_ms > kWorkerIdleTimeoutMs) {
|
||||
ShutdownWorker(candidate);
|
||||
local_pool->erase(local_pool->begin() + i);
|
||||
continue;
|
||||
}
|
||||
if (best_index < 0 ||
|
||||
(!candidate->loaded_graph_path.isEmpty() &&
|
||||
candidate->loaded_graph_path == graph_path &&
|
||||
((*local_pool)[size_t(best_index)]->loaded_graph_path != graph_path))) {
|
||||
best_index = int(i);
|
||||
}
|
||||
++i;
|
||||
}
|
||||
|
||||
if (best_index >= 0) {
|
||||
std::unique_ptr<PooledWorker> worker =
|
||||
std::move((*local_pool)[size_t(best_index)]);
|
||||
local_pool->erase(local_pool->begin() + best_index);
|
||||
worker->last_used_ms = now;
|
||||
++worker->use_count;
|
||||
return worker;
|
||||
}
|
||||
|
||||
|
||||
// No idle worker available: start a new one.
|
||||
auto *process = new QProcess();
|
||||
process->setProgram(WorkerProgramPath());
|
||||
process->setArguments({QStringLiteral("--backend"), gpu_backend_});
|
||||
|
||||
const QString worker_stderr_path = QDir(QDir::tempPath()).filePath(
|
||||
QStringLiteral("oak-render-worker-%1-%2.stderr.log")
|
||||
.arg(QCoreApplication::applicationPid())
|
||||
.arg(QDateTime::currentMSecsSinceEpoch()));
|
||||
process->setStandardErrorFile(worker_stderr_path);
|
||||
|
||||
process->start();
|
||||
if (!process->waitForStarted(10000)) {
|
||||
qWarning() << "RenderWorkerPool failed to start worker"
|
||||
<< process->errorString();
|
||||
delete process;
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
QString error;
|
||||
QJsonObject response;
|
||||
if (!ReadControlMessage(process, &response, &error)) {
|
||||
qWarning() << "RenderWorkerPool did not receive startup handshake"
|
||||
<< error << process->readAllStandardError();
|
||||
process->kill();
|
||||
process->waitForFinished();
|
||||
delete process;
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
auto worker = std::make_unique<PooledWorker>();
|
||||
worker->process = process;
|
||||
worker->last_used_ms = now;
|
||||
worker->use_count = 1;
|
||||
return worker;
|
||||
}
|
||||
|
||||
void RenderWorkerPool::ReturnWorker(
|
||||
std::vector<std::unique_ptr<PooledWorker>> *local_pool,
|
||||
std::unique_ptr<PooledWorker> worker,
|
||||
bool keep_alive)
|
||||
{
|
||||
if (!worker || !worker->process) {
|
||||
return;
|
||||
}
|
||||
|
||||
const bool pool_full = worker->use_count >= kWorkerMaxUses;
|
||||
if (!keep_alive || stopping_ || pool_full) {
|
||||
ShutdownWorker(worker.get());
|
||||
return;
|
||||
}
|
||||
|
||||
worker->last_used_ms = QDateTime::currentMSecsSinceEpoch();
|
||||
local_pool->push_back(std::move(worker));
|
||||
}
|
||||
|
||||
void RenderWorkerPool::ShutdownWorker(PooledWorker *worker)
|
||||
{
|
||||
if (!worker || !worker->process) {
|
||||
return;
|
||||
}
|
||||
|
||||
QProcess *process = worker->process;
|
||||
worker->process = nullptr;
|
||||
worker->loaded_graph_path.clear();
|
||||
worker->use_count = 0;
|
||||
|
||||
if (process->state() == QProcess::Running) {
|
||||
QJsonObject shutdown;
|
||||
shutdown[QStringLiteral("type")] = ipc::msgtype::kShutdown;
|
||||
TryWriteControlMessage(process, shutdown);
|
||||
process->closeWriteChannel();
|
||||
if (!process->waitForFinished(5000)) {
|
||||
process->kill();
|
||||
process->waitForFinished();
|
||||
}
|
||||
}
|
||||
delete process;
|
||||
}
|
||||
|
||||
void RenderWorkerPool::ShutdownLocalPool(
|
||||
std::vector<std::unique_ptr<PooledWorker>> *local_pool)
|
||||
{
|
||||
if (!local_pool) {
|
||||
return;
|
||||
}
|
||||
for (std::unique_ptr<PooledWorker> &worker : *local_pool) {
|
||||
ShutdownWorker(worker.get());
|
||||
}
|
||||
local_pool->clear();
|
||||
}
|
||||
|
||||
void RenderWorkerPool::ClearGraphCache()
|
||||
{
|
||||
QMutexLocker locker(&mutex_);
|
||||
for (auto it = graph_cache_.begin(); it != graph_cache_.end(); ++it) {
|
||||
CleanupGraphFile(it->path);
|
||||
}
|
||||
graph_cache_.clear();
|
||||
}
|
||||
|
||||
void RenderWorkerPool::FinishWithFrame(RenderTicketPtr ticket,
|
||||
|
||||
Reference in New Issue
Block a user