diff --git a/.beads/interactions.jsonl b/.beads/interactions.jsonl index 8af12a8..bf96f73 100644 --- a/.beads/interactions.jsonl +++ b/.beads/interactions.jsonl @@ -89,3 +89,4 @@ {"id":"int-5c780dfb","kind":"field_change","created_at":"2026-06-30T22:34:06.198966004Z","actor":"dawn","issue_id":"klbr-f8w","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Split Discord formatting, inbox storage, and tool schema builders into modules and deduplicated batch row rendering prep"}} {"id":"int-1f1c3a07","kind":"field_change","created_at":"2026-07-01T08:52:06.275727261Z","actor":"dawn","issue_id":"klbr-5h3","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Re-exported core DTO types from klbr-ipc and removed daemon field-by-field mappers"}} {"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."}} diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index 9076b99..2c9a22d 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -50,7 +50,7 @@ {"_type":"issue","id":"klbr-5h3","title":"ipc/daemon: eliminate duplicate DTO structs and field mapping boilerplate","description":"HistoryEntry, ToolCall, CompactionRecord, and ResolutionEventDto are identical duplicates between klbr-ipc and klbr-core. The daemon has extensive boilerplate mapping them field-by-field. Share or re-export these types.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:40:01Z","created_by":"dawn","updated_at":"2026-07-01T08:52:06Z","started_at":"2026-06-30T22:34:30Z","closed_at":"2026-07-01T08:52:06Z","close_reason":"Re-exported core DTO types from klbr-ipc and removed daemon field-by-field mappers","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-2kc","title":"core: reduce KDL config parser boilerplate in parser.rs","description":"parser.rs contains ~20 copy-pasted optional_*_node helpers that share identical structure. Simplify with macros or generic helpers. Also resolve duplicate deserialize_usize.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:55Z","created_by":"dawn","updated_at":"2026-07-01T08:54:41Z","started_at":"2026-07-01T08:52:25Z","closed_at":"2026-07-01T08:54:41Z","close_reason":"Collapsed KDL optional-node parser helpers and centralized serde integer deserializers","dependency_count":0,"dependent_count":0,"comment_count":0} {"_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":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:42Z","created_by":"dawn","updated_at":"2026-06-30T21:39:42Z","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":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:36Z","created_by":"dawn","updated_at":"2026-06-30T21:39:36Z","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} diff --git a/klbr-core/src/models.rs b/klbr-core/src/models.rs index 1b24382..0151d80 100644 --- a/klbr-core/src/models.rs +++ b/klbr-core/src/models.rs @@ -677,12 +677,7 @@ impl LlmClient { tracing::debug!("stream line read: {}", line); - // Standard SSE uses "data: " but some providers (e.g. Kimi) - // omit the space and send "data:". Accept both. - let Some(data) = line - .strip_prefix("data: ") - .or_else(|| line.strip_prefix("data:")) - else { + let Some(data) = sse_data_payload(&line) else { continue; }; if data == "[DONE]" { @@ -903,55 +898,7 @@ impl LlmClient { let text = Self::error_for_status_with_body(v).await?.text().await?; if text.trim().starts_with("data:") { - let mut content = String::new(); - let mut completion_tokens = 0; - let mut prompt_tokens = 0; - let mut total_tokens = 0; - for line in text.lines() { - let line = line.trim(); - if line.starts_with("data:") { - let data_part = line[5..].trim(); - if data_part == "[DONE]" { - continue; - } - if let Ok(val) = serde_json::from_str::(data_part) { - if let Some(choices) = val["choices"].as_array() { - if let Some(choice) = choices.get(0) { - if let Some(delta) = choice.get("delta") { - if let Some(chunk) = delta["content"].as_str() { - content.push_str(chunk); - } - } - if let Some(message) = choice.get("message") { - if let Some(text_content) = message["content"].as_str() { - content.push_str(text_content); - } - } - } - } - if let Some(usage) = val.get("usage") { - if let Some(ct) = usage["completion_tokens"].as_u64() { - completion_tokens = ct as usize; - } - if let Some(pt) = usage["prompt_tokens"].as_u64() { - prompt_tokens = pt as usize; - } - if let Some(tt) = usage["total_tokens"].as_u64() { - total_tokens = tt as usize; - } - } - } - } - } - Ok(( - content, - Usage { - prompt_tokens, - completion_tokens, - total_tokens, - ..Default::default() - }, - )) + Ok(complete_from_sse_text(&text)) } else { let v: Value = serde_json::from_str(&text)?; let content = v["choices"][0]["message"]["content"] @@ -1019,36 +966,7 @@ impl LlmClient { } fn strip_media_urls(text: &str) -> String { - let mut result = String::new(); - let mut current_idx = 0; - while let Some(start_idx) = text[current_idx..].find("data:") { - let abs_start = current_idx + start_idx; - if let Some(comma_offset) = text[abs_start..].find(";base64,") { - if comma_offset < 100 { - let in_between = &text[abs_start + 5..abs_start + comma_offset]; - if !in_between.contains("data:") { - let base64_start = abs_start + comma_offset + 8; - let mut base64_end = base64_start; - for c in text[base64_start..].chars() { - if c.is_ascii_alphanumeric() || c == '+' || c == '/' || c == '=' { - base64_end += c.len_utf8(); - } else { - break; - } - } - if base64_end > base64_start { - result.push_str(&text[current_idx..abs_start]); - current_idx = base64_end; - continue; - } - } - } - } - result.push_str(&text[current_idx..abs_start + 5]); - current_idx = abs_start + 5; - } - result.push_str(&text[current_idx..]); - result + rewrite_media_data_urls(text, |_| Some(String::new())) } async fn embed_with_config( @@ -1503,47 +1421,137 @@ impl LlmClient { } } -fn replace_media_urls_with_placeholders(text: &str, limit_bytes: usize) -> String { - let mut result = String::new(); - let mut current_idx = 0; - while let Some(start_idx) = text[current_idx..].find("data:") { - let abs_start = current_idx + start_idx; - if let Some(comma_offset) = text[abs_start..].find(";base64,") { - if comma_offset < 100 { - let mime = &text[abs_start + 5..abs_start + comma_offset]; - if !mime.contains("data:") { - let base64_start = abs_start + comma_offset + 8; - let mut base64_end = base64_start; - for c in text[base64_start..].chars() { - if c.is_ascii_alphanumeric() || c == '+' || c == '/' || c == '=' { - base64_end += c.len_utf8(); - } else { - break; - } +fn sse_data_payload(line: &str) -> Option<&str> { + let line = line.trim(); + // Standard SSE uses "data: " but some providers omit the space. + line.strip_prefix("data: ") + .or_else(|| line.strip_prefix("data:")) + .map(str::trim) +} + +fn complete_from_sse_text(text: &str) -> (String, Usage) { + let mut content = String::new(); + let mut usage = Usage::default(); + + for data in text.lines().filter_map(sse_data_payload) { + if data == "[DONE]" { + continue; + } + let Ok(val) = serde_json::from_str::(data) else { + continue; + }; + if let Some(choices) = val["choices"].as_array() { + if let Some(choice) = choices.first() { + if let Some(delta) = choice.get("delta") { + if let Some(chunk) = delta["content"].as_str() { + content.push_str(chunk); } - if base64_end > base64_start { - let b64_len = base64_end - base64_start; - let approx_size = (b64_len * 3) / 4; - if approx_size > limit_bytes { - result.push_str(&text[current_idx..abs_start]); - result.push_str(&format!( - "[Media attachment: {}, approximate size: {} bytes]", - mime, approx_size - )); - current_idx = base64_end; - continue; - } + } + if let Some(message) = choice.get("message") { + if let Some(text_content) = message["content"].as_str() { + content.push_str(text_content); } } } } - result.push_str(&text[current_idx..abs_start + 5]); - current_idx = abs_start + 5; + if let Some(value) = val.get("usage") { + usage.prompt_tokens = value["prompt_tokens"].as_u64().unwrap_or(0) as usize; + usage.completion_tokens = value["completion_tokens"].as_u64().unwrap_or(0) as usize; + usage.total_tokens = value["total_tokens"].as_u64().unwrap_or(0) as usize; + } + } + + (content, usage) +} + +struct MediaDataUrl<'a> { + start: usize, + end: usize, + mime: &'a str, + approx_size: usize, +} + +struct MediaScan<'a> { + data_start: usize, + media: Option>, +} + +fn next_media_scan(text: &str, current_idx: usize) -> Option> { + let data_start = current_idx + text[current_idx..].find("data:")?; + Some(MediaScan { + data_start, + media: parse_media_data_url_at(text, data_start), + }) +} + +fn parse_media_data_url_at(text: &str, start: usize) -> Option> { + let comma_offset = text[start..].find(";base64,")?; + if comma_offset >= 100 { + return None; + } + let mime = &text[start + "data:".len()..start + comma_offset]; + if mime.contains("data:") { + return None; + } + + let base64_start = start + comma_offset + ";base64,".len(); + let mut base64_end = base64_start; + for ch in text[base64_start..].chars() { + if ch.is_ascii_alphanumeric() || ch == '+' || ch == '/' || ch == '=' { + base64_end += ch.len_utf8(); + } else { + break; + } + } + if base64_end == base64_start { + return None; + } + + let base64_len = base64_end - base64_start; + Some(MediaDataUrl { + start, + end: base64_end, + mime, + approx_size: (base64_len * 3) / 4, + }) +} + +fn rewrite_media_data_urls( + text: &str, + mut replacement: impl FnMut(&MediaDataUrl<'_>) -> Option, +) -> String { + let mut result = String::new(); + let mut current_idx = 0; + while let Some(scan) = next_media_scan(text, current_idx) { + if let Some(media) = scan.media { + if let Some(value) = replacement(&media) { + result.push_str(&text[current_idx..media.start]); + result.push_str(&value); + current_idx = media.end; + continue; + } + } + + result.push_str(&text[current_idx..scan.data_start + "data:".len()]); + current_idx = scan.data_start + "data:".len(); } result.push_str(&text[current_idx..]); result } +fn replace_media_urls_with_placeholders(text: &str, limit_bytes: usize) -> String { + rewrite_media_data_urls(text, |media| { + if media.approx_size > limit_bytes { + Some(format!( + "[Media attachment: {}, approximate size: {} bytes]", + media.mime, media.approx_size + )) + } else { + None + } + }) +} + fn prepare_messages_for_api(messages: &[Message]) -> Vec { let mut prepared = Vec::new(); for mut message in messages.iter().cloned() { @@ -1625,6 +1633,39 @@ mod tests { ); } + #[test] + fn media_placeholders_preserve_small_attachments() { + let input = concat!( + "small data:image/png;base64,YWJj ", + "large data:audio/mp3;base64,YWJjZGVm" + ); + + assert_eq!( + replace_media_urls_with_placeholders(input, 4), + concat!( + "small data:image/png;base64,YWJj ", + "large [Media attachment: audio/mp3, approximate size: 6 bytes]" + ) + ); + } + + #[test] + fn complete_from_sse_text_accepts_sse_spacing_variants() { + let text = concat!( + "data:{\"choices\":[{\"delta\":{\"content\":\"hel\"}}]}\n", + "data: {\"choices\":[{\"message\":{\"content\":\"lo\"}}],", + "\"usage\":{\"prompt_tokens\":2,\"completion_tokens\":3,\"total_tokens\":5}}\n", + "data: [DONE]\n" + ); + + let (content, usage) = complete_from_sse_text(text); + + assert_eq!(content, "hello"); + assert_eq!(usage.prompt_tokens, 2); + assert_eq!(usage.completion_tokens, 3); + assert_eq!(usage.total_tokens, 5); + } + async fn spawn_http_response(status: &str, content_type: &str, body: &str) -> String { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap();