diff --git a/.beads/interactions.jsonl b/.beads/interactions.jsonl index a7b179f..4d362bc 100644 --- a/.beads/interactions.jsonl +++ b/.beads/interactions.jsonl @@ -91,3 +91,4 @@ {"id":"int-623e8776","kind":"field_change","created_at":"2026-07-01T08:54:40.646821049Z","actor":"dawn","issue_id":"klbr-2kc","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Collapsed KDL optional-node parser helpers and centralized serde integer deserializers"}} {"id":"int-4cb60a18","kind":"field_change","created_at":"2026-07-01T09:02:05.887375856Z","actor":"dawn","issue_id":"klbr-3ha","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Refactored shared SSE and media data-url parsing in models.rs; added focused coverage and verified klbr-core tests."}} {"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."}} diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index edccb3e..5ad0f8e 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -53,7 +53,7 @@ {"_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-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":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:23Z","created_by":"dawn","updated_at":"2026-06-30T21:39:23Z","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":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:10Z","created_by":"dawn","updated_at":"2026-06-30T21:39:10Z","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-6m2","title":"web: refactor and split App.svelte monolith","description":"App.svelte is currently a 3000+ line file holding all UI, WS client, and styles. Needs to be split into modular components like Sidebar, ChatPane, Metrics, and XmlVisualizer.","status":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:03Z","created_by":"dawn","updated_at":"2026-06-30T21:39:03Z","dependency_count":0,"dependent_count":0,"comment_count":0} diff --git a/klbr-core/src/agent.rs b/klbr-core/src/agent.rs index 6d5027e..4c50b9d 100644 --- a/klbr-core/src/agent.rs +++ b/klbr-core/src/agent.rs @@ -19,6 +19,8 @@ use crate::{ AgentEvent, AgentMetrics, CompactionRecord, }; +const MAX_CONSECUTIVE_STREAM_FAILURES: usize = 5; + pub struct Agent { config: Config, memory: MemoryStore, @@ -198,7 +200,7 @@ impl Agent { } else { last_event_at = std::time::Instant::now(); skip_idle_nudge = true; - if failure_count >= 5 { + if failure_count >= MAX_CONSECUTIVE_STREAM_FAILURES { tracing::info!( "cleaning up failed/problematic turns due to new message after errors" ); @@ -239,32 +241,14 @@ impl Agent { ) .await?; - let stream_failed = ctx - .turns - .last() - .map(|m| { - m.role == Role::Assistant && m.content_str().starts_with("error:") - }) - .unwrap_or(false); - if stream_failed { - failure_count += 1; - if failure_count < 5 { - let backoff_secs = 2u64.pow(failure_count as u32); - tracing::warn!("llm stream failed; backoff sleeping for {} seconds (failure {}/5)", backoff_secs, failure_count); - tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)) - .await; - } else { - tracing::error!("llm stream failed 5 times consecutively; stopping automatic retries and waiting for new message/event"); - break; - } - } else { - failure_count = 0; + if handle_stream_failure_backoff(&ctx, &mut failure_count).await { + break; } } } } else { // Continuous execution: run a turn without any new input if there is a pending message needing response - if failure_count >= 5 { + if failure_count >= MAX_CONSECUTIVE_STREAM_FAILURES { tokio::time::sleep(std::time::Duration::from_millis(500)).await; continue; } @@ -284,27 +268,7 @@ impl Agent { ) .await?; - let stream_failed = ctx - .turns - .last() - .map(|m| m.role == Role::Assistant && m.content_str().starts_with("error:")) - .unwrap_or(false); - if stream_failed { - failure_count += 1; - if failure_count < 5 { - let backoff_secs = 2u64.pow(failure_count as u32); - tracing::warn!( - "llm stream failed; backoff sleeping for {} seconds (failure {}/5)", - backoff_secs, - failure_count - ); - tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)).await; - } else { - tracing::error!("llm stream failed 5 times consecutively; stopping automatic retries and waiting for new message/event"); - } - } else { - failure_count = 0; - } + handle_stream_failure_backoff(&ctx, &mut failure_count).await; while let Some(curr_int) = next_interrupt { if !matches!( @@ -337,26 +301,8 @@ impl Agent { ) .await?; - let stream_failed = ctx - .turns - .last() - .map(|m| { - m.role == Role::Assistant && m.content_str().starts_with("error:") - }) - .unwrap_or(false); - if stream_failed { - failure_count += 1; - if failure_count < 5 { - let backoff_secs = 2u64.pow(failure_count as u32); - tracing::warn!("llm stream failed; backoff sleeping for {} seconds (failure {}/5)", backoff_secs, failure_count); - tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)) - .await; - } else { - tracing::error!("llm stream failed 5 times consecutively; stopping automatic retries and waiting for new message/event"); - break; - } - } else { - failure_count = 0; + if handle_stream_failure_backoff(&ctx, &mut failure_count).await { + break; } } } else { @@ -384,27 +330,7 @@ impl Agent { ) .await?; - let stream_failed = ctx - .turns - .last() - .map(|m| m.role == Role::Assistant && m.content_str().starts_with("error:")) - .unwrap_or(false); - if stream_failed { - failure_count += 1; - if failure_count < 5 { - let backoff_secs = 2u64.pow(failure_count as u32); - tracing::warn!( - "llm stream failed; backoff sleeping for {} seconds (failure {}/5)", - backoff_secs, - failure_count - ); - tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)).await; - } else { - tracing::error!("llm stream failed 5 times consecutively; stopping automatic retries and waiting for new message/event"); - } - } else { - failure_count = 0; - } + handle_stream_failure_backoff(&ctx, &mut failure_count).await; while let Some(curr_int) = next_interrupt { if !matches!( @@ -437,26 +363,8 @@ impl Agent { ) .await?; - let stream_failed = ctx - .turns - .last() - .map(|m| { - m.role == Role::Assistant && m.content_str().starts_with("error:") - }) - .unwrap_or(false); - if stream_failed { - failure_count += 1; - if failure_count < 5 { - let backoff_secs = 2u64.pow(failure_count as u32); - tracing::warn!("llm stream failed; backoff sleeping for {} seconds (failure {}/5)", backoff_secs, failure_count); - tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)) - .await; - } else { - tracing::error!("llm stream failed 5 times consecutively; stopping automatic retries and waiting for new message/event"); - break; - } - } else { - failure_count = 0; + if handle_stream_failure_backoff(&ctx, &mut failure_count).await { + break; } } } @@ -1262,6 +1170,38 @@ fn idle_nudge_message(since_last_event: std::time::Duration) -> String { ) } +fn last_turn_is_stream_error(ctx: &Context) -> bool { + ctx.turns + .last() + .is_some_and(|m| m.role == Role::Assistant && m.content_str().starts_with("error:")) +} + +async fn handle_stream_failure_backoff(ctx: &Context, failure_count: &mut usize) -> bool { + if !last_turn_is_stream_error(ctx) { + *failure_count = 0; + return false; + } + + *failure_count += 1; + if *failure_count < MAX_CONSECUTIVE_STREAM_FAILURES { + let backoff_secs = 2u64.pow(*failure_count as u32); + tracing::warn!( + "llm stream failed; backoff sleeping for {} seconds (failure {}/{})", + backoff_secs, + *failure_count, + MAX_CONSECUTIVE_STREAM_FAILURES + ); + tokio::time::sleep(std::time::Duration::from_secs(backoff_secs)).await; + false + } else { + tracing::error!( + "llm stream failed {} times consecutively; stopping automatic retries and waiting for new message/event", + MAX_CONSECUTIVE_STREAM_FAILURES + ); + true + } +} + fn clean_failed_turns(ctx: &mut Context) { while let Some(last) = ctx.turns.last() { let content = last.content_str();