diff --git a/flake.lock b/flake.lock index b910ac4..baa3c9f 100644 --- a/flake.lock +++ b/flake.lock @@ -1,17 +1,113 @@ { "nodes": { + "crane": { + "flake": false, + "locked": { + "lastModified": 1758758545, + "narHash": "sha256-NU5WaEdfwF6i8faJ2Yh+jcK9vVFrofLcwlD/mP65JrI=", + "owner": "ipetkov", + "repo": "crane", + "rev": "95d528a5f54eaba0d12102249ce42f4d01f4e364", + "type": "github" + }, + "original": { + "owner": "ipetkov", + "ref": "v0.21.1", + "repo": "crane", + "type": "github" + } + }, + "dream2nix": { + "inputs": { + "nixpkgs": [ + "nci", + "nixpkgs" + ], + "purescript-overlay": "purescript-overlay", + "pyproject-nix": "pyproject-nix" + }, + "locked": { + "lastModified": 1765953015, + "narHash": "sha256-5FBZbbWR1Csp3Y2icfRkxMJw/a/5FGg8hCXej2//bbI=", + "owner": "nix-community", + "repo": "dream2nix", + "rev": "69eb01fa0995e1e90add49d8ca5bcba213b0416f", + "type": "github" + }, + "original": { + "owner": "nix-community", + "repo": "dream2nix", + "type": "github" + } + }, + "flake-compat": { + "flake": false, + "locked": { + "lastModified": 1696426674, + "narHash": "sha256-kvjfFW7WAETZlt09AgDn1MrtKzP7t90Vf7vypd3OL1U=", + "owner": "edolstra", + "repo": "flake-compat", + "rev": "0f9255e01c2351cc7d116c072cb317785dd33b33", + "type": "github" + }, + "original": { + "owner": "edolstra", + "repo": "flake-compat", + "type": "github" + } + }, + "mk-naked-shell": { + "flake": false, + "locked": { + "lastModified": 1681286841, + "narHash": "sha256-3XlJrwlR0nBiREnuogoa5i1b4+w/XPe0z8bbrJASw0g=", + "owner": "90-008", + "repo": "mk-naked-shell", + "rev": "7612f828dd6f22b7fb332cc69440e839d7ffe6bd", + "type": "github" + }, + "original": { + "owner": "90-008", + "repo": "mk-naked-shell", + "type": "github" + } + }, + "nci": { + "inputs": { + "crane": "crane", + "dream2nix": "dream2nix", + "mk-naked-shell": "mk-naked-shell", + "nixpkgs": "nixpkgs", + "parts": "parts", + "rust-overlay": "rust-overlay", + "treefmt": "treefmt" + }, + "locked": { + "lastModified": 1775459945, + "narHash": "sha256-52p03BHauMmZyYYhx1QFtPooDwL6YxDtFMFNDtqB5Oo=", + "owner": "90-008", + "repo": "nix-cargo-integration", + "rev": "e53db757c693e2f1265a3b680f5360c3f528b7a7", + "type": "github" + }, + "original": { + "owner": "90-008", + "repo": "nix-cargo-integration", + "type": "github" + } + }, "nixpkgs": { "locked": { - "lastModified": 1775403759, - "narHash": "sha256-cGyKiTspHEUx3QwAnV3RfyT+VOXhHLs+NEr17HU34Wo=", - "owner": "nixos", + "lastModified": 1775036866, + "narHash": "sha256-ZojAnPuCdy657PbTq5V0Y+AHKhZAIwSIT2cb8UgAz/U=", + "owner": "NixOS", "repo": "nixpkgs", - "rev": "5e11f7acce6c3469bef9df154d78534fa7ae8b6c", + "rev": "6201e203d09599479a3b3450ed24fa81537ebc4e", "type": "github" }, "original": { - "owner": "nixos", - "ref": "nixpkgs-unstable", + "owner": "NixOS", + "ref": "nixos-unstable", "repo": "nixpkgs", "type": "github" } @@ -31,7 +127,44 @@ "type": "github" } }, + "nixpkgs_2": { + "locked": { + "lastModified": 1775403759, + "narHash": "sha256-cGyKiTspHEUx3QwAnV3RfyT+VOXhHLs+NEr17HU34Wo=", + "owner": "nixos", + "repo": "nixpkgs", + "rev": "5e11f7acce6c3469bef9df154d78534fa7ae8b6c", + "type": "github" + }, + "original": { + "owner": "nixos", + "ref": "nixpkgs-unstable", + "repo": "nixpkgs", + "type": "github" + } + }, "parts": { + "inputs": { + "nixpkgs-lib": [ + "nci", + "nixpkgs" + ] + }, + "locked": { + "lastModified": 1775087534, + "narHash": "sha256-91qqW8lhL7TLwgQWijoGBbiD4t7/q75KTi8NxjVmSmA=", + "owner": "hercules-ci", + "repo": "flake-parts", + "rev": "3107b77cd68437b9a76194f0f7f9c55f2329ca5b", + "type": "github" + }, + "original": { + "owner": "hercules-ci", + "repo": "flake-parts", + "type": "github" + } + }, + "parts_2": { "inputs": { "nixpkgs-lib": "nixpkgs-lib" }, @@ -49,10 +182,122 @@ "type": "github" } }, + "purescript-overlay": { + "inputs": { + "flake-compat": "flake-compat", + "nixpkgs": [ + "nci", + "dream2nix", + "nixpkgs" + ], + "slimlock": "slimlock" + }, + "locked": { + "lastModified": 1728546539, + "narHash": "sha256-Sws7w0tlnjD+Bjck1nv29NjC5DbL6nH5auL9Ex9Iz2A=", + "owner": "thomashoneyman", + "repo": "purescript-overlay", + "rev": "4ad4c15d07bd899d7346b331f377606631eb0ee4", + "type": "github" + }, + "original": { + "owner": "thomashoneyman", + "repo": "purescript-overlay", + "type": "github" + } + }, + "pyproject-nix": { + "inputs": { + "nixpkgs": [ + "nci", + "dream2nix", + "nixpkgs" + ] + }, + "locked": { + "lastModified": 1763017646, + "narHash": "sha256-Z+R2lveIp6Skn1VPH3taQIuMhABg1IizJd8oVdmdHsQ=", + "owner": "pyproject-nix", + "repo": "pyproject.nix", + "rev": "47bd6f296502842643078d66128f7b5e5370790c", + "type": "github" + }, + "original": { + "owner": "pyproject-nix", + "repo": "pyproject.nix", + "type": "github" + } + }, "root": { "inputs": { - "nixpkgs": "nixpkgs", - "parts": "parts" + "nci": "nci", + "nixpkgs": "nixpkgs_2", + "parts": "parts_2" + } + }, + "rust-overlay": { + "inputs": { + "nixpkgs": [ + "nci", + "nixpkgs" + ] + }, + "locked": { + "lastModified": 1775445266, + "narHash": "sha256-3fgIj85WHQbOamrpIw9WY3ZL1PoEvjPOjmzMYNsEQJo=", + "owner": "oxalica", + "repo": "rust-overlay", + "rev": "61747bc3cf2da179b9af356ae9e70f3b895b0c24", + "type": "github" + }, + "original": { + "owner": "oxalica", + "repo": "rust-overlay", + "type": "github" + } + }, + "slimlock": { + "inputs": { + "nixpkgs": [ + "nci", + "dream2nix", + "purescript-overlay", + "nixpkgs" + ] + }, + "locked": { + "lastModified": 1688756706, + "narHash": "sha256-xzkkMv3neJJJ89zo3o2ojp7nFeaZc2G0fYwNXNJRFlo=", + "owner": "thomashoneyman", + "repo": "slimlock", + "rev": "cf72723f59e2340d24881fd7bf61cb113b4c407c", + "type": "github" + }, + "original": { + "owner": "thomashoneyman", + "repo": "slimlock", + "type": "github" + } + }, + "treefmt": { + "inputs": { + "nixpkgs": [ + "nci", + "nixpkgs" + ] + }, + "locked": { + "lastModified": 1775125835, + "narHash": "sha256-2qYcPgzFhnQWchHo0SlqLHrXpux5i6ay6UHA+v2iH4U=", + "owner": "numtide", + "repo": "treefmt-nix", + "rev": "75925962939880974e3ab417879daffcba36c4a3", + "type": "github" + }, + "original": { + "owner": "numtide", + "repo": "treefmt-nix", + "type": "github" } } }, diff --git a/flake.nix b/flake.nix index b4c496e..261d450 100644 --- a/flake.nix +++ b/flake.nix @@ -1,11 +1,13 @@ { inputs.parts.url = "github:hercules-ci/flake-parts"; inputs.nixpkgs.url = "github:nixos/nixpkgs/nixpkgs-unstable"; + inputs.nci.url = "github:90-008/nix-cargo-integration"; outputs = inp: inp.parts.lib.mkFlake { inputs = inp; } { systems = [ "x86_64-linux" ]; + imports = [inp.nci.flakeModule]; perSystem = { pkgs, @@ -13,21 +15,15 @@ ... }: { + nci.projects."klbr".path = ./.; packages.default = pkgs.callPackage ./default.nix {}; - devShells = { - default = pkgs.mkShell { - packages = with pkgs; [ - rustPlatform.rustLibSrc - rust-analyzer - cargo - cargo-outdated - rustc - rustfmt - clang - wild - ]; - }; - }; + devShells.default = config.nci.outputs."klbr".devShell.overrideAttrs (old: { + packages = (old.packages or []) ++ (with pkgs; [ + cargo-outdated + clang + wild + ]); + }); }; }; } diff --git a/klbr-core/src/agent.rs b/klbr-core/src/agent.rs index ea8cbb8..6c22dd2 100644 --- a/klbr-core/src/agent.rs +++ b/klbr-core/src/agent.rs @@ -26,7 +26,10 @@ pub async fn run( // resume: replay the sliding window from the last run if let Ok(prior) = memory.recent_turns(config.compaction_keep) { - let pairs: Vec<(String, String)> = prior.into_iter().map(|(_, r, c)| (r, c)).collect(); + let pairs: Vec<(String, String)> = prior + .into_iter() + .map(|entry| (entry.role, entry.content)) + .collect(); ctx.load_turns(&pairs); turn_count = pairs.len(); } @@ -66,6 +69,7 @@ pub async fn run( tokio::spawn(async move { let _ = llm2.stream(&msgs, tok_tx).await; }); + let _ = output.send(ServerMsg::Started); let mut response = String::new(); let mut thinking = String::new(); diff --git a/klbr-core/src/daemon.rs b/klbr-core/src/daemon.rs index 4a25deb..e1b0b45 100644 --- a/klbr-core/src/daemon.rs +++ b/klbr-core/src/daemon.rs @@ -89,7 +89,7 @@ async fn handle( ClientMsg::FetchHistory { before_id, limit } => { let limit = limit.min(HISTORY_PAGE); let turns = memory.turns_before(before_id, limit).unwrap_or_default(); - send_msg(&mut sock_tx, &ServerMsg::HistoryPage { turns }).await?; + send_msg(&mut sock_tx, &ServerMsg::History { turns }).await?; } } } diff --git a/klbr-core/src/memory.rs b/klbr-core/src/memory.rs index 939aebc..9164069 100644 --- a/klbr-core/src/memory.rs +++ b/klbr-core/src/memory.rs @@ -1,4 +1,5 @@ use anyhow::Result; +use klbr_ipc::HistoryEntry; use rusqlite::{ffi::sqlite3_auto_extension, params, Connection}; use sqlite_vec::sqlite3_vec_init; use std::sync::{Arc, Mutex}; @@ -103,18 +104,20 @@ impl MemoryStore { } /// last `n` turns in chronological order (oldest first, ready to replay into context). - pub fn recent_turns(&self, n: usize) -> Result> { + pub fn recent_turns(&self, n: usize) -> Result> { let conn = self.0.lock().unwrap(); let mut stmt = conn.prepare( - "SELECT id, role, content FROM (SELECT id, role, content, ts FROM turns ORDER BY ts DESC LIMIT ?1) ORDER BY ts ASC", + "SELECT id, role, content, thinking, ts FROM (SELECT id, role, content, thinking, ts FROM turns ORDER BY ts DESC LIMIT ?1) ORDER BY ts ASC", )?; let results = stmt .query_map(params![n as i64], |row| { - Ok(( - row.get::<_, i64>(0)?, - row.get::<_, String>(1)?, - row.get::<_, String>(2)?, - )) + Ok(HistoryEntry { + id: row.get(0)?, + role: row.get(1)?, + content: row.get(2)?, + reasoning: row.get(3)?, + timestamp: row.get(4)?, + }) })? .filter_map(|r| r.ok()) .collect(); @@ -122,18 +125,20 @@ impl MemoryStore { } /// turns older than `before_id`, newest-first then reversed, for scroll-back paging. - pub fn turns_before(&self, before_id: i64, limit: usize) -> Result> { + pub fn turns_before(&self, before_id: i64, limit: usize) -> Result> { let conn = self.0.lock().unwrap(); let mut stmt = conn.prepare( - "SELECT id, role, content FROM (SELECT id, role, content FROM turns WHERE id < ?1 ORDER BY id DESC LIMIT ?2) ORDER BY id ASC", + "SELECT id, role, content, thinking, ts FROM (SELECT id, role, content, thinking, ts FROM turns WHERE id < ?1 ORDER BY id DESC LIMIT ?2) ORDER BY id ASC", )?; let results = stmt .query_map(params![before_id, limit as i64], |row| { - Ok(( - row.get::<_, i64>(0)?, - row.get::<_, String>(1)?, - row.get::<_, String>(2)?, - )) + Ok(HistoryEntry { + id: row.get(0)?, + role: row.get(1)?, + content: row.get(2)?, + reasoning: row.get(3)?, + timestamp: row.get(4)?, + }) })? .filter_map(|r| r.ok()) .collect(); diff --git a/klbr-ipc/src/lib.rs b/klbr-ipc/src/lib.rs index 97fd19c..123fe33 100644 --- a/klbr-ipc/src/lib.rs +++ b/klbr-ipc/src/lib.rs @@ -9,10 +9,21 @@ pub enum ClientMsg { FetchHistory { before_id: i64, limit: usize }, } +#[derive(Debug, Serialize, Deserialize, Clone)] +#[serde(rename_all = "snake_case")] +pub struct HistoryEntry { + pub id: i64, + pub timestamp: i64, + pub role: String, + pub content: String, + pub reasoning: Option, +} + /// daemon → client #[derive(Debug, Serialize, Deserialize, Clone)] #[serde(tag = "type", rename_all = "snake_case")] pub enum ServerMsg { + Started, Token { content: String, }, @@ -29,13 +40,8 @@ pub enum ServerMsg { context_tokens: usize, watermark: usize, }, - /// pushed on connect: current context window turns History { - turns: Vec<(i64, String, String)>, - }, - /// response to FetchHistory - HistoryPage { - turns: Vec<(i64, String, String)>, + turns: Vec, }, } diff --git a/klbr-tui/src/main.rs b/klbr-tui/src/main.rs index 8b4047f..61c9c09 100644 --- a/klbr-tui/src/main.rs +++ b/klbr-tui/src/main.rs @@ -34,11 +34,22 @@ struct Reason { done: bool, } +#[derive(Clone)] +enum AssistantStep { + PromptProcessing, + Reasoning, + Response, + Done, +} + #[derive(Clone)] enum Role { System, User, - Assistant { reason: Option, done: bool }, + Assistant { + reason: Option, + step: AssistantStep, + }, } impl Role { @@ -48,10 +59,10 @@ impl Role { fn if_assistant(&mut self, f: F) -> Option where - F: FnOnce(&mut Option, &mut bool) -> T, + F: FnOnce(&mut Option, &mut AssistantStep) -> T, { match self { - Role::Assistant { reason, done } => Some(f(reason, done)), + Role::Assistant { reason, step } => Some(f(reason, step)), _ => None, } } @@ -84,7 +95,7 @@ impl ChatMsg { fn assistant() -> Self { Self { role: Role::Assistant { - done: false, + step: AssistantStep::PromptProcessing, reason: None, }, content: String::new(), @@ -189,25 +200,30 @@ impl App { self.at_bottom = false; } - fn prepend_turns(&mut self, turns: Vec<(i64, String, String)>) -> bool { + fn prepend_turns(&mut self, turns: Vec) -> bool { if turns.is_empty() { return false; } let msgs: Vec = turns - .iter() - .filter_map(|(id, role, content)| { + .into_iter() + .filter_map(|entry| { // track oldest id for next page request - if self.oldest_turn_id.map_or(true, |cur| *id < cur) { - self.oldest_turn_id = Some(*id); + if self.oldest_turn_id.map_or(true, |cur| entry.id < cur) { + self.oldest_turn_id = Some(entry.id); } - match role.as_str() { - "user" => Some(ChatMsg::user(content.clone())), + match entry.role.as_str() { + "user" => Some(ChatMsg::user(entry.content)), + "system" => Some(ChatMsg::system(entry.content)), "assistant" => { let mut m = ChatMsg::assistant(); - m.content = content.clone(); + m.content = entry.content; m.role = Role::Assistant { - reason: None, - done: true, + reason: entry.reasoning.map(|content| Reason { + content, + expanded: false, + done: true, + }), + step: AssistantStep::Done, }; Some(m) } @@ -226,11 +242,19 @@ impl App { self.at_bottom = true; } - // gets the last assistant message + // gets the last assistant message that we are streaming fn streaming_mut(&mut self) -> Option<&mut ChatMsg> { - self.history - .last_mut() - .filter(|m| matches!(m.role, Role::Assistant { done: false, .. })) + self.history.last_mut().filter(|m| { + matches!( + m.role, + Role::Assistant { + step: AssistantStep::PromptProcessing + | AssistantStep::Reasoning + | AssistantStep::Response, + .. + } + ) + }) } fn toggle_last_think(&mut self) { @@ -250,10 +274,12 @@ impl App { }); } - fn get_assistant_entry(&mut self) -> &mut ChatMsg { + fn token_entry(&mut self) -> &mut ChatMsg { + if self.stream_start.is_none() { + self.stream_start = Some(Instant::now()); + } if self.streaming_mut().is_none() { self.history.push(ChatMsg::assistant()); - self.stream_start = Some(Instant::now()); self.stream_tokens = 0; } self.stream_tokens += 1; @@ -329,18 +355,31 @@ fn render_history<'a>(history: &'a [ChatMsg]) -> Vec> { } was_assistant = false; } - Role::Assistant { reason, done } => { + Role::Assistant { reason, step } => { + was_assistant = true; + let klbr_span = Span::styled( + "klbr", + Style::default() + .fg(Color::Green) + .add_modifier(Modifier::DIM), + ); + if let AssistantStep::PromptProcessing = step { + lines.push(Line::from(vec![ + klbr_span, + Span::styled( + " is processing prompt...", + Style::default() + .fg(Color::DarkGray) + .add_modifier(Modifier::DIM), + ), + ])); + continue; + } let think_expanded = if let Some(think) = reason { let think_token = think .done .then_some("has reasoned") .unwrap_or("is reasoning..."); - let klbr_span = Span::styled( - "klbr", - Style::default() - .fg(Color::Green) - .add_modifier(Modifier::DIM), - ); if think.expanded { lines.push(Line::from(vec![ klbr_span, @@ -382,9 +421,6 @@ fn render_history<'a>(history: &'a [ChatMsg]) -> Vec> { Span::raw(line), ])); } - was_assistant = true; - // suppress unused warning — the field exists for future use - let _ = done; } } } @@ -595,8 +631,8 @@ async fn handle_event( _ => {} }, Event::Mouse(m) => match m.kind { - MouseEventKind::ScrollUp => app.scroll_up(1), - MouseEventKind::ScrollDown => app.scroll_down(1), + MouseEventKind::ScrollUp => app.scroll_up(2), + MouseEventKind::ScrollDown => app.scroll_down(2), _ => {} }, _ => {} @@ -607,18 +643,26 @@ async fn handle_event( fn handle_message(app: &mut App, line: String) -> Result<()> { let msg: ServerMsg = serde_json::from_str(&line)?; match msg { + ServerMsg::Started => { + if app.streaming_mut().is_none() { + app.history.push(ChatMsg::assistant()); + app.stream_tokens = 0; + } + } ServerMsg::Token { content } => { - let entry = app.get_assistant_entry(); + let entry = app.token_entry(); entry.content.push_str(&content); // mark reasoning as done - entry.role.if_assistant(|r, _| { + entry.role.if_assistant(|r, step| { + *step = AssistantStep::Response; if let Some(r) = r { r.done = true; } }); } ServerMsg::ThinkToken { content } => { - app.get_assistant_entry().role.if_assistant(|r, _| { + app.token_entry().role.if_assistant(|r, step| { + *step = AssistantStep::Reasoning; r.get_or_insert_with(Reason::default) .content .push_str(&content) @@ -626,9 +670,9 @@ fn handle_message(app: &mut App, line: String) -> Result<()> { } ServerMsg::Done => { // mark assistant message as done - app.get_assistant_entry() + app.token_entry() .role - .if_assistant(|_, done| *done = true); + .if_assistant(|_, step| *step = AssistantStep::Done); app.status.clear(); app.stream_start = None; } @@ -645,9 +689,6 @@ fn handle_message(app: &mut App, line: String) -> Result<()> { app.watermark = watermark; } ServerMsg::History { turns } => { - app.prepend_turns(turns); - } - ServerMsg::HistoryPage { turns } => { app.loading_history = false; if !app.prepend_turns(turns) { app.history_exhausted = true; diff --git a/rust-toolchain.toml b/rust-toolchain.toml new file mode 100644 index 0000000..0440a5d --- /dev/null +++ b/rust-toolchain.toml @@ -0,0 +1,3 @@ +[toolchain] +channel = "stable" +components = ["rust-analyzer", "rust-src", "rustfmt"]