From 453b0144d9cb83949c344c99fce9ad0c845ff47d Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 1 Jul 2026 12:25:54 +0300 Subject: [PATCH] split agent turn phases --- .beads/interactions.jsonl | 1 + .beads/issues.jsonl | 2 +- klbr-core/src/agent.rs | 796 ++++++++++++++++++++------------------ 3 files changed, 427 insertions(+), 372 deletions(-) diff --git a/.beads/interactions.jsonl b/.beads/interactions.jsonl index f618edc..827e5f6 100644 --- a/.beads/interactions.jsonl +++ b/.beads/interactions.jsonl @@ -93,3 +93,4 @@ {"id":"int-d2a36a78","kind":"field_change","created_at":"2026-07-01T09:11:34.33603338Z","actor":"dawn","issue_id":"klbr-647","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Migrated klbr-core Message.role to a Role enum with lowercase serde compatibility and updated core role checks; verified workspace tests."}} {"id":"int-e8560d77","kind":"field_change","created_at":"2026-07-01T09:14:49.216625941Z","actor":"dawn","issue_id":"klbr-ow7","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Extracted duplicated agent run-loop stream-failure backoff handling into a shared helper with a named retry limit; verified klbr-core tests."}} {"id":"int-24fd625d","kind":"field_change","created_at":"2026-07-01T09:18:07.32698818Z","actor":"dawn","issue_id":"klbr-q53","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Moved external crate versions into root workspace.dependencies and switched member manifests to workspace=true; verified cargo metadata, workspace check, and workspace tests."}} +{"id":"int-04183019","kind":"field_change","created_at":"2026-07-01T09:25:30.275787695Z","actor":"dawn","issue_id":"klbr-so2","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Split agent run_turn into stream_iteration and execute_tool_iteration phase helpers, reducing run_turn to a smaller dispatcher while preserving tool/preemption behavior; verified core and workspace tests."}} diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index 5df88de..8841254 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -52,7 +52,7 @@ {"_type":"issue","id":"klbr-zue","title":"core: fix shell.rs unsafe byte-slice truncation","description":"shell.rs truncates command output using byte offsets on a lossy UTF-8 string (e.g. stdout[..20_000]), which can panic if slicing in the middle of a multi-byte character. Truncate by char count or indices instead.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:48Z","created_by":"dawn","updated_at":"2026-06-30T22:03:16Z","started_at":"2026-06-30T22:02:13Z","closed_at":"2026-06-30T22:03:16Z","close_reason":"Replaced shell stdout/stderr byte slicing with char-boundary truncation helper and added UTF-8 regression tests. Verified cargo test -p klbr-core shell_truncation, full cargo test -p klbr-core (168 passed, 1 ignored), cargo fmt --check, git diff --check, and rg confirms the unsafe stdout/stderr slices are gone.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-3ha","title":"core: refactor models.rs media stripping and SSE parser duplication","description":"strip_media_urls and replace_media_urls_with_placeholders share near-identical regex parsing logic. Also, complete() and stream() duplicate manual SSE chunk parsing. Clean these up.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:42Z","created_by":"dawn","updated_at":"2026-07-01T09:02:06Z","started_at":"2026-07-01T08:55:00Z","closed_at":"2026-07-01T09:02:06Z","close_reason":"Refactored shared SSE and media data-url parsing in models.rs; added focused coverage and verified klbr-core tests.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-647","title":"core: migrate Message.role from String to Role enum","description":"Message.role is currently represented as a raw String, which leads to fragile string comparisons ('assistant', 'tool', etc.) all over the codebase. Migrate to a proper Role enum.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:36Z","created_by":"dawn","updated_at":"2026-07-01T09:11:34Z","started_at":"2026-07-01T09:03:58Z","closed_at":"2026-07-01T09:11:34Z","close_reason":"Migrated klbr-core Message.role to a Role enum with lowercase serde compatibility and updated core role checks; verified workspace tests.","dependency_count":0,"dependent_count":0,"comment_count":0} -{"_type":"issue","id":"klbr-so2","title":"core: decompose agent.rs run_turn() god method","description":"agent.rs run_turn() is 500+ lines, managing execution, tool loops, preemption, and discord integration. Break down tool execution loops and phase handlers into separate functions.","status":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:29Z","created_by":"dawn","updated_at":"2026-06-30T21:39:29Z","dependency_count":0,"dependent_count":0,"comment_count":0} +{"_type":"issue","id":"klbr-so2","title":"core: decompose agent.rs run_turn() god method","description":"agent.rs run_turn() is 500+ lines, managing execution, tool loops, preemption, and discord integration. Break down tool execution loops and phase handlers into separate functions.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:29Z","created_by":"dawn","updated_at":"2026-07-01T09:25:30Z","started_at":"2026-07-01T09:18:51Z","closed_at":"2026-07-01T09:25:30Z","close_reason":"Split agent run_turn into stream_iteration and execute_tool_iteration phase helpers, reducing run_turn to a smaller dispatcher while preserving tool/preemption behavior; verified core and workspace tests.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-ow7","title":"core: refactor agent.rs run() loop failure-handling duplication","description":"agent.rs run() loop copy-pastes the exact same ~30-line backoff/failure-handling parsing block 4 times across different event stream branches. Extract to a helper function.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:23Z","created_by":"dawn","updated_at":"2026-07-01T09:14:49Z","started_at":"2026-07-01T09:12:29Z","closed_at":"2026-07-01T09:14:49Z","close_reason":"Extracted duplicated agent run-loop stream-failure backoff handling into a shared helper with a named retry limit; verified klbr-core tests.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-7e0","title":"core: clean up dead tool files in src/tools/","description":"There are 12 unused tool files in klbr-core/src/tools/ (remember.rs, recall.rs, tag_memory.rs, etc.) that are not declared as submodules or registered. They should be deleted.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:16Z","created_by":"dawn","updated_at":"2026-06-30T22:08:24Z","started_at":"2026-06-30T22:05:01Z","closed_at":"2026-06-30T22:08:24Z","close_reason":"Removed unregistered old memory tool modules and updated memory docs","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-q53","title":"workspace: consolidate dependency versions via workspace.dependencies","description":"Workspace is not utilizing cargo workspace dependencies, causing version drift between core, daemon, discord, and ipc. Consolidate them in the root Cargo.toml.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:10Z","created_by":"dawn","updated_at":"2026-07-01T09:18:07Z","started_at":"2026-07-01T09:15:45Z","closed_at":"2026-07-01T09:18:07Z","close_reason":"Moved external crate versions into root workspace.dependencies and switched member manifests to workspace=true; verified cargo metadata, workspace check, and workspace tests.","dependency_count":0,"dependent_count":0,"comment_count":0} diff --git a/klbr-core/src/agent.rs b/klbr-core/src/agent.rs index 4c50b9d..0bd9da6 100644 --- a/klbr-core/src/agent.rs +++ b/klbr-core/src/agent.rs @@ -12,7 +12,7 @@ use crate::{ harness_block::{self, HarnessBlock, DEFAULT_FORMAT_STYLE}, interrupt::Interrupt, memory::MemoryStore, - models::{LlmClient, LlmEvent, Message, Role}, + models::{LlmClient, LlmEvent, Message, Role, ToolCall}, pipeline::{BenchQuery, ContextBudget, MemoryPipeline}, router::{RouteDecision, Router}, tools::{self, ToolContext}, @@ -21,6 +21,20 @@ use crate::{ const MAX_CONSECUTIVE_STREAM_FAILURES: usize = 5; +struct StreamIteration { + response: String, + thinking: String, + tool_calls: Vec, + stream_error: Option, + cancelled_by: Option, +} + +enum ToolIterationOutcome { + Continue, + Suspend(Option), + Restart, +} + pub struct Agent { config: Config, memory: MemoryStore, @@ -541,425 +555,465 @@ impl Agent { .await } - async fn run_turn( + async fn stream_iteration( &mut self, llm: &LlmClient, - tool_ctx: &ToolContext, ctx: &mut Context, - turn_count: &mut usize, - runtime_soul: &mut String, - recalled_index: Option, - is_local_message: bool, - is_discord_event: bool, - consecutive_waits: &mut usize, - ) -> Result> { - let mut tool_iterations = 0usize; - let mut discord_send_called = false; - let mut discord_nudge_injected = false; - let mut local_send_called = false; - let mut local_send_nudge_injected = false; - let mut should_restart = false; - // non-preemptive interrupt received during streaming/tool execution; returned at end of turn - let mut stashed_interrupt: Option = None; - const MAX_TOOL_ITERATIONS: usize = 20; - - // emit Started once — all tokens/thinking for the whole turn go into one bubble - let _ = self.output.send(AgentEvent::Started); + stashed_interrupt: &mut Option, + ) -> StreamIteration { + let (tok_tx, mut tok_rx) = mpsc::channel(256); + let llm2 = llm.clone(); + let msgs = ctx.as_messages_with_refs(&self.memory); + let defs = self.registry.definitions(); + let stream_task = tokio::spawn(async move { llm2.stream(&msgs, &defs, tok_tx).await }); + + let mut response = String::new(); + let mut thinking = String::new(); + let mut tool_calls = vec![]; + let mut cancelled_by = None; loop { - let (tok_tx, mut tok_rx) = mpsc::channel(256); - let llm2 = llm.clone(); - let msgs = ctx.as_messages_with_refs(&self.memory); - let defs = self.registry.definitions(); - let stream_task = tokio::spawn(async move { llm2.stream(&msgs, &defs, tok_tx).await }); - - let mut iteration_response = String::new(); - let mut iteration_thinking = String::new(); - let mut tool_calls = vec![]; - let mut stream_error = None::; - - let mut stream_cancelled_by: Option = None; - loop { - tokio::select! { - biased; - interrupt = self.rx.recv() => { - match interrupt { - Some(int) if int.is_preemptive() => { - tracing::info!(source = int.source_tag(), "preemptive interrupt during stream; aborting"); - stream_task.abort(); - stream_cancelled_by = Some(int); - break; - } - Some(int) => { - // non-preemptive (e.g. ExternalEvent) — stash for after turn - if stashed_interrupt.is_none() { - stashed_interrupt = Some(int); - } + tokio::select! { + biased; + interrupt = self.rx.recv() => { + match interrupt { + Some(int) if int.is_preemptive() => { + tracing::info!(source = int.source_tag(), "preemptive interrupt during stream; aborting"); + stream_task.abort(); + cancelled_by = Some(int); + break; + } + Some(int) => { + // non-preemptive (e.g. ExternalEvent) — stash for after turn + if stashed_interrupt.is_none() { + *stashed_interrupt = Some(int); } - None => {} // channel closed = shutdown } + None => {} // channel closed = shutdown } - ev = tok_rx.recv() => { - match ev { - None => break, - Some(LlmEvent::ThinkToken(tok)) => { - iteration_thinking.push_str(&tok); - let _ = self.output.send(AgentEvent::ThinkToken(tok)); - } - Some(LlmEvent::Token(tok)) => { - iteration_response.push_str(&tok); - let _ = self.output.send(AgentEvent::ScratchToken(tok)); - } - Some(LlmEvent::Usage(usage)) => { - ctx.update_tokens(&usage); - } - Some(LlmEvent::ToolCalls(calls)) => { - tool_calls = calls; - } - Some(LlmEvent::ToolCallStarted) => { - let _ = self.output.send(AgentEvent::ToolCallStarted); - } + } + ev = tok_rx.recv() => { + match ev { + None => break, + Some(LlmEvent::ThinkToken(tok)) => { + thinking.push_str(&tok); + let _ = self.output.send(AgentEvent::ThinkToken(tok)); + } + Some(LlmEvent::Token(tok)) => { + response.push_str(&tok); + let _ = self.output.send(AgentEvent::ScratchToken(tok)); + } + Some(LlmEvent::Usage(usage)) => { + ctx.update_tokens(&usage); + } + Some(LlmEvent::ToolCalls(calls)) => { + tool_calls = calls; + } + Some(LlmEvent::ToolCallStarted) => { + let _ = self.output.send(AgentEvent::ToolCallStarted); } } } } + } - if let Some(interrupt) = stream_cancelled_by { - // turn was preempted mid-stream — finish cleanly and hand off - let _ = self.output.send(AgentEvent::Done); - cleanup_recalled_context(ctx, recalled_index); - save_context_snapshot(&self.memory, ctx); - return Ok(Some(interrupt)); + let stream_error = match stream_task.await { + Ok(Ok(())) => None, + Ok(Err(e)) => { + let msg = format!("llm stream failed: {e}"); + tracing::error!(%msg); + let _ = self.output.send(AgentEvent::Error(msg.clone())); + let _ = self + .output + .send(AgentEvent::Status("llm stream failed".into())); + Some(msg) } - - match stream_task.await { - Ok(Ok(())) => {} - Ok(Err(e)) => { - let msg = format!("llm stream failed: {e}"); - tracing::error!(%msg); - let _ = self.output.send(AgentEvent::Error(msg.clone())); - let _ = self - .output - .send(AgentEvent::Status("llm stream failed".into())); - stream_error = Some(msg); - } - Err(e) if e.is_cancelled() => { - // intentional abort — already handled above - } - Err(e) => { - let msg = format!("llm stream task panicked/cancelled: {e}"); - tracing::error!(%msg); - let _ = self.output.send(AgentEvent::Error(msg.clone())); - let _ = self - .output - .send(AgentEvent::Status("llm stream task failed".into())); - stream_error = Some(msg); - } + Err(e) if e.is_cancelled() => { + // intentional abort — already handled by cancelled_by + None + } + Err(e) => { + let msg = format!("llm stream task panicked/cancelled: {e}"); + tracing::error!(%msg); + let _ = self.output.send(AgentEvent::Error(msg.clone())); + let _ = self + .output + .send(AgentEvent::Status("llm stream task failed".into())); + Some(msg) } + }; - if !tool_calls.is_empty() && tool_iterations < MAX_TOOL_ITERATIONS { - tool_iterations += 1; - let reasoning_content = if iteration_thinking.is_empty() { - None - } else { - Some(iteration_thinking.clone()) - }; - ctx.push_assistant_tool_calls(tool_calls.clone(), None, reasoning_content); - let _ = self.memory.log_turn_with_tools( - "assistant", - "", - (!iteration_thinking.is_empty()).then_some(iteration_thinking.as_str()), - Some(&tool_calls), - None, - ); + StreamIteration { + response, + thinking, + tool_calls, + stream_error, + cancelled_by, + } + } - let mut pending_local_sends = Vec::new(); - for call in &tool_calls { - let name = call.function.name.clone(); - if name == "restart_harness" { - should_restart = true; + async fn execute_tool_iteration( + &mut self, + tool_ctx: &ToolContext, + ctx: &mut Context, + turn_count: &mut usize, + runtime_soul: &mut String, + tool_calls: &[ToolCall], + iteration_thinking: &str, + stashed_interrupt: &mut Option, + consecutive_waits: &mut usize, + discord_send_called: &mut bool, + local_send_called: &mut bool, + ) -> Result { + let reasoning_content = if iteration_thinking.is_empty() { + None + } else { + Some(iteration_thinking.to_string()) + }; + ctx.push_assistant_tool_calls(tool_calls.to_vec(), None, reasoning_content); + let _ = self.memory.log_turn_with_tools( + "assistant", + "", + (!iteration_thinking.is_empty()).then_some(iteration_thinking), + Some(tool_calls), + None, + ); + + let mut pending_local_sends = Vec::new(); + let mut should_restart = false; + for call in tool_calls { + let name = call.function.name.clone(); + if name == "restart_harness" { + should_restart = true; + } + if name != "wait_and_continue" { + *consecutive_waits = 0; + } + if name == "discord_send" { + *discord_send_called = true; + } + if name == "local_send" { + *local_send_called = true; + } + let args = call.function.arguments.clone(); + if name != "local_send" { + let _ = self.output.send(AgentEvent::ToolCall { + name: name.clone(), + args: args.clone(), + }); + } + + let wait_result = if name == "wait_and_continue" { + Some(execute_wait_and_continue(&args, &mut self.rx).await) + } else { + None + }; + let mut incoming_after_tool = None; + let mut suspend_after_wait = false; + let mut sent_local_content = None; + let result = if name == "local_send" { + match tools::local_send_content(&args) { + Ok(content) => { + let result = + format!("sent {} chars to local operator", content.chars().count()); + sent_local_content = Some(content); + result } - if name != "wait_and_continue" { - *consecutive_waits = 0; + Err(err) => { + tracing::warn!(err = %err, "local_send call had invalid arguments"); + format!("error: {err}") } - if name == "discord_send" { - discord_send_called = true; + } + } else { + match wait_result { + Some(WaitAndContinueResult::Timeout(mut result)) => { + suspend_after_wait = is_wait_timeout_without_event(&result); + if suspend_after_wait { + *consecutive_waits += 1; + if *consecutive_waits >= 5 && *consecutive_waits % 5 == 0 { + let nudge = format!( + "\n\nnudge: you have called wait_and_continue {} consecutive times without taking other actions. to conserve local resources, please increase the wait duration in your next wait_and_continue call (e.g. wait_and_continue(duration=\"10m\")).", + *consecutive_waits + ); + result.push_str(&nudge); + } + } else { + *consecutive_waits = 0; + } + result } - if name == "local_send" { - local_send_called = true; + Some(WaitAndContinueResult::Incoming { result, interrupt }) => { + incoming_after_tool = Some(interrupt); + result } - let args = call.function.arguments.clone(); - if name != "local_send" { - let _ = self.output.send(AgentEvent::ToolCall { - name: name.clone(), - args: args.clone(), - }); + Some(WaitAndContinueResult::Reset) => { + ctx.clear(); + *turn_count = 0; + *runtime_soul = self.sys_prompt()?; + ctx.update_soul(runtime_soul, &[]); + save_context_snapshot(&self.memory, ctx); + let _ = self.output.send(AgentEvent::Reset); + "interrupted by reset; context reset".to_string() } - - let wait_result = if name == "wait_and_continue" { - Some(execute_wait_and_continue(&args, &mut self.rx).await) - } else { - None - }; - let mut incoming_after_tool = None; - let mut suspend_after_wait = false; - let mut sent_local_content = None; - let result = if name == "local_send" { - match tools::local_send_content(&args) { - Ok(content) => { - let result = format!( - "sent {} chars to local operator", - content.chars().count() - ); - sent_local_content = Some(content); - result - } - Err(err) => { - tracing::warn!(err = %err, "local_send call had invalid arguments"); - format!("error: {err}") - } + Some(WaitAndContinueResult::UpdateSoul(content)) => { + self.memory.set_soul_text(&content)?; + *runtime_soul = self.sys_prompt()?; + let pinned = self.memory.pinned_memories().unwrap_or_default(); + ctx.update_soul(runtime_soul, &pinned); + save_context_snapshot(&self.memory, ctx); + let _ = self.output.send(AgentEvent::Status("soul updated".into())); + "interrupted by soul update; soul updated and continuing".to_string() + } + Some(WaitAndContinueResult::Compact) => { + let _ = self.output.send(AgentEvent::Status("compacting...".into())); + if let Err(e) = compact(tool_ctx, ctx, 0, 0, &self.output).await { + tracing::error!(err = %e, "compaction while waiting failed"); + format!("compaction failed while waiting: {e}") + } else { + *runtime_soul = self.sys_prompt()?; + refresh_context_soul(&self.memory, runtime_soul, ctx); + *turn_count = ctx.turn_count(); + "interrupted by compaction; compacted and continuing".to_string() } - } else { - match wait_result { - Some(WaitAndContinueResult::Timeout(mut result)) => { - suspend_after_wait = is_wait_timeout_without_event(&result); - if suspend_after_wait { - *consecutive_waits += 1; - if *consecutive_waits >= 5 && *consecutive_waits % 5 == 0 { - let nudge = format!( - "\n\nnudge: you have called wait_and_continue {} consecutive times without taking other actions. to conserve local resources, please increase the wait duration in your next wait_and_continue call (e.g. wait_and_continue(duration=\"10m\")).", - consecutive_waits - ); - result.push_str(&nudge); - } - } else { - *consecutive_waits = 0; + } + None => { + // Check for a preemptive interrupt before running the tool. + // We use try_recv here (not select!) so the tool future is never + // created and dropped — that would risk double-execution on retry. + match self.rx.try_recv() { + Ok(int) if int.is_preemptive() => { + tracing::info!( + source = int.source_tag(), + tool = name, + "preemptive interrupt before tool execution" + ); + let cancel_result = format!( + "tool execution cancelled by incoming {} message", + int.source_tag() + ); + if name != "local_send" { + let _ = self.output.send(AgentEvent::ToolResult { + name: name.clone(), + content: cancel_result.clone(), + }); } - result - } - Some(WaitAndContinueResult::Incoming { result, interrupt }) => { - incoming_after_tool = Some(interrupt); - result - } - Some(WaitAndContinueResult::Reset) => { - ctx.clear(); - *turn_count = 0; - *runtime_soul = self.sys_prompt()?; - ctx.update_soul(runtime_soul, &[]); - save_context_snapshot(&self.memory, ctx); - let _ = self.output.send(AgentEvent::Reset); - "interrupted by reset; context reset".to_string() - } - Some(WaitAndContinueResult::UpdateSoul(content)) => { - self.memory.set_soul_text(&content)?; - *runtime_soul = self.sys_prompt()?; - let pinned = self.memory.pinned_memories().unwrap_or_default(); - ctx.update_soul(runtime_soul, &pinned); - save_context_snapshot(&self.memory, ctx); - let _ = self.output.send(AgentEvent::Status("soul updated".into())); - "interrupted by soul update; soul updated and continuing" - .to_string() + ctx.push_tool_result(&call.id, &cancel_result); + let _ = self.memory.log_turn_with_tools( + "tool", + &cancel_result, + None, + None, + Some(&call.id), + ); + flush_pending_local_sends( + &self.memory, + ctx, + turn_count, + &mut pending_local_sends, + ); + return Ok(ToolIterationOutcome::Suspend(Some(int))); } - Some(WaitAndContinueResult::Compact) => { - let _ = - self.output.send(AgentEvent::Status("compacting...".into())); - if let Err(e) = compact(tool_ctx, ctx, 0, 0, &self.output).await { - tracing::error!(err = %e, "compaction while waiting failed"); - format!("compaction failed while waiting: {e}") - } else { - *runtime_soul = self.sys_prompt()?; - refresh_context_soul(&self.memory, runtime_soul, ctx); - *turn_count = ctx.turn_count(); - "interrupted by compaction; compacted and continuing" - .to_string() + Ok(int) => { + // non-preemptive — stash + if stashed_interrupt.is_none() { + *stashed_interrupt = Some(int); } } - None => { - // Check for a preemptive interrupt before running the tool. - // We use try_recv here (not select!) so the tool future is never - // created and dropped — that would risk double-execution on retry. - match self.rx.try_recv() { - Ok(int) if int.is_preemptive() => { - tracing::info!( - source = int.source_tag(), - tool = name, - "preemptive interrupt before tool execution" - ); - let cancel_result = format!( - "tool execution cancelled by incoming {} message", - int.source_tag() - ); - if name != "local_send" { - let _ = self.output.send(AgentEvent::ToolResult { - name: name.clone(), - content: cancel_result.clone(), - }); - } - ctx.push_tool_result(&call.id, &cancel_result); - let _ = self.memory.log_turn_with_tools( - "tool", - &cancel_result, - None, - None, - Some(&call.id), - ); - flush_pending_local_sends( - &self.memory, - ctx, - turn_count, - &mut pending_local_sends, - ); - self.finish_suspended_turn(ctx, turn_count, recalled_index) - .await; - return Ok(Some(int)); - } - Ok(int) => { - // non-preemptive — stash - if stashed_interrupt.is_none() { - stashed_interrupt = Some(int); - } - } - Err(_) => {} - } + Err(_) => {} + } - // Now run the tool, allowing preemptive interrupts during its execution. - let mut tool_output = None; - let mut interrupted_by = None; - let mut exec_fut = Box::pin(self.registry.execute(call, tool_ctx)); - loop { - tokio::select! { - res = &mut exec_fut => { - tool_output = Some(res); + // Now run the tool, allowing preemptive interrupts during its execution. + let mut tool_output = None; + let mut interrupted_by = None; + let mut exec_fut = Box::pin(self.registry.execute(call, tool_ctx)); + loop { + tokio::select! { + res = &mut exec_fut => { + tool_output = Some(res); + break; + } + int_opt = self.rx.recv() => { + match int_opt { + Some(int) if int.is_preemptive() => { + interrupted_by = Some(int); break; } - int_opt = self.rx.recv() => { - match int_opt { - Some(int) if int.is_preemptive() => { - interrupted_by = Some(int); - break; - } - Some(int) => { - if stashed_interrupt.is_none() { - stashed_interrupt = Some(int); - } - } - None => {} + Some(int) => { + if stashed_interrupt.is_none() { + *stashed_interrupt = Some(int); } } + None => {} } } + } + } - if let Some(int) = interrupted_by { - tracing::info!( - source = int.source_tag(), - tool = name, - "tool execution interrupted by incoming preemptive event" - ); - let cancel_result = format!( - "tool execution interrupted by incoming {} message", - int.source_tag() - ); - if name != "local_send" { - let _ = self.output.send(AgentEvent::ToolResult { - name: name.clone(), - content: cancel_result.clone(), - }); - } - ctx.push_tool_result(&call.id, &cancel_result); - let _ = self.memory.log_turn_with_tools( - "tool", - &cancel_result, - None, - None, - Some(&call.id), - ); - flush_pending_local_sends( - &self.memory, - ctx, - turn_count, - &mut pending_local_sends, - ); - self.finish_suspended_turn(ctx, turn_count, recalled_index) - .await; - return Ok(Some(int)); - } - - tool_output.unwrap_or_default() + if let Some(int) = interrupted_by { + tracing::info!( + source = int.source_tag(), + tool = name, + "tool execution interrupted by incoming preemptive event" + ); + let cancel_result = format!( + "tool execution interrupted by incoming {} message", + int.source_tag() + ); + if name != "local_send" { + let _ = self.output.send(AgentEvent::ToolResult { + name: name.clone(), + content: cancel_result.clone(), + }); } + ctx.push_tool_result(&call.id, &cancel_result); + let _ = self.memory.log_turn_with_tools( + "tool", + &cancel_result, + None, + None, + Some(&call.id), + ); + flush_pending_local_sends( + &self.memory, + ctx, + turn_count, + &mut pending_local_sends, + ); + return Ok(ToolIterationOutcome::Suspend(Some(int))); } - }; - if name != "local_send" { - let _ = self.output.send(AgentEvent::ToolResult { - name: name.clone(), - content: result.clone(), - }); + tool_output.unwrap_or_default() } + } + }; + + if name != "local_send" { + let _ = self.output.send(AgentEvent::ToolResult { + name: name.clone(), + content: result.clone(), + }); + } - ctx.push_tool_result(&call.id, &result); - let _ = self.memory.log_turn_with_tools( - "tool", - &result, - None, - None, - Some(&call.id), + ctx.push_tool_result(&call.id, &result); + let _ = self + .memory + .log_turn_with_tools("tool", &result, None, None, Some(&call.id)); + if let Some(content) = sent_local_content { + let _ = self.output.send(AgentEvent::Token(content.clone())); + pending_local_sends.push(content); + } + if let Some(interrupt) = incoming_after_tool { + flush_pending_local_sends(&self.memory, ctx, turn_count, &mut pending_local_sends); + return Ok(ToolIterationOutcome::Suspend(Some(interrupt))); + } + if suspend_after_wait { + flush_pending_local_sends(&self.memory, ctx, turn_count, &mut pending_local_sends); + return Ok(ToolIterationOutcome::Suspend(None)); + } + // Check for interrupts that arrived *during* tool execution. + // Preemptive ones cut the turn short; non-preemptive are stashed. + match self.rx.try_recv() { + Ok(int) if int.is_preemptive() => { + flush_pending_local_sends( + &self.memory, + ctx, + turn_count, + &mut pending_local_sends, ); - if let Some(content) = sent_local_content { - let _ = self.output.send(AgentEvent::Token(content.clone())); - pending_local_sends.push(content); + return Ok(ToolIterationOutcome::Suspend(Some(int))); + } + Ok(int) => { + if stashed_interrupt.is_none() { + *stashed_interrupt = Some(int); } - if let Some(interrupt) = incoming_after_tool { - flush_pending_local_sends( - &self.memory, - ctx, - turn_count, - &mut pending_local_sends, - ); + } + Err(_) => {} + } + } + + flush_pending_local_sends(&self.memory, ctx, turn_count, &mut pending_local_sends); + if should_restart { + Ok(ToolIterationOutcome::Restart) + } else { + Ok(ToolIterationOutcome::Continue) + } + } + + async fn run_turn( + &mut self, + llm: &LlmClient, + tool_ctx: &ToolContext, + ctx: &mut Context, + turn_count: &mut usize, + runtime_soul: &mut String, + recalled_index: Option, + is_local_message: bool, + is_discord_event: bool, + consecutive_waits: &mut usize, + ) -> Result> { + let mut tool_iterations = 0usize; + let mut discord_send_called = false; + let mut discord_nudge_injected = false; + let mut local_send_called = false; + let mut local_send_nudge_injected = false; + // non-preemptive interrupt received during streaming/tool execution; returned at end of turn + let mut stashed_interrupt: Option = None; + const MAX_TOOL_ITERATIONS: usize = 20; + + // emit Started once — all tokens/thinking for the whole turn go into one bubble + let _ = self.output.send(AgentEvent::Started); + + loop { + let StreamIteration { + response: iteration_response, + thinking: iteration_thinking, + tool_calls, + stream_error, + cancelled_by: stream_cancelled_by, + } = self + .stream_iteration(llm, ctx, &mut stashed_interrupt) + .await; + + if let Some(interrupt) = stream_cancelled_by { + // turn was preempted mid-stream — finish cleanly and hand off + let _ = self.output.send(AgentEvent::Done); + cleanup_recalled_context(ctx, recalled_index); + save_context_snapshot(&self.memory, ctx); + return Ok(Some(interrupt)); + } + + if !tool_calls.is_empty() && tool_iterations < MAX_TOOL_ITERATIONS { + tool_iterations += 1; + let outcome = self + .execute_tool_iteration( + tool_ctx, + ctx, + turn_count, + runtime_soul, + &tool_calls, + &iteration_thinking, + &mut stashed_interrupt, + consecutive_waits, + &mut discord_send_called, + &mut local_send_called, + ) + .await?; + match outcome { + ToolIterationOutcome::Continue => {} + ToolIterationOutcome::Suspend(interrupt) => { self.finish_suspended_turn(ctx, turn_count, recalled_index) .await; - return Ok(Some(interrupt)); + return Ok(interrupt); } - if suspend_after_wait { - flush_pending_local_sends( - &self.memory, - ctx, - turn_count, - &mut pending_local_sends, - ); + ToolIterationOutcome::Restart => { self.finish_suspended_turn(ctx, turn_count, recalled_index) .await; - return Ok(None); + tracing::info!("restart harness requested, performing restart"); + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + perform_restart(); } - // Check for interrupts that arrived *during* tool execution. - // Preemptive ones cut the turn short; non-preemptive are stashed. - match self.rx.try_recv() { - Ok(int) if int.is_preemptive() => { - flush_pending_local_sends( - &self.memory, - ctx, - turn_count, - &mut pending_local_sends, - ); - self.finish_suspended_turn(ctx, turn_count, recalled_index) - .await; - return Ok(Some(int)); - } - Ok(int) => { - if stashed_interrupt.is_none() { - stashed_interrupt = Some(int); - } - } - Err(_) => {} - } - } - flush_pending_local_sends(&self.memory, ctx, turn_count, &mut pending_local_sends); - if should_restart { - self.finish_suspended_turn(ctx, turn_count, recalled_index) - .await; - tracing::info!("restart harness requested, performing restart"); - tokio::time::sleep(std::time::Duration::from_millis(500)).await; - perform_restart(); } + *runtime_soul = self.sys_prompt()?; refresh_context_soul(&self.memory, runtime_soul, ctx); continue; -- 2.51.2