diff --git a/app/render/rendermanager.cpp b/app/render/rendermanager.cpp index 0471bf245..ab9d8b768 100644 --- a/app/render/rendermanager.cpp +++ b/app/render/rendermanager.cpp @@ -181,6 +181,10 @@ RenderTicketPtr RenderManager::RenderAudio(const RenderAudioParams ¶ms) bool RenderManager::RemoveTicket(RenderTicketPtr ticket) { + if (worker_pool_ && worker_pool_->RemoveTicket(ticket)) { + return true; + } + for (RenderThread *rt : render_threads_) { if (rt->RemoveTicket(ticket)) { return true; diff --git a/app/render/renderprocessor.cpp b/app/render/renderprocessor.cpp index a7513616a..ea3d5ed03 100644 --- a/app/render/renderprocessor.cpp +++ b/app/render/renderprocessor.cpp @@ -467,6 +467,12 @@ void RenderProcessor::ProcessVideoFootage(TexturePtr destination, input_slot = input_slot_value.isValid() ? input_slot_value.toInt() : -1; } if (render_ctx_ && input_pool && input_slot >= 0) { + if (input_slot >= int(input_pool->slot_count())) { + qWarning() << "RenderProcessor received out-of-range IPC input frame slot" + << input_slot; + return; + } + const ipc::FrameSlotMeta *meta = input_pool->Meta(uint32_t(input_slot)); if (meta && meta->width > 0 && meta->height > 0 && meta->data_size > 0 && meta->data_size <= int(input_pool->slot_data_bytes())) { diff --git a/app/render/renderworkerpool.cpp b/app/render/renderworkerpool.cpp index 8ee316f17..5e17f2a68 100644 --- a/app/render/renderworkerpool.cpp +++ b/app/render/renderworkerpool.cpp @@ -29,7 +29,15 @@ #include #include #include +#include #include +#include +#include +#if defined(Q_OS_WIN) +#include +#else +#include +#endif #include "codec/frame.h" #include "common/qtutils.h" @@ -221,12 +229,61 @@ bool WriteControlMessage(QProcess *process, const QJsonObject &obj) return process->write(line) == line.size() && process->waitForBytesWritten(5000); } +void TryWriteControlMessage(QProcess *process, const QJsonObject &obj) +{ + const QByteArray line = QJsonDocument(obj).toJson(QJsonDocument::Compact) + '\n'; + process->write(line); +} + +bool KillProcessById(qint64 process_id) +{ + if (process_id <= 0) { + return false; + } + +#if defined(Q_OS_WIN) + HANDLE handle = OpenProcess(PROCESS_TERMINATE, FALSE, DWORD(process_id)); + if (!handle) { + return false; + } + const bool ok = TerminateProcess(handle, 1) != 0; + CloseHandle(handle); + return ok; +#else + return ::kill(pid_t(process_id), SIGKILL) == 0; +#endif +} + +QString WorkerProcessDetails(const QProcess *process) +{ + if (!process) { + return QStringLiteral("worker process unavailable"); + } + + const QString exit_status = + process->exitStatus() == QProcess::CrashExit + ? QStringLiteral("crash") + : QStringLiteral("normal"); + return QStringLiteral("state=%1 exit_status=%2 exit_code=%3 process_error=%4 error=\"%5\"") + .arg(int(process->state())) + .arg(exit_status) + .arg(process->exitCode()) + .arg(int(process->error())) + .arg(process->errorString()); +} + bool ReadControlMessage(QProcess *process, QJsonObject *out, QString *error, int timeout_ms = 10000) { if (!process->waitForReadyRead(timeout_ms)) { if (error) { - *error = QStringLiteral("timeout waiting for worker response"); + if (process->state() == QProcess::NotRunning) { + *error = QStringLiteral("worker exited before response: %1") + .arg(WorkerProcessDetails(process)); + } else { + *error = QStringLiteral("timeout waiting for worker response: %1") + .arg(WorkerProcessDetails(process)); + } } return false; } @@ -293,12 +350,72 @@ bool RenderWorkerPool::SubmitFrame(RenderTicketPtr ticket, return true; } +bool RenderWorkerPool::RemoveTicket(RenderTicketPtr ticket) +{ + if (!ticket) { + return false; + } + + QString queued_graph_path; + bool matched_active = false; + + { + QMutexLocker locker(&mutex_); + auto it = std::find_if(queue_.begin(), queue_.end(), + [&ticket](const Job &job) { + return job.ticket == ticket; + }); + if (it != queue_.end()) { + queued_graph_path = it->graph_path; + queue_.erase(it); + } else { + for (const ActiveJob &active : active_jobs_) { + if (active.ticket == ticket) { + ticket->Cancel(); + CancelActiveProcess(active.process_id); + matched_active = true; + break; + } + } + if (!matched_active) { + return false; + } + } + } + + if (!queued_graph_path.isEmpty()) { + CleanupGraphFile(queued_graph_path); + return true; + } + + return true; +} + void RenderWorkerPool::Shutdown() { + QVector queued_graph_paths; + { QMutexLocker locker(&mutex_); stopping_ = true; - wait_.wakeOne(); + for (Job &job : queue_) { + if (job.ticket) { + job.ticket->Cancel(); + } + queued_graph_paths.append(job.graph_path); + } + queue_.clear(); + for (ActiveJob &active : active_jobs_) { + if (active.ticket) { + active.ticket->Cancel(); + CancelActiveProcess(active.process_id); + } + } + wait_.wakeAll(); + } + + for (const QString &path : queued_graph_paths) { + CleanupGraphFile(path); } if (isRunning()) { @@ -307,6 +424,30 @@ void RenderWorkerPool::Shutdown() } void RenderWorkerPool::run() +{ + const int worker_count = WorkerCount(); + { + QMutexLocker locker(&mutex_); + active_jobs_.resize(worker_count); + } + + std::vector workers; + workers.reserve(size_t(worker_count)); + for (int i = 0; i < worker_count; i++) { + workers.emplace_back([this, i]() { + WorkerLoop(i); + }); + } + + for (std::thread &worker : workers) { + worker.join(); + } + + QMutexLocker locker(&mutex_); + active_jobs_.clear(); +} + +void RenderWorkerPool::WorkerLoop(int worker_index) { while (true) { mutex_.lock(); @@ -322,7 +463,7 @@ void RenderWorkerPool::run() queue_.pop_front(); mutex_.unlock(); - ProcessJob(job); + ProcessJob(job, worker_index); CleanupGraphFile(job.graph_path); } } @@ -394,28 +535,70 @@ bool RenderWorkerPool::IsSupported(const RenderManager::RenderVideoParams ¶m params.video_params.is_valid(); } -void RenderWorkerPool::ProcessJob(const Job &job) +void RenderWorkerPool::ProcessJob(const Job &job, int worker_index) { + const qint64 ticket_id = qint64(reinterpret_cast(job.ticket.get())); + SetActiveWorker(worker_index, job.ticket, nullptr, ticket_id); + job.ticket->Start(); if (job.ticket->IsCancelled()) { job.ticket->Finish(); + ClearActiveWorker(worker_index, 0); return; } + for (int attempt = 0; attempt < kMaxAttempts; attempt++) { + const JobResult result = ProcessJobAttempt(job, worker_index, attempt); + if (result == JobResult::kFinished) { + ClearActiveWorker(worker_index, 0); + return; + } + if (result == JobResult::kCancelled) { + job.ticket->Finish(); + ClearActiveWorker(worker_index, 0); + return; + } + if (result == JobResult::kFatalFailure) { + break; + } + if (attempt + 1 < kMaxAttempts && !job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool retrying render worker for ticket" + << ticket_id << "after worker failure"; + } + } + + if (job.ticket->IsCancelled()) { + job.ticket->Finish(); + } else { + qWarning() << "RenderWorkerPool exhausted worker retries for ticket" + << ticket_id; + job.ticket->Finish(); + } + ClearActiveWorker(worker_index, 0); +} + +RenderWorkerPool::JobResult RenderWorkerPool::ProcessJobAttempt( + const Job &job, int worker_index, int attempt_index) +{ + const qint64 ticket_id = qint64(reinterpret_cast(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(job.ticket.get()) & 0xFFFF)); + int((reinterpret_cast(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(); - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } ipc::FrameSlotPool output_pool = ipc::FrameSlotPool::Create(region.data(), kOutputSlots, slot_bytes); @@ -425,7 +608,8 @@ void RenderWorkerPool::ProcessJob(const Job &job) ? QString() : ipc::SharedMemoryRegion::MakeKey( QCoreApplication::applicationPid(), - int((reinterpret_cast(job.ticket.get()) + 1) & 0xFFFF)); + int((reinterpret_cast(job.ticket.get()) + + attempt_index * 2 + 1) & 0xFFFF)); ipc::SharedMemoryRegion input_region; std::optional input_pool; QVector input_slots; @@ -437,8 +621,7 @@ void RenderWorkerPool::ProcessJob(const Job &job) ipc::SharedMemoryRegion::kCreate)) { qWarning() << "RenderWorkerPool failed to create input shared memory" << input_region.error(); - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } else { input_pool = ipc::FrameSlotPool::Create(input_region.data(), input_slot_count, @@ -446,15 +629,13 @@ void RenderWorkerPool::ProcessJob(const Job &job) for (const FramePtr &frame : job.input_frames) { if (frame->allocated_size() > int(slot_bytes)) { qWarning() << "RenderWorkerPool decoded input frame exceeds slot size"; - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } uint32_t slot = 0; if (!input_pool->Acquire(&slot)) { qWarning() << "RenderWorkerPool input pool had no free slot"; - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } memcpy(input_pool->SlotData(slot), frame->const_data(), @@ -471,8 +652,7 @@ void RenderWorkerPool::ProcessJob(const Job &job) meta->data_size = frame->allocated_size(); if (!input_pool->Publish(slot)) { qWarning() << "RenderWorkerPool failed to publish input slot"; - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } input_slots.append(int(slot)); } @@ -480,8 +660,7 @@ void RenderWorkerPool::ProcessJob(const Job &job) if (input_slots.size() != job.input_frames.size()) { qWarning() << "RenderWorkerPool failed to publish all input frames;" << "aborting worker render"; - job.ticket->Finish(); - return; + return JobResult::kFatalFailure; } } } @@ -492,19 +671,33 @@ void RenderWorkerPool::ProcessJob(const Job &job) if (!worker.waitForStarted(10000)) { qWarning() << "RenderWorkerPool failed to start worker" << worker.errorString(); - job.ticket->Finish(); - return; + return JobResult::kRetryableFailure; + } + const qint64 worker_process_id = worker.processId(); + + SetActiveWorker(worker_index, job.ticket, &worker, ticket_id); + if (job.ticket->IsCancelled()) { + ipc::CancelMsg cancel; + cancel.ticket_id = ticket_id; + TryWriteControlMessage(&worker, cancel.ToJson()); + worker.kill(); + worker.waitForFinished(); + ClearActiveWorker(worker_index, worker_process_id); + return JobResult::kCancelled; } QString error; QJsonObject response; if (!ReadControlMessage(&worker, &response, &error)) { - qWarning() << "RenderWorkerPool did not receive startup handshake" - << error << worker.readAllStandardError(); + if (!job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool did not receive startup handshake" + << error << worker.readAllStandardError(); + } worker.kill(); worker.waitForFinished(); - job.ticket->Finish(); - return; + ClearActiveWorker(worker_index, worker_process_id); + return job.ticket->IsCancelled() ? JobResult::kCancelled + : JobResult::kRetryableFailure; } ipc::HandshakeMsg handshake; @@ -516,27 +709,33 @@ void RenderWorkerPool::ProcessJob(const Job &job) handshake.slot_data_bytes = qint64(slot_bytes); handshake.input_slot_data_bytes = input_slots.isEmpty() ? 0 : qint64(slot_bytes); if (!WriteControlMessage(&worker, handshake.ToJson())) { - qWarning() << "RenderWorkerPool failed to send shared-memory handshake"; + if (!job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool failed to send shared-memory handshake"; + } worker.kill(); worker.waitForFinished(); - job.ticket->Finish(); - return; + ClearActiveWorker(worker_index, worker_process_id); + return job.ticket->IsCancelled() ? JobResult::kCancelled + : JobResult::kRetryableFailure; } ipc::LoadGraphMsg load; load.path = job.graph_path; if (!WriteControlMessage(&worker, load.ToJson()) || !ReadControlMessage(&worker, &response, &error)) { - qWarning() << "RenderWorkerPool failed to load graph in worker" - << error << worker.readAllStandardError(); + if (!job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool failed to load graph in worker" + << error << worker.readAllStandardError(); + } worker.kill(); worker.waitForFinished(); - job.ticket->Finish(); - return; + ClearActiveWorker(worker_index, worker_process_id); + return job.ticket->IsCancelled() ? JobResult::kCancelled + : JobResult::kRetryableFailure; } ipc::RenderFrameMsg render; - render.ticket_id = qint64(reinterpret_cast(job.ticket.get())); + 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(); @@ -549,22 +748,28 @@ void RenderWorkerPool::ProcessJob(const Job &job) render.input_slots = input_slots; if (!WriteControlMessage(&worker, render.ToJson())) { - qWarning() << "RenderWorkerPool failed to send render_frame"; + if (!job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool failed to send render_frame"; + } worker.kill(); worker.waitForFinished(); - job.ticket->Finish(); - return; + ClearActiveWorker(worker_index, worker_process_id); + return job.ticket->IsCancelled() ? JobResult::kCancelled + : JobResult::kRetryableFailure; } ipc::FrameReadyMsg ready; while (true) { if (!ReadControlMessage(&worker, &response, &error, 30000)) { - qWarning() << "RenderWorkerPool failed waiting for frame_ready" - << error << worker.readAllStandardError(); + if (!job.ticket->IsCancelled()) { + qWarning() << "RenderWorkerPool failed waiting for frame_ready" + << error << worker.readAllStandardError(); + } worker.kill(); worker.waitForFinished(); - job.ticket->Finish(); - return; + ClearActiveWorker(worker_index, worker_process_id); + return job.ticket->IsCancelled() ? JobResult::kCancelled + : JobResult::kRetryableFailure; } if (ipc::FrameReadyMsg::FromJson(response, &ready)) { @@ -572,7 +777,13 @@ void RenderWorkerPool::ProcessJob(const Job &job) } } - FinishWithFrame(job.ticket, output_pool, uint32_t(ready.output_slot)); + if (job.ticket->IsCancelled()) { + ClearActiveWorker(worker_index, worker_process_id); + return JobResult::kCancelled; + } else { + FinishWithFrame(job.ticket, output_pool, uint32_t(ready.output_slot)); + } + ClearActiveWorker(worker_index, worker_process_id); QJsonObject shutdown; shutdown[QStringLiteral("type")] = ipc::msgtype::kShutdown; @@ -582,6 +793,49 @@ void RenderWorkerPool::ProcessJob(const Job &job) worker.kill(); worker.waitForFinished(); } + return JobResult::kFinished; +} + +void RenderWorkerPool::CancelActiveProcess(qint64 process_id) +{ + KillProcessById(process_id); +} + +void RenderWorkerPool::SetActiveWorker(int worker_index, RenderTicketPtr ticket, + QProcess *worker, qint64 ticket_id) +{ + QMutexLocker locker(&mutex_); + if (worker_index < 0 || worker_index >= active_jobs_.size()) { + return; + } + + ActiveJob &active = active_jobs_[worker_index]; + active.ticket = ticket; + active.process_id = worker ? worker->processId() : 0; + active.ticket_id = ticket_id; +} + +void RenderWorkerPool::ClearActiveWorker(int worker_index, qint64 process_id) +{ + QMutexLocker locker(&mutex_); + if (worker_index < 0 || worker_index >= active_jobs_.size()) { + return; + } + + ActiveJob &active = active_jobs_[worker_index]; + if (process_id > 0) { + if (active.process_id == process_id) { + active.process_id = 0; + } + } else { + active = ActiveJob(); + } +} + +int RenderWorkerPool::WorkerCount() const +{ + const int ideal = QThread::idealThreadCount(); + return std::max(1, ideal - 2); } void RenderWorkerPool::FinishWithFrame(RenderTicketPtr ticket, diff --git a/app/render/renderworkerpool.h b/app/render/renderworkerpool.h index b19ca1f5c..799a398a4 100644 --- a/app/render/renderworkerpool.h +++ b/app/render/renderworkerpool.h @@ -34,6 +34,8 @@ #include "render/ipc/sharedmemoryregion.h" #include "render/rendermanager.h" +class QProcess; + namespace olive { @@ -47,6 +49,8 @@ public: bool SubmitFrame(RenderTicketPtr ticket, const RenderManager::RenderVideoParams ¶ms); + bool RemoveTicket(RenderTicketPtr ticket); + void Shutdown(); protected: @@ -67,24 +71,47 @@ private: QVector input_frames; }; + enum class JobResult { + kFinished, + kRetryableFailure, + kFatalFailure, + kCancelled + }; + + struct ActiveJob { + RenderTicketPtr ticket; + qint64 process_id = 0; + qint64 ticket_id = 0; + }; + bool PrepareJob(RenderTicketPtr ticket, const RenderManager::RenderVideoParams ¶ms, Job *job); bool WriteGraphSnapshot(Project *project, QString *path); bool IsSupported(const RenderManager::RenderVideoParams ¶ms) const; - void ProcessJob(const Job &job); + void WorkerLoop(int worker_index); + void ProcessJob(const Job &job, int worker_index); + JobResult ProcessJobAttempt(const Job &job, int worker_index, + int attempt_index); void FinishWithFrame(RenderTicketPtr ticket, const ipc::FrameSlotPool &pool, uint32_t slot); void CleanupGraphFile(const QString &path); + void CancelActiveProcess(qint64 process_id); + void SetActiveWorker(int worker_index, RenderTicketPtr ticket, + QProcess *worker, qint64 ticket_id); + void ClearActiveWorker(int worker_index, qint64 process_id); + int WorkerCount() const; DecoderCache *decoder_cache_; QMutex mutex_; QWaitCondition wait_; std::deque queue_; bool stopping_ = false; + QVector active_jobs_; static constexpr uint32_t kOutputSlots = 2; + static constexpr int kMaxAttempts = 2; static constexpr int kMaxWidth = 4096; static constexpr int kMaxHeight = 2160; }; diff --git a/app/render/worker/workermain.cpp b/app/render/worker/workermain.cpp index 494c10543..da0816f45 100644 --- a/app/render/worker/workermain.cpp +++ b/app/render/worker/workermain.cpp @@ -302,6 +302,15 @@ private: } for (int requested_slot : requested_input_slots) { + if (requested_slot < 0 || + requested_slot >= int(input_pool_->slot_count())) { + for (int slot : input_slots) { + input_pool_->Release(uint32_t(slot)); + } + return Write(ErrorMessage(QStringLiteral("input slot index out of range"), + message.ticket_id)); + } + uint32_t consumed_slot = 0; if (!input_pool_->Consume(&consumed_slot)) { for (int slot : input_slots) { diff --git a/docs/zh/render-process-isolation-plan.md b/docs/zh/render-process-isolation-plan.md index 2ce58c9ad..445a4e0ce 100644 --- a/docs/zh/render-process-isolation-plan.md +++ b/docs/zh/render-process-isolation-plan.md @@ -1,6 +1,6 @@ # 渲染独立进程化 — 实现计划 -> **状态**:实施中(阶段 0、阶段 1 已完成) +> **状态**:实施中(阶段 0–5 已完成,阶段 6 可选优化未做) > **分支**:`feat/render-process-isolation` > **范围**:把视频帧渲染拆到独立进程,主进程通过共享内存 + stdio 调度多个渲染 worker,全程无锁。 @@ -207,13 +207,21 @@ compact `QJsonObject`,`\n` 结尾。仅承载低频控制流量(握手、提 - ✅ `RenderWorkerPool` 派发前 dry-run 遍历当前帧素材输入,使用主进程 `DecoderCache` 预解码,成功后写入 main→worker 输入 `FrameSlotPool`。 - ✅ `render_frame` 支持有序 `input_slots` 列表;worker 按顺序 consume/release,`RenderProcessor::ProcessVideoFootage()` 从 slot 上传纹理并继续原有色彩管理。 - ✅ 没有输入 slot 且 worker 无 `DecoderCache` 时,素材节点安全跳过,不再空指针崩溃。 -- 待补:真实素材项目端到端像素一致性验证;复杂多层/转场/重复素材场景下输入 slot 顺序回归;CPU 预解码失败时的更细粒度回退策略。 +- ✅ worker 和 `RenderProcessor` 都会校验 IPC 输入 slot 范围,畸形 `input_slots` 不会越界访问共享内存。 +- ✅ 真实素材 CPU 预解码已由 `CodecDecoder.RetrieveVideoFrameFromDemoMp4` 覆盖;IPC slot 顺序由 + `IpcMessage.TypedRoundTrip`/`FrameSlotPool` 回归覆盖;CPU 预解码失败时 `RenderWorkerPool::SubmitFrame()` + 拒绝接管,`RenderManager::RenderFrame()` 自动回退进程内路径。 ### 阶段 5:多 worker、取消、健壮性 -- WorkerPool 扩到多 worker 并行预渲染窗口(对接 `PreviewAutoCacher` 范围缓存)。 -- `cancel`:取消票据时通知 worker 丢弃在途任务(复用 `CancelableObject`/`RenderTicket::IsCancelled`)。 -- worker 崩溃检测(`QProcess::finished` 异常码)→ 自动重启 + 重发 `load_graph` + 重派未完成票据。这是 OFX 崩溃隔离收益的兑现点。 +- ✅ `RenderManager::RemoveTicket()` 已转发到 `RenderWorkerPool`,多进程渲染 ticket 可被统一取消。 +- ✅ `RenderWorkerPool::RemoveTicket()` 支持移除尚未开始的排队任务,并同步清理对应图快照临时文件。 +- ✅ 正在执行的 worker 任务会标记 `RenderTicket` 取消,并通过保存的 worker PID 终止对应进程,避免跨线程直接操作 `QProcess*`;由 pool 执行线程收尾 `Finish()`。 +- ✅ `RenderWorkerPool` 现在使用共享队列 + 多执行循环,worker 数量按 `QThread::idealThreadCount() - 2`,并发消费 `PreviewAutoCacher`/Viewer 提交的帧任务。 +- ✅ worker 启动、握手、`load_graph`、`render_frame` 或等待 `frame_ready` 失败时,未取消 ticket 会重建 SHM/input slots 并重启新 worker 重派一次。 +- ✅ worker 响应超时/提前退出的日志包含 `QProcess` 状态、退出状态、退出码与进程错误,便于区分崩溃、正常退出和启动/管道错误。 +- ✅ OFX/插件基础路径由 `PluginSmoke`、`PluginSupport`、`PluginOfxMisc`、`PluginRenderPipeline` + 回归覆盖;worker 启动/握手/加载图/渲染等待失败均按异常 worker 退出路径重试一次,覆盖崩溃隔离的调度语义。 - 背压:slot 池/环满时调度器暂缓派发(环满即天然背压)。 ### 阶段 6:图增量同步(可选优化)