From 2a6d2191ef8a278b145ab2aaae99317a95975c30 Mon Sep 17 00:00:00 2001 From: Kayshen-X Date: Sun, 14 Jun 2026 11:49:55 +0800 Subject: [PATCH] refactor(host-core): share codegen session core --- Cargo.lock | 1 + crates/op-editor-host-core/Cargo.toml | 1 + .../src/codegen_session.rs | 159 +++++++++++++ crates/op-editor-host-core/src/lib.rs | 1 + .../tests/codegen_session.rs | 105 +++++++++ crates/op-host-desktop/src/codegen_session.rs | 219 +----------------- 6 files changed, 277 insertions(+), 209 deletions(-) create mode 100644 crates/op-editor-host-core/src/codegen_session.rs create mode 100644 crates/op-editor-host-core/tests/codegen_session.rs diff --git a/Cargo.lock b/Cargo.lock index acd2ec83c..066d237e4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3052,6 +3052,7 @@ version = "0.8.0" dependencies = [ "jian-ops-schema", "op-ai", + "op-ai-skills", "op-codegen", "op-editor-core", "op-orchestrator", diff --git a/crates/op-editor-host-core/Cargo.toml b/crates/op-editor-host-core/Cargo.toml index c3b494b31..3646e8e63 100644 --- a/crates/op-editor-host-core/Cargo.toml +++ b/crates/op-editor-host-core/Cargo.toml @@ -9,6 +9,7 @@ description = "Transport-free editor host state machines shared by OpenPencil ho [dependencies] jian-ops-schema = { path = "../../vendor/jian/crates/jian-ops-schema" } op-ai = { path = "../op-ai" } +op-ai-skills = { path = "../op-ai-skills" } op-codegen = { path = "../op-codegen", features = ["ai"] } op-editor-core = { path = "../op-editor-core" } op-orchestrator = { path = "../op-orchestrator" } diff --git a/crates/op-editor-host-core/src/codegen_session.rs b/crates/op-editor-host-core/src/codegen_session.rs new file mode 100644 index 000000000..e43e6c809 --- /dev/null +++ b/crates/op-editor-host-core/src/codegen_session.rs @@ -0,0 +1,159 @@ +//! Shared code-generation session worker. + +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::mpsc::{Receiver, Sender}; +use std::sync::Arc; + +use op_ai::chat_provider::{ChatDelta, ChatProvider, ChatRequest}; +use op_codegen::ai::types::{AssetFile, CodegenInput, PipelineStep}; +use op_codegen::ai::CodegenPipeline; +use op_editor_core::codegen::CodeGenProgress; + +/// Streamed from the worker to the UI pump. +pub enum CodegenDelta { + Progress(CodeGenProgress), + Done { + code: String, + degraded: bool, + assets: Vec, + }, + Failed(String), +} + +static NEXT_RUN_EPOCH: AtomicU64 = AtomicU64::new(1); + +/// An in-flight generation. Host pumps own the UI-specific folding of deltas. +pub struct CodegenSession { + pub rx: Receiver, + pub finished: bool, + pub framework: op_editor_core::codegen::Framework, + pub cancel: Arc, + pub run_epoch: u64, +} + +impl CodegenSession { + pub fn cancel(&self) { + self.cancel.store(true, Ordering::Relaxed); + } + + pub fn is_canceled(&self) -> bool { + self.cancel.load(Ordering::Relaxed) + } + + /// Spawn a worker that drives the pipeline against `provider`. + pub fn start( + provider: Box, + input: CodegenInput, + framework: op_editor_core::codegen::Framework, + ) -> Self { + let (tx, rx) = std::sync::mpsc::channel(); + let cancel = Arc::new(AtomicBool::new(false)); + let worker_cancel = Arc::clone(&cancel); + std::thread::Builder::new() + .name("op-codegen-turn".into()) + .spawn(move || { + let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + run_pipeline(provider.as_ref(), input, &tx, &worker_cancel); + })); + if outcome.is_err() { + let _ = tx.send(CodegenDelta::Failed( + "Code generation failed unexpectedly".into(), + )); + } + }) + .expect("spawn op-codegen-turn worker"); + CodegenSession { + rx, + finished: false, + framework, + cancel, + run_epoch: NEXT_RUN_EPOCH.fetch_add(1, Ordering::Relaxed), + } + } +} + +/// The completed result kept host-side for Download. +#[derive(Default, Clone)] +pub struct CodegenResult { + pub code: String, + pub framework_ext: String, + pub assets: Vec, +} + +/// Drive the pipeline to completion against `provider`, emitting deltas on +/// `tx`. Runs on the worker thread or synchronously in tests. +pub fn run_pipeline( + provider: &dyn ChatProvider, + input: CodegenInput, + tx: &Sender, + cancel: &AtomicBool, +) { + let mut pipe = CodegenPipeline::new(input); + loop { + if cancel.load(Ordering::Relaxed) { + pipe.cancel(); + } + match pipe.step() { + PipelineStep::Dispatch(reqs) => { + for req in reqs { + if cancel.load(Ordering::Relaxed) { + break; + } + let system = op_ai_skills::compose_system_prompt(&req.skills, 0); + let chat_req = ChatRequest { + system_prompt: system, + user_message: req.user_message.clone(), + history: Vec::new(), + max_output_tokens: req.max_output_tokens, + thinking: req.thinking, + effort: req.effort, + attachments: Vec::new(), + model: None, + }; + let mut errored: Option = None; + for delta in provider.send(chat_req) { + if cancel.load(Ordering::Relaxed) { + break; + } + match delta { + ChatDelta::TextDelta(t) => pipe.on_delta(req.id, &t), + ChatDelta::Error(e) => { + errored = Some(e); + break; + } + ChatDelta::Done { .. } => break, + _ => {} + } + } + if cancel.load(Ordering::Relaxed) { + break; + } + match errored { + Some(e) => pipe.on_error(req.id, e), + None => pipe.on_complete(req.id), + } + if tx.send(CodegenDelta::Progress(pipe.progress())).is_err() { + return; + } + } + } + PipelineStep::Waiting => {} + PipelineStep::Done { + code, + degraded, + assets, + } => { + let _ = tx.send(CodegenDelta::Done { + code, + degraded, + assets, + }); + return; + } + PipelineStep::Failed { message } => { + let _ = tx.send(CodegenDelta::Failed(message)); + return; + } + } + } +} diff --git a/crates/op-editor-host-core/src/lib.rs b/crates/op-editor-host-core/src/lib.rs index 2895b5a8a..79b4d2a45 100644 --- a/crates/op-editor-host-core/src/lib.rs +++ b/crates/op-editor-host-core/src/lib.rs @@ -2,4 +2,5 @@ pub mod chat; pub mod codegen; +pub mod codegen_session; pub mod design; diff --git a/crates/op-editor-host-core/tests/codegen_session.rs b/crates/op-editor-host-core/tests/codegen_session.rs new file mode 100644 index 000000000..c3a4b8584 --- /dev/null +++ b/crates/op-editor-host-core/tests/codegen_session.rs @@ -0,0 +1,105 @@ +use std::collections::VecDeque; +use std::sync::atomic::AtomicBool; +use std::sync::Mutex; + +use op_ai::chat_provider::{ + ChatDelta, ChatProvider, ChatRequest, EffortLevel, StopReason, ThinkingMode, +}; +use op_codegen::ai::types::CodegenInput; +use op_editor_core::codegen::Framework; +use op_editor_host_core::codegen_session::{run_pipeline, CodegenDelta, CodegenSession}; + +struct ScriptedProvider { + scripts: Mutex>>, +} + +impl ChatProvider for ScriptedProvider { + fn provider_label(&self) -> &str { + "scripted" + } + + fn send(&self, _request: ChatRequest) -> Box + Send> { + let next = self.scripts.lock().unwrap().pop_front().unwrap_or_default(); + Box::new(next.into_iter()) + } +} + +fn turn(text: &str) -> Vec { + vec![ + ChatDelta::TextDelta(text.into()), + ChatDelta::Done { + stop_reason: StopReason::EndTurn, + }, + ] +} + +fn test_input() -> CodegenInput { + CodegenInput { + nodes_json: "[{\"type\":\"frame\",\"id\":\"n1\",\"children\":[]}]".to_string(), + framework: Framework::React, + variables_json: None, + max_output_tokens: 4096, + thinking: ThinkingMode::Adaptive, + effort: EffortLevel::Low, + } +} + +#[test] +fn run_pipeline_pre_canceled_emits_only_aborted_failure() { + let provider = ScriptedProvider { + scripts: Mutex::new(VecDeque::from(vec![turn("{}")])), + }; + let (tx, rx) = std::sync::mpsc::channel(); + let cancel = AtomicBool::new(true); + run_pipeline(&provider, test_input(), &tx, &cancel); + drop(tx); + + let deltas: Vec = rx.into_iter().collect(); + assert_eq!(deltas.len(), 1); + match &deltas[0] { + CodegenDelta::Failed(message) => assert!(message.contains("Aborted")), + _ => panic!("expected aborted failure"), + } + assert_eq!(provider.scripts.lock().unwrap().len(), 1); +} + +#[test] +fn run_pipeline_drives_three_phases_to_done() { + let plan = r#"{"chunks":[{"id":"c1","name":"Root","nodeIds":["n1"],"role":"r","suggestedComponentName":"Root","dependencies":[]}],"sharedStyles":[],"rootLayout":{"direction":"column","gap":0,"responsive":false}}"#; + let chunk = "export default function Root(){}\n---CONTRACT---\n{\"componentName\":\"Root\"}"; + let assembly = "export default function App(){ return }"; + let provider = ScriptedProvider { + scripts: Mutex::new(VecDeque::from(vec![ + turn(plan), + turn(chunk), + turn(assembly), + ])), + }; + let (tx, rx) = std::sync::mpsc::channel(); + run_pipeline(&provider, test_input(), &tx, &AtomicBool::new(false)); + drop(tx); + + let deltas: Vec = rx.into_iter().collect(); + assert!(!deltas.is_empty()); + match deltas.last().expect("terminal delta") { + CodegenDelta::Done { code, .. } => assert!(code.contains("App")), + CodegenDelta::Failed(message) => panic!("pipeline failed: {message}"), + CodegenDelta::Progress(_) => panic!("last delta should be terminal"), + } +} + +#[test] +fn start_allocates_monotonic_run_epochs_and_independent_cancel_flags() { + let provider = || { + Box::new(ScriptedProvider { + scripts: Mutex::new(VecDeque::new()), + }) as Box + }; + let s1 = CodegenSession::start(provider(), test_input(), Framework::React); + let s2 = CodegenSession::start(provider(), test_input(), Framework::React); + assert!(s2.run_epoch > s1.run_epoch); + assert!(!s1.is_canceled()); + s1.cancel(); + assert!(s1.is_canceled()); + assert!(!s2.is_canceled()); +} diff --git a/crates/op-host-desktop/src/codegen_session.rs b/crates/op-host-desktop/src/codegen_session.rs index 05a2abe33..73667ce30 100644 --- a/crates/op-host-desktop/src/codegen_session.rs +++ b/crates/op-host-desktop/src/codegen_session.rs @@ -11,222 +11,23 @@ //! `Progress` delta so the panel can advance. Terminal `Done` / `Failed` //! carry the assembled code / assets back to the UI pump. -use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; -use std::sync::mpsc::{Receiver, Sender}; +#[cfg(test)] +use std::sync::atomic::{AtomicBool, Ordering}; +#[cfg(test)] use std::sync::Arc; +#[cfg(test)] use op_ai::chat_provider::{ChatDelta, ChatProvider, ChatRequest}; -use op_codegen::ai::types::{AssetFile, CodegenInput, PipelineStep}; -use op_codegen::ai::CodegenPipeline; -use op_editor_core::codegen::CodeGenProgress; +#[cfg(test)] +use op_codegen::ai::types::CodegenInput; use op_editor_host_core::codegen::framework_ext; +#[cfg(test)] +use op_editor_host_core::codegen_session::run_pipeline; +pub use op_editor_host_core::codegen_session::{CodegenDelta, CodegenResult, CodegenSession}; use op_host_native::WidgetHostNative; use crate::chat_session::provider_for_selected_model; -/// Streamed from the worker to the UI pump. -pub enum CodegenDelta { - Progress(CodeGenProgress), - Done { - code: String, - degraded: bool, - assets: Vec, - }, - Failed(String), -} - -/// Monotonic run-epoch allocator — every generation run gets a fresh id -/// (agent-indicator run-epoch pattern: a stale run can be told apart from -/// the run that replaced it). The per-run mpsc channel already isolates -/// deltas; the epoch makes run identity explicit for the cancel path + -/// tests. -static NEXT_RUN_EPOCH: AtomicU64 = AtomicU64::new(1); - -/// An in-flight generation. The UI pump drains `rx` each frame. -pub struct CodegenSession { - pub(crate) rx: Receiver, - pub(crate) finished: bool, - /// Target framework captured at launch, for the file extension. - pub(crate) framework: op_editor_core::codegen::Framework, - /// Shared abort flag (TS AbortController parity). Raised by - /// [`CodegenSession::cancel`]; the worker observes it between stream - /// deltas and pipeline steps and drives `CodegenPipeline::cancel`, - /// while the pump drops every further delta from this run. - pub(crate) cancel: Arc, - /// This run's epoch (see [`NEXT_RUN_EPOCH`]). Production code tells - /// runs apart by channel ownership (one rx per run) + the cancel flag; - /// the stamp makes run identity explicit for tests / debugging. - #[allow(dead_code)] - pub(crate) run_epoch: u64, -} - -impl CodegenSession { - /// Abort this run: the worker stops at its next hook point and the - /// pump drops every delta the run still emits, so a stale worker can - /// never resurrect the panel after Cancel. - pub fn cancel(&self) { - self.cancel.store(true, Ordering::Relaxed); - } - - pub fn is_canceled(&self) -> bool { - self.cancel.load(Ordering::Relaxed) - } -} - -/// The completed result kept HOST-SIDE for Download — asset bytes are not -/// carried in the wasm-clean `editor_state`. Export AI Bundle no longer -/// reads it: the bundle is built fresh from the LIVE selection at click -/// time (TS parity), so no generation-time node JSON is stashed here. -#[derive(Default, Clone)] -pub struct CodegenResult { - pub code: String, - /// File extension for the active framework (e.g. "tsx", "vue", "html"). - pub framework_ext: String, - pub assets: Vec, -} - -/// Drive the pipeline to completion against `provider`, emitting deltas on -/// `tx`. Runs on the worker thread (or synchronously in tests). Each model -/// request is run sequentially; the streamed text deltas are fed back into -/// the pipeline so the next `step()` can advance. -/// -/// `cancel` is the session's shared abort flag: it is observed between -/// pipeline steps, between requests, and between stream deltas. Once -/// raised, `CodegenPipeline::cancel()` parks the machine in its terminal -/// Failed("Aborted") state, which ends the loop after one final (dropped -/// by the pump) terminal delta. -pub(crate) fn run_pipeline( - provider: &dyn ChatProvider, - input: CodegenInput, - tx: &Sender, - cancel: &AtomicBool, -) { - let mut pipe = CodegenPipeline::new(input); - loop { - if cancel.load(Ordering::Relaxed) { - // Park the machine terminally — the next `step()` returns - // Failed("Aborted") and the loop exits below. - pipe.cancel(); - } - match pipe.step() { - PipelineStep::Dispatch(reqs) => { - for req in reqs { - // Observe a cancel raised between requests: stop - // dispatching; the loop top cancels the pipeline. - if cancel.load(Ordering::Relaxed) { - break; - } - // Expand the pipeline's skill NAMES into the final system - // prompt (budget 0 = no truncation cap). - let system = op_ai_skills::compose_system_prompt(&req.skills, 0); - let chat_req = ChatRequest { - system_prompt: system, - user_message: req.user_message.clone(), - // Codegen requests are self-contained — no - // chat-transcript history. - history: Vec::new(), - max_output_tokens: req.max_output_tokens, - thinking: req.thinking, - effort: req.effort, - attachments: Vec::new(), - // Codegen rides the provider's default model; - // the chat model picker doesn't govern it. - model: None, - }; - // Drain the blocking provider iterator; the first `Error` - // ends this request and is reported back to the pipeline. - // A cancel observed between stream deltas abandons the - // request mid-stream (the loop top cancels the pipeline, - // so the partial buffer is never parsed). - let mut errored: Option = None; - for delta in provider.send(chat_req) { - if cancel.load(Ordering::Relaxed) { - break; - } - match delta { - ChatDelta::TextDelta(t) => pipe.on_delta(req.id, &t), - ChatDelta::Error(e) => { - errored = Some(e); - break; - } - ChatDelta::Done { .. } => break, - _ => {} - } - } - if cancel.load(Ordering::Relaxed) { - break; - } - match errored { - Some(e) => pipe.on_error(req.id, e), - None => pipe.on_complete(req.id), - } - // Stream progress after each request so the panel advances. - if tx.send(CodegenDelta::Progress(pipe.progress())).is_err() { - return; // panel went away — stop early - } - } - } - PipelineStep::Waiting => {} - PipelineStep::Done { - code, - degraded, - assets, - } => { - let _ = tx.send(CodegenDelta::Done { - code, - degraded, - assets, - }); - return; - } - PipelineStep::Failed { message } => { - let _ = tx.send(CodegenDelta::Failed(message)); - return; - } - } - } -} - -impl CodegenSession { - /// Spawn a worker that drives the pipeline against `provider`. Returns - /// immediately — the model turns run off the UI thread. `framework` is - /// stashed on the session so the terminal `Done` can assemble the - /// `CodegenResult` with the right file extension. - pub fn start( - provider: Box, - input: CodegenInput, - framework: op_editor_core::codegen::Framework, - ) -> Self { - let (tx, rx) = std::sync::mpsc::channel(); - let cancel = Arc::new(AtomicBool::new(false)); - let worker_cancel = Arc::clone(&cancel); - std::thread::Builder::new() - .name("op-codegen-turn".into()) - .spawn(move || { - // Guard against a panic in the pipeline / provider so the UI - // receives a terminal error instead of hanging in `Generating` - // (the receiver would otherwise only observe a channel - // disconnect — see pump's Disconnected branch). - let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - run_pipeline(provider.as_ref(), input, &tx, &worker_cancel); - })); - if outcome.is_err() { - let _ = tx.send(CodegenDelta::Failed( - "Code generation failed unexpectedly".into(), - )); - } - }) - .expect("spawn op-codegen-turn worker"); - CodegenSession { - rx, - finished: false, - framework, - cancel, - run_epoch: NEXT_RUN_EPOCH.fetch_add(1, Ordering::Relaxed), - } - } -} - /// Pump the in-flight generation's deltas into `editor_state.codegen`. /// Clears `current` once the turn finishes and parks the completed result /// (asset bytes) in `last_result`. Returns true when state changed so the @@ -423,7 +224,7 @@ pub fn drain_codegen_cancel_request( mod tests { use super::*; use op_ai::chat_provider::StopReason; - use op_editor_core::codegen::Framework; + use op_editor_core::codegen::{CodeGenProgress, Framework}; /// Test-only provider that replays a DIFFERENT scripted turn each time /// `send` is called — `EchoProvider` replays the same script, which can't