From fec6e9dba70d2e867d4a3656a2e7c833faca0a9e Mon Sep 17 00:00:00 2001 From: Mike Solar Date: Sat, 12 Sep 2026 23:10:43 +0800 Subject: [PATCH] =?UTF-8?q?render:=20the=20M3=20OFX=20host=20=E2=80=94=20o?= =?UTF-8?q?ne=20oak-worker=20--ofx-host=20process=20for=20every=20plugin?= =?UTF-8?q?=20job?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- crates/oak-node/src/nodes/plugin.rs | 4 + crates/oak-plugin/cbits/oak_test_plugin.c | 101 ++- crates/oak-plugin/src/lib.rs | 21 + crates/oak-plugin/src/node_factory.rs | 2 + crates/oak-render/src/eval.rs | 101 ++- crates/oak-render/src/ipc.rs | 163 +++++ crates/oak-render/src/lib.rs | 1 + crates/oak-render/src/manager.rs | 39 +- crates/oak-render/src/ofxhost.rs | 829 ++++++++++++++++++++++ crates/oak-render/src/procpool.rs | 6 +- crates/oak-worker/src/main.rs | 6 + crates/oak-worker/src/ofx_host.rs | 517 ++++++++++++++ crates/oak-worker/tests/ofx_host.rs | 261 +++++++ docs/zh/plans/render-pipeline-threads.md | 30 +- 14 files changed, 2054 insertions(+), 27 deletions(-) create mode 100644 crates/oak-render/src/ofxhost.rs create mode 100644 crates/oak-worker/src/ofx_host.rs create mode 100644 crates/oak-worker/tests/ofx_host.rs diff --git a/crates/oak-node/src/nodes/plugin.rs b/crates/oak-node/src/nodes/plugin.rs index 07cf9d5be..9ad972047 100644 --- a/crates/oak-node/src/nodes/plugin.rs +++ b/crates/oak-node/src/nodes/plugin.rs @@ -80,6 +80,9 @@ impl PluginInstanceHandle { pub struct PluginJobPayload { /// The instance identity (oakplugin registry key). pub instance: PluginInstanceHandle, + /// The OFX plugin identifier (cross-process stable; the single OFX + /// host process resolves its own instance from it). + pub type_id: String, /// The request time (C++ `globals.time().in()`). pub time: Rational, /// The effect input id the main source texture arrives on (C++ @@ -238,6 +241,7 @@ impl NodeBehavior for PluginNode { if tex.is_some() && !self.instance.is_null() { let payload = Job::PluginJob(PluginJobPayload { instance: self.instance, + type_id: self.type_id.clone(), time, effect_input_id: core.effect_input.clone(), values: inputs.clone(), diff --git a/crates/oak-plugin/cbits/oak_test_plugin.c b/crates/oak-plugin/cbits/oak_test_plugin.c index d4ed9ca37..62a391a75 100644 --- a/crates/oak-plugin/cbits/oak_test_plugin.c +++ b/crates/oak-plugin/cbits/oak_test_plugin.c @@ -37,6 +37,8 @@ #include #ifdef _WIN32 #include +#else +#include #endif /* Windows 上插件 DLL 带自己的 CRT 环境块:宿主 Rust 侧的 @@ -373,6 +375,90 @@ static OfxStatus actionRender(OfxImageEffectHandle inst, OfxPropertySetHandle in return kOfxStatOK; } +/* ---------- 慢速变体(M3):取消 e2e 测试 ---------- */ + +/* progressUpdate 之间的睡眠(宿主取消在渲染中置位,下一次 + * progressUpdate 返回非 OK,插件立即中止)。 */ +static void slow_sleep_ms(int ms) +{ +#ifdef _WIN32 + Sleep((DWORD)ms); +#else + usleep((useconds_t)ms * 1000); +#endif +} + +/* 100 次 progressUpdate x 20ms(约 2 秒):把输出填成常量 0.25。 + * 任一次 progressUpdate 返回非 OK(宿主取消)→ kOfxStatFailed。 */ +static OfxStatus actionRenderSlow(OfxImageEffectHandle inst, OfxPropertySetHandle inArgs) +{ + double time = 0.0; + propGetDouble(inArgs, kOfxPropTime, 0, &time); + + OfxImageClipHandle clip = NULL; + OfxPropertySetHandle clipProps = NULL; + OfxStatus st = g_imageEffectSuite->clipGetHandle(inst, "Output", &clip, &clipProps); + if (st != kOfxStatOK) + return st; + + OfxPropertySetHandle image = NULL; + st = g_imageEffectSuite->clipGetImage(clip, time, NULL, &image); + if (st != kOfxStatOK) + return st; + + void *data = NULL; + int rowBytes = 0; + int bounds[4] = { 0, 0, 0, 0 }; + propGetPointer(image, kOfxImagePropData, 0, &data); + propGetInt(image, kOfxImagePropRowBytes, 0, &rowBytes); + propGetIntN(image, kOfxImagePropBounds, 4, bounds); + + if (g_progressSuite) { + g_progressSuite->progressStart((OfxImageEffectHandle)inst, "slow-render"); + for (int i = 0; i < 100; i++) { + OfxStatus ps = g_progressSuite->progressUpdate((OfxImageEffectHandle)inst, + (double)(i + 1) / 100.0); + if (ps != kOfxStatOK) { + g_progressSuite->progressEnd((OfxImageEffectHandle)inst); + return kOfxStatFailed; + } + slow_sleep_ms(20); + } + } + + int w = bounds[2] - bounds[0]; + int h = bounds[3] - bounds[1]; + if (data && w > 0 && h > 0) { + for (int y = 0; y < h; y++) { + float *row = (float *)((char *)data + (size_t)y * (size_t)rowBytes); + for (int x = 0; x < w; x++) { + row[x * 4 + 0] = 0.25f; + row[x * 4 + 1] = 0.25f; + row[x * 4 + 2] = 0.25f; + row[x * 4 + 3] = 1.0f; + } + } + } + + g_imageEffectSuite->clipReleaseImage(image); + if (g_progressSuite) { + g_progressSuite->progressEnd((OfxImageEffectHandle)inst); + } + return kOfxStatOK; +} + +/* 慢速变体的效果入口:render 走慢速路径,其余 action 委托基础实现。 */ +static OfxStatus mainEntry(const char *action, const void *handle, + OfxPropertySetHandle inArgs, OfxPropertySetHandle outArgs); + +static OfxStatus mainEntrySlow(const char *action, const void *handle, + OfxPropertySetHandle inArgs, OfxPropertySetHandle outArgs) +{ + if (strcmp(action, kOfxImageEffectActionRender) == 0) { + return actionRenderSlow((OfxImageEffectHandle)handle, inArgs); + } + return mainEntry(action, handle, inArgs, outArgs); +} /* ---------- ofxColour(M11 §4):GetOutputColourspace ---------- */ static OfxStatus actionGetOutputColourspace(OfxPropertySetHandle inArgs, @@ -827,9 +913,20 @@ static const OfxPlugin test_plugin_interact = { /* mainEntry */ mainEntryInteractEffect, }; +/* M3 取消 e2e:慢速 filter(progressUpdate x20ms 循环,取消即中止)。 */ +static const OfxPlugin test_plugin_slow = { + /* pluginApi */ kOfxImageEffectPluginApi, + /* apiVersion */ kOfxImageEffectPluginApiVersion, + /* pluginIdentifier */ "org.oak.test-plugin.slow", + /* pluginVersionMajor */ 1, + /* pluginVersionMinor */ 0, + /* setHost */ setHost, + /* mainEntry */ mainEntrySlow, +}; + OfxExport int OfxGetNumberOfPlugins(void) { - return 4; + return 5; } OfxExport OfxPlugin *OfxGetPlugin(int nth) @@ -842,5 +939,7 @@ OfxExport OfxPlugin *OfxGetPlugin(int nth) return (OfxPlugin *)&test_plugin_id; if (nth == 3) return (OfxPlugin *)&test_plugin_interact; + if (nth == 4) + return (OfxPlugin *)&test_plugin_slow; return NULL; } diff --git a/crates/oak-plugin/src/lib.rs b/crates/oak-plugin/src/lib.rs index febf6db68..897f59c64 100644 --- a/crates/oak-plugin/src/lib.rs +++ b/crates/oak-plugin/src/lib.rs @@ -72,3 +72,24 @@ pub mod property; pub mod render; pub mod render_driver; pub mod suites; + +/// The minimal test plugin shared library built by this crate's build +/// script (`$OUT_DIR/oak_test_plugin.{so,dylib}`), for integration tests +/// that exercise the real plugin path — including the single OFX host +/// process (M3) — without a system plugin. The build script always +/// compiles it, so this is `None` only if the build directory was +/// tampered with (tests treat that as a hard failure, not a skip). +pub fn bundled_test_plugin() -> Option { + let dir = std::path::Path::new(env!("OUT_DIR")); + for name in [ + "oak_test_plugin.dylib", + "oak_test_plugin.so", + "oak_test_plugin.dll", + ] { + let path = dir.join(name); + if path.exists() { + return Some(path); + } + } + None +} diff --git a/crates/oak-plugin/src/node_factory.rs b/crates/oak-plugin/src/node_factory.rs index eedd63769..f8338557f 100644 --- a/crates/oak-plugin/src/node_factory.rs +++ b/crates/oak-plugin/src/node_factory.rs @@ -837,6 +837,7 @@ fn execute_plugin_job( let oak_render::eval::JobSpec::Plugin { instance, + type_id: _, time, effect_input_id, inputs, @@ -1167,6 +1168,7 @@ mod tests { } let spec = oak_render::eval::JobSpec::Plugin { instance, + type_id: "org.oak.test-plugin".to_string(), time: 0.0, effect_input_id: Some("Source".to_string()), inputs: Vec::new(), diff --git a/crates/oak-render/src/eval.rs b/crates/oak-render/src/eval.rs index efad7d8b7..ed87a7a30 100644 --- a/crates/oak-render/src/eval.rs +++ b/crates/oak-render/src/eval.rs @@ -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 { 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, diff --git a/crates/oak-render/src/ipc.rs b/crates/oak-render/src/ipc.rs index 751070605..5012b6177 100644 --- a/crates/oak-render/src/ipc.rs +++ b/crates/oak-render/src/ipc.rs @@ -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, + /// Param overrides (input id -> value). + pub values: Vec, + /// Input-pool slot of the main source frame. + pub src_slot: Option, +} + +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":""}"#; diff --git a/crates/oak-render/src/lib.rs b/crates/oak-render/src/lib.rs index b790363cc..8fac8c997 100644 --- a/crates/oak-render/src/lib.rs +++ b/crates/oak-render/src/lib.rs @@ -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; diff --git a/crates/oak-render/src/manager.rs b/crates/oak-render/src/manager.rs index 31b9a6d06..291c21712 100644 --- a/crates/oak-render/src/manager.rs +++ b/crates/oak-render/src/manager.rs @@ -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>, + /// 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>, } 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, Arc, Option>, Option>, Option>, + Option>, ) = 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> { + self.ofx_host.clone() + } + /// Global access; `None` before init. pub fn global() -> Option> { 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(); diff --git a/crates/oak-render/src/ofxhost.rs b/crates/oak-render/src/ofxhost.rs new file mode 100644 index 000000000..83a5584e5 --- /dev/null +++ b/crates/oak-render/src/ofxhost.rs @@ -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 . + +//! 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, + /// 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, + /// 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(m: &Mutex) -> 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, + stdin: Option, + input: Option>, + output: Option>, + /// 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>, + /// Total host respawns (tests/observability). + restarts: u64, +} + +struct WaitState { + results: HashMap>, + dead: bool, +} + +/// The single-host client. +pub struct OfxHost { + config: OfxHostConfig, + inner: Mutex, + wait: Arc<(Mutex, 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, + values: Vec, +} + +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> { + 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 { + 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 { + 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 { + 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::>>()?; + // 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, 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>>> = OnceLock::new(); + +fn client_slot() -> &'static Mutex>> { + 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>) { + *lock(client_slot()) = client; +} + +/// The installed client, if any. +pub fn client() -> Option> { + 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 { + 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 { + 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 { + 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, + wait: Arc<(Mutex, 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::(&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()); + } +} diff --git a/crates/oak-render/src/procpool.rs b/crates/oak-render/src/procpool.rs index 117f38de9..f8b0c6637 100644 --- a/crates/oak-render/src/procpool.rs +++ b/crates/oak-render/src/procpool.rs @@ -148,7 +148,7 @@ pub fn set_plugin_progress_cb(cb: Option) { .unwrap_or_else(|e| e.into_inner()) = cb; } -fn plugin_progress_cb() -> Option { +pub(crate) fn plugin_progress_cb() -> Option { 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> { + pub(crate) fn create(key: &str, slots: u32, slot_bytes: usize) -> Result> { let mut region = SharedMemoryRegion::new(); let bytes = FrameSlotPool::bytes_needed(slots, slot_bytes); if !region.open(key, bytes, ShmMode::Create) { diff --git a/crates/oak-worker/src/main.rs b/crates/oak-worker/src/main.rs index f294402dc..73d308c82 100644 --- a/crates/oak-worker/src/main.rs +++ b/crates/oak-worker/src/main.rs @@ -32,6 +32,7 @@ mod framecache; mod ipc; +mod ofx_host; mod worker; use std::process::exit; @@ -64,6 +65,11 @@ fn parse_backend(args: &[String]) -> String { fn main() { let args: Vec = std::env::args().collect(); + // M3: `--ofx-host` runs the single OpenFX host loop instead of the + // ticket worker (same binary, no extra target to package). + if args.iter().any(|arg| arg == "--ofx-host") { + exit(ofx_host::ofx_host_main(&args)); + } let backend = parse_backend(&args); exit(worker::worker_main(&backend)); } diff --git a/crates/oak-worker/src/ofx_host.rs b/crates/oak-worker/src/ofx_host.rs new file mode 100644 index 000000000..aea4b169f --- /dev/null +++ b/crates/oak-worker/src/ofx_host.rs @@ -0,0 +1,517 @@ +// 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 . + +//! `oak-worker --ofx-host`: the single OpenFX host process (M3, design +//! §3.2). +//! +//! The host loads every OFX plugin once (the worker binary's normal plugin +//! runtime) and then serves `ofx_job` messages from the main process over +//! NDJSON, with frames moving through the input/output +//! [`FrameSlotPool`]s announced by the handshake. Each job is resolved to +//! a host-local instance by the plugin **identifier** (the cross-process +//! stable key) and rendered through the same in-process executor the +//! workers used to install ([`oak_plugin::node_factory::install_render_executor`]). +//! +//! Progress is flushed to stdout immediately (the main process's reader +//! forwards it to the plugin-progress dialog), and `plugin_cancel` sets +//! the same sticky flag protocol the worker used: the next `progressStart` +//! resets it and every `progressUpdate` after a cancel answers false, so +//! the plugin aborts at its next progress call. +//! +//! The host is deliberately single-threaded and synchronous, matching the +//! worker model; crash isolation comes from the parent respawning it. + +use std::io::{BufRead, BufReader, Write}; +use std::path::PathBuf; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::mpsc; +use std::sync::{Arc, Mutex, MutexGuard}; + +use oak_core::frame::VideoParamsPod; +use oak_core::texture::{Frame, Texture}; +use oak_core::PixelFormat; +use oak_render::eval::{self, JobSpec, PluginJobRequest}; +use oak_render::ipc::{ + error_message, write_message, FrameSlotPool, HandshakeMsg, OfxJobMsg, OfxResultMsg, + PluginProgressMsg, SharedMemoryRegion, ShmMode, TYPE_HANDSHAKE, TYPE_OFX_JOB, + TYPE_PLUGIN_CANCEL, TYPE_SHUTDOWN, +}; + +/// Sticky plugin cancel (protocol parity with `worker.rs`): set by +/// `plugin_cancel`, reset by the next `progressStart`. +static OFX_CANCEL: AtomicBool = AtomicBool::new(false); + +/// stdout is shared by progress events (emitted from inside the plugin +/// render) and control responses. +static OUT_LOCK: Mutex<()> = Mutex::new(()); + +fn lock(m: &Mutex) -> MutexGuard<'_, T> { + m.lock().unwrap_or_else(|e| e.into_inner()) +} + +/// Write one NDJSON line to stdout, flushed (progress must stream while +/// the job renders). No-op under unit tests (they exercise the reporter +/// return values, not the pipe). +fn emit(value: &serde_json::Value) { + if cfg!(test) { + return; + } + let _guard = lock(&OUT_LOCK); + let stdout = std::io::stdout(); + let mut out = stdout.lock(); + let _ = write_message(&mut out, value); + let _ = out.flush(); +} + +/// The host's progress reporter: emits `plugin_progress` immediately and +/// reports cancellation to the plugin. +struct HostProgressReporter { + label: String, + message: String, +} + +impl oak_plugin::progress::UiProgressReporter for HostProgressReporter { + fn update(&mut self, progress: f64) -> bool { + emit( + &PluginProgressMsg { + label: self.label.clone(), + message: self.message.clone(), + fraction: progress.clamp(0.0, 1.0), + } + .to_json(), + ); + !OFX_CANCEL.load(Ordering::Relaxed) + } + + fn end(&mut self) { + emit( + &PluginProgressMsg { + label: self.label.clone(), + message: self.message.clone(), + fraction: 1.0, + } + .to_json(), + ); + } +} + +/// Build one progress reporter for `progressStart`: clears the sticky +/// cancel (the protocol's "a fresh render starts uncancelled") and emits +/// the fraction-0 start event. The reporter factory installed with the +/// progress suite calls this; tests call it directly. +fn host_progress_reporter(label: &str, message: &str) -> Box { + OFX_CANCEL.store(false, Ordering::Relaxed); + emit( + &PluginProgressMsg { + label: label.to_string(), + message: message.to_string(), + fraction: 0.0, + } + .to_json(), + ); + Box::new(HostProgressReporter { + label: label.to_string(), + message: message.to_string(), + }) +} + +/// Install the reporter factory (`progressStart` → [`host_progress_reporter`]). +fn install_progress_factory() { + oak_plugin::progress::set_reporter_factory(Some(Arc::new(host_progress_reporter))); +} + +/// Read stdin on its own thread. `plugin_cancel` is handled inline (sets +/// the sticky flag) so a cancel is observed while a plugin render is in +/// flight; every other line goes to the main loop through the channel. +/// Dropping the sender on EOF closes the channel and ends the loop. +fn spawn_stdin_reader() -> mpsc::Receiver { + let (tx, rx) = mpsc::channel(); + std::thread::Builder::new() + .name("oak-ofx-host-stdin".into()) + .spawn(move || { + let stdin = std::io::stdin(); + let mut reader = BufReader::new(stdin.lock()); + let mut line = String::new(); + loop { + line.clear(); + match reader.read_line(&mut line) { + Ok(0) | Err(_) => break, + Ok(_) => {} + } + if line.trim().is_empty() { + continue; + } + if let Ok(value) = serde_json::from_str::(&line) { + if value.get("type").and_then(|t| t.as_str()) == Some(TYPE_PLUGIN_CANCEL) { + OFX_CANCEL.store(true, Ordering::Relaxed); + continue; + } + } + if tx.send(line.clone()).is_err() { + break; + } + } + }) + .expect("spawn OFX host stdin reader"); + rx +} + +/// The attached pools. The `SharedMemoryRegion`s must outlive the pools +/// (the pool views point into the mappings). +struct HostPools { + _input_region: SharedMemoryRegion, + input: FrameSlotPool, + _output_region: SharedMemoryRegion, + output: FrameSlotPool, +} + +/// Attach the pools announced by a `handshake` (the worker's +/// `handle_handshake`, minus the render session). +fn attach_pools(msg: &serde_json::Value) -> Result { + let hs: HandshakeMsg = + serde_json::from_value(msg.clone()).map_err(|e| format!("invalid handshake: {e}"))?; + if hs.shm_key.is_empty() || hs.output_slots <= 0 || hs.slot_data_bytes <= 0 { + return Err("handshake missing output shared-memory geometry".to_string()); + } + let output_bytes = + FrameSlotPool::bytes_needed(hs.output_slots as u32, hs.slot_data_bytes as usize); + let mut output_region = SharedMemoryRegion::new(); + if !output_region.open(&hs.shm_key, output_bytes, ShmMode::Attach) { + return Err(format!( + "failed to attach output shared memory: {}", + output_region.error() + )); + } + // SAFETY: the mapping is live and sized for the pool. + let output = unsafe { FrameSlotPool::attach(output_region.data()) }; + if !output.is_valid() { + return Err("output shared memory does not contain a frame slot pool".to_string()); + } + + if hs.input_shm_key.is_empty() || hs.input_slots <= 0 || hs.input_slot_data_bytes <= 0 { + return Err("handshake missing input shared-memory geometry".to_string()); + } + let input_bytes = + FrameSlotPool::bytes_needed(hs.input_slots as u32, hs.input_slot_data_bytes as usize); + let mut input_region = SharedMemoryRegion::new(); + if !input_region.open(&hs.input_shm_key, input_bytes, ShmMode::Attach) { + return Err(format!( + "failed to attach input shared memory: {}", + input_region.error() + )); + } + // SAFETY: the mapping is live and sized for the pool. + let input = unsafe { FrameSlotPool::attach(input_region.data()) }; + if !input.is_valid() { + return Err("input shared memory does not contain a frame slot pool".to_string()); + } + + Ok(HostPools { + _input_region: input_region, + input, + _output_region: output_region, + output, + }) +} + +/// Consume one input frame (the producer is the main process). +fn read_input_frame(pool: &FrameSlotPool, slot: u32) -> Result { + let mut consumed = 0u32; + // SAFETY: consumer side of the SPSC protocol, single-threaded host. + if !unsafe { pool.consume(&mut consumed) } { + return Err("input slot missing".to_string()); + } + if consumed != slot { + unsafe { + pool.release(consumed); + } + return Err(format!( + "input slot mismatch: expected {slot}, got {consumed}" + )); + } + // SAFETY: `slot` was just consumed; the meta POD is initialized by the + // producer before publish. + let meta = unsafe { &*pool.meta_const(slot) }; + if meta.width <= 0 || meta.height <= 0 || meta.data_size < 0 { + unsafe { + pool.release(slot); + } + return Err("input frame has invalid metadata".to_string()); + } + let len = (meta.data_size as usize).min(pool.slot_data_bytes()); + // SAFETY: `len` is within the slot block. + let data = unsafe { std::slice::from_raw_parts(pool.slot_data_const(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) +} + +/// Publish one output frame (the consumer is the main process). +fn write_output_frame(pool: &FrameSlotPool, frame: &Frame) -> Result { + let mut slot = 0u32; + // SAFETY: producer side of the SPSC protocol, single-threaded host. + if !unsafe { pool.acquire(&mut slot) } { + return Err("output pool is full".to_string()); + } + let bytes = frame.data.len(); + if bytes > pool.slot_data_bytes() { + unsafe { + pool.release(slot); + } + return Err(format!( + "output frame is {bytes} bytes, larger than the output slot ({})", + pool.slot_data_bytes() + )); + } + // SAFETY: `slot` was just acquired; fill then publish. + 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("output publish failed".to_string()); + } + } + Ok(slot as i32) +} + +/// Render one `ofx_job` and return its `ofx_result` response. +fn handle_job(msg: serde_json::Value, pools: &HostPools) -> serde_json::Value { + let job: OfxJobMsg = match serde_json::from_value(msg) { + Ok(job) => job, + Err(err) => return error_message(&format!("invalid ofx_job: {err}"), None), + }; + let fail = |message: String| { + OfxResultMsg { + job: job.job, + slot: -1, + error: message, + } + .to_json() + }; + + let mut inputs = Vec::with_capacity(job.inputs.len()); + for input in &job.inputs { + match read_input_frame(&pools.input, input.slot) { + Ok(frame) => inputs.push((input.name.clone(), Texture::wrap_frame(frame))), + Err(err) => return fail(err), + } + } + let src = match job.src_slot { + Some(slot) => match read_input_frame(&pools.input, slot) { + Ok(frame) => Texture::wrap_frame(frame), + Err(err) => return fail(err), + }, + None => { + // No explicit source: mirror the evaluator's fallback (the + // declared effect input, else the first clip, else dummy). + let named = inputs + .iter() + .find(|(name, _)| name == &job.effect_input_id) + .or_else(|| inputs.first()) + .map(|(_, texture)| texture.clone()); + named.unwrap_or_else(Texture::dummy) + } + }; + + let Some(factory) = eval::plugin_instance_factory() else { + return fail("no plugin instance factory installed in the OFX host".to_string()); + }; + let Some(instance) = factory(&job.type_id) else { + return fail(format!("unknown or unavailable OFX plugin: {}", job.type_id)); + }; + let Some(executor) = eval::plugin_executor() else { + return fail("no plugin executor installed in the OFX host".to_string()); + }; + let spec = JobSpec::Plugin { + instance, + type_id: job.type_id.clone(), + time: job.time, + effect_input_id: if job.effect_input_id.is_empty() { + None + } else { + Some(job.effect_input_id.clone()) + }, + inputs, + values: job + .values + .iter() + .map(|param| (param.input.clone(), param.value.to_node_value())) + .collect(), + }; + match executor(&PluginJobRequest { spec: &spec, src }) { + Ok(texture) => match texture.to_frame() { + Ok(frame) => match write_output_frame(&pools.output, &frame) { + Ok(slot) => OfxResultMsg { + job: job.job, + slot, + error: String::new(), + } + .to_json(), + Err(err) => fail(err), + }, + Err(err) => fail(format!("plugin output readback failed: {err:?}")), + }, + Err(err) => fail(format!("plugin render failed: {err:?}")), + } +} + +/// Test-only crash hooks (deterministic crash/restart acceptance tests). +struct CrashHooks { + /// Crash on every job. + always: bool, + /// Crash on the first job of each process unless the marker file + /// exists (create it before crashing, so a respawned host renders). + once_marker: Option, +} + +impl CrashHooks { + fn from_args(args: &[String]) -> Self { + let mut hooks = CrashHooks { + always: false, + once_marker: None, + }; + let mut i = 0usize; + while i < args.len() { + match args[i].as_str() { + "--ofx-crash-always" => hooks.always = true, + "--ofx-crash-once" if i + 1 < args.len() => { + hooks.once_marker = Some(PathBuf::from(&args[i + 1])); + i += 1; + } + _ => {} + } + i += 1; + } + hooks + } + + fn maybe_crash(&self) { + if self.always { + std::process::abort(); + } + if let Some(path) = &self.once_marker { + if !path.exists() { + let _ = std::fs::write(path, b"ofx-host-crashed"); + std::process::abort(); + } + } + } +} + +/// The `--ofx-host` main loop. Returns the process exit code. +pub fn ofx_host_main(args: &[String]) -> i32 { + // The same plugin runtime the render workers install: the executor + // (so `plugin_executor` is callable) and the identifier-keyed instance + // factory (so jobs resolve their own instances). + oak_plugin::node_factory::install_render_executor(); + if let Err(err) = oak_plugin::host::Host::global().cache.scan() { + eprintln!("ofx-host: plugin scan failed: {err}"); + } + install_progress_factory(); + let crash = CrashHooks::from_args(args); + + // stdin runs on its own thread (cancel must be observed mid-render); + // the main loop consumes the forwarded control messages. + let control = spawn_stdin_reader(); + let mut pools: Option = None; + while let Ok(line) = control.recv() { + let Ok(msg) = serde_json::from_str::(&line) else { + emit(&error_message("invalid JSON message", None)); + continue; + }; + match msg.get("type").and_then(|t| t.as_str()) { + Some(TYPE_HANDSHAKE) => match attach_pools(&msg) { + Ok(attached) => pools = Some(attached), + Err(err) => emit(&error_message(&err, None)), + }, + Some(TYPE_OFX_JOB) => match &pools { + Some(pools) => { + crash.maybe_crash(); + let response = handle_job(msg, pools); + emit(&response); + } + None => emit(&error_message("ofx_job before handshake", None)), + }, + Some(TYPE_PLUGIN_CANCEL) => OFX_CANCEL.store(true, Ordering::Relaxed), + Some(TYPE_SHUTDOWN) => break, + other => emit(&error_message( + &format!( + "unknown message type: {}", + other.unwrap_or("") + ), + None, + )), + } + } + 0 +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cancel_flag_is_reset_by_progress_start() { + // Cancel semantics (protocol parity with the worker): a cancelled + // reporter tells the plugin to abort at its next progressUpdate; + // the next progressStart (the factory path below) clears the + // sticky flag. + OFX_CANCEL.store(true, Ordering::Relaxed); + let mut reporter = host_progress_reporter("render", "msg"); + assert!( + reporter.update(0.5), + "the progressStart factory path reset the sticky cancel" + ); + // A cancel after the start makes the next update answer false. + OFX_CANCEL.store(true, Ordering::Relaxed); + assert!(!reporter.update(0.6), "a cancelled reporter answers false"); + OFX_CANCEL.store(false, Ordering::Relaxed); + } + + #[test] + fn crash_hooks_parse_args() { + let hooks = CrashHooks::from_args(&[ + "oak-worker".to_string(), + "--ofx-host".to_string(), + "--ofx-crash-once".to_string(), + "/tmp/marker".to_string(), + ]); + assert!(!hooks.always); + assert_eq!(hooks.once_marker.as_deref(), Some(std::path::Path::new("/tmp/marker"))); + + let hooks = CrashHooks::from_args(&["oak-worker".to_string(), "--ofx-crash-always".to_string()]); + assert!(hooks.always); + assert!(hooks.once_marker.is_none()); + } +} diff --git a/crates/oak-worker/tests/ofx_host.rs b/crates/oak-worker/tests/ofx_host.rs new file mode 100644 index 000000000..a93df3d04 --- /dev/null +++ b/crates/oak-worker/tests/ofx_host.rs @@ -0,0 +1,261 @@ +// 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 . + +//! M3: the single OpenFX host process (`oak-worker --ofx-host`). +//! +//! These tests spawn the real worker binary in host mode against the +//! crate's bundled test plugin and assert the milestone's acceptance +//! points: a plugin job renders and its progress events reach the app +//! callback, a crashed host is respawned and the in-flight job is +//! re-posted, and three consecutive crashes exhaust the budget (the +//! evaluator's purple fallback is asserted by the eval unit test). + +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use oak_core::texture::Texture; +use oak_core::{PixelFormat, Rational}; +use oak_render::eval::{generate_frame, JobSpec}; +use oak_render::ofxhost::{OfxHost, OfxHostConfig}; + +/// The plugin runtime and the progress callback are process-wide. +static LOCK: Mutex<()> = Mutex::new(()); + +/// Assemble the bundled test plugin into an OFX bundle and return the +/// directory that contains it (`OFX_PLUGIN_PATH`). `None` when the test +/// plugin was not built (release builds). +fn bundle_parent(tag: &str) -> Option { + let lib = oak_plugin::bundled_test_plugin()?; + let root = std::env::temp_dir().join(format!( + "oak-ofx-host-{tag}-{}", + std::process::id() + )); + let bundle = root.join("OakTest.ofx.bundle"); + let platform = if cfg!(target_os = "macos") { + "MacOS" + } else if cfg!(target_os = "windows") { + "Win64" + } else if cfg!(target_arch = "aarch64") { + "Linux-aarch64" + } else { + "Linux-x86-64" + }; + let dest = bundle.join("Contents").join(platform); + std::fs::create_dir_all(&dest).ok()?; + let name = lib.file_name()?; + std::fs::copy(&lib, dest.join(name)).ok()?; + Some(root) +} + +fn host(plugin_dir: &Path, args: Vec, max_failures: u32) -> Arc { + OfxHost::new(OfxHostConfig { + host_bin: Some(PathBuf::from(env!("CARGO_BIN_EXE_oak-worker"))), + max_failures, + host_args: args, + env: vec![( + "OFX_PLUGIN_PATH".to_string(), + plugin_dir.to_string_lossy().into_owned(), + )], + ..Default::default() + }) + .expect("host client") +} + +fn source_texture(w: i32, h: i32) -> Texture { + let mut frame = generate_frame(Rational::new(0, 1), (w, h), PixelFormat::F32).unwrap(); + for px in frame.data.chunks_exact_mut(16) { + for (c, v) in px.chunks_exact_mut(4).zip([0.1f32, 0.2, 0.3, 1.0]) { + c.copy_from_slice(&v.to_le_bytes()); + } + } + Texture::wrap_frame(frame) +} + +/// The bundled test plugin's filter (`org.oak.test-plugin`) paints a +/// constant 0.5 grey with alpha 1. +fn plugin_spec(src: &Texture) -> JobSpec { + JobSpec::Plugin { + instance: 0, + type_id: "org.oak.test-plugin".to_string(), + time: 0.0, + effect_input_id: Some("Source".to_string()), + inputs: vec![("Source".to_string(), src.clone())], + values: Vec::new(), + } +} + +fn first_pixel(frame: &oak_core::texture::Frame) -> [f32; 4] { + let mut out = [0f32; 4]; + for (i, value) in out.iter_mut().enumerate() { + *value = f32::from_le_bytes(frame.data[i * 4..i * 4 + 4].try_into().unwrap()); + } + out +} + +#[test] +fn host_renders_plugin_job_and_forwards_progress() { + let _lock = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let dir = bundle_parent("render") + .expect("the bundled test plugin must exist (build.rs compiles it)"); + let events: Arc>> = Arc::new(Mutex::new(Vec::new())); + let sink = events.clone(); + oak_render::procpool::set_plugin_progress_cb(Some(Arc::new(move |label, message, fraction| { + sink.lock() + .unwrap_or_else(|e| e.into_inner()) + .push((label, message, fraction)); + }))); + + let client = host(&dir, Vec::new(), 3); + let src = source_texture(16, 16); + let out = client + .submit(&plugin_spec(&src), &src) + .expect("the host renders the plugin job"); + assert_eq!((out.width, out.height), (16, 16)); + assert_eq!(out.format, PixelFormat::F32); + let px = first_pixel(&out); + assert!( + (px[0] - 0.5).abs() < 1e-3 && (px[3] - 1.0).abs() < 1e-6, + "the test plugin's constant grey must come back: {px:?}" + ); + + let seen = events.lock().unwrap_or_else(|e| e.into_inner()).clone(); + assert!( + seen.iter() + .any(|(label, _, fraction)| label == "render" && (*fraction - 0.5).abs() < 1e-6), + "the plugin's progressUpdate must reach the app callback: {seen:?}" + ); + assert!( + seen.iter().any(|(_, _, fraction)| (*fraction - 1.0).abs() < 1e-6), + "progressEnd closes the dialog: {seen:?}" + ); + + oak_render::procpool::set_plugin_progress_cb(None); + client.shutdown(); + let _ = std::fs::remove_dir_all(&dir); +} + +#[test] +fn host_respawns_and_reposts_after_a_crash() { + let _lock = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let dir = bundle_parent("respawn") + .expect("the bundled test plugin must exist (build.rs compiles it)"); + let marker = dir.join("crash-once.marker"); + let client = host( + &dir, + vec![ + "--ofx-crash-once".to_string(), + marker.to_string_lossy().into_owned(), + ], + 3, + ); + let src = source_texture(16, 16); + let out = client + .submit(&plugin_spec(&src), &src) + .expect("the in-flight job is re-posted to the respawned host"); + assert_eq!(out.width, 16); + assert!(client.restarts() >= 1, "the crashed host was respawned"); + assert_eq!(client.failures(), 0, "a successful re-post clears the crash count"); + assert_eq!(first_pixel(&out)[0], 0.5, "the re-posted job still renders"); + + client.shutdown(); + let _ = std::fs::remove_dir_all(&dir); +} + +#[test] +fn host_cancels_an_in_flight_slow_render() { + let _lock = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let dir = bundle_parent("cancel").expect("the bundled test plugin must exist"); + let client = host(&dir, Vec::new(), 3); + let src = source_texture(16, 16); + let spec = JobSpec::Plugin { + instance: 0, + type_id: "org.oak.test-plugin.slow".to_string(), + time: 0.0, + effect_input_id: Some("Source".to_string()), + inputs: vec![("Source".to_string(), src.clone())], + values: Vec::new(), + }; + let submitter = client.clone(); + let src_for_job = src.clone(); + let job = std::thread::spawn(move || submitter.submit(&spec, &src_for_job)); + // The slow plugin renders for ~2 s in 20 ms progress steps; cancel + // lands while it is in flight and the next progressUpdate aborts it. + std::thread::sleep(Duration::from_millis(200)); + client.cancel(); + let result = job.join().expect("the submit thread must not panic"); + assert!( + result.is_err(), + "a cancelled slow render must fail instead of returning a frame" + ); + assert_eq!( + client.failures(), + 0, + "a cancel is not a crash; the budget stays untouched" + ); + + client.shutdown(); + let _ = std::fs::remove_dir_all(&dir); +} + +#[test] +fn host_serializes_concurrent_submits() { + let _lock = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let dir = bundle_parent("concurrent").expect("the bundled test plugin must exist"); + let client = host(&dir, Vec::new(), 3); + let src = source_texture(16, 16); + + let a = client.clone(); + let spec_a = plugin_spec(&src); + let src_a = src.clone(); + let first = std::thread::spawn(move || a.submit(&spec_a, &src_a)); + let b = client.clone(); + let spec_b = plugin_spec(&src); + let src_b = src.clone(); + let second = std::thread::spawn(move || b.submit(&spec_b, &src_b)); + + // The submit lock serializes the jobs (each writes the shared input + // pool); both must come back with the plugin's constant grey. + for handle in [first, second] { + let frame = handle.join().expect("submit thread").expect("job renders"); + assert_eq!(first_pixel(&frame)[0], 0.5); + } + + client.shutdown(); + let _ = std::fs::remove_dir_all(&dir); +} + +#[test] +fn host_gives_up_after_three_consecutive_crashes() { + let _lock = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let dir = bundle_parent("giveup") + .expect("the bundled test plugin must exist (build.rs compiles it)"); + let client = host(&dir, vec!["--ofx-crash-always".to_string()], 3); + let src = source_texture(8, 8); + let err = client + .submit(&plugin_spec(&src), &src) + .expect_err("three consecutive crashes must give up"); + let _ = err; + assert_eq!(client.failures(), 3); + assert!(client.is_permanently_dead()); + // Fast-fail afterwards: no more spawns. + let restarts = client.restarts(); + assert!(client.submit(&plugin_spec(&src), &src).is_err()); + assert_eq!(client.restarts(), restarts); + + client.shutdown(); + let _ = std::fs::remove_dir_all(&dir); +} diff --git a/docs/zh/plans/render-pipeline-threads.md b/docs/zh/plans/render-pipeline-threads.md index ab8406af4..57bb1b32a 100644 --- a/docs/zh/plans/render-pipeline-threads.md +++ b/docs/zh/plans/render-pipeline-threads.md @@ -200,6 +200,34 @@ - GPU 型 OFX 插件(OpenGL/CUDA 上下文的)本期不支持直通,按 CPU 插件处理 (文档明示;真有需要时按 §3.6 的平台分支再开 GPU 句柄通道)。 +> **M3 落地回填(2026-09-12)**: +> +> - 采用 **`oak-worker --ofx-host` 模式**(新增 `oak-worker/src/ofx_host.rs`), +> 不新增可执行目标,省掉打包/路径解析的分支。宿主启动即 +> `install_render_executor` + `Host::cache.scan()`,按**插件标识符** +> (`shared_plugin_instance`,跨进程稳定)解析宿主本地实例。 +> - 客户端在 `oak-render/src/ofxhost.rs`:`RenderManager` 在 **Pipeline +> 后端**创建并安装进程级客户端(首个个 PluginJob 才 spawn,空闲不占进程), +> `eval::process_plugin_job` 优先走它,失败/缺席再回退进程内 executor(进程池 +> 后端在 M4 之前仍按现状在每个 worker 内渲染,见 §5 回退策略)。 +> 数据面为输入/输出两个 `FrameSlotPool`(handshake 的 `input_*` 字段首次启用): +> 命名 clip + src 各写一个输入槽,插件输出写输出槽;槽容量/数量按 job 需求 +> 自动扩容(重启一次宿主,无在途任务时安全)。 +> - 崩溃:reader 线程 EOF 即判死,`submit` 在同一调用内**重生宿主并重投同一 +> job**(输入帧只回读一次,重投不重复回读);连续 +> `HOST_MAX_FAILURES=3` 次崩溃后熔断,`process_plugin_job` 落紫帧 +> (`plugin_job_host_failure_yields_purple_frame` 直接断言)。 +> - 进度/取消:宿主 reporter 即时 flush `plugin_progress`(不再是 worker 的 +> 批末缓冲,进度条可实时更新);宿主用独立 stdin 线程,`plugin_cancel` +> 在该线程内直接置黏性 flag(渲染中也能生效),下次 `progressStart` 复位; +> `request_plugin_cancel_all` 同时广播给 worker 池与宿主。取消的生效粒度 +> 取决于插件是否回调 progressUpdate(与现状契约一致)。 +> - 验收测试(`oak-worker/tests/ofx_host.rs`,真实宿主 + 内置测试插件): +> 渲染与进度转发;`--ofx-crash-once` 崩溃后重投成功;`--ofx-crash-always` +> 连续三次崩溃后熔断(紫帧由 eval 单测覆盖);慢速测试插件 +> (`org.oak.test-plugin.slow`,progressUpdate x20ms)上取消在途渲染; +> 并发 submit 由客户端互斥串行化。 + ### 3.3 队列与 ticket - 对上层(oak-app/oak-cli)**ticket API 不变**:RenderManager 仍是唯一入口, @@ -386,7 +414,7 @@ fallback。** 解码上传与上屏共用一层 `gpuinteop` 抽象,按后端 | **M0b Job 图 + 虚拟端点 + BFS** | §3.8 全量:图固定 GraphInput/GraphOutput 虚拟节点(默认相连、禁删禁复制、序列化往返)、节点编辑器显示两节点、resolve 改为从输入节点的 Kahn 形态 BFS | 新增测试:多输入汇合等齐全部输入、多输出分叉各自成帧、非全连通图不可达节点不执行、环报错断支、虚拟节点删除/复制被拒、序列化往返后端点仍在;节点编辑器 UI 测试(端点可见、入线/出线规则);既有测试全绿 | | **M1 线程管线骨架** | 解码/渲染/上屏三线程+三队列进 oak-render(`pipeline` 模块);RenderManager 增加线程后端,进程池后端保留,`OAK_PIPELINE=processes` 可回退 | 同一套渲染测试在两个后端下都绿(测试矩阵化);播放/seek/导出 smoke 等价 | | **M2 GPU 零拷贝** | 图内全程 `Texture::Gpu`(合成/转场/调整层不再逐帧回读);wgpu 29 统一 + 采用 gpui device(§3.5 攻关已回填);GPU 色彩管理(工作空间→输出规格→显示器 ICC 烘焙 3D LUT,GPU 执行);导出/缓存/OFX 三处边界显式回读;内置 YUV→RGB GPU pass(M5 解码导入的依赖项,解码接线随 M5) | 图播放路径 **GPU→CPU 回读为 0**(`oak_core::backend::gpu_transfer_counters` 计数断言,M1 帧缓存范式);`RenderedFrame::Gpu` + `to_display` 上屏在 adopted device 上零拷贝(app 测试);YUV→RGB pass 与 `colormath::yuv444p16_to_rgb_f32` 对拍;全 workspace 测试绿 | -| **M3 OFX 独立进程** | oak-ofx-host 单进程宿主;PluginJob 经 IPC;崩溃重生+紫帧回退;进度/取消协议搬运 | 杀掉 ofx-host 进程 → 在途 job 重投成功;连续三次崩溃 → 紫帧;进度条/取消行为与现状一致 | +| **M3 OFX 独立进程** | `oak-worker --ofx-host` 单进程宿主 + `oak-render/ofxhost` 客户端(Pipeline 后端安装,首 job 惰性 spawn);PluginJob 经 NDJSON + 输入/输出 shm 槽;崩溃重生+在途 job 重投+三次熔断紫帧;`plugin_progress`/`plugin_cancel` 搬运(宿主即时 flush) | 植入确定性崩溃钩子:`--ofx-crash-once` 杀掉宿主 → 在途 job 重投成功;`--ofx-crash-always` 连续三次崩溃 → 客户端熔断、eval 紫帧;进度事件(含 0.5/1.0)到达 app 回调;取消 flag 语义单测(`oak-worker/tests/ofx_host.rs` + eval/ofx_host 单测) | | **M4 流水线预取** | 调度层按 §3.4 投依赖窗口;背压策略 | 1080p 播放 CPU 占用不升、fps 不低于进程池后端;首帧延迟不劣化(基准对比留档) | | **M5 GPU 解码零拷贝** | §3.6 表逐行落地:staging fallback 基线 → Linux NVDEC/VAAPI 导入 → Windows D3D11VA 导入 → macOS VideoToolbox 导入;FFmpeg 无 hwaccel 的组合才评估手写 GPU 解码 | 硬解路径 `HW_TRANSFERS` 计数归零(不再下载);逐平台导入开/关对比测试;每行独立 PR 可回退 |