feat: oakengine facade, oaknode/oakrender impls, worker+CLI, app skeleton
- oaknode Rust crate: full implementation (core engine, sequence/ track/block/footage, traverser, serializer, 43 node behaviors; 493 tests green) - oakrender Rust crate: full implementation incl. wgpu backend skeleton, ticket arena, worker pool (136 tests green; fixed lost-wakeup and ticket ordering races) - src/facade/rust (oakfacade): 222 oakengine_* exports over the module C ABIs (61 tests green); worker_main + real POSIX shm frame-slot transport (SpscRingBuffer/FrameSlotPool, wire-compatible with engine/render/ipc) - cli/rust + worker/rust binaries (29 + 29 tests green) - oakotio: FCPXML import/export (49 tests green) - oaktask: OTIO/FCPXML format dispatch (90 tests green) - app/rust: gpui app skeleton — dock panels (viewers/timeline/ explorer/inspector/node editor), transport, olive themes, i18n (en/zh), 37 tests green - gpui submodule: menu checkmarks, dock ratios, vertical meter, CPU-frame viewer surface, drop-frame timecode
This commit is contained in:
Generated
+1470
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,50 @@
|
||||
# Oak - 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/>.
|
||||
|
||||
[package]
|
||||
name = "oak-worker"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
description = "Oak Video Editor headless render worker process (Rust)"
|
||||
license = "GPL-3.0-or-later"
|
||||
|
||||
[[bin]]
|
||||
name = "oak-worker"
|
||||
path = "src/main.rs"
|
||||
|
||||
[dependencies]
|
||||
clap = { version = "4", features = ["derive"] }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
|
||||
# The liboakengine facade (Rust): its worker module owns the whole worker
|
||||
# runtime — render backend selection through the oakrender module C ABI,
|
||||
# the startup handshake and the NDJSON control loop
|
||||
# (`oakfacade::worker::worker_main`, the port of engine/src/capi/worker.cpp).
|
||||
# This binary is a thin shell over it, like worker/workermain.cpp. Its ipc
|
||||
# module provides the real shared-memory frame-slot transport
|
||||
# (`oakfacade::ipc`) that src/session.rs + src/transport.rs attach through.
|
||||
oakfacade = { path = "../../src/facade/rust" }
|
||||
|
||||
# The oakrender module crate (the Rust rewrite of the oakrender module):
|
||||
# its C ABI (include/render/renderer.h) is how the facade's worker
|
||||
# initializes the render backend — oakrender_display_renderer_create_dynamic()
|
||||
# /_init() — the Rust equivalent of the C++ worker's DynamicRenderer/
|
||||
# OpenGLRenderer. Linked here so those facade imports resolve in the binary.
|
||||
oakrender = { path = "../../src/render/rust" }
|
||||
|
||||
[profile.release]
|
||||
panic = "unwind"
|
||||
@@ -0,0 +1,106 @@
|
||||
# oak-worker (Rust)
|
||||
|
||||
Headless render worker process — the Rust rewrite of `worker/workermain.cpp`
|
||||
(contract: `engine/include/oakengine/worker.h` and
|
||||
`engine/include/oakengine/ipc.h`).
|
||||
|
||||
## Build and test
|
||||
|
||||
```sh
|
||||
cargo build --release # binary: target/release/oak-worker
|
||||
cargo test # unit + integration tests (29 tests)
|
||||
```
|
||||
|
||||
The worker is a **thin shell over the facade**, exactly like the C++
|
||||
`worker/workermain.cpp` is a thin shell over `liboakengine`:
|
||||
|
||||
- `oakfacade::worker::worker_main` (the port of
|
||||
`engine/src/capi/worker.cpp` `oakengine_worker_main()`) owns the whole
|
||||
runtime: render backend selection through the oakrender module C ABI
|
||||
(dynamic → OpenGL fallback), the startup handshake and the NDJSON
|
||||
control loop. `src/main.rs` only parses `--backend` (clap) and forwards.
|
||||
- `oakfacade::ipc` owns the shared-memory frame-slot transport (the real
|
||||
`SpscRingBuffer` + `FrameSlotPool` over POSIX `shm_open`/`mmap`);
|
||||
`src/transport.rs` attaches through it.
|
||||
|
||||
The oakrender module crate (`../../src/render/rust`) is linked so the
|
||||
facade's renderer imports resolve; oakrender depends on `ocio-rs` with the
|
||||
`bundled` feature, whose first-time build fetches a vendored OpenColorIO
|
||||
dependency (`sse2neon`) from github.com. On networks without github access,
|
||||
build with a shared target directory that already contains a completed
|
||||
oakrender build tree, e.g.:
|
||||
|
||||
```sh
|
||||
CARGO_TARGET_DIR=/path/to/oak/src/render/rust/target cargo build --release
|
||||
```
|
||||
|
||||
## What the worker does
|
||||
|
||||
Same flow as the C++ main, in the same order:
|
||||
|
||||
1. **parse `--backend <name>`** (clap; default `opengl`; `none` skips
|
||||
renderer creation and the process exits 1, like the C++ main).
|
||||
2. **initialize the render backend** (inside `oakfacade::worker`): the
|
||||
oakrender module C ABI `oakrender_display_renderer_create_dynamic` +
|
||||
`_init`, falling back to the direct OpenGL renderer exactly like the
|
||||
C++ `create_renderer()` fallback chain.
|
||||
3. **write the startup handshake** (protocol version 1, empty shared-memory
|
||||
geometry — same as the C++ worker's startup handshake; the parent
|
||||
creates the segments and announces their geometry in its reply).
|
||||
4. **serve the NDJSON control loop** on stdin/stdout until a `shutdown`
|
||||
message or EOF: `handshake` attaches the announced shared-memory
|
||||
frame-slot pools through the real transport; `load_graph` /
|
||||
`render_frame` / `cancel` / `shutdown` are dispatched by the facade
|
||||
session. Responses are one compact JSON line per message.
|
||||
|
||||
## Implemented vs stubbed (nothing is faked)
|
||||
|
||||
**Real:** argument parsing, render backend initialization (real wgpu
|
||||
renderer, dynamic → OpenGL fallback), startup handshake, NDJSON framing,
|
||||
message validation (protocol version, handshake geometry, `load_graph`
|
||||
file existence/size — the same messages the C++ worker emits), the
|
||||
**shared-memory frame-slot transport** (`oakfacade::ipc` — POSIX
|
||||
`shm_open`/`mmap`/`munmap`/`shm_unlink`, the SPSC ring buffer and the
|
||||
frame-slot pool with the exact version-1 shared layout; a `handshake`
|
||||
genuinely attaches the output and input pools), unknown-type/
|
||||
malformed-message errors, shutdown/EOF termination.
|
||||
|
||||
**Stubbed (documented in `src/transport.rs`):**
|
||||
|
||||
| area | reason |
|
||||
|---|---|
|
||||
| node-graph deserialization (`load_graph` beyond the file checks) | the oaknode crate is a `todo!()` skeleton |
|
||||
| frame rendering (`render_frame`) | no graph/render-pipeline backing (the shm frame-slot transport is attached, but there is no graph to render) |
|
||||
|
||||
Stubbed requests answer with a clear `{"type":"error","message":…}` that
|
||||
names the missing piece (a `render_frame` error also carries the ticket,
|
||||
mirroring the C++ `error_message()` shape). A real `load_graph` on a
|
||||
non-existent/empty file produces the C++-identical error before reaching
|
||||
the stub.
|
||||
|
||||
**Deviation from the C++:** the startup handshake omits `gl_major`/
|
||||
`gl_minor` — the oakrender module C ABI exposes no GL context version (the
|
||||
C++ worker reads them off its `QOpenGLContext`).
|
||||
|
||||
## Layout
|
||||
|
||||
```
|
||||
src/
|
||||
main.rs clap entry; thin shell forwarding to oakfacade::worker
|
||||
(renderer init, handshake, NDJSON loop all live there)
|
||||
ipc.rs control-plane message structs + NDJSON framing (serde)
|
||||
session.rs in-process session mirror (message dispatch + real shm
|
||||
handshake attach), exercised by the unit tests
|
||||
transport.rs real shared-memory frame-slot transport over oakfacade::ipc
|
||||
tests/worker.rs binary-level tests (help, clap errors, --backend none exit 1)
|
||||
```
|
||||
|
||||
The NDJSON control-loop behavior is exercised in-process in `src/session.rs`
|
||||
against the facade's real shared memory (no GPU needed via `--backend none`
|
||||
sessions); a binary-level loop test would require a working GPU backend and
|
||||
is deliberately not part of the unit suite. Run the binary against a
|
||||
created segment to see the real attach path:
|
||||
|
||||
```sh
|
||||
target/release/oak-worker --backend opengl <<< '{"type":"shutdown"}'
|
||||
```
|
||||
@@ -0,0 +1,257 @@
|
||||
// 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/>.
|
||||
|
||||
//! Control-plane NDJSON protocol (contract: `engine/include/oakengine/ipc.h`
|
||||
//! and `engine/include/oakengine/worker.h`).
|
||||
//!
|
||||
//! The wire format is **one compact JSON object per line** on the stdio
|
||||
//! pipes (worker.cpp / ipcmessage.cpp `write_message`/`read_message`).
|
||||
//! Every message carries a `"type"` string; the field names below are the
|
||||
//! ones the C++ serializers actually emit (`engine/render/ipc/ipcmessage.cpp`):
|
||||
//! note `ticket` / `node` / `channels` / `slot` — the longer names
|
||||
//! (`ticket_id`, `node_uuid`, `channel_count`, `output_slot`) exist only on
|
||||
//! the C POD structs in `ipc.h`.
|
||||
//!
|
||||
//! Message types (M = main/editor, W = worker):
|
||||
//! handshake M<->W negotiate protocol version + announce shm geometry
|
||||
//! load_graph M ->W path to a temp file holding the serialized graph
|
||||
//! render_frame M ->W request a frame render (ticket, node, time, params)
|
||||
//! frame_ready W ->M a rendered frame is published (slot + ticket)
|
||||
//! cancel M ->W abandon an in-flight ticket
|
||||
//! graph_update M ->W reserved (no payload struct yet)
|
||||
//! shutdown M ->W finish current work and exit cleanly
|
||||
//! error W ->M worker-side failure report ("message" field)
|
||||
//!
|
||||
//! Items the worker does not emit yet (frame_ready, graph_update,
|
||||
//! `FrameReadyMsg`) and message ids it ignores (`cancel`) are kept as the
|
||||
//! documented protocol surface; `dead_code` until the frame-slot transport
|
||||
//! lands (see crate::transport).
|
||||
|
||||
#![allow(dead_code)]
|
||||
|
||||
use std::io::{self, Write};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::{json, Value};
|
||||
|
||||
/// `"handshake"`.
|
||||
pub const TYPE_HANDSHAKE: &str = "handshake";
|
||||
/// `"load_graph"`.
|
||||
pub const TYPE_LOAD_GRAPH: &str = "load_graph";
|
||||
/// `"render_frame"`.
|
||||
pub const TYPE_RENDER_FRAME: &str = "render_frame";
|
||||
/// `"frame_ready"`.
|
||||
pub const TYPE_FRAME_READY: &str = "frame_ready";
|
||||
/// `"cancel"`.
|
||||
pub const TYPE_CANCEL: &str = "cancel";
|
||||
/// `"graph_update"`.
|
||||
pub const TYPE_GRAPH_UPDATE: &str = "graph_update";
|
||||
/// `"shutdown"`.
|
||||
pub const TYPE_SHUTDOWN: &str = "shutdown";
|
||||
/// `"error"`.
|
||||
pub const TYPE_ERROR: &str = "error";
|
||||
|
||||
/// `handshake` — field-for-field equivalent of `oak_ipc_handshake`
|
||||
/// (ipc.h). Wire field names match the C++ serializer.
|
||||
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
|
||||
#[serde(default)]
|
||||
pub struct HandshakeMsg {
|
||||
/// Protocol version.
|
||||
pub protocol_version: i32,
|
||||
/// Worker->main output shared-memory segment key.
|
||||
pub shm_key: String,
|
||||
/// Main->worker input shared-memory segment key (optional).
|
||||
pub input_shm_key: String,
|
||||
/// Number of main->worker input frame slots.
|
||||
pub input_slots: i32,
|
||||
/// Number of worker->main output frame slots.
|
||||
pub output_slots: i32,
|
||||
/// Per-output-slot pixel block size.
|
||||
pub slot_data_bytes: i64,
|
||||
/// Per-input-slot pixel block size.
|
||||
pub input_slot_data_bytes: i64,
|
||||
}
|
||||
|
||||
impl HandshakeMsg {
|
||||
/// The worker's startup handshake (`worker.cpp startup_handshake()`).
|
||||
pub fn to_json(&self) -> Value {
|
||||
json!({
|
||||
"type": TYPE_HANDSHAKE,
|
||||
"protocol_version": self.protocol_version,
|
||||
"shm_key": self.shm_key,
|
||||
"input_shm_key": self.input_shm_key,
|
||||
"input_slots": self.input_slots,
|
||||
"output_slots": self.output_slots,
|
||||
"slot_data_bytes": self.slot_data_bytes,
|
||||
"input_slot_data_bytes": self.input_slot_data_bytes,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// `render_frame` — request a frame render. Wire names per ipcmessage.cpp:
|
||||
/// `ticket`, `node`, `channels` (not the ipc.h POD names).
|
||||
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
|
||||
#[serde(default)]
|
||||
pub struct RenderFrameMsg {
|
||||
/// Correlates with the eventual frame_ready.
|
||||
pub ticket: i64,
|
||||
/// Viewer node stable uuid in the loaded graph.
|
||||
pub node: String,
|
||||
pub time_num: i64,
|
||||
pub time_den: i64,
|
||||
/// Forced output size (0 = graph default).
|
||||
pub width: i32,
|
||||
pub height: i32,
|
||||
/// Forced PixelFormat (-1 = default).
|
||||
pub format: i32,
|
||||
/// Channel count (0 = default).
|
||||
pub channels: i32,
|
||||
/// RenderMode.
|
||||
pub mode: i32,
|
||||
/// Optional decoded input slot (-1 = none).
|
||||
pub input_slot: i32,
|
||||
/// Ordered decoded input slots.
|
||||
pub input_slots: Vec<i32>,
|
||||
/// Output color transform present?
|
||||
pub has_color_transform: bool,
|
||||
pub color_is_display: bool,
|
||||
pub color_output: String,
|
||||
pub color_view: String,
|
||||
pub color_look: String,
|
||||
}
|
||||
|
||||
/// `frame_ready` — a rendered frame is published (wire names `ticket`/
|
||||
/// `slot`).
|
||||
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
|
||||
#[serde(default)]
|
||||
pub struct FrameReadyMsg {
|
||||
pub ticket: i64,
|
||||
/// Index into the worker->main output FrameSlotPool.
|
||||
pub slot: i32,
|
||||
}
|
||||
|
||||
/// `cancel` — abandon an in-flight ticket by id.
|
||||
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
|
||||
#[serde(default)]
|
||||
pub struct CancelMsg {
|
||||
pub ticket: i64,
|
||||
}
|
||||
|
||||
/// `load_graph` — path to a temporary file holding the serialized graph.
|
||||
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
|
||||
#[serde(default)]
|
||||
pub struct LoadGraphMsg {
|
||||
pub path: String,
|
||||
}
|
||||
|
||||
/// Build a worker-side error report, mirroring `error_message()` in
|
||||
/// worker.cpp: `{"type":"error","message":...}` plus `"ticket"` when
|
||||
/// non-zero.
|
||||
pub fn error_message(message: &str, ticket: Option<i64>) -> Value {
|
||||
match ticket.filter(|t| *t != 0) {
|
||||
Some(t) => json!({ "type": TYPE_ERROR, "message": message, "ticket": t }),
|
||||
None => json!({ "type": TYPE_ERROR, "message": message }),
|
||||
}
|
||||
}
|
||||
|
||||
/// Write one NDJSON message line (compact JSON + `\n`), the Rust port of
|
||||
/// `ipcmessage.cpp write_message()`.
|
||||
pub fn write_message(w: &mut impl Write, msg: &Value) -> io::Result<()> {
|
||||
let line = serde_json::to_string(msg)
|
||||
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
|
||||
w.write_all(line.as_bytes())?;
|
||||
w.write_all(b"\n")
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn handshake_wire_format_matches_cpp_field_names() {
|
||||
let hs = HandshakeMsg {
|
||||
protocol_version: 1,
|
||||
shm_key: "olive-rw-1234-0-out".into(),
|
||||
input_shm_key: "".into(),
|
||||
input_slots: 0,
|
||||
output_slots: 6,
|
||||
slot_data_bytes: 4096,
|
||||
input_slot_data_bytes: 0,
|
||||
};
|
||||
let value = hs.to_json();
|
||||
// Key order is not part of the contract (JSON objects; the C++
|
||||
// QJsonObject is hash-ordered too), but the names must match the
|
||||
// C++ serializer exactly.
|
||||
assert_eq!(value["type"], "handshake");
|
||||
assert_eq!(value["protocol_version"], 1);
|
||||
assert_eq!(value["shm_key"], "olive-rw-1234-0-out");
|
||||
assert_eq!(value["input_shm_key"], "");
|
||||
assert_eq!(value["input_slots"], 0);
|
||||
assert_eq!(value["output_slots"], 6);
|
||||
assert_eq!(value["slot_data_bytes"], 4096);
|
||||
assert_eq!(value["input_slot_data_bytes"], 0);
|
||||
// And the serialized line must parse back to the same object.
|
||||
let round: serde_json::Value =
|
||||
serde_json::from_str(&serde_json::to_string(&value).unwrap()).unwrap();
|
||||
assert_eq!(round, value);
|
||||
}
|
||||
|
||||
#[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":""}"#;
|
||||
let m: RenderFrameMsg = serde_json::from_str(json).unwrap();
|
||||
assert_eq!(m.ticket, 42);
|
||||
assert_eq!(m.node, "abcd");
|
||||
assert_eq!(m.time_num, 1);
|
||||
assert_eq!(m.time_den, 24);
|
||||
assert_eq!(m.width, 1920);
|
||||
assert_eq!(m.input_slot, -1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn render_frame_defaults_on_missing_fields() {
|
||||
// The C++ parser defaults missing fields (QJsonValue defaults);
|
||||
// serde(default) mirrors that.
|
||||
let m: RenderFrameMsg = serde_json::from_str(r#"{"type":"render_frame","ticket":7}"#).unwrap();
|
||||
assert_eq!(m.ticket, 7);
|
||||
assert_eq!(m.time_den, 0);
|
||||
assert!(m.node.is_empty());
|
||||
assert!(!m.has_color_transform);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn error_message_carries_ticket_only_when_nonzero() {
|
||||
assert_eq!(
|
||||
error_message("boom", None),
|
||||
json!({ "type": "error", "message": "boom" })
|
||||
);
|
||||
assert_eq!(
|
||||
error_message("boom", Some(0)),
|
||||
json!({ "type": "error", "message": "boom" })
|
||||
);
|
||||
assert_eq!(
|
||||
error_message("boom", Some(9)),
|
||||
json!({ "type": "error", "message": "boom", "ticket": 9 })
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_message_emits_one_json_line() {
|
||||
let mut buf = Vec::new();
|
||||
write_message(&mut buf, &json!({ "type": "shutdown" })).unwrap();
|
||||
assert_eq!(String::from_utf8(buf).unwrap(), "{\"type\":\"shutdown\"}\n");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
// 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/>.
|
||||
|
||||
//! oak-worker: headless render worker process (Rust).
|
||||
//!
|
||||
//! A thin shell over the facade, mirroring `worker/workermain.cpp`: all
|
||||
//! runtime logic — render backend selection (dynamic -> OpenGL fallback
|
||||
//! through the oakrender module C ABI), the startup handshake and the
|
||||
//! NDJSON control loop — lives in `oakfacade::worker` (the Rust port of
|
||||
//! `engine/src/capi/worker.cpp`, contract in
|
||||
//! `engine/include/oakengine/worker.h`). This crate keeps only the CLI
|
||||
//! surface (arg parsing) and the in-process session mirror
|
||||
//! ([`session`], [`transport`]) that its unit tests exercise against the
|
||||
//! facade's real shared-memory transport.
|
||||
//!
|
||||
//! See README.md for the full status.
|
||||
|
||||
mod ipc;
|
||||
mod session;
|
||||
mod transport;
|
||||
|
||||
use std::process::exit;
|
||||
|
||||
use clap::Parser;
|
||||
|
||||
// Force-link the oakrender module crate: the facade's worker module
|
||||
// (oakfacade::worker) initializes the render backend through the oakrender
|
||||
// C ABI, but its bridge imports are `extern "C"` declarations — nothing in
|
||||
// the worker source names the crate, so without this the oakrender rlib
|
||||
// would not be added to the link and those imports would stay undefined.
|
||||
#[allow(unused_imports)]
|
||||
use oakrender as _;
|
||||
|
||||
/// Protocol version announced in the startup handshake
|
||||
/// (`k_protocol_version` in worker.cpp). Mirrors
|
||||
/// `oakfacade::worker::PROTOCOL_VERSION`.
|
||||
pub const PROTOCOL_VERSION: i32 = 1;
|
||||
|
||||
/// CLI surface (the C++ worker scans argv for `--backend`; clap formalizes
|
||||
/// that single option).
|
||||
#[derive(Parser, Debug)]
|
||||
#[command(
|
||||
name = "oak-worker",
|
||||
about = "Oak render worker: headless render process for the editor's worker pool",
|
||||
disable_version_flag = true
|
||||
)]
|
||||
struct Args {
|
||||
/// Render backend to initialize: "opengl", "vulkan", "metal", "auto",
|
||||
/// or "none" (no renderer; the process exits 1 like the C++ worker).
|
||||
#[arg(long, default_value = "opengl")]
|
||||
backend: String,
|
||||
}
|
||||
|
||||
fn main() {
|
||||
let args = Args::parse();
|
||||
// The facade's worker_main is the C++ oakengine_worker_main() — the
|
||||
// whole worker flow. Like workermain.cpp, this main only forwards.
|
||||
exit(oakfacade::worker::worker_main(&args.backend.to_ascii_lowercase()));
|
||||
}
|
||||
|
||||
/// Log a worker-side message to stderr, mirroring worker.cpp `log_error()`
|
||||
/// (the `worker: ` prefix).
|
||||
pub fn log_error(message: &str) {
|
||||
eprintln!("worker: {message}");
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn clap_parses_backend_default() {
|
||||
use clap::Parser;
|
||||
let args = Args::try_parse_from(["oak-worker"]).unwrap();
|
||||
assert_eq!(args.backend, "opengl");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clap_parses_backend_flag() {
|
||||
use clap::Parser;
|
||||
let args = Args::try_parse_from(["oak-worker", "--backend", "none"]).unwrap();
|
||||
assert_eq!(args.backend, "none");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clap_rejects_unknown_flags() {
|
||||
use clap::Parser;
|
||||
assert!(Args::try_parse_from(["oak-worker", "--frobnicate"]).is_err());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,487 @@
|
||||
// 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 worker-side session state machine — the in-process mirror of
|
||||
//! `OakWorkerSession` in `engine/src/capi/worker.cpp` (whose production
|
||||
//! Rust port lives in `oakfacade::worker`).
|
||||
//!
|
||||
//! The session holds the attached shared-memory frame-slot pools
|
||||
//! ([`crate::transport::AttachedPools`]) and the shutdown flag, and
|
||||
//! answers one NDJSON control message at a time. Renderer creation and the
|
||||
//! main loop are the facade's job (see `crate::main`); this module keeps
|
||||
//! the message handling testable in-process against the facade's real
|
||||
//! shared-memory transport. Where the C++ session has machinery the Rust
|
||||
//! mirror lacks, the handler reproduces the *validation* faithfully and
|
||||
//! then reports the documented stub ([`crate::transport`]) — it never
|
||||
//! fakes a result.
|
||||
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::ipc::{self, HandshakeMsg, LoadGraphMsg, RenderFrameMsg};
|
||||
use crate::transport::{self, AttachedPools};
|
||||
|
||||
/// Worker-side session: attached frame-slot pools + message-handling state.
|
||||
pub struct WorkerSession {
|
||||
/// Attached shared-memory frame-slot pools (output + optional input).
|
||||
pools: Option<AttachedPools>,
|
||||
shutdown_requested: bool,
|
||||
}
|
||||
|
||||
impl WorkerSession {
|
||||
/// A fresh session with no attached pools.
|
||||
pub fn new() -> WorkerSession {
|
||||
WorkerSession {
|
||||
pools: None,
|
||||
shutdown_requested: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// 1 once the handshake has attached the shared-memory frame-slot
|
||||
/// pools.
|
||||
pub fn has_pools(&self) -> bool {
|
||||
self.pools.is_some()
|
||||
}
|
||||
|
||||
/// 1 once a shutdown control message has been received.
|
||||
pub fn shutdown_requested(&self) -> bool {
|
||||
self.shutdown_requested
|
||||
}
|
||||
|
||||
/// The attached output pool (the worker->main frame-slot pool).
|
||||
pub fn output_pool(&self) -> Option<&oakfacade::ipc::FrameSlotPool> {
|
||||
self.pools.as_ref().map(|p| &p.output_pool)
|
||||
}
|
||||
|
||||
/// The startup handshake the worker sends to its parent
|
||||
/// (`worker.cpp startup_handshake()`): protocol version 1 and empty
|
||||
/// shared-memory geometry — the parent creates the segments and
|
||||
/// announces their geometry in its handshake reply.
|
||||
pub fn startup_handshake(&self) -> Value {
|
||||
HandshakeMsg {
|
||||
protocol_version: crate::PROTOCOL_VERSION,
|
||||
shm_key: String::new(),
|
||||
input_shm_key: String::new(),
|
||||
input_slots: 0,
|
||||
output_slots: 0,
|
||||
slot_data_bytes: 0,
|
||||
input_slot_data_bytes: 0,
|
||||
}
|
||||
.to_json()
|
||||
}
|
||||
|
||||
/// Handle one complete NDJSON control line and produce the response, if
|
||||
/// any — the port of worker.cpp `handle()`. A malformed line yields an
|
||||
/// error response (the loop continues), never a failure.
|
||||
pub fn handle_line(&mut self, line: &str) -> Option<Value> {
|
||||
let msg: Value = match serde_json::from_str::<Value>(line) {
|
||||
Ok(v) if v.is_object() => v,
|
||||
_ => return Some(ipc::error_message("malformed control message", None)),
|
||||
};
|
||||
let typ = msg.get("type").and_then(Value::as_str).unwrap_or("");
|
||||
match typ {
|
||||
ipc::TYPE_HANDSHAKE => self.handle_handshake(&msg),
|
||||
ipc::TYPE_LOAD_GRAPH => self.handle_load_graph(&msg),
|
||||
ipc::TYPE_RENDER_FRAME => self.handle_render_frame(&msg),
|
||||
// cancel: the C++ worker does synchronous single-frame work
|
||||
// (nothing in flight), so a cancel produces no response.
|
||||
ipc::TYPE_CANCEL => None,
|
||||
ipc::TYPE_SHUTDOWN => {
|
||||
self.shutdown_requested = true;
|
||||
None
|
||||
}
|
||||
other => Some(ipc::error_message(
|
||||
&format!("unknown message type: {other}"),
|
||||
None,
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
/// `handshake`: validate and attach the shared-memory frame-slot pools.
|
||||
/// Validation mirrors worker.cpp `attach_output_pool()`; the attachment
|
||||
/// itself goes through the real [`crate::transport`].
|
||||
fn handle_handshake(&mut self, msg: &Value) -> Option<Value> {
|
||||
let hs: HandshakeMsg = match serde_json::from_value(msg.clone()) {
|
||||
Ok(hs) => hs,
|
||||
Err(_) => return Some(ipc::error_message("invalid handshake message", None)),
|
||||
};
|
||||
if hs.protocol_version != crate::PROTOCOL_VERSION {
|
||||
return Some(ipc::error_message(
|
||||
&format!("unsupported protocol version {}", hs.protocol_version),
|
||||
None,
|
||||
));
|
||||
}
|
||||
if hs.shm_key.is_empty() || hs.output_slots <= 0 || hs.slot_data_bytes <= 0 {
|
||||
return Some(ipc::error_message(
|
||||
"handshake missing output shared-memory geometry",
|
||||
None,
|
||||
));
|
||||
}
|
||||
if hs.input_slots > 0 && (hs.input_shm_key.is_empty() || hs.input_slot_data_bytes <= 0) {
|
||||
return Some(ipc::error_message(
|
||||
"handshake missing input shared-memory geometry",
|
||||
None,
|
||||
));
|
||||
}
|
||||
match transport::attach_pools(&hs) {
|
||||
Ok(pools) => {
|
||||
self.pools = Some(pools);
|
||||
None
|
||||
}
|
||||
Err(msg) => Some(ipc::error_message(&msg, None)),
|
||||
}
|
||||
}
|
||||
|
||||
/// `load_graph`: the file checks are real (mirror worker.cpp
|
||||
/// `load_graph()`); the deserialization is the documented stub.
|
||||
fn handle_load_graph(&mut self, msg: &Value) -> Option<Value> {
|
||||
let load: LoadGraphMsg = match serde_json::from_value(msg.clone()) {
|
||||
Ok(l) => l,
|
||||
Err(_) => return Some(ipc::error_message("invalid load_graph message", None)),
|
||||
};
|
||||
match std::fs::metadata(&load.path) {
|
||||
Err(_) => Some(ipc::error_message(
|
||||
&format!("graph file does not exist: {}", load.path),
|
||||
None,
|
||||
)),
|
||||
Ok(md) if md.len() == 0 => Some(ipc::error_message(
|
||||
&format!("graph file is empty: {}", load.path),
|
||||
None,
|
||||
)),
|
||||
Ok(md) => {
|
||||
crate::log_error(&format!(
|
||||
"LoadGraph: loading {} ({} bytes)",
|
||||
load.path,
|
||||
md.len()
|
||||
));
|
||||
Some(ipc::error_message(transport::GRAPH_STUB, None))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `render_frame`: the graph/render pipeline has no Rust backing, so a
|
||||
/// render request is answered with a clear error carrying the ticket —
|
||||
/// the same `error_message()` shape the C++ worker uses for its own
|
||||
/// failures.
|
||||
fn handle_render_frame(&mut self, msg: &Value) -> Option<Value> {
|
||||
let render: RenderFrameMsg = match serde_json::from_value(msg.clone()) {
|
||||
Ok(r) => r,
|
||||
Err(_) => return Some(ipc::error_message("invalid render_frame message", None)),
|
||||
};
|
||||
Some(ipc::error_message(
|
||||
transport::RENDER_STUB,
|
||||
Some(render.ticket),
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for WorkerSession {
|
||||
fn default() -> Self {
|
||||
WorkerSession::new()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use oakfacade::ipc::{FrameSlotPool, SharedMemoryRegion, ShmMode};
|
||||
use serde_json::json;
|
||||
|
||||
/// A unique, temporary POSIX segment key for a test.
|
||||
fn test_key(name: &str) -> String {
|
||||
static COUNTER: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0);
|
||||
let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||
SharedMemoryRegion::make_key(i64::from(std::process::id()), (n & 0x7FFF) as i32)
|
||||
+ &format!("-t-{name}")
|
||||
}
|
||||
|
||||
/// The "parent" side of a handshake: create an output segment holding a
|
||||
/// pool, optionally an input segment, and return the handshake message
|
||||
/// plus the owner regions (kept alive by the caller).
|
||||
fn parent_side(slots: i32, slot_bytes: i64, input: bool) -> (Value, SharedMemoryRegion, Option<SharedMemoryRegion>) {
|
||||
let out_key = test_key("out");
|
||||
let out_bytes = FrameSlotPool::bytes_needed(slots as u32, slot_bytes as usize);
|
||||
let mut out_region = SharedMemoryRegion::new();
|
||||
assert!(
|
||||
out_region.open(&out_key, out_bytes, ShmMode::Create),
|
||||
"{}",
|
||||
out_region.error()
|
||||
);
|
||||
// SAFETY: live mapping sized by bytes_needed.
|
||||
let _ = unsafe { FrameSlotPool::create(out_region.data(), slots as u32, slot_bytes as usize) };
|
||||
|
||||
let (in_key, in_bytes, in_region) = if input {
|
||||
let in_key = test_key("in");
|
||||
let in_bytes = FrameSlotPool::bytes_needed(slots as u32, slot_bytes as usize);
|
||||
let mut in_region = SharedMemoryRegion::new();
|
||||
assert!(in_region.open(&in_key, in_bytes, ShmMode::Create));
|
||||
// SAFETY: live mapping.
|
||||
let _ = unsafe { FrameSlotPool::create(in_region.data(), slots as u32, slot_bytes as usize) };
|
||||
(Some(in_key), Some(in_bytes), Some(in_region))
|
||||
} else {
|
||||
(None, None, None)
|
||||
};
|
||||
|
||||
let hs = json!({
|
||||
"type": "handshake",
|
||||
"protocol_version": crate::PROTOCOL_VERSION,
|
||||
"shm_key": out_key,
|
||||
"input_shm_key": in_key.unwrap_or_default(),
|
||||
"input_slots": if input { slots } else { 0 },
|
||||
"output_slots": slots,
|
||||
"slot_data_bytes": slot_bytes,
|
||||
"input_slot_data_bytes": in_bytes.unwrap_or(0),
|
||||
});
|
||||
(hs, out_region, in_region)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn session_starts_without_pools() {
|
||||
let s = WorkerSession::new();
|
||||
assert!(!s.has_pools());
|
||||
assert!(!s.shutdown_requested());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn startup_handshake_is_protocol_version_1_with_empty_geometry() {
|
||||
let s = WorkerSession::new();
|
||||
let hs = s.startup_handshake();
|
||||
assert_eq!(
|
||||
hs,
|
||||
json!({
|
||||
"type": "handshake",
|
||||
"protocol_version": 1,
|
||||
"shm_key": "",
|
||||
"input_shm_key": "",
|
||||
"input_slots": 0,
|
||||
"output_slots": 0,
|
||||
"slot_data_bytes": 0,
|
||||
"input_slot_data_bytes": 0,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn malformed_line_yields_error_response() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s.handle_line("this is not json").unwrap();
|
||||
assert_eq!(resp["type"], "error");
|
||||
assert_eq!(resp["message"], "malformed control message");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_object_json_yields_error_response() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s.handle_line("[1,2,3]").unwrap();
|
||||
assert_eq!(resp["message"], "malformed control message");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_message_type_yields_error_response() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s.handle_line(r#"{"type":"frobnicate"}"#).unwrap();
|
||||
assert_eq!(resp["type"], "error");
|
||||
assert_eq!(resp["message"], "unknown message type: frobnicate");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_type_field_yields_unknown_error() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s.handle_line(r#"{"hello":1}"#).unwrap();
|
||||
assert_eq!(resp["message"], "unknown message type: ");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cancel_produces_no_response() {
|
||||
let mut s = WorkerSession::new();
|
||||
assert!(s.handle_line(r#"{"type":"cancel","ticket":5}"#).is_none());
|
||||
assert!(!s.shutdown_requested());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn shutdown_sets_flag_and_has_no_response() {
|
||||
let mut s = WorkerSession::new();
|
||||
assert!(s.handle_line(r#"{"type":"shutdown"}"#).is_none());
|
||||
assert!(s.shutdown_requested());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_wrong_protocol_version() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s
|
||||
.handle_line(r#"{"type":"handshake","protocol_version":99,"shm_key":"k","output_slots":1,"slot_data_bytes":16}"#)
|
||||
.unwrap();
|
||||
assert_eq!(resp["message"], "unsupported protocol version 99");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_missing_geometry() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s
|
||||
.handle_line(r#"{"type":"handshake","protocol_version":1}"#)
|
||||
.unwrap();
|
||||
assert_eq!(resp["message"], "handshake missing output shared-memory geometry");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_missing_input_geometry_is_an_error() {
|
||||
let mut s = WorkerSession::new();
|
||||
let (mut hs, _out, _in) = parent_side(2, 256, false);
|
||||
// Ask for input slots without announcing their geometry.
|
||||
hs["input_slots"] = json!(2);
|
||||
let resp = s.handle_line(&hs.to_string()).unwrap();
|
||||
assert_eq!(resp["message"], "handshake missing input shared-memory geometry");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_attaches_real_output_pool() {
|
||||
let mut s = WorkerSession::new();
|
||||
let (hs, out_region, _in) = parent_side(4, 4096, false);
|
||||
let resp = s.handle_line(&hs.to_string());
|
||||
assert!(resp.is_none(), "unexpected error: {resp:?}");
|
||||
assert!(s.has_pools());
|
||||
let out_pool = s.output_pool().unwrap();
|
||||
assert_eq!(out_pool.slot_count(), 4);
|
||||
assert_eq!(out_pool.slot_data_bytes(), 4096);
|
||||
// The pool is real shared state: the parent's publish lands in the
|
||||
// worker's ready ring.
|
||||
// SAFETY: `out_region` is a live mapping of the pool the session
|
||||
// attached to.
|
||||
let parent_pool = unsafe { FrameSlotPool::attach(out_region.data()) };
|
||||
let mut slot = 0;
|
||||
assert!(unsafe { parent_pool.acquire(&mut slot) });
|
||||
assert_eq!(slot, 0);
|
||||
// SAFETY: acquired slot.
|
||||
let meta = unsafe { &mut *parent_pool.meta(slot) };
|
||||
meta.id = 7;
|
||||
assert!(unsafe { parent_pool.publish(slot) });
|
||||
let mut consumed = 0;
|
||||
assert!(unsafe { out_pool.consume(&mut consumed) });
|
||||
assert_eq!(consumed, 0);
|
||||
// SAFETY: consumed slot.
|
||||
assert_eq!(unsafe { (*out_pool.meta_const(consumed)).id }, 7);
|
||||
unsafe { out_pool.release(consumed) };
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_attaches_input_pool_too() {
|
||||
let mut s = WorkerSession::new();
|
||||
let (hs, _out, _in) = parent_side(2, 256, true);
|
||||
let resp = s.handle_line(&hs.to_string());
|
||||
assert!(resp.is_none(), "unexpected error: {resp:?}");
|
||||
assert!(s.has_pools());
|
||||
let pools = s.pools.as_ref().unwrap();
|
||||
assert!(pools.input_pool.is_some());
|
||||
let in_pool = pools.input_pool.as_ref().unwrap();
|
||||
assert_eq!(in_pool.slot_count(), 2);
|
||||
assert_eq!(in_pool.slot_data_bytes(), 256);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_attach_failure_reports_error() {
|
||||
let mut s = WorkerSession::new();
|
||||
// A key that was never created.
|
||||
let resp = s
|
||||
.handle_line(&json!({
|
||||
"type": "handshake",
|
||||
"protocol_version": 1,
|
||||
"shm_key": format!("olive-rw-{}-missing", std::process::id()),
|
||||
"output_slots": 4,
|
||||
"slot_data_bytes": 4096,
|
||||
})
|
||||
.to_string())
|
||||
.unwrap();
|
||||
assert_eq!(resp["type"], "error");
|
||||
assert!(resp["message"]
|
||||
.as_str()
|
||||
.unwrap()
|
||||
.starts_with("failed to attach shared memory: "));
|
||||
assert!(!s.has_pools());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_rejects_non_pool_segment() {
|
||||
let mut s = WorkerSession::new();
|
||||
// A real segment of the right size that does not contain a pool
|
||||
// (zeroed memory -> wrong magic).
|
||||
let key = test_key("nopool");
|
||||
let bytes = FrameSlotPool::bytes_needed(4, 4096);
|
||||
let mut region = SharedMemoryRegion::new();
|
||||
assert!(region.open(&key, bytes, ShmMode::Create));
|
||||
let resp = s
|
||||
.handle_line(&json!({
|
||||
"type": "handshake",
|
||||
"protocol_version": 1,
|
||||
"shm_key": key,
|
||||
"output_slots": 4,
|
||||
"slot_data_bytes": 4096,
|
||||
})
|
||||
.to_string())
|
||||
.unwrap();
|
||||
assert_eq!(resp["message"], "shared memory does not contain a frame slot pool");
|
||||
assert!(!s.has_pools());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handshake_bad_json_shape_is_invalid_handshake() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s.handle_line(r#"{"type":"handshake","protocol_version":"x"}"#).unwrap();
|
||||
assert_eq!(resp["message"], "invalid handshake message");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn load_graph_file_checks_are_real_then_stub() {
|
||||
let mut s = WorkerSession::new();
|
||||
|
||||
let missing = "/definitely/not/a/real/graph.ove";
|
||||
let resp = s
|
||||
.handle_line(&json!({ "type": "load_graph", "path": missing }).to_string())
|
||||
.unwrap();
|
||||
assert_eq!(resp["message"], format!("graph file does not exist: {missing}"));
|
||||
|
||||
let empty = std::env::temp_dir().join("oak_worker_test_empty_graph.ove");
|
||||
std::fs::write(&empty, b"").unwrap();
|
||||
let resp = s
|
||||
.handle_line(&json!({ "type": "load_graph", "path": empty.display().to_string() }).to_string())
|
||||
.unwrap();
|
||||
assert_eq!(resp["message"], format!("graph file is empty: {}", empty.display()));
|
||||
let _ = std::fs::remove_file(&empty);
|
||||
|
||||
let real = std::env::temp_dir().join("oak_worker_test_graph.ove");
|
||||
std::fs::write(&real, b"<root/>").unwrap();
|
||||
let resp = s
|
||||
.handle_line(&json!({ "type": "load_graph", "path": real.display().to_string() }).to_string())
|
||||
.unwrap();
|
||||
assert!(resp["message"]
|
||||
.as_str()
|
||||
.unwrap()
|
||||
.contains("node-graph deserialization is not yet available"));
|
||||
let _ = std::fs::remove_file(&real);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn render_frame_reports_stub_with_ticket() {
|
||||
let mut s = WorkerSession::new();
|
||||
let resp = s
|
||||
.handle_line(r#"{"type":"render_frame","ticket":123,"node":"abc"}"#)
|
||||
.unwrap();
|
||||
assert_eq!(resp["type"], "error");
|
||||
assert_eq!(resp["ticket"], 123);
|
||||
assert!(resp["message"]
|
||||
.as_str()
|
||||
.unwrap()
|
||||
.contains("frame rendering is not yet available"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
// 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/>.
|
||||
|
||||
//! Shared-memory frame-slot transport.
|
||||
//!
|
||||
//! The C++ worker exchanges bulk pixel data with the editor through named
|
||||
//! shared-memory segments holding an `olive::ipc::FrameSlotPool` — a fixed
|
||||
//! pool of frame slots whose free/ready queues are synchronized by the
|
||||
//! lock-free single-producer/single-consumer `SpscRingBuffer`. That
|
||||
//! machinery is implemented in the facade crate (`oakfacade::ipc`, the
|
||||
//! Rust port of `engine/render/ipc/` behind
|
||||
//! `engine/include/oakengine/ipc.h`) — this module is the worker-side
|
||||
//! transport over it.
|
||||
//!
|
||||
//! [`attach_pools`] mirrors the C++ `attach_output_pool()`: attach the
|
||||
//! output segment in [`ShmMode::Attach`], map the [`FrameSlotPool`] it
|
||||
//! contains (rejecting segments without the pool magic), and attach the
|
||||
//! input pool when the handshake announces one. The returned
|
||||
//! [`AttachedPools`] is owned by the session and dropped (unmapped) with
|
||||
//! it.
|
||||
//!
|
||||
//! The node-graph and render-pipeline stubs below carry the same rationale
|
||||
//! as before: `oaknode` is a `todo!()` skeleton and the oakrender crate
|
||||
//! does not yet evaluate an arbitrary loaded graph to a frame.
|
||||
|
||||
use oakfacade::ipc::{FrameSlotPool, SharedMemoryRegion, ShmMode};
|
||||
|
||||
use crate::ipc::HandshakeMsg;
|
||||
|
||||
/// Why `load_graph` answers "not yet available" (after the real file checks).
|
||||
pub const GRAPH_STUB: &str = "load_graph: node-graph deserialization is not yet available in the \
|
||||
Rust worker (the oaknode crate is a todo!() skeleton; see worker/rust/README.md)";
|
||||
|
||||
/// Why `render_frame` answers "not yet available".
|
||||
pub const RENDER_STUB: &str = "render_frame: frame rendering is not yet available in the Rust \
|
||||
worker (no node-graph or render-pipeline backing; the shm frame-slot transport is \
|
||||
attached but there is no graph to render; see worker/rust/README.md)";
|
||||
|
||||
/// The shared-memory frame-slot pools attached by a successful handshake,
|
||||
/// kept alive for the session's lifetime.
|
||||
pub struct AttachedPools {
|
||||
/// Worker->main output segment (unmapped on drop).
|
||||
pub output_region: SharedMemoryRegion,
|
||||
/// Output frame-slot pool view.
|
||||
pub output_pool: FrameSlotPool,
|
||||
/// Main->worker input segment, when the handshake announced one.
|
||||
pub input_region: Option<SharedMemoryRegion>,
|
||||
/// Input frame-slot pool view.
|
||||
pub input_pool: Option<FrameSlotPool>,
|
||||
}
|
||||
|
||||
/// Attach the handshake's shared-memory frame-slot pools — the real port of
|
||||
/// worker.cpp `attach_output_pool()`.
|
||||
///
|
||||
/// The output segment must exist and contain a valid [`FrameSlotPool`]
|
||||
/// (magic check); the input pool is attached when `input_slots > 0`.
|
||||
/// Returns `Err(message)` describing the failure, matching the C++ error
|
||||
/// strings.
|
||||
pub fn attach_pools(hs: &HandshakeMsg) -> Result<AttachedPools, String> {
|
||||
// Attach the worker->main output pool.
|
||||
let 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, bytes, ShmMode::Attach) {
|
||||
return Err(format!(
|
||||
"failed to attach shared memory: {}",
|
||||
output_region.error()
|
||||
));
|
||||
}
|
||||
// SAFETY: `output_region` is a live mapping of at least `bytes` bytes
|
||||
// (checked inside `open`).
|
||||
let output_pool = unsafe { FrameSlotPool::attach(output_region.data()) };
|
||||
if !output_pool.is_valid() {
|
||||
return Err("shared memory does not contain a frame slot pool".to_string());
|
||||
}
|
||||
|
||||
// Attach the main->worker input pool when the handshake announced one.
|
||||
let mut input_region = None;
|
||||
let mut input_pool = None;
|
||||
if hs.input_slots > 0 {
|
||||
let input_bytes = FrameSlotPool::bytes_needed(
|
||||
hs.input_slots as u32,
|
||||
hs.input_slot_data_bytes as usize,
|
||||
);
|
||||
let mut region = SharedMemoryRegion::new();
|
||||
if !region.open(&hs.input_shm_key, input_bytes, ShmMode::Attach) {
|
||||
return Err(format!(
|
||||
"failed to attach input shared memory: {}",
|
||||
region.error()
|
||||
));
|
||||
}
|
||||
// SAFETY: `region` is a live mapping of at least `input_bytes`
|
||||
// bytes (checked inside `open`).
|
||||
let pool = unsafe { FrameSlotPool::attach(region.data()) };
|
||||
if !pool.is_valid() {
|
||||
return Err("input shared memory does not contain a frame slot pool".to_string());
|
||||
}
|
||||
input_region = Some(region);
|
||||
input_pool = Some(pool);
|
||||
}
|
||||
|
||||
Ok(AttachedPools {
|
||||
output_region,
|
||||
output_pool,
|
||||
input_region,
|
||||
input_pool,
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
// 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/>.
|
||||
|
||||
//! Binary-level tests for `oak-worker`: argument handling and the process
|
||||
//! exit contract. The NDJSON control-loop behavior itself is exercised
|
||||
//! in-process in `src/session.rs` (a real loop test would require a working
|
||||
//! GPU backend, so it stays out of the unit suite).
|
||||
|
||||
use std::process::Command;
|
||||
|
||||
fn bin() -> &'static str {
|
||||
env!("CARGO_BIN_EXE_oak-worker")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn help_exits_zero() {
|
||||
let out = Command::new(bin()).arg("--help").output().expect("spawn oak-worker");
|
||||
assert_eq!(out.status.code(), Some(0));
|
||||
let stdout = String::from_utf8_lossy(&out.stdout);
|
||||
assert!(stdout.contains("oak-worker"));
|
||||
assert!(stdout.contains("--backend"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_flag_is_a_clap_usage_error() {
|
||||
let out = Command::new(bin()).arg("--frobnicate").output().expect("spawn oak-worker");
|
||||
// clap's usage-error exit code.
|
||||
assert_eq!(out.status.code(), Some(2));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn backend_none_exits_one_like_the_cpp_main() {
|
||||
// Mirrors oakengine_worker_main(): without a renderer the worker cannot
|
||||
// do anything and exits 1.
|
||||
let out = Command::new(bin())
|
||||
.args(["--backend", "none"])
|
||||
.output()
|
||||
.expect("spawn oak-worker");
|
||||
assert_eq!(out.status.code(), Some(1));
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(stderr.contains("no renderer initialized"), "stderr: {stderr}");
|
||||
}
|
||||
Reference in New Issue
Block a user