refactor(host-core): share chat turn core
This commit is contained in:
parent
207560b7a3
commit
fcaf902849
245
crates/op-editor-host-core/src/chat.rs
Normal file
245
crates/op-editor-host-core/src/chat.rs
Normal file
|
|
@ -0,0 +1,245 @@
|
|||
//! Shared chat turn worker and transcript folding logic.
|
||||
|
||||
use std::sync::mpsc::{self, Receiver, Sender, SyncSender, TryRecvError};
|
||||
use std::sync::Mutex;
|
||||
use std::thread;
|
||||
|
||||
use op_ai::chat_provider::{
|
||||
ChatDelta, ChatProvider, ChatRequest, ChatToolExecutor, ChatToolResult,
|
||||
};
|
||||
use op_editor_core::{ChatMessage, ChatToolCall};
|
||||
|
||||
/// One tool call forwarded from an agent-loop worker to a host UI thread.
|
||||
pub struct ChatToolRequest {
|
||||
pub name: String,
|
||||
pub args_json: String,
|
||||
pub ack: SyncSender<ChatToolResult>,
|
||||
}
|
||||
|
||||
/// Worker-side [`ChatToolExecutor`] that forwards calls over a channel and
|
||||
/// blocks until the host acks with the real tool result.
|
||||
pub struct UiChatToolExecutor {
|
||||
tx: Mutex<Sender<ChatToolRequest>>,
|
||||
}
|
||||
|
||||
impl UiChatToolExecutor {
|
||||
pub fn new(tx: Sender<ChatToolRequest>) -> Self {
|
||||
Self { tx: Mutex::new(tx) }
|
||||
}
|
||||
}
|
||||
|
||||
impl ChatToolExecutor for UiChatToolExecutor {
|
||||
fn execute(&self, name: &str, args_json: &str) -> ChatToolResult {
|
||||
let (ack_tx, ack_rx) = std::sync::mpsc::sync_channel::<ChatToolResult>(1);
|
||||
let req = ChatToolRequest {
|
||||
name: name.to_string(),
|
||||
args_json: args_json.to_string(),
|
||||
ack: ack_tx,
|
||||
};
|
||||
let sent = match self.tx.lock() {
|
||||
Ok(tx) => tx.send(req).is_ok(),
|
||||
Err(_) => false,
|
||||
};
|
||||
if !sent {
|
||||
return aborted_result();
|
||||
}
|
||||
match ack_rx.recv_timeout(std::time::Duration::from_secs(60)) {
|
||||
Ok(result) => result,
|
||||
Err(_) => timeout_result(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn timeout_result() -> ChatToolResult {
|
||||
ChatToolResult {
|
||||
content: r#"{"success":false,"error":"tool execution timed out waiting for the editor"}"#
|
||||
.into(),
|
||||
is_error: true,
|
||||
}
|
||||
}
|
||||
|
||||
fn aborted_result() -> ChatToolResult {
|
||||
ChatToolResult {
|
||||
content: r#"{"success":false,"error":"chat turn aborted before the tool ran"}"#.into(),
|
||||
is_error: true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Create the worker-to-host tool channel for one chat turn.
|
||||
pub fn chat_tool_channel() -> (UiChatToolExecutor, Receiver<ChatToolRequest>) {
|
||||
let (tx, rx) = std::sync::mpsc::channel::<ChatToolRequest>();
|
||||
(UiChatToolExecutor::new(tx), rx)
|
||||
}
|
||||
|
||||
/// One in-flight chat turn.
|
||||
pub struct ChatSession {
|
||||
rx: Receiver<ChatDelta>,
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
finished: bool,
|
||||
}
|
||||
|
||||
/// Result of a single non-blocking [`ChatSession::poll`].
|
||||
pub struct ChatPoll {
|
||||
pub text: String,
|
||||
pub thinking: String,
|
||||
pub tool_calls: Vec<ChatToolCall>,
|
||||
pub error: Option<String>,
|
||||
pub finished: bool,
|
||||
}
|
||||
|
||||
impl ChatPoll {
|
||||
/// True when this poll carried no new content and the turn has not ended.
|
||||
pub fn is_idle(&self) -> bool {
|
||||
self.text.is_empty()
|
||||
&& self.thinking.is_empty()
|
||||
&& self.tool_calls.is_empty()
|
||||
&& self.error.is_none()
|
||||
&& !self.finished
|
||||
}
|
||||
}
|
||||
|
||||
/// Fold one [`ChatPoll`] into the trailing assistant `message`.
|
||||
pub fn apply_poll_to_message(message: &mut ChatMessage, poll: &ChatPoll) {
|
||||
if let Some(err) = &poll.error {
|
||||
message.content = format!("error: {err}");
|
||||
} else {
|
||||
message.content.push_str(&poll.text);
|
||||
}
|
||||
message.thinking.push_str(&poll.thinking);
|
||||
if poll.tool_calls.iter().any(tool_call_defaults_open) {
|
||||
message.tools_collapsed = false;
|
||||
}
|
||||
message.tool_calls.extend(poll.tool_calls.iter().cloned());
|
||||
if poll.finished {
|
||||
message.streaming = false;
|
||||
}
|
||||
}
|
||||
|
||||
fn tool_call_defaults_open(call: &ChatToolCall) -> bool {
|
||||
if let Some(level) = tool_level_from_args(&call.args) {
|
||||
return matches!(level.as_str(), "modify" | "delete" | "orchestrate");
|
||||
}
|
||||
matches!(
|
||||
call.name.as_str(),
|
||||
"update_node"
|
||||
| "replace_node"
|
||||
| "move_node"
|
||||
| "set_variables"
|
||||
| "set_themes"
|
||||
| "load_theme_preset"
|
||||
| "rename_page"
|
||||
| "reorder_page"
|
||||
| "batch_design"
|
||||
| "set_design_md"
|
||||
| "export_design_md"
|
||||
| "delete_node"
|
||||
| "remove_page"
|
||||
)
|
||||
}
|
||||
|
||||
fn tool_level_from_args(args: &str) -> Option<String> {
|
||||
serde_json::from_str::<serde_json::Value>(args)
|
||||
.ok()?
|
||||
.get("level")?
|
||||
.as_str()
|
||||
.map(str::to_string)
|
||||
}
|
||||
|
||||
impl ChatSession {
|
||||
/// Spawn a worker that drains `provider.send(req)` into a channel.
|
||||
pub fn start(provider: Box<dyn ChatProvider>, req: ChatRequest) -> Self {
|
||||
Self::start_with_tools(provider, req, None)
|
||||
}
|
||||
|
||||
/// [`start`](Self::start) plus a tool-request receiver for providers that
|
||||
/// execute host tools.
|
||||
pub fn start_with_tools(
|
||||
provider: Box<dyn ChatProvider>,
|
||||
req: ChatRequest,
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
) -> Self {
|
||||
let (tx, rx) = mpsc::channel();
|
||||
thread::Builder::new()
|
||||
.name("op-chat-turn".into())
|
||||
.spawn(move || {
|
||||
for delta in provider.send(req) {
|
||||
if tx.send(delta).is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
})
|
||||
.expect("spawn op-chat-turn thread");
|
||||
Self {
|
||||
rx,
|
||||
tool_rx,
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Wrap externally supplied channels, used when a host owns its own
|
||||
/// routing worker but wants the shared poll/finish behavior.
|
||||
pub fn from_channels(
|
||||
rx: Receiver<ChatDelta>,
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
) -> Self {
|
||||
Self {
|
||||
rx,
|
||||
tool_rx,
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain every delta ready right now without blocking.
|
||||
pub fn poll(&mut self) -> ChatPoll {
|
||||
let mut text = String::new();
|
||||
let mut thinking = String::new();
|
||||
let mut tool_calls = Vec::new();
|
||||
let mut error = None;
|
||||
loop {
|
||||
match self.rx.try_recv() {
|
||||
Ok(ChatDelta::TextDelta(s)) => text.push_str(&s),
|
||||
Ok(ChatDelta::Thinking(s)) => thinking.push_str(&s),
|
||||
Ok(ChatDelta::ToolUse { name, args }) => {
|
||||
tool_calls.push(ChatToolCall { name, args });
|
||||
}
|
||||
Ok(ChatDelta::Error(msg)) => {
|
||||
if error.is_none() {
|
||||
error = Some(msg);
|
||||
}
|
||||
}
|
||||
Ok(ChatDelta::Done { .. }) => self.finished = true,
|
||||
Err(TryRecvError::Empty) => break,
|
||||
Err(TryRecvError::Disconnected) => {
|
||||
self.finished = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
ChatPoll {
|
||||
text,
|
||||
thinking,
|
||||
tool_calls,
|
||||
error,
|
||||
finished: self.finished,
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain pending host tool requests without blocking.
|
||||
pub fn drain_tool_requests(&mut self) -> Vec<ChatToolRequest> {
|
||||
let Some(tool_rx) = self.tool_rx.as_ref() else {
|
||||
return Vec::new();
|
||||
};
|
||||
let mut requests = Vec::new();
|
||||
loop {
|
||||
match tool_rx.try_recv() {
|
||||
Ok(req) => requests.push(req),
|
||||
Err(TryRecvError::Empty) | Err(TryRecvError::Disconnected) => break,
|
||||
}
|
||||
}
|
||||
requests
|
||||
}
|
||||
|
||||
pub fn finished(&self) -> bool {
|
||||
self.finished
|
||||
}
|
||||
}
|
||||
|
|
@ -1,3 +1,4 @@
|
|||
//! Transport-free editor host state machines shared by native and web hosts.
|
||||
|
||||
pub mod chat;
|
||||
pub mod codegen;
|
||||
|
|
|
|||
196
crates/op-editor-host-core/tests/chat.rs
Normal file
196
crates/op-editor-host-core/tests/chat.rs
Normal file
|
|
@ -0,0 +1,196 @@
|
|||
use op_ai::chat_provider::{
|
||||
ChatDelta, ChatRequest, ChatToolExecutor, ChatToolResult, EchoProvider, StopReason,
|
||||
};
|
||||
use op_editor_core::{ChatMessage, ChatToolCall};
|
||||
use op_editor_host_core::chat::{apply_poll_to_message, chat_tool_channel, ChatPoll, ChatSession};
|
||||
|
||||
fn drain_session(session: &mut ChatSession) -> (String, String, Vec<ChatToolCall>, Option<String>) {
|
||||
let mut text = String::new();
|
||||
let mut thinking = String::new();
|
||||
let mut tools = Vec::new();
|
||||
let mut error = None;
|
||||
for _ in 0..1000 {
|
||||
let poll = session.poll();
|
||||
text.push_str(&poll.text);
|
||||
thinking.push_str(&poll.thinking);
|
||||
tools.extend(poll.tool_calls);
|
||||
if poll.error.is_some() {
|
||||
error = poll.error;
|
||||
}
|
||||
if poll.finished {
|
||||
break;
|
||||
}
|
||||
std::thread::sleep(std::time::Duration::from_millis(1));
|
||||
}
|
||||
(text, thinking, tools, error)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn session_streams_provider_deltas_to_completion() {
|
||||
let provider = Box::new(EchoProvider {
|
||||
script: vec![
|
||||
ChatDelta::TextDelta("Hel".into()),
|
||||
ChatDelta::TextDelta("lo".into()),
|
||||
ChatDelta::Done {
|
||||
stop_reason: StopReason::EndTurn,
|
||||
},
|
||||
],
|
||||
});
|
||||
let mut session = ChatSession::start(
|
||||
provider,
|
||||
ChatRequest {
|
||||
system_prompt: String::new(),
|
||||
user_message: "hi".into(),
|
||||
max_output_tokens: 256,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
|
||||
let (text, _, _, _) = drain_session(&mut session);
|
||||
assert!(session.finished());
|
||||
assert_eq!(text, "Hello");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn poll_splits_thinking_tools_errors_and_answer_text() {
|
||||
let provider = Box::new(EchoProvider {
|
||||
script: vec![
|
||||
ChatDelta::Thinking("let me think".into()),
|
||||
ChatDelta::ToolUse {
|
||||
name: "insert_node".into(),
|
||||
args: "{\"kind\":\"rect\"}".into(),
|
||||
},
|
||||
ChatDelta::TextDelta("answer".into()),
|
||||
ChatDelta::Error("boom".into()),
|
||||
ChatDelta::Done {
|
||||
stop_reason: StopReason::EndTurn,
|
||||
},
|
||||
],
|
||||
});
|
||||
let mut session = ChatSession::start(
|
||||
provider,
|
||||
ChatRequest {
|
||||
user_message: "x".into(),
|
||||
max_output_tokens: 64,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
|
||||
let (text, thinking, tools, error) = drain_session(&mut session);
|
||||
assert_eq!(text, "answer");
|
||||
assert_eq!(thinking, "let me think");
|
||||
assert_eq!(tools.len(), 1);
|
||||
assert_eq!(tools[0].name, "insert_node");
|
||||
assert_eq!(tools[0].args, "{\"kind\":\"rect\"}");
|
||||
assert_eq!(error.as_deref(), Some("boom"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn apply_poll_appends_content_and_clears_streaming_on_finish() {
|
||||
let mut msg = ChatMessage::assistant_streaming();
|
||||
apply_poll_to_message(
|
||||
&mut msg,
|
||||
&ChatPoll {
|
||||
text: "hi".into(),
|
||||
thinking: "reasoning".into(),
|
||||
tool_calls: vec![ChatToolCall {
|
||||
name: "t".into(),
|
||||
args: "{}".into(),
|
||||
}],
|
||||
error: None,
|
||||
finished: false,
|
||||
},
|
||||
);
|
||||
assert_eq!(msg.content, "hi");
|
||||
assert_eq!(msg.thinking, "reasoning");
|
||||
assert_eq!(msg.tool_calls.len(), 1);
|
||||
assert!(msg.streaming);
|
||||
|
||||
apply_poll_to_message(
|
||||
&mut msg,
|
||||
&ChatPoll {
|
||||
text: "!".into(),
|
||||
thinking: String::new(),
|
||||
tool_calls: vec![],
|
||||
error: None,
|
||||
finished: true,
|
||||
},
|
||||
);
|
||||
assert_eq!(msg.content, "hi!");
|
||||
assert!(!msg.streaming);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn apply_poll_opens_modify_tools_but_keeps_read_tools_collapsed() {
|
||||
let mut modify = ChatMessage::assistant_streaming();
|
||||
apply_poll_to_message(
|
||||
&mut modify,
|
||||
&ChatPoll {
|
||||
text: String::new(),
|
||||
thinking: String::new(),
|
||||
tool_calls: vec![ChatToolCall {
|
||||
name: "batch_design".into(),
|
||||
args: "{}".into(),
|
||||
}],
|
||||
error: None,
|
||||
finished: false,
|
||||
},
|
||||
);
|
||||
assert!(!modify.tools_collapsed);
|
||||
|
||||
let mut read = ChatMessage::assistant_streaming();
|
||||
apply_poll_to_message(
|
||||
&mut read,
|
||||
&ChatPoll {
|
||||
text: String::new(),
|
||||
thinking: String::new(),
|
||||
tool_calls: vec![ChatToolCall {
|
||||
name: "snapshot_layout".into(),
|
||||
args: "{}".into(),
|
||||
}],
|
||||
error: None,
|
||||
finished: false,
|
||||
},
|
||||
);
|
||||
assert!(read.tools_collapsed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn apply_poll_error_replaces_content_and_ends_stream() {
|
||||
let mut msg = ChatMessage::assistant_streaming();
|
||||
msg.content = "partial answer".into();
|
||||
apply_poll_to_message(
|
||||
&mut msg,
|
||||
&ChatPoll {
|
||||
text: String::new(),
|
||||
thinking: String::new(),
|
||||
tool_calls: vec![],
|
||||
error: Some("rate limited".into()),
|
||||
finished: true,
|
||||
},
|
||||
);
|
||||
assert_eq!(msg.content, "error: rate limited");
|
||||
assert!(!msg.streaming);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tool_executor_forwards_request_and_waits_for_ack() {
|
||||
let (executor, rx) = chat_tool_channel();
|
||||
let worker = std::thread::spawn(move || executor.execute("insert_node", "{\"kind\":\"rect\"}"));
|
||||
|
||||
let req = rx
|
||||
.recv_timeout(std::time::Duration::from_secs(1))
|
||||
.expect("tool request");
|
||||
assert_eq!(req.name, "insert_node");
|
||||
assert_eq!(req.args_json, "{\"kind\":\"rect\"}");
|
||||
req.ack
|
||||
.send(ChatToolResult {
|
||||
content: r#"{"success":true}"#.into(),
|
||||
is_error: false,
|
||||
})
|
||||
.expect("ack");
|
||||
|
||||
let result = worker.join().expect("worker joins");
|
||||
assert_eq!(result.content, r#"{"success":true}"#);
|
||||
assert!(!result.is_error);
|
||||
}
|
||||
|
|
@ -16,11 +16,11 @@
|
|||
//! event loop drains requests each frame (`chat_session::pump`) and
|
||||
//! executes via [`execute_chat_tool`] against the canonical state.
|
||||
|
||||
use std::sync::mpsc::{Receiver, Sender, SyncSender};
|
||||
use std::sync::Mutex;
|
||||
|
||||
use op_ai::chat_provider::{ChatToolDef, ChatToolExecutor, ChatToolResult};
|
||||
use op_ai::chat_provider::{ChatToolDef, ChatToolResult};
|
||||
use op_editor_core::EditorState;
|
||||
pub(crate) use op_editor_host_core::chat::{
|
||||
chat_tool_channel, ChatToolRequest, UiChatToolExecutor,
|
||||
};
|
||||
use op_mcp::{ToolRegistry, ToolResponse};
|
||||
|
||||
/// TS `maxTurns` for the chat agent loop (`ai-chat-handlers.ts:254`).
|
||||
|
|
@ -99,78 +99,6 @@ pub(crate) fn chat_tool_defs() -> Vec<ChatToolDef> {
|
|||
]
|
||||
}
|
||||
|
||||
/// One tool call forwarded from the agent-loop worker to the UI
|
||||
/// thread. The worker blocks on `ack` until the host executed the
|
||||
/// call against the live `EditorState`.
|
||||
pub(crate) struct ChatToolRequest {
|
||||
pub name: String,
|
||||
pub args_json: String,
|
||||
pub ack: SyncSender<ChatToolResult>,
|
||||
}
|
||||
|
||||
/// Worker-side [`ChatToolExecutor`] — forwards each call over the
|
||||
/// session's tool channel and blocks until the UI thread acks. The
|
||||
/// sender rides a `Mutex` because `std::sync::mpsc::Sender` is `Send`
|
||||
/// but not `Sync`, while `ChatToolExecutor` objects are shared behind
|
||||
/// an `Arc` (`Send + Sync` bound).
|
||||
pub(crate) struct UiChatToolExecutor {
|
||||
tx: Mutex<Sender<ChatToolRequest>>,
|
||||
}
|
||||
|
||||
impl UiChatToolExecutor {
|
||||
pub(crate) fn new(tx: Sender<ChatToolRequest>) -> Self {
|
||||
Self { tx: Mutex::new(tx) }
|
||||
}
|
||||
}
|
||||
|
||||
impl ChatToolExecutor for UiChatToolExecutor {
|
||||
fn execute(&self, name: &str, args_json: &str) -> ChatToolResult {
|
||||
let (ack_tx, ack_rx) = std::sync::mpsc::sync_channel::<ChatToolResult>(1);
|
||||
let req = ChatToolRequest {
|
||||
name: name.to_string(),
|
||||
args_json: args_json.to_string(),
|
||||
ack: ack_tx,
|
||||
};
|
||||
let sent = match self.tx.lock() {
|
||||
Ok(tx) => tx.send(req).is_ok(),
|
||||
Err(_) => false,
|
||||
};
|
||||
if !sent {
|
||||
return aborted_result();
|
||||
}
|
||||
// The ack sender lives inside the request; if the UI drops the
|
||||
// session (turn aborted), the sender drops and recv errors.
|
||||
// The timeout is a backstop for the pump stalling while the
|
||||
// session still owns the receiver (app shutdown / frozen UI):
|
||||
// the worker must never block a tool turn forever.
|
||||
match ack_rx.recv_timeout(std::time::Duration::from_secs(60)) {
|
||||
Ok(result) => result,
|
||||
Err(_) => timeout_result(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn timeout_result() -> ChatToolResult {
|
||||
ChatToolResult {
|
||||
content: r#"{"success":false,"error":"tool execution timed out waiting for the editor"}"#
|
||||
.into(),
|
||||
is_error: true,
|
||||
}
|
||||
}
|
||||
|
||||
fn aborted_result() -> ChatToolResult {
|
||||
ChatToolResult {
|
||||
content: r#"{"success":false,"error":"chat turn aborted before the tool ran"}"#.into(),
|
||||
is_error: true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Create the worker↔UI tool channel for one chat turn.
|
||||
pub(crate) fn chat_tool_channel() -> (UiChatToolExecutor, Receiver<ChatToolRequest>) {
|
||||
let (tx, rx) = std::sync::mpsc::channel::<ChatToolRequest>();
|
||||
(UiChatToolExecutor::new(tx), rx)
|
||||
}
|
||||
|
||||
/// Execute one chat tool call against the live editor state. Returns
|
||||
/// the TS-shaped tool result (`{"success":…}`) plus whether the call
|
||||
/// mutated the document (caller marks the redraw dirty).
|
||||
|
|
@ -341,6 +269,7 @@ fn chat_tool_registry(state: &EditorState, requested: &str) -> ToolRegistry {
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use op_ai::chat_provider::ChatToolExecutor;
|
||||
|
||||
#[test]
|
||||
fn chat_tool_defs_match_ts_crud_subset_and_auth_levels() {
|
||||
|
|
|
|||
|
|
@ -1,20 +1,17 @@
|
|||
//! Background chat-turn runner.
|
||||
//! Desktop chat-session host glue.
|
||||
//!
|
||||
//! `ChatProvider::send()` hands back a *blocking* iterator — draining
|
||||
//! it on the UI thread would freeze the window for the whole LLM
|
||||
//! turn. `ChatSession` drains the turn on a dedicated worker thread
|
||||
//! and exposes a non-blocking [`ChatSession::poll`] the winit event
|
||||
//! loop pumps each frame, appending deltas to the in-flight assistant
|
||||
//! message.
|
||||
//! The transport-free turn worker, poll result, transcript folding, and tool
|
||||
//! channel live in `op-editor-host-core::chat`. This module keeps desktop
|
||||
//! provider routing plus UI-thread tool execution against `WidgetHostNative`.
|
||||
|
||||
use std::sync::mpsc::{self, Receiver, TryRecvError};
|
||||
use std::thread;
|
||||
|
||||
use op_ai::chat_provider::{ChatDelta, ChatProvider, ChatRequest, ChatToolResult};
|
||||
use op_editor_core::{ChatMessage, ChatState, ChatToolCall, EditorState};
|
||||
use op_ai::chat_provider::ChatToolResult;
|
||||
use op_editor_core::{ChatState, EditorState};
|
||||
#[cfg(test)]
|
||||
pub use op_editor_host_core::chat::ChatPoll;
|
||||
pub use op_editor_host_core::chat::{apply_poll_to_message, ChatSession};
|
||||
use op_host_native::WidgetHostNative;
|
||||
|
||||
use crate::chat_canvas_tools::{execute_chat_tool, ChatToolRequest};
|
||||
use crate::chat_canvas_tools::execute_chat_tool;
|
||||
|
||||
// Turn launch + provider routing (split out at the 800-line cap).
|
||||
// `launch_if_pending` and friends live in the sibling file; the
|
||||
|
|
@ -28,199 +25,8 @@ pub(crate) use launch::{
|
|||
};
|
||||
pub use launch::{drain_new_chat_request, drain_stop_request, launch_if_pending};
|
||||
|
||||
/// One in-flight chat turn. The worker thread owns the provider and
|
||||
/// drains `provider.send()` into the channel; [`poll`] consumes
|
||||
/// whatever is ready without blocking.
|
||||
pub struct ChatSession {
|
||||
rx: Receiver<ChatDelta>,
|
||||
/// Canvas tool-call requests from the builtin agent loop. `None`
|
||||
/// for providers without tool execution. Drained by [`pump`] each
|
||||
/// frame — the worker blocks on each request's ack, mirroring the
|
||||
/// design session's command channel.
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
finished: bool,
|
||||
}
|
||||
|
||||
/// Result of a single non-blocking [`ChatSession::poll`].
|
||||
pub struct ChatPoll {
|
||||
/// Answer-text fragments (`TextDelta`) accumulated since the last
|
||||
/// poll. Empty when nothing new arrived.
|
||||
pub text: String,
|
||||
/// Reasoning fragments (`Thinking`) accumulated since the last
|
||||
/// poll — kept separate from `text` so the chat panel can render
|
||||
/// them in their own collapsible block.
|
||||
pub thinking: String,
|
||||
/// Tool invocations (`ToolUse`) seen this poll.
|
||||
pub tool_calls: Vec<ChatToolCall>,
|
||||
/// First error seen this poll, if any. When set the caller
|
||||
/// should surface it as the assistant message body.
|
||||
pub error: Option<String>,
|
||||
/// True once the turn's terminal `Done` arrived, or the worker
|
||||
/// thread / channel closed.
|
||||
pub finished: bool,
|
||||
}
|
||||
|
||||
impl ChatPoll {
|
||||
/// True when this poll carried no new content and the turn has
|
||||
/// not ended — the caller can skip touching the transcript.
|
||||
fn is_idle(&self) -> bool {
|
||||
self.text.is_empty()
|
||||
&& self.thinking.is_empty()
|
||||
&& self.tool_calls.is_empty()
|
||||
&& self.error.is_none()
|
||||
&& !self.finished
|
||||
}
|
||||
}
|
||||
|
||||
/// Fold one [`ChatPoll`] into the trailing assistant `message`. An
|
||||
/// error replaces the visible body; otherwise answer text + thinking
|
||||
/// accumulate and tool calls append. A finished poll clears the
|
||||
/// `streaming` flag so the panel stops the streaming animation.
|
||||
pub fn apply_poll_to_message(message: &mut ChatMessage, poll: &ChatPoll) {
|
||||
if let Some(err) = &poll.error {
|
||||
message.content = format!("error: {err}");
|
||||
} else {
|
||||
message.content.push_str(&poll.text);
|
||||
}
|
||||
message.thinking.push_str(&poll.thinking);
|
||||
if poll.tool_calls.iter().any(tool_call_defaults_open) {
|
||||
message.tools_collapsed = false;
|
||||
}
|
||||
message.tool_calls.extend(poll.tool_calls.iter().cloned());
|
||||
if poll.finished {
|
||||
message.streaming = false;
|
||||
}
|
||||
}
|
||||
|
||||
fn tool_call_defaults_open(call: &ChatToolCall) -> bool {
|
||||
if let Some(level) = tool_level_from_args(&call.args) {
|
||||
return matches!(level.as_str(), "modify" | "delete" | "orchestrate");
|
||||
}
|
||||
matches!(
|
||||
call.name.as_str(),
|
||||
"update_node"
|
||||
| "replace_node"
|
||||
| "move_node"
|
||||
| "set_variables"
|
||||
| "set_themes"
|
||||
| "load_theme_preset"
|
||||
| "rename_page"
|
||||
| "reorder_page"
|
||||
| "batch_design"
|
||||
| "set_design_md"
|
||||
| "export_design_md"
|
||||
| "delete_node"
|
||||
| "remove_page"
|
||||
)
|
||||
}
|
||||
|
||||
fn tool_level_from_args(args: &str) -> Option<String> {
|
||||
serde_json::from_str::<serde_json::Value>(args)
|
||||
.ok()?
|
||||
.get("level")?
|
||||
.as_str()
|
||||
.map(str::to_string)
|
||||
}
|
||||
|
||||
impl ChatSession {
|
||||
/// Spawn a worker that drains `provider.send(req)` into a
|
||||
/// channel. Returns immediately — the LLM turn runs off-thread.
|
||||
pub fn start(provider: Box<dyn ChatProvider>, req: ChatRequest) -> Self {
|
||||
Self::start_with_tools(provider, req, None)
|
||||
}
|
||||
|
||||
/// [`start`](Self::start) plus a canvas tool-request channel for
|
||||
/// tool-executing providers (the builtin agent loop). [`pump`]
|
||||
/// drains `tool_rx` each frame and executes the calls against the
|
||||
/// live editor state.
|
||||
pub fn start_with_tools(
|
||||
provider: Box<dyn ChatProvider>,
|
||||
req: ChatRequest,
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
) -> Self {
|
||||
let (tx, rx) = mpsc::channel();
|
||||
thread::Builder::new()
|
||||
.name("op-chat-turn".into())
|
||||
.spawn(move || {
|
||||
// `provider.send` itself returns a blocking iterator;
|
||||
// draining it here keeps the block off the UI thread.
|
||||
for delta in provider.send(req) {
|
||||
if tx.send(delta).is_err() {
|
||||
return; // chat panel went away — stop early
|
||||
}
|
||||
}
|
||||
})
|
||||
.expect("spawn op-chat-turn thread");
|
||||
Self {
|
||||
rx,
|
||||
tool_rx,
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Wrap externally-supplied channels — the CLI intent router
|
||||
/// (GAP #33) owns its own worker thread and feeds these directly.
|
||||
pub(crate) fn from_channels(
|
||||
rx: Receiver<ChatDelta>,
|
||||
tool_rx: Option<Receiver<ChatToolRequest>>,
|
||||
) -> Self {
|
||||
Self {
|
||||
rx,
|
||||
tool_rx,
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain every delta ready right now without blocking.
|
||||
pub fn poll(&mut self) -> ChatPoll {
|
||||
let mut text = String::new();
|
||||
let mut thinking = String::new();
|
||||
let mut tool_calls = Vec::new();
|
||||
let mut error = None;
|
||||
loop {
|
||||
match self.rx.try_recv() {
|
||||
Ok(ChatDelta::TextDelta(s)) => text.push_str(&s),
|
||||
Ok(ChatDelta::Thinking(s)) => thinking.push_str(&s),
|
||||
// Tool dispatch is the agent runtime's job; the panel
|
||||
// only surfaces the call in its collapsible tool view.
|
||||
Ok(ChatDelta::ToolUse { name, args }) => {
|
||||
tool_calls.push(ChatToolCall { name, args });
|
||||
}
|
||||
Ok(ChatDelta::Error(msg)) => {
|
||||
if error.is_none() {
|
||||
error = Some(msg);
|
||||
}
|
||||
}
|
||||
Ok(ChatDelta::Done { .. }) => self.finished = true,
|
||||
Err(TryRecvError::Empty) => break,
|
||||
Err(TryRecvError::Disconnected) => {
|
||||
self.finished = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
ChatPoll {
|
||||
text,
|
||||
thinking,
|
||||
tool_calls,
|
||||
error,
|
||||
finished: self.finished,
|
||||
}
|
||||
}
|
||||
|
||||
/// True once the turn has fully completed. Test-only accessor —
|
||||
/// the event-loop glue keys off `ChatPoll::finished` instead.
|
||||
#[cfg(test)]
|
||||
pub fn finished(&self) -> bool {
|
||||
self.finished
|
||||
}
|
||||
}
|
||||
|
||||
/// Pump the in-flight turn's deltas into the trailing (assistant)
|
||||
/// message, then execute any pending canvas tool calls against the
|
||||
/// live editor state. Clears `current` once the turn finishes.
|
||||
/// Returns true when the transcript changed so the caller can dirty
|
||||
/// the redraw.
|
||||
/// Pump the in-flight turn's deltas into the trailing assistant message, then
|
||||
/// execute any pending canvas tool calls against the live editor state.
|
||||
pub fn pump(host: &mut WidgetHostNative, current: &mut Option<ChatSession>) -> bool {
|
||||
let Some(session) = current.as_mut() else {
|
||||
return false;
|
||||
|
|
@ -233,10 +39,6 @@ pub fn pump(host: &mut WidgetHostNative, current: &mut Option<ChatSession>) -> b
|
|||
changed = true;
|
||||
}
|
||||
}
|
||||
// Execute pending canvas tool calls AFTER folding the deltas in —
|
||||
// the agent loop emits the `ToolUse` delta (which creates the
|
||||
// transcript card) before it forwards the request, so by the time
|
||||
// a request is visible its card already exists.
|
||||
if drain_tool_requests(host.editor_state_mut(), session) {
|
||||
changed = true;
|
||||
}
|
||||
|
|
@ -249,31 +51,15 @@ pub fn pump(host: &mut WidgetHostNative, current: &mut Option<ChatSession>) -> b
|
|||
changed
|
||||
}
|
||||
|
||||
/// Drain every pending canvas tool request from the in-flight turn
|
||||
/// and execute it against the live `EditorState` — the chat-loop
|
||||
/// mirror of `design_session::pump_commands`. Each request is acked
|
||||
/// with its result so the blocked worker resumes; the matching
|
||||
/// transcript tool card is updated with the real result. Returns true
|
||||
/// when state or transcript changed.
|
||||
/// Drain every pending canvas tool request from the in-flight turn and execute
|
||||
/// it against the live `EditorState`.
|
||||
fn drain_tool_requests(state: &mut EditorState, session: &mut ChatSession) -> bool {
|
||||
let Some(tool_rx) = session.tool_rx.as_ref() else {
|
||||
return false;
|
||||
};
|
||||
let mut requests = Vec::new();
|
||||
loop {
|
||||
match tool_rx.try_recv() {
|
||||
Ok(req) => requests.push(req),
|
||||
Err(TryRecvError::Empty) | Err(TryRecvError::Disconnected) => break,
|
||||
}
|
||||
}
|
||||
let requests = session.drain_tool_requests();
|
||||
if requests.is_empty() {
|
||||
return false;
|
||||
}
|
||||
let mut changed = false;
|
||||
for req in requests {
|
||||
// Internal host op from the DESIGN_MODIFY route (GAP #33) —
|
||||
// applies the parsed modification nodes against the live
|
||||
// state. Never advertised to a model; no transcript card.
|
||||
if req.name == crate::chat_intent::APPLY_MODIFICATION_OP {
|
||||
let nodes = serde_json::from_str::<serde_json::Value>(&req.args_json)
|
||||
.ok()
|
||||
|
|
@ -297,17 +83,12 @@ fn drain_tool_requests(state: &mut EditorState, session: &mut ChatSession) -> bo
|
|||
if attach_tool_result_to_transcript(&mut state.chat, &req.name, &result) {
|
||||
changed = true;
|
||||
}
|
||||
// If the ack fails the worker already dropped its receiver
|
||||
// (turn aborted) — nothing to do.
|
||||
let _ = req.ack.send(result);
|
||||
}
|
||||
changed
|
||||
}
|
||||
|
||||
/// Record an executed tool call's result on its transcript card: the
|
||||
/// last matching `status:"running"` envelope gains `result` +
|
||||
/// `status: done|error`, which the chat panel's tool card renders as
|
||||
/// its Result line (same envelope shape the TS cards use).
|
||||
/// Record an executed tool call's result on its transcript card.
|
||||
fn attach_tool_result_to_transcript(
|
||||
chat: &mut ChatState,
|
||||
name: &str,
|
||||
|
|
|
|||
|
|
@ -1,7 +1,8 @@
|
|||
use super::*;
|
||||
use crate::chat_system_prompt::{build_chat_system_prompt, chat_history_from_transcript};
|
||||
use op_ai::chat_history::{trim_chat_history, DEFAULT_MAX_CHARS, DEFAULT_MAX_MESSAGES};
|
||||
use op_ai::chat_provider::{EchoProvider, StopReason};
|
||||
use op_ai::chat_provider::{ChatDelta, ChatRequest, EchoProvider, StopReason};
|
||||
use op_editor_core::{ChatMessage, ChatToolCall};
|
||||
|
||||
#[test]
|
||||
fn session_streams_echo_provider_deltas_to_completion() {
|
||||
|
|
|
|||
Loading…
Reference in a new issue