render: the M3 OFX host — one oak-worker --ofx-host process for every plugin job

docs/zh/plans/render-pipeline-threads.md M3 (design 3.2): OpenFX crash
isolation moves from "every worker hosts plugins" to a single dedicated
host process, served over NDJSON + shared memory.

- oak-worker --ofx-host mode (src/ofx_host.rs): loads every plugin once,
  resolves jobs by the cross-process-stable OFX identifier, and renders
  through the same in-process executor the workers used to install.
- oak-render/ofxhost.rs: the single-host client. The render manager
  creates and installs it for the Pipeline backend (lazy spawn on the
  first plugin job); eval::process_plugin_job prefers it and falls back
  to the in-process executor otherwise, so the process backend keeps its
  current behavior until M4.
- Data plane: input/output FrameSlotPool pairs (the handshake's input_*
  fields are used for the first time). Named clips and the source frame
  are written to input slots after the explicit CPU readback; the plugin
  output returns through an output slot. Pool size/capacity grow by a
  host restart when a job needs more (safe: submissions are serialized
  and one job is in flight).
- Crash loop: reader EOF fails the in-flight submit, which respawns the
  host and re-posts the same job (frames are read back once); after three
  consecutive crashes the client is permanently dead and the evaluator
  falls back to a purple frame. The dead child is reaped immediately, and
  a submit mutex enforces the one-job-in-flight contract.
- Progress/cancel: the host flushes plugin_progress immediately (live
  progress), and reads stdin on its own thread so plugin_cancel takes
  effect mid-render at the plugin's next progressUpdate; the sticky flag
  resets at progressStart and request_plugin_cancel_all broadcasts to
  both the worker pool and the host.
- JobSpec::Plugin / PluginJobPayload carry the plugin type_id (stable
  across processes); `--ofx-crash-once` / `--ofx-crash-always` are the
  deterministic crash hooks, matching the worker's env hooks.
- Tests: wire round-trips; host unit tests (crash budget, cancel-flag
  reset through the factory, source mapping); oak-worker integration
  tests against the real host + bundled test plugin (render + progress,
  crash respawn and re-post, three-crash give-up, mid-render cancel on
  the new slow variant, concurrent submits); eval's purple fallback.
This commit is contained in:
2026-09-12 23:10:43 +08:00
parent 4337559ed0
commit fec6e9dba7
14 changed files with 2054 additions and 27 deletions
+82 -19
View File
@@ -142,6 +142,10 @@ pub enum JobSpec {
Plugin {
/// Plugin instance identity (oakplugin instance registry key).
instance: u64,
/// OFX plugin identifier (cross-process stable). The single OFX
/// host process (M3) resolves its own instance from this; the
/// in-process executor ignores it and uses `instance`.
type_id: String,
/// Request time in seconds (C++ `PluginJob` time).
time: f64,
/// Clip name the main source texture arrives on (C++
@@ -401,13 +405,16 @@ impl RenderEvalHooks {
Ok(())
}
/// C++ process_plugin_job: dispatch through the installed plugin
/// executor (the oakplugin render driver; dependency inversion). A
/// missing executor or a failed render yields a purple failure frame
/// instead of aborting the graph, matching pluginjob.cpp's fallback.
/// C++ process_plugin_job: dispatch through the single OFX host
/// process when one is installed (M3), else through the in-process
/// plugin executor (the oakplugin render driver; dependency
/// inversion). A missing host/executor or a failed render yields a
/// purple failure frame instead of aborting the graph, matching
/// pluginjob.cpp's fallback.
fn process_plugin_job(&mut self, src: Texture, spec: &JobSpec) -> Result<Texture> {
let JobSpec::Plugin {
instance,
type_id,
time,
effect_input_id,
inputs,
@@ -417,10 +424,21 @@ impl RenderEvalHooks {
return Err(Error::Invalid);
};
let size = src.size();
if let Some(host) = crate::ofxhost::client() {
return match host.submit(spec, &src) {
Ok(frame) => Ok(Texture::wrap_frame(frame)),
Err(err) => {
eprintln!(
"OFX host job {type_id} (instance {instance}) at t={time}s failed: {err:#}"
);
Ok(purple_frame(Rational::from_double(*time), size))
}
};
}
let Some(executor) = plugin_executor() else {
return Ok(purple_frame(Rational::from_double(*time), size));
};
let _ = (instance, effect_input_id, inputs, values);
let _ = (instance, type_id, effect_input_id, inputs, values);
match executor(&PluginJobRequest { spec, src }) {
Ok(texture) => Ok(texture),
Err(err) => {
@@ -602,6 +620,7 @@ impl RenderEvalHooks {
let spec = JobSpec::Plugin {
instance: payload.instance.0,
type_id: payload.type_id.clone(),
time: payload.time.to_f64(),
effect_input_id: if payload.effect_input_id.is_empty() {
None
@@ -2867,23 +2886,31 @@ fn apply_montage_effect(
// plugin identifier; the rendering process resolves it to a live
// instance through the oakplugin-installed factory (lazily created
// and cached per identifier), then dispatches through the plugin
// executor exactly like the graph path's plugin jobs.
let Some(factory) = plugin_instance_factory() else {
warn_unsupported_once(
&effect.type_id,
"no plugin instance factory installed (oakplugin init missing in this process)",
);
return src;
};
let Some(instance) = factory(&effect.type_id) else {
warn_unsupported_once(
&effect.type_id,
"no evaluator: unknown built-in effect or OFX plugin unavailable in this process",
);
return src;
// executor exactly like the graph path's plugin jobs. When the single
// OFX host client is installed the local instance is not needed: the
// host resolves the identifier itself.
let instance = if crate::ofxhost::client().is_some() {
0
} else {
let Some(factory) = plugin_instance_factory() else {
warn_unsupported_once(
&effect.type_id,
"no plugin instance factory installed (oakplugin init missing in this process)",
);
return src;
};
let Some(instance) = factory(&effect.type_id) else {
warn_unsupported_once(
&effect.type_id,
"no evaluator: unknown built-in effect or OFX plugin unavailable in this process",
);
return src;
};
instance
};
let spec = JobSpec::Plugin {
instance,
type_id: effect.type_id.clone(),
time: time.to_f64(),
effect_input_id: effect.effect_input_id.clone(),
inputs: Vec::new(),
@@ -3113,6 +3140,7 @@ mod tests {
fn plugin_spec() -> JobSpec {
JobSpec::Plugin {
instance: 7,
type_id: "org.oak.test-plugin".into(),
time: 0.5,
effect_input_id: Some("Source".into()),
inputs: Vec::new(),
@@ -3134,6 +3162,7 @@ mod tests {
#[test]
fn plugin_job_without_executor_yields_purple_frame() {
let _guard = PLUGIN_TEST_LOCK.lock().unwrap();
crate::ofxhost::install_client(None);
set_plugin_executor(None);
let mut hooks = RenderEvalHooks::new();
let src = Texture::wrap_frame(
@@ -3147,6 +3176,7 @@ mod tests {
#[test]
fn plugin_job_executor_error_falls_back_to_purple() {
let _guard = PLUGIN_TEST_LOCK.lock().unwrap();
crate::ofxhost::install_client(None);
set_plugin_executor(Some(Arc::new(|_req: &PluginJobRequest<'_>| {
Err(Error::Failed("boom".into()))
})));
@@ -3159,9 +3189,40 @@ mod tests {
set_plugin_executor(None);
}
/// M3 acceptance: a host that cannot come up (or crashes past its
/// budget) yields the purple failure frame, exactly like a failing
/// in-process executor.
#[test]
#[cfg(unix)]
fn plugin_job_host_failure_yields_purple_frame() {
let _guard = PLUGIN_TEST_LOCK.lock().unwrap();
// An executor that would succeed, to prove the host result wins.
set_plugin_executor(Some(Arc::new(|_req: &PluginJobRequest<'_>| {
Ok(Texture::wrap_frame(
generate_frame(Rational::new(0, 1), (2, 2), PixelFormat::F32).unwrap(),
))
})));
let host = crate::ofxhost::OfxHost::new(crate::ofxhost::OfxHostConfig {
host_bin: Some(std::path::PathBuf::from("/bin/false")),
max_failures: 1,
..Default::default()
})
.unwrap();
crate::ofxhost::install_client(Some(host));
let mut hooks = RenderEvalHooks::new();
let src = Texture::wrap_frame(
generate_frame(Rational::new(0, 1), (2, 2), PixelFormat::F32).unwrap(),
);
let out = hooks.process_plugin_job(src, &plugin_spec()).unwrap();
assert_eq!(first_pixel(&out), [1.0, 0.0, 1.0, 1.0]);
crate::ofxhost::install_client(None);
set_plugin_executor(None);
}
#[test]
fn plugin_job_dispatches_through_installed_executor() {
let _guard = PLUGIN_TEST_LOCK.lock().unwrap();
crate::ofxhost::install_client(None);
set_plugin_executor(Some(Arc::new(|req: &PluginJobRequest<'_>| {
// Echo: paint the source size with the instance id.
let JobSpec::Plugin { instance, .. } = req.spec else {
@@ -3190,6 +3251,7 @@ mod tests {
use oak_node::nodes::plugin::{PluginInstanceHandle, PluginJobPayload};
let _guard = PLUGIN_TEST_LOCK.lock().unwrap();
crate::ofxhost::install_client(None);
set_plugin_executor(Some(Arc::new(|req: &PluginJobRequest<'_>| {
let JobSpec::Plugin {
instance,
@@ -3227,6 +3289,7 @@ mod tests {
values.insert("gain".into(), NodeValue::Float(0.25));
let payload = PluginJobPayload {
instance: PluginInstanceHandle(7),
type_id: "org.oak.test-plugin".into(),
time: Rational::new(1, 2),
effect_input_id: "Source".into(),
values,
+163
View File
@@ -132,6 +132,13 @@ pub const TYPE_PLUGIN_PROGRESS: &str = "plugin_progress";
/// progress reporter then answers false (the plugin aborts at its next
/// progressUpdate).
pub const TYPE_PLUGIN_CANCEL: &str = "plugin_cancel";
/// `"ofx_job"` (M3): main->ofx-host — one OpenFX `PluginJob` to render
/// against frames in the dedicated input pool; the result returns through
/// the output pool with [`TYPE_OFX_RESULT`].
pub const TYPE_OFX_JOB: &str = "ofx_job";
/// `"ofx_result"` (M3): ofx-host->main — the rendered job's output-pool
/// slot, or `slot: -1` plus an error string.
pub const TYPE_OFX_RESULT: &str = "ofx_result";
/// Wire-format slot format for 8-bit BGRA frames (M15 S1). The viewer
/// preview path requests BGRA8 so the worker converts its F32 pipeline
@@ -652,6 +659,85 @@ pub fn plugin_cancel_json() -> Value {
json!({ "type": TYPE_PLUGIN_CANCEL })
}
/// One clip input on an [`OfxJobMsg`]: the OFX clip name and the
/// input-pool slot holding the frame.
#[derive(Serialize, Deserialize, Default, Debug, Clone, PartialEq)]
#[serde(default)]
pub struct OfxInputRef {
/// OFX clip name (`"Source"`, …).
pub name: String,
/// Input-pool slot index.
pub slot: u32,
}
/// `ofx_job` (main->ofx-host, M3) — one OpenFX `PluginJob` render request.
///
/// The host owns the plugin instances (resolved from `type_id` through the
/// process-local instance factory) and the frames live in the dedicated
/// input/output [`FrameSlotPool`]s announced by the handshake. `src_slot`
/// is the input-pool slot the effect's main source arrived on (`None` when
/// the job has no source).
#[derive(Serialize, Deserialize, Default, Debug, Clone, PartialEq)]
#[serde(default)]
pub struct OfxJobMsg {
/// Client-assigned job id (the result echoes it).
pub job: u64,
/// OFX plugin identifier (`org.oak.test-plugin`, …).
pub type_id: String,
/// Request time in seconds.
pub time: f64,
/// Effect input id the main source arrives on (`""` = none).
pub effect_input_id: String,
/// Clip inputs by name.
pub inputs: Vec<OfxInputRef>,
/// Param overrides (input id -> value).
pub values: Vec<WireEffectParam>,
/// Input-pool slot of the main source frame.
pub src_slot: Option<u32>,
}
impl OfxJobMsg {
/// The wire `ofx_job` value.
pub fn to_json(&self) -> Value {
json!({
"type": TYPE_OFX_JOB,
"job": self.job,
"type_id": self.type_id,
"time": self.time,
"effect_input_id": self.effect_input_id,
"inputs": self.inputs,
"values": self.values,
"src_slot": self.src_slot,
})
}
}
/// `ofx_result` (ofx-host->main, M3) — the job's output-pool slot on
/// success, or `slot: -1` with the failure reason (the client then reports
/// the job as failed and the evaluator falls back to a purple frame).
#[derive(Serialize, Deserialize, Default, Debug, Clone, PartialEq)]
#[serde(default)]
pub struct OfxResultMsg {
/// The job id being answered.
pub job: u64,
/// Output-pool slot, or -1 on failure.
pub slot: i32,
/// Failure reason (empty on success).
pub error: String,
}
impl OfxResultMsg {
/// The wire `ofx_result` value.
pub fn to_json(&self) -> Value {
json!({
"type": TYPE_OFX_RESULT,
"job": self.job,
"slot": self.slot,
"error": self.error,
})
}
}
/// Build a worker-side error report, mirroring `error_message()` in
/// worker.cpp: `{"type":"error","message":...}` plus `"ticket"` when
/// non-zero.
@@ -1938,6 +2024,83 @@ mod tests {
assert_eq!(value["type"], TYPE_PLUGIN_CANCEL);
}
// ---- M3 OFX-host job protocol ---------------------------------------
#[test]
fn ofx_job_wire_roundtrip() {
let job = OfxJobMsg {
job: 7,
type_id: "org.oak.test-plugin".into(),
time: 0.5,
effect_input_id: "Source".into(),
inputs: vec![
OfxInputRef {
name: "Source".into(),
slot: 3,
},
OfxInputRef {
name: "Background".into(),
slot: 4,
},
],
values: vec![WireEffectParam {
input: "gain".into(),
value: WireNodeValue::Float(0.25),
}],
src_slot: Some(3),
};
let value = job.to_json();
assert_eq!(value["type"], TYPE_OFX_JOB);
assert_eq!(value["job"], 7);
assert_eq!(value["type_id"], "org.oak.test-plugin");
assert_eq!(value["inputs"][0]["name"], "Source");
assert_eq!(value["inputs"][0]["slot"], 3);
assert_eq!(value["values"][0]["input"], "gain");
assert_eq!(value["src_slot"], 3);
// The typed form parses the same line (the type tag is ignored by
// serde's default unknown-field handling).
let parsed: OfxJobMsg = serde_json::from_value(value).unwrap();
assert_eq!(parsed, job);
// A src-less job omits `src_slot` (defaults to None).
let bare: OfxJobMsg = serde_json::from_str(r#"{"job":1,"type_id":"x"}"#).unwrap();
assert_eq!(bare.src_slot, None);
assert!(bare.inputs.is_empty());
assert!(bare.values.is_empty());
}
#[test]
fn ofx_result_wire_roundtrip() {
let ok = OfxResultMsg {
job: 7,
slot: 2,
error: String::new(),
};
let value = ok.to_json();
assert_eq!(value["type"], TYPE_OFX_RESULT);
assert_eq!(value["job"], 7);
assert_eq!(value["slot"], 2);
let parsed: OfxResultMsg = serde_json::from_value(value).unwrap();
assert_eq!(parsed, ok);
let failed = OfxResultMsg {
job: 8,
slot: -1,
error: "plugin not found".into(),
};
let parsed: OfxResultMsg = serde_json::from_value(failed.to_json()).unwrap();
assert_eq!(parsed, failed);
}
#[test]
fn ofx_protocol_type_constants_are_stable() {
assert_eq!(TYPE_OFX_JOB, "ofx_job");
assert_eq!(TYPE_OFX_RESULT, "ofx_result");
assert_eq!(TYPE_PLUGIN_PROGRESS, "plugin_progress");
assert_eq!(TYPE_PLUGIN_CANCEL, "plugin_cancel");
}
#[test]
fn render_frame_parse_accepts_cpp_field_names() {
let json = r#"{"type":"render_frame","ticket":42,"node":"abcd","time_num":1,"time_den":24,"width":1920,"height":1080,"format":-1,"channels":0,"mode":0,"input_slot":-1,"input_slots":[],"has_color_transform":false,"color_output":"","color_view":"","color_look":""}"#;
+1
View File
@@ -53,6 +53,7 @@ pub mod frameio;
pub mod handle;
pub mod ipc;
pub mod manager;
pub mod ofxhost;
pub mod pipeline;
pub mod procpool;
pub mod scheduler;
+35 -4
View File
@@ -117,6 +117,10 @@ pub struct RenderManager {
/// the inline and process backends). Holding the `Arc` keeps the
/// render/decode threads alive for the manager's lifetime.
pipeline: Option<Arc<crate::pipeline::PipelineBackend>>,
/// The single OpenFX host client (M3), when the Pipeline backend is
/// live: the render thread's plugin jobs go through it instead of the
/// in-process executor (design §3.2).
ofx_host: Option<Arc<crate::ofxhost::OfxHost>>,
}
impl RenderManager {
@@ -171,18 +175,19 @@ impl RenderManager {
eval::render_produced_frame(time, params)
.map(crate::ticket::TicketPayload::Video)
});
let (dispatch, audio_dispatch, audio_fallback, process_pool, pipeline): (
let (dispatch, audio_dispatch, audio_fallback, process_pool, pipeline, ofx_host): (
Arc<dyn JobDispatch>,
Arc<dyn JobDispatch>,
Option<Arc<dyn JobDispatch>>,
Option<Arc<ProcessDispatcher>>,
Option<Arc<crate::pipeline::PipelineBackend>>,
Option<Arc<crate::ofxhost::OfxHost>>,
) = match choice {
RenderBackendChoice::Threads => {
// Test-only inline backend: synchronous execution on the
// calling thread, shared by video and audio.
let inline = InlineDispatcher::sync();
(inline.clone(), inline, None, None, None)
(inline.clone(), inline, None, None, None, None)
}
RenderBackendChoice::Processes(config) => {
let dispatcher = ProcessDispatcher::new(config)?;
@@ -200,16 +205,29 @@ impl RenderManager {
// audio-side plugin crash now takes down the main process,
// and the mix cost lands on the UI tick.
let inline = InlineDispatcher::sync();
(dispatcher.clone(), inline, None, Some(dispatcher), None)
(dispatcher.clone(), inline, None, Some(dispatcher), None, None)
}
RenderBackendChoice::Pipeline => {
// M1 thread pipeline: the render thread executes the same
// producer (so graph mode included), and the decode thread
// the producer's footage decodes rendezvous with. Audio
// stays inline, as on the process backend.
//
// M3: one OFX host process serves every plugin job the
// in-process render thread evaluates (isolation without a
// pool). The client starts lazily on the first plugin job.
let host = crate::ofxhost::OfxHost::new(crate::ofxhost::OfxHostConfig::default())?;
crate::ofxhost::install_client(Some(host.clone()));
let pipeline = crate::pipeline::PipelineBackend::new()?;
let inline = InlineDispatcher::sync();
(pipeline.clone(), inline, None, None, Some(pipeline))
(
pipeline.clone(),
inline,
None,
None,
Some(pipeline),
Some(host),
)
}
};
let tickets = Arc::new(TicketArena::new_with_audio_fallback(
@@ -233,6 +251,7 @@ impl RenderManager {
stopping: AtomicBool::new(false),
process_pool,
pipeline,
ofx_host,
}));
Ok(())
}
@@ -252,6 +271,12 @@ impl RenderManager {
self.pipeline.clone()
}
/// The single OpenFX host client, when the Pipeline backend is live
/// (tests/observability). `None` on the inline and process backends.
pub fn ofx_host(&self) -> Option<Arc<crate::ofxhost::OfxHost>> {
self.ofx_host.clone()
}
/// Global access; `None` before init.
pub fn global() -> Option<Arc<RenderManager>> {
lock(&MANAGER).clone()
@@ -276,6 +301,12 @@ impl RenderManager {
// dispatcher after it is drained.
manager.stopping.store(true, Ordering::Release);
manager.tickets.cancel_all();
// Stop routing plugin jobs to the OFX host before tearing it
// down (an in-flight submit fails out to a purple frame).
if let Some(host) = &manager.ofx_host {
crate::ofxhost::install_client(None);
host.shutdown();
}
// Drain after the cancels so queued completions fire. Both
// dispatches are idempotent.
manager.dispatch.shutdown();
+829
View File
@@ -0,0 +1,829 @@
// Oak Video Editor - Non-Linear Video Editor
// Copyright (C) 2026 Oak Team
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU General Public License for more details.
//
// You should have received a copy of the GNU General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
//! The single OpenFX host client (M3, design §3.2).
//!
//! One `oak-worker --ofx-host` child process hosts every OFX plugin
//! instance; the render eval thread submits [`JobSpec::Plugin`] jobs to it
//! over the NDJSON control plane and the shm frame-slot data plane:
//!
//! - Input textures are read back to CPU and written into the host's
//! **input pool** (the explicit CPU boundary of the plugin path);
//! the host renders and publishes its output frame into the **output
//! pool**, which this client copies out and re-uploads into the
//! evaluator's texture value.
//! - `plugin_progress` lines from the host are forwarded to the
//! process-wide progress callback (`procpool::set_plugin_progress_cb`);
//! `plugin_cancel` is sent on [`OfxHost::cancel`].
//! - A dying host (EOF on stdout) fails the in-flight submit, which
//! respawns the host and **re-posts the same job**; after
//! [`OfxHostConfig::max_failures`] consecutive crashes the client stays
//! permanently dead and the evaluator falls back to a purple frame.
//!
//! The client is single-slot by design (design §3.2): submissions are
//! synchronous, one job in flight, because the evaluator itself is
//! synchronous. The process-wide client is installed by the render
//! manager when the thread pipeline is selected ([`install_client`]);
//! with no client installed the evaluator keeps using the in-process
//! oakplugin executor.
use std::collections::HashMap;
use std::io::{BufRead, BufReader, Write};
use std::path::PathBuf;
use std::process::{Child, ChildStdin, Command, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock};
use std::thread::JoinHandle;
use oak_core::frame::VideoParamsPod;
use oak_core::texture::{Frame, Texture};
use oak_core::PixelFormat;
use crate::error::{Error, Result};
use crate::eval::JobSpec;
use crate::ipc::{
plugin_cancel_json, write_message, HandshakeMsg, OfxInputRef, OfxJobMsg, OfxResultMsg,
PluginProgressMsg, WireEffectParam, WireNodeValue, TYPE_ERROR, TYPE_OFX_RESULT,
TYPE_PLUGIN_PROGRESS, TYPE_SHUTDOWN,
};
use crate::procpool::{plugin_progress_cb, ShmRegionView};
/// Slots per direction (input/output pools). Two slots per direction is
/// the floor; the client restarts the host with a bigger pool when a job
/// needs more (plugins with many clip inputs), so the idle cost stays
/// small.
pub const HOST_SLOTS: u32 = 2;
/// Initial per-slot capacity. Plugin frames are render-size F32; a 1080p
/// frame is ~8 MiB, and the pool grows (with a host restart) when a job
/// needs more. Small enough that an idle host costs little memory.
pub const HOST_SLOT_BYTES: usize = 8 * 1024 * 1024;
/// Consecutive crash budget before the client is permanently dead
/// (acceptance: three consecutive crashes fall back to purple frames).
pub const HOST_MAX_FAILURES: u32 = 3;
/// Configuration for [`OfxHost`].
#[derive(Clone, Debug)]
pub struct OfxHostConfig {
/// Host executable; `None` resolves `OAK_WORKER_BIN` / the
/// `oak-worker` binary next to the current executable.
pub host_bin: Option<PathBuf>,
/// Slots per direction.
pub slots: u32,
/// Initial per-slot capacity in bytes.
pub slot_bytes: usize,
/// Consecutive crash budget (a successful job resets it).
pub max_failures: u32,
/// Extra host argv (test crash hooks).
pub host_args: Vec<String>,
/// Extra host environment (test plugin discovery).
pub env: Vec<(String, String)>,
}
impl Default for OfxHostConfig {
fn default() -> Self {
Self {
host_bin: None,
slots: HOST_SLOTS,
slot_bytes: HOST_SLOT_BYTES,
max_failures: HOST_MAX_FAILURES,
host_args: Vec::new(),
env: Vec::new(),
}
}
}
fn lock<T>(m: &Mutex<T>) -> MutexGuard<'_, T> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
/// Why a submit did not produce a frame: the host died (retryable) or it
/// reported a render failure (not a crash).
enum SubmitError {
/// EOF/pipe error: the host is gone and the job can be re-posted.
Died(String),
/// A real plugin/render failure (the evaluator falls back).
Failed(String),
}
struct HostInner {
child: Option<Child>,
stdin: Option<ChildStdin>,
input: Option<Arc<ShmRegionView>>,
output: Option<Arc<ShmRegionView>>,
/// Current per-slot capacity (grows on demand).
slot_bytes: usize,
/// Current slots per direction (grows on demand).
slots: u32,
generation: u64,
consecutive_failures: u32,
permanently_dead: bool,
reader: Option<JoinHandle<()>>,
/// Total host respawns (tests/observability).
restarts: u64,
}
struct WaitState {
results: HashMap<u64, std::result::Result<Frame, String>>,
dead: bool,
}
/// The single-host client.
pub struct OfxHost {
config: OfxHostConfig,
inner: Mutex<HostInner>,
wait: Arc<(Mutex<WaitState>, Condvar)>,
next_job: AtomicU64,
/// Serializes submissions: the one-job-in-flight contract is enforced
/// here, not left to the render thread being the only caller (a
/// concurrent second submit would interleave input-pool writes).
submit_lock: Mutex<()>,
}
/// One job's pre-read frames (read once, re-posted unchanged across
/// crashes).
struct HostJobFrames {
type_id: String,
time: f64,
effect_input_id: String,
inputs: Vec<(String, Frame)>,
/// The main source frame (sent separately from the named inputs; the
/// evaluator's `src` may not correspond to any named clip).
src: Option<Frame>,
values: Vec<WireEffectParam>,
}
impl HostJobFrames {
fn max_bytes(&self) -> usize {
self.inputs
.iter()
.map(|(_, f)| f.data.len())
.chain(self.src.iter().map(|f| f.data.len()))
.max()
.unwrap_or(0)
}
/// Input-pool slots the job needs (named clips + the source).
fn input_slots(&self) -> u32 {
self.inputs.len() as u32 + u32::from(self.src.is_some())
}
}
impl OfxHost {
/// A client for `config`; the host process starts lazily on the first
/// submit (or via [`OfxHost::start`]).
pub fn new(config: OfxHostConfig) -> Result<Arc<Self>> {
Ok(Arc::new(Self {
config,
inner: Mutex::new(HostInner {
child: None,
stdin: None,
input: None,
output: None,
slot_bytes: 0,
slots: 0,
generation: 0,
consecutive_failures: 0,
permanently_dead: false,
reader: None,
restarts: 0,
}),
wait: Arc::new((
Mutex::new(WaitState {
results: HashMap::new(),
dead: false,
}),
Condvar::new(),
)),
next_job: AtomicU64::new(1),
submit_lock: Mutex::new(()),
}))
}
/// Start (or restart) the host now instead of on the first submit.
pub fn start(&self) -> Result<()> {
let mut inner = lock(&self.inner);
Self::ensure_started(&mut inner, &self.config, &self.wait)
}
/// True once the crash budget is exhausted.
pub fn is_permanently_dead(&self) -> bool {
lock(&self.inner).permanently_dead
}
/// Consecutive crashes since the last successful job.
pub fn failures(&self) -> u32 {
lock(&self.inner).consecutive_failures
}
/// Total host respawns (tests).
pub fn restarts(&self) -> u64 {
lock(&self.inner).restarts
}
/// Whether a host process is currently running.
pub fn is_running(&self) -> bool {
lock(&self.inner).child.is_some()
}
/// Submit one plugin job for `spec` (the evaluator's
/// [`JobSpec::Plugin`]) against `src`, returning the rendered frame.
/// Transparently re-posts the job across host crashes until the
/// crash budget is exhausted.
pub fn submit(&self, spec: &JobSpec, src: &Texture) -> Result<Frame> {
let _submit = lock(&self.submit_lock);
let frames = Self::read_job_frames(spec, src)?;
let job = self.next_job.fetch_add(1, Ordering::Relaxed);
loop {
match self.attempt(job, &frames) {
Ok(frame) => {
lock(&self.inner).consecutive_failures = 0;
return Ok(frame);
}
Err(SubmitError::Failed(err)) => return Err(Error::Failed(err)),
Err(SubmitError::Died(why)) => {
let mut inner = lock(&self.inner);
// Reap immediately: the crash budget may be exhausted
// here, and a dead child must not linger (is_running
// stays truthful) until shutdown.
Self::reap(&mut inner);
inner.consecutive_failures += 1;
if inner.consecutive_failures >= self.config.max_failures {
inner.permanently_dead = true;
return Err(Error::Failed(format!(
"OFX host crashed {} times in a row ({why}); giving up",
inner.consecutive_failures
)));
}
lock(&self.wait.0).dead = false;
// Loop: ensure_started spawns a fresh host and the job
// is sent again unchanged.
}
}
}
}
/// Send `plugin_cancel` to the host (sticky flag; the host's progress
/// reporter then answers false at the plugin's next progressUpdate).
pub fn cancel(&self) {
let mut inner = lock(&self.inner);
if let Some(stdin) = inner.stdin.as_mut() {
let _ = write_message(stdin, &plugin_cancel_json());
let _ = stdin.flush();
}
}
/// Stop the host (idempotent). The client is permanently dead
/// afterwards, mirroring the evaluation fallback.
pub fn shutdown(&self) {
let mut inner = lock(&self.inner);
if let Some(stdin) = inner.stdin.as_mut() {
let _ = write_message(stdin, &serde_json::json!({ "type": TYPE_SHUTDOWN }));
let _ = stdin.flush();
}
Self::reap(&mut inner);
inner.permanently_dead = true;
}
/// One attempt: ensure a live host, write the inputs, send the job and
/// wait for its result (or the host's death).
fn attempt(&self, job: u64, frames: &HostJobFrames) -> std::result::Result<Frame, SubmitError> {
let mut inner = lock(&self.inner);
if inner.permanently_dead {
return Err(SubmitError::Failed(
"OFX host is permanently dead".to_string(),
));
}
// Size the pools before the first spawn (and restart the host when
// a job outgrows them; no in-flight job exists — submits are
// synchronous).
let needed_bytes = frames.max_bytes();
let needed_slots = frames.input_slots().max(1);
if needed_bytes > inner.slot_bytes || needed_slots > inner.slots {
if inner.child.is_some() {
Self::reap(&mut inner);
}
inner.slot_bytes = needed_bytes.next_power_of_two().max(self.config.slot_bytes);
inner.slots = needed_slots.next_power_of_two().max(self.config.slots);
}
if let Err(err) = Self::ensure_started(&mut inner, &self.config, &self.wait) {
return Err(SubmitError::Died(err.to_string()));
}
let input = inner
.input
.clone()
.ok_or_else(|| SubmitError::Died("OFX host has no input pool".to_string()))?;
let stdin = inner
.stdin
.as_mut()
.ok_or_else(|| SubmitError::Died("OFX host has no stdin".to_string()))?;
let mut inputs = Vec::with_capacity(frames.inputs.len());
for (name, frame) in &frames.inputs {
let slot = match write_input_frame(&input, frame) {
Ok(slot) => slot,
Err(err) => return Err(SubmitError::Failed(err)),
};
inputs.push(OfxInputRef {
name: name.clone(),
slot,
});
}
let src_slot = match &frames.src {
Some(frame) => match write_input_frame(&input, frame) {
Ok(slot) => Some(slot),
Err(err) => return Err(SubmitError::Failed(err)),
},
None => None,
};
let msg = OfxJobMsg {
job,
type_id: frames.type_id.clone(),
time: frames.time,
effect_input_id: frames.effect_input_id.clone(),
inputs,
values: frames.values.clone(),
src_slot,
};
if let Err(err) = write_message(stdin, &msg.to_json()).and_then(|_| stdin.flush()) {
return Err(SubmitError::Died(format!("OFX host stdin: {err}")));
}
drop(inner);
let (wait_lock, cv) = &*self.wait;
let mut wait = lock(wait_lock);
loop {
if let Some(result) = wait.results.remove(&job) {
return result.map_err(SubmitError::Failed);
}
if wait.dead {
return Err(SubmitError::Died("OFX host exited".to_string()));
}
wait = cv.wait(wait).unwrap_or_else(|e| e.into_inner());
}
}
/// Read the job's textures back to CPU frames once (the explicit OFX
/// boundary) and convert the scalar params to their wire form. The
/// `src` frame is kept separate: plugin jobs may name no clip for it
/// (montage) or the evaluator's `src` may be a clone the host cannot
/// re-identify by pointer.
fn read_job_frames(spec: &JobSpec, src: &Texture) -> Result<HostJobFrames> {
let JobSpec::Plugin {
type_id,
time,
effect_input_id,
inputs,
values,
..
} = spec
else {
return Err(Error::Invalid);
};
let frames = inputs
.iter()
.map(|(name, texture)| Ok((name.clone(), texture.to_frame()?)))
.collect::<Result<Vec<_>>>()?;
// A dummy source mirrors the in-process executor's rejection (the
// host receives no src and fails the job the same way).
let src = if src.is_dummy() {
None
} else {
Some(src.to_frame()?)
};
let values = values
.iter()
.filter_map(|(input, value)| {
WireNodeValue::from_node_value(value).map(|value| WireEffectParam {
input: input.clone(),
value,
})
})
.collect();
Ok(HostJobFrames {
type_id: type_id.clone(),
time: *time,
effect_input_id: effect_input_id.clone().unwrap_or_default(),
inputs: frames,
src,
values,
})
}
/// Create the pools and spawn the host. No-op when it is already up.
fn ensure_started(
inner: &mut HostInner,
config: &OfxHostConfig,
wait: &Arc<(Mutex<WaitState>, Condvar)>,
) -> Result<()> {
if inner.permanently_dead {
return Err(Error::Failed("OFX host is permanently dead".into()));
}
if inner.child.is_some() {
return Ok(());
}
if inner.slot_bytes == 0 {
inner.slot_bytes = config.slot_bytes;
}
if inner.slots == 0 {
inner.slots = config.slots.max(1);
}
inner.generation += 1;
let generation = inner.generation;
let slots = inner.slots.max(1);
let input = ShmRegionView::create(
&format!("oak-ofx-in-{}-{generation}", std::process::id()),
slots,
inner.slot_bytes,
)?;
let output = ShmRegionView::create(
&format!("oak-ofx-out-{}-{generation}", std::process::id()),
slots,
inner.slot_bytes,
)?;
let bin = resolve_host_bin(config)?;
let mut command = Command::new(&bin);
command
.arg("--ofx-host")
.args(&config.host_args)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::inherit());
for (key, value) in &config.env {
command.env(key, value);
}
let mut child = command
.spawn()
.map_err(|e| Error::Failed(format!("spawn OFX host {}: {e}", bin.display())))?;
let mut stdin = child
.stdin
.take()
.ok_or_else(|| Error::Failed("OFX host stdin not piped".into()))?;
let stdout = child
.stdout
.take()
.ok_or_else(|| Error::Failed("OFX host stdout not piped".into()))?;
let handshake = HandshakeMsg {
protocol_version: 1,
shm_key: output.key().to_string(),
input_shm_key: input.key().to_string(),
input_slots: slots as i32,
output_slots: slots as i32,
slot_data_bytes: inner.slot_bytes as i64,
input_slot_data_bytes: inner.slot_bytes as i64,
};
write_message(&mut stdin, &handshake.to_json())
.and_then(|_| stdin.flush())
.map_err(|e| Error::Failed(format!("OFX host handshake: {e}")))?;
let reader = spawn_reader(stdout, output.clone(), wait.clone());
inner.child = Some(child);
inner.stdin = Some(stdin);
inner.input = Some(input);
inner.output = Some(output);
inner.reader = Some(reader);
Ok(())
}
/// Kill/join the current generation and drop its pools.
fn reap(inner: &mut HostInner) {
if let Some(mut child) = inner.child.take() {
let _ = child.kill();
let _ = child.wait();
}
inner.stdin = None;
if let Some(reader) = inner.reader.take() {
let _ = reader.join();
}
inner.input = None;
inner.output = None;
inner.restarts += 1;
}
}
// ---------------------------------------------------------------------------
// Process-wide client slot
// ---------------------------------------------------------------------------
static CLIENT: OnceLock<Mutex<Option<Arc<OfxHost>>>> = OnceLock::new();
fn client_slot() -> &'static Mutex<Option<Arc<OfxHost>>> {
CLIENT.get_or_init(|| Mutex::new(None))
}
/// Install (or clear) the process-wide host client. Installed by the
/// render manager when the thread pipeline is selected; the evaluator
/// prefers the client over the in-process plugin executor while it is
/// installed.
pub fn install_client(client: Option<Arc<OfxHost>>) {
*lock(client_slot()) = client;
}
/// The installed client, if any.
pub fn client() -> Option<Arc<OfxHost>> {
lock(client_slot()).clone()
}
/// Broadcast `plugin_cancel` to the host (the app's cancel button; the
/// worker pool gets its own broadcast in `procpool`).
pub fn request_cancel_all() {
if let Some(host) = client() {
host.cancel();
}
}
// ---------------------------------------------------------------------------
// Helpers
// ---------------------------------------------------------------------------
/// Resolve the host executable: the configured path, `OAK_WORKER_BIN`, or
/// an `oak-worker` binary next to the current executable (the host is a
/// mode of the worker binary, not a separate target).
fn resolve_host_bin(config: &OfxHostConfig) -> Result<PathBuf> {
if let Some(path) = &config.host_bin {
return Ok(path.clone());
}
if let Ok(path) = std::env::var("OAK_WORKER_BIN") {
return Ok(PathBuf::from(path));
}
let exe = std::env::current_exe()
.map_err(|e| Error::Failed(format!("resolve oak-worker: current exe: {e}")))?;
let candidate = exe
.parent()
.ok_or_else(|| Error::Failed("resolve oak-worker: no exe parent".into()))?
.join(format!("oak-worker{}", std::env::consts::EXE_SUFFIX));
if candidate.exists() {
return Ok(candidate);
}
Err(Error::Failed(format!(
"oak-worker binary not found at {}; set OfxHostConfig::host_bin or OAK_WORKER_BIN",
candidate.display()
)))
}
/// Write one input frame into the host's input pool, returning its slot.
fn write_input_frame(region: &ShmRegionView, frame: &Frame) -> std::result::Result<u32, String> {
let pool = region.pool();
let mut slot = 0u32;
if !unsafe { pool.acquire(&mut slot) } {
return Err("OFX host input pool is full".to_string());
}
let bytes = frame.data.len();
if bytes > region.slot_data_bytes() {
unsafe {
pool.release(slot);
}
return Err(format!(
"input frame is {bytes} bytes, larger than the OFX host slot ({})",
region.slot_data_bytes()
));
}
// SAFETY: `slot` was just acquired from this live pool; the copy and
// meta fill are the producer side of the SPSC protocol.
unsafe {
std::ptr::copy_nonoverlapping(frame.data.as_ptr(), pool.slot_data(slot), bytes);
let meta = &mut *pool.meta(slot);
*meta = Default::default();
meta.width = frame.width;
meta.height = frame.height;
meta.format = PixelFormat::F32 as i32;
meta.channel_count = 4;
meta.linesize = frame.linesize_bytes() as i32;
meta.data_size = bytes as i32;
meta.time_num = frame.timestamp.numerator();
meta.time_den = frame.timestamp.denominator().max(1);
if !pool.publish(slot) {
pool.release(slot);
return Err("OFX host input publish failed".to_string());
}
}
Ok(slot)
}
/// Copy one output frame out of the host's output pool and release the
/// slot (the consumer side of the SPSC protocol).
fn read_output_frame(
region: &ShmRegionView,
slot: u32,
) -> std::result::Result<Frame, String> {
let pool = region.pool();
let mut consumed = 0u32;
if !unsafe { pool.consume(&mut consumed) } {
return Err("OFX host output slot missing".to_string());
}
if consumed != slot {
// Keep the pool consistent: recycle the unexpected entry.
unsafe {
pool.release(consumed);
}
return Err(format!(
"OFX host output slot mismatch: expected {slot}, got {consumed}"
));
}
let meta = region.meta_copy(slot);
if meta.width <= 0 || meta.height <= 0 || meta.data_size < 0 {
unsafe {
pool.release(slot);
}
return Err("OFX host returned invalid frame metadata".to_string());
}
let len = (meta.data_size as usize).min(region.slot_data_bytes());
let data = region.slot_bytes(slot)[..len].to_vec();
unsafe {
pool.release(slot);
}
let mut frame = Frame::new();
let mut pod = VideoParamsPod::default();
pod.width = meta.width;
pod.height = meta.height;
pod.format = PixelFormat::F32 as i32;
frame.set_video_params(pod);
frame.data = data;
Ok(frame)
}
/// The host's stdout reader: forwards progress, delivers results and
/// flags death for the blocked submitter.
fn spawn_reader(
stdout: std::process::ChildStdout,
output: Arc<ShmRegionView>,
wait: Arc<(Mutex<WaitState>, Condvar)>,
) -> JoinHandle<()> {
std::thread::Builder::new()
.name("oak-ofx-host-reader".into())
.spawn(move || {
let mut reader = BufReader::new(stdout);
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) | Err(_) => break,
Ok(_) => {}
}
let Ok(value) = serde_json::from_str::<serde_json::Value>(&line) else {
continue;
};
match value.get("type").and_then(|t| t.as_str()) {
Some(TYPE_OFX_RESULT) => {
let msg: OfxResultMsg = serde_json::from_value(value).unwrap_or_default();
let result = if msg.slot >= 0 {
read_output_frame(&output, msg.slot as u32)
} else {
Err(msg.error)
};
let (state, cv) = &*wait;
lock(state).results.insert(msg.job, result);
cv.notify_all();
}
Some(TYPE_PLUGIN_PROGRESS) => {
let msg: PluginProgressMsg = serde_json::from_value(value).unwrap_or_default();
if let Some(cb) = plugin_progress_cb() {
cb(msg.label, msg.message, msg.fraction);
}
}
Some(TYPE_ERROR) => {
let message = value
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("unknown OFX host error");
eprintln!("OFX host error: {message}");
}
_ => {}
}
}
let (state, cv) = &*wait;
lock(state).dead = true;
cv.notify_all();
})
.expect("spawn OFX host reader")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::eval::generate_frame;
use oak_core::Rational;
fn plugin_spec(type_id: &str) -> JobSpec {
JobSpec::Plugin {
instance: 0,
type_id: type_id.to_string(),
time: 0.0,
effect_input_id: None,
inputs: Vec::new(),
values: Vec::new(),
}
}
fn small_texture() -> Texture {
Texture::wrap_frame(
generate_frame(Rational::new(0, 1), (4, 4), PixelFormat::F32).unwrap(),
)
}
#[test]
#[cfg(unix)]
fn submit_exhausts_the_crash_budget() {
let host = OfxHost::new(OfxHostConfig {
host_bin: Some(PathBuf::from("/bin/false")),
max_failures: 3,
..Default::default()
})
.unwrap();
let spec = plugin_spec("org.oak.missing");
let src = small_texture();
assert!(
host.submit(&spec, &src).is_err(),
"a host that never comes up must fail the submit"
);
assert_eq!(host.failures(), 3, "three consecutive crashes");
assert!(host.is_permanently_dead());
// Further submissions fail fast without spawning again.
let restarts = host.restarts();
assert!(host.submit(&spec, &src).is_err());
assert_eq!(host.restarts(), restarts, "no more spawns after the budget");
}
#[test]
#[cfg(unix)]
fn shutdown_is_idempotent_and_marks_dead() {
let host = OfxHost::new(OfxHostConfig {
host_bin: Some(PathBuf::from("/bin/false")),
..Default::default()
})
.unwrap();
host.shutdown();
host.shutdown();
assert!(host.is_permanently_dead());
assert!(!host.is_running());
}
#[test]
fn read_job_frames_keeps_the_source_mapping() {
let mut source = generate_frame(Rational::new(0, 1), (2, 2), PixelFormat::F32).unwrap();
source.data[0] = 0x7F;
let source = Texture::wrap_frame(source);
// The same texture value appears under the effect input name.
let spec = JobSpec::Plugin {
instance: 0,
type_id: "org.oak.test-plugin".into(),
time: 1.0,
effect_input_id: Some("Source".into()),
inputs: vec![("Source".to_string(), source.clone())],
values: vec![(
"gain".to_string(),
oak_node::value::NodeValue::Float(0.5),
)],
};
let frames = OfxHost::read_job_frames(&spec, &source).unwrap();
assert!(frames.src.is_some(), "a real source is sent separately");
assert_eq!(frames.src.as_ref().unwrap().data[0], 0x7F);
assert_eq!(frames.effect_input_id, "Source");
assert_eq!(frames.inputs.len(), 1);
assert_eq!(frames.inputs[0].1.data[0], 0x7F);
assert_eq!(frames.input_slots(), 2, "named clip + source");
assert_eq!(
frames.values,
vec![WireEffectParam {
input: "gain".into(),
value: WireNodeValue::Float(0.5),
}]
);
}
#[test]
fn wait_state_wakes_on_result() {
// The wait protocol itself (independent of a real host).
let host = OfxHost::new(OfxHostConfig::default()).unwrap();
let (state, cv) = &*host.wait;
let mut wait = lock(state);
// No result yet and not dead: a timed wait returns empty-handed.
let (guard, timeout) = cv
.wait_timeout(wait, std::time::Duration::from_millis(10))
.unwrap();
wait = guard;
assert!(timeout.timed_out());
assert!(!wait.dead);
// A result wakes a waiter.
let frame = Frame::dummy();
wait.results.insert(7, Ok(frame));
let (guard, _) = cv
.wait_timeout(wait, std::time::Duration::from_millis(10))
.unwrap();
wait = guard;
assert!(wait.results.remove(&7).is_some());
}
}
+4 -2
View File
@@ -148,7 +148,7 @@ pub fn set_plugin_progress_cb(cb: Option<PluginProgressCb>) {
.unwrap_or_else(|e| e.into_inner()) = cb;
}
fn plugin_progress_cb() -> Option<PluginProgressCb> {
pub(crate) fn plugin_progress_cb() -> Option<PluginProgressCb> {
PLUGIN_PROGRESS_CB
.get_or_init(|| Mutex::new(None))
.lock()
@@ -178,6 +178,8 @@ pub fn request_plugin_cancel_all() {
if let Some(dispatcher) = dispatcher {
dispatcher.broadcast_plugin_cancel();
}
// M3: the single OFX host (thread pipeline) gets the same broadcast.
crate::ofxhost::request_cancel_all();
}
// ---------------------------------------------------------------------------
@@ -249,7 +251,7 @@ impl ShmRegionView {
/// Create (and initialize) a segment of `slots` x `slot_bytes` under
/// `key`. A stale segment under the same name (left by a crashed
/// previous owner) is unlinked and the create retried once.
fn create(key: &str, slots: u32, slot_bytes: usize) -> Result<Arc<ShmRegionView>> {
pub(crate) fn create(key: &str, slots: u32, slot_bytes: usize) -> Result<Arc<ShmRegionView>> {
let mut region = SharedMemoryRegion::new();
let bytes = FrameSlotPool::bytes_needed(slots, slot_bytes);
if !region.open(key, bytes, ShmMode::Create) {