From 5ff0b2501cf00e1c324749da14ed4a0eef697721 Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Mon, 3 Aug 2026 21:21:40 -0400 Subject: [PATCH 1/5] feat(ui): persist plugin session state --- n00n-lua/tests/plugin_host.rs | 41 +++++++++++++++++ n00n-ui/src/app/mod.rs | 1 + n00n-ui/src/app/session.rs | 83 ++++++++++++++++++++++++++++++++++- n00n-ui/src/app/tests.rs | 82 +++++++++++++++++++++++++++++++++- n00n-ui/src/event_loop.rs | 4 +- plugins/todo_write/init.lua | 25 ++++++++--- 6 files changed, 226 insertions(+), 10 deletions(-) diff --git a/n00n-lua/tests/plugin_host.rs b/n00n-lua/tests/plugin_host.rs index 197a0f302..8e6ddf18a 100644 --- a/n00n-lua/tests/plugin_host.rs +++ b/n00n-lua/tests/plugin_host.rs @@ -23,6 +23,7 @@ use n00n_providers::{ StreamResponse, System, TokenUsage, }; use n00n_storage::id::SessionRef; +use n00n_storage::sessions::StoredStateScope; const TOOL_DEFINITIONS_BYTE_BUDGET: usize = 46_000; @@ -5337,6 +5338,46 @@ fn bundled_todo_panel_keeps_current_todo_stable_in_hint() { "ToolDone", serde_json::json!({ "id": "cmd-1", "tool": "bash", "is_error": false }), ); + handle.fire_autocmd("TurnEnd", serde_json::json!({})); + barrier(&host); + let hints = host.hint_reader().load(); + let text = hints + .entries + .iter() + .flat_map(|(_, spans)| spans.iter().map(|(text, _)| text.as_str())) + .collect::(); + assert!(text.contains("Run tests"), "turn end cleared todos: {text}"); +} + +#[test] +fn bundled_todo_persists_root_scoped_state() { + let (reg, host) = builtins_host(); + let identity = SessionIdentity::root(SessionRef::generate()); + let entry = reg.get("todo_write").unwrap(); + let invocation = entry + .tool + .parse(&serde_json::json!({ + "todos": [{ "content": "Resume work", "status": "in_progress", "priority": "high" }] + })) + .unwrap(); + let mut ctx = n00n_agent::tools::test_support::stub_ctx(&n00n_agent::AgentMode::Build); + ctx.identity = Some(identity.clone()); + + smol::block_on(invocation.execute(&ctx)).output.unwrap(); + + let snapshot = host + .event_handle() + .unwrap() + .capture_state(&identity, 1) + .unwrap(); + assert_eq!( + snapshot + .plugin_payload_for_apply("todo_write", 1, StoredStateScope::Root) + .unwrap(), + Some(&serde_json::json!({ + "todos": [{ "content": "Resume work", "status": "in_progress", "priority": "high" }] + })) + ); } #[test] diff --git a/n00n-ui/src/app/mod.rs b/n00n-ui/src/app/mod.rs index b338680ab..6f163f14d 100644 --- a/n00n-ui/src/app/mod.rs +++ b/n00n-ui/src/app/mod.rs @@ -1596,6 +1596,7 @@ impl App { self.subagent_answers.remove(&e.id); self.subagent_prompts.remove(&e.id); } + self.save_session(); } if let AgentEvent::Retry { diff --git a/n00n-ui/src/app/session.rs b/n00n-ui/src/app/session.rs index e89f045ec..def62b944 100644 --- a/n00n-ui/src/app/session.rs +++ b/n00n-ui/src/app/session.rs @@ -6,9 +6,10 @@ use crate::chat::{Chat, DONE_TEXT, RESTORE_BATCH_SIZE, history_to_display, trans use crate::components::DisplayRole; use crate::components::rewind_picker::RewindEntry; use crate::components::{Action, LoadedSession}; +use n00n_agent::tools::SessionIdentity; use n00n_agent::{AgentInput, AgentMode, McpPromptRef}; use n00n_providers::{Model, TokenUsage}; -use n00n_storage::id::n00nId; +use n00n_storage::id::{SessionRef, n00nId}; use n00n_storage::sessions::{ StoredDelivery, StoredImageMediaType, StoredImageSource, StoredMcpPrompt, StoredMode, StoredQueuedMessage, StoredSubagent, StoredThinking, @@ -150,7 +151,7 @@ impl App { } pub(crate) fn save_session(&mut self) { - let snapshot = self.session_snapshot(); + let snapshot = self.session_snapshot_with_plugin_state(); if !session_has_content(&snapshot) { return; } @@ -168,6 +169,77 @@ impl App { self.state.session.clone() } + fn session_snapshot_with_plugin_state(&mut self) -> AppSession { + let mut snapshot = self.session_snapshot(); + self.capture_plugin_state(); + snapshot.meta.state_snapshot = self.state.session.meta.state_snapshot.clone(); + snapshot + } + + pub(crate) fn checkpoint_session(&mut self) { + let snapshot = self.session_snapshot(); + if session_has_content(&snapshot) { + self.storage_writer.send(Box::new(snapshot)); + } + } + + pub(crate) fn hydrate_plugin_state(&mut self) { + let Some(handle) = &self.lua_event_handle else { + return; + }; + let session_id = self.state.session.id; + let identity = SessionIdentity::root(SessionRef::from_id(session_id)); + if let Err(error) = handle.drop_state_owner(session_id) { + tracing::warn!(%session_id, %error, "failed to clear stale plugin session state"); + return; + } + if let Err(error) = + handle.hydrate_state(&identity, self.state.session.meta.state_snapshot.clone()) + { + tracing::warn!(%session_id, %error, "failed to restore plugin session state"); + } + } + + fn capture_plugin_state(&mut self) { + let Some(handle) = &self.lua_event_handle else { + return; + }; + let session_id = self.state.session.id; + let identity = SessionIdentity::root(SessionRef::from_id(session_id)); + let persisted_revision = match self + .state + .session + .meta + .state_snapshot + .as_ref() + .and_then(|snapshot| snapshot.state_revision()) + { + Some(revision) => revision, + None => 0, + }; + let revision = self + .state + .session + .meta + .revision + .max(persisted_revision.saturating_add(1)); + match handle.capture_state(&identity, revision) { + Ok(snapshot) => self.state.session.meta.state_snapshot = Some(snapshot), + Err(error) => { + tracing::warn!(%session_id, %error, "failed to capture plugin session state"); + } + } + } + + pub(crate) fn drop_plugin_state(&self, session_id: n00nId) { + let Some(handle) = &self.lua_event_handle else { + return; + }; + if let Err(error) = handle.drop_state_owner(session_id) { + tracing::warn!(%session_id, %error, "failed to drop plugin session state"); + } + } + fn sync_ephemeral_state(&mut self) { let draft = self.input_box.buffer.value(); self.state.session.meta.input_draft = if draft.is_empty() { None } else { Some(draft) }; @@ -332,6 +404,7 @@ impl App { pub(super) fn reset_session(&mut self) -> Vec { self.save_session(); + self.drop_plugin_state(self.state.session.id); self.reset_ui_chrome(); self.state.token_usage = TokenUsage::default(); self.state.context_size = 0; @@ -340,6 +413,7 @@ impl App { self.enter_plan(); } self.state.session = AppSession::new(&self.state.session.model, &self.state.session.cwd); + self.hydrate_plugin_state(); self.fire_session_autocmd("SessionReset", serde_json::json!({})); vec![Action::NewSession] } @@ -392,7 +466,12 @@ impl App { ) -> LoadedSession { self.permissions .load_session_rules(stored_to_rules(&session.meta.session_rules)); + let previous_session_id = self.state.session.id; self.state = SessionState::from_session(session, fallback_model, &self.storage); + if previous_session_id != self.state.session.id { + self.drop_plugin_state(previous_session_id); + } + self.hydrate_plugin_state(); self.state .session .prune_orphans(|m| m.tool_uses().map(|(id, _, _)| id.to_owned()).collect()); diff --git a/n00n-ui/src/app/tests.rs b/n00n-ui/src/app/tests.rs index 2e892e6cd..8cde501f9 100644 --- a/n00n-ui/src/app/tests.rs +++ b/n00n-ui/src/app/tests.rs @@ -8,15 +8,19 @@ use crate::selection::{SelectableZone, SelectionState, SelectionZone}; use arc_swap::ArcSwap; use crossterm::event::{KeyCode, KeyEvent, KeyModifiers, MouseButton, MouseEventKind}; use n00n_agent::permissions::PermissionManager; +use n00n_agent::tools::{SessionIdentity, ToolRegistry}; use n00n_agent::{ ExtractedCommand, FusionLane, FusionUsageStats, ImageMediaType, InterruptPoint, InterruptSource, McpConfigErrors, McpPromptArg, McpServerInfo, McpServerStatus, McpSnapshot, McpSnapshotReader, ToolDoneEvent, ToolOutput, ToolStartEvent, TurnCompleteEvent, }; use n00n_config::{PermissionsConfig, UiConfig}; -use n00n_lua::{HintReader, KeymapReader, LuaCommandReader}; +use n00n_lua::{HintReader, KeymapReader, LuaCommandReader, PluginHost}; use n00n_providers::{ContentBlock, Effort, Role, TokenUsage}; -use n00n_storage::sessions::{StoredMode, StoredThinking, TranscriptEntry}; +use n00n_storage::id::SessionRef; +use n00n_storage::sessions::{ + StoredMode, StoredSessionStateSnapshot, StoredStateScope, StoredThinking, TranscriptEntry, +}; use ratatui::{Terminal, backend::TestBackend, layout::Rect}; use ratatui_image::picker::Picker; use std::path::{Path, PathBuf}; @@ -2570,6 +2574,80 @@ fn drain_writer(app: App, writer: Arc) { .shutdown(WRITER_DRAIN_TIMEOUT); } +#[test] +fn save_session_captures_plugin_state_snapshot() { + let (_tmp, dir, writer, mut app) = tempdir_app(); + let host = PluginHost::new(Arc::new(ToolRegistry::new())).unwrap(); + let handle = host.event_handle().unwrap(); + let session_id = app.state.session.id; + let mut snapshot = StoredSessionStateSnapshot::new(3); + snapshot + .set_plugin_state( + "todo_write", + 1, + StoredStateScope::Root, + serde_json::json!({"todos": [{"content": "ship", "status": "in_progress"}]}), + ) + .unwrap(); + app.state.session.meta.state_snapshot = Some(snapshot); + app.lua_event_handle = Some(handle); + app.hydrate_plugin_state(); + app.state + .session + .messages + .push(Message::user("persist state".into())); + + app.save_session(); + assert!(app.state.session.meta.state_snapshot.is_some()); + drain_writer(app, writer); + + let loaded = AppSession::load(session_id, &dir).unwrap(); + let payload = loaded + .meta + .state_snapshot + .as_ref() + .unwrap() + .plugin_payload_for_apply("todo_write", 1, StoredStateScope::Root) + .unwrap(); + assert_eq!( + payload, + Some(&serde_json::json!({"todos": [{"content": "ship", "status": "in_progress"}]})) + ); +} + +#[test] +fn apply_loaded_session_hydrates_plugin_state_snapshot() { + let mut app = test_app(); + let host = PluginHost::new(Arc::new(ToolRegistry::new())).unwrap(); + let handle = host.event_handle().unwrap(); + app.lua_event_handle = Some(handle.clone()); + let mut loaded = AppSession::new("test-model", "/tmp/test"); + loaded.messages.push(Message::user("restore state".into())); + let loaded_id = loaded.id; + let mut snapshot = StoredSessionStateSnapshot::new(9); + snapshot + .set_plugin_state( + "todo_write", + 1, + StoredStateScope::Root, + serde_json::json!({"todos": [{"content": "resume", "status": "pending"}]}), + ) + .unwrap(); + loaded.meta.state_snapshot = Some(snapshot); + let model = app.state.model.clone(); + + app.apply_loaded_session(loaded, &model); + + let identity = SessionIdentity::root(SessionRef::from_id(loaded_id)); + let captured = handle.capture_state(&identity, 10).unwrap(); + assert_eq!( + captured + .plugin_payload_for_apply("todo_write", 1, StoredStateScope::Root) + .unwrap(), + Some(&serde_json::json!({"todos": [{"content": "resume", "status": "pending"}]})) + ); +} + #[test] fn reload_persists_session_with_content_to_disk() { let (_tmp, dir, writer, mut app) = tempdir_app(); diff --git a/n00n-ui/src/event_loop.rs b/n00n-ui/src/event_loop.rs index 5a541d424..fdaafd92a 100644 --- a/n00n-ui/src/event_loop.rs +++ b/n00n-ui/src/event_loop.rs @@ -251,6 +251,7 @@ impl SpawnCtx { picker: Arc::clone(&self.picker), }); app.lua_event_handle.clone_from(&self.lua_event_handle); + app.hydrate_plugin_state(); handles.apply_to_app(&mut app); if resumed { restore_session(&mut app, &handles); @@ -710,7 +711,7 @@ impl<'t> EventLoop<'t> { } for rt in &mut self.sessions { if should_save_periodically(&rt.app.status) { - rt.app.save_session(); + rt.app.checkpoint_session(); } } self.last_save = Instant::now(); @@ -866,6 +867,7 @@ impl<'t> EventLoop<'t> { } let rt = self.remove_runtime(i); rt.handles.cancel(); + rt.app.drop_plugin_state(id); } self.ctx.storage_writer.delete(id, move |res| { let reply = match res { diff --git a/plugins/todo_write/init.lua b/plugins/todo_write/init.lua index b4b2becee..4c1962c0c 100644 --- a/plugins/todo_write/init.lua +++ b/plugins/todo_write/init.lua @@ -222,8 +222,17 @@ n00n.api.register_tool({ return ToolView.restore_lines(build_lines(), { max_lines = DEFAULT_PREVIEW_LINES, keep = "head" }) end, - handler = function(input) + handler = function(input, ctx) items = input.todos or {} + local _, state_err + if #items == 0 then + _, state_err = ctx:state_remove("root") + else + _, state_err = ctx:state_replace("root", { todos = items }) + end + if state_err then + error(state_err) + end if #items == 0 then if win and win:is_open() then win:hide() @@ -308,16 +317,22 @@ n00n.api.create_autocmd("ToolDone", { end, }) -local function clear_todos() - items = {} - seen_first = false +local function clear_activity() running = {} running_order = {} activity_expanded = false + refresh_activity() +end + +local function clear_todos() + items = {} + seen_first = false + clear_activity() if win and win:is_open() then win:hide() end n00n.ui.set_status_hint(nil) end -n00n.api.create_autocmd({ "TurnEnd", "TurnError", "SessionReset" }, { callback = clear_todos }) +n00n.api.create_autocmd({ "TurnEnd", "TurnError" }, { callback = clear_activity }) +n00n.api.create_autocmd("SessionReset", { callback = clear_todos }) From 564ece909b7b7d4aaf9f3788cd7990939ae303ce Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Mon, 3 Aug 2026 21:55:08 -0400 Subject: [PATCH 2/5] feat(sessions): persist plugin state across runtimes --- n00n-acp/src/lib.rs | 3 +- n00n-acp/src/server.rs | 1 + n00n-agent/src/headless.rs | 283 +++++++++++++++++++++++++++++++++++-- n00n-lua/src/loader.rs | 28 ++++ n00n-ui/src/app/session.rs | 32 +++-- src/cmd/acp.rs | 12 +- src/cmd/agent.rs | 12 +- src/cmd/tui.rs | 12 +- src/print.rs | 3 + src/sdk_mode.rs | 5 +- 10 files changed, 357 insertions(+), 34 deletions(-) diff --git a/n00n-acp/src/lib.rs b/n00n-acp/src/lib.rs index c6912411f..bfeb7583b 100644 --- a/n00n-acp/src/lib.rs +++ b/n00n-acp/src/lib.rs @@ -6,7 +6,7 @@ pub mod translate; use std::path::PathBuf; use std::sync::Arc; -use n00n_agent::headless::InteractiveHandle; +use n00n_agent::headless::{InteractiveHandle, SessionStatePersistence}; use n00n_agent::prompt::ResolvedSlots; use n00n_agent::{AgentConfig, PermissionsConfig}; use n00n_providers::model::Model; @@ -27,6 +27,7 @@ pub struct AcpParams { pub initial_wd: PathBuf, pub mcp_handle: Option, pub prompt_slots: Arc, + pub state_persistence: Option>, pub yolo: bool, pub session_daemon_register: Option, } diff --git a/n00n-acp/src/server.rs b/n00n-acp/src/server.rs index 251abf103..a8b1d3591 100644 --- a/n00n-acp/src/server.rs +++ b/n00n-acp/src/server.rs @@ -193,6 +193,7 @@ fn spawn_session( timeouts: params.timeouts, openai_options: params.openai_options, prompt_slots: Arc::clone(¶ms.prompt_slots), + state_persistence: params.state_persistence.clone(), excluded_tools: Vec::new(), mcp_handle: params.mcp_handle.clone(), initial_wd: cwd, diff --git a/n00n-agent/src/headless.rs b/n00n-agent/src/headless.rs index 862a8e591..45c350abc 100644 --- a/n00n-agent/src/headless.rs +++ b/n00n-agent/src/headless.rs @@ -12,7 +12,7 @@ use n00n_providers::model::Model; use n00n_providers::provider::{self, Provider}; use n00n_storage::StateDir; use n00n_storage::id::{SessionRef, n00nId}; -use n00n_storage::sessions::{Session, StoredMode}; +use n00n_storage::sessions::{Session, StoredMode, StoredSessionStateSnapshot}; use serde_json::Value; use tracing::{error, warn}; @@ -33,20 +33,68 @@ use crate::{ type StoredSession = Session; const NON_UTF8_PLAN_PATH_ERR: &str = "plan path must be valid UTF-8"; +const INITIAL_STATE_REVISION: u64 = 0; + +pub trait SessionStatePersistence: Send + Sync { + /// # Errors + /// Returns an error when the runtime cannot restore the snapshot. + fn hydrate( + &self, + identity: &SessionIdentity, + snapshot: Option, + ) -> Result<(), String>; + + /// # Errors + /// Returns an error when the runtime cannot capture its current state. + fn capture( + &self, + identity: &SessionIdentity, + revision: u64, + ) -> Result; + + /// # Errors + /// Returns an error when the runtime cannot remove the owner state. + fn drop_owner(&self, owner: n00nId) -> Result<(), String>; +} +fn state_revision_or_initial(snapshot: Option<&StoredSessionStateSnapshot>) -> u64 { + let Some(snapshot) = snapshot else { + return INITIAL_STATE_REVISION; + }; + let Some(revision) = snapshot.state_revision() else { + return INITIAL_STATE_REVISION; + }; + revision +} struct SessionStore { dir: StateDir, session: StoredSession, + state_persistence: Option>, + identity: SessionIdentity, } impl SessionStore { - fn open(session_id: n00nId, cwd: &str, model_spec: &str, mode: &AgentMode) -> Option { + fn open( + session_id: n00nId, + cwd: &str, + model_spec: &str, + mode: &AgentMode, + state_persistence: Option>, + ) -> Option { let dir = StateDir::resolve() .map_err(|e| warn!(error = %e, "state dir unavailable; session will not be persisted")) .ok()?; - Some(Self::open_in(dir, session_id, cwd, model_spec, mode)) + Some(Self::open_in_with_state( + dir, + session_id, + cwd, + model_spec, + mode, + state_persistence, + )) } + #[cfg(test)] fn open_in( dir: StateDir, session_id: n00nId, @@ -54,22 +102,76 @@ impl SessionStore { model_spec: &str, mode: &AgentMode, ) -> Self { - if let Ok(session) = StoredSession::load(session_id, &dir) { - Self { dir, session } + Self::open_in_with_state(dir, session_id, cwd, model_spec, mode, None) + } + + fn open_in_with_state( + dir: StateDir, + session_id: n00nId, + cwd: &str, + model_spec: &str, + mode: &AgentMode, + state_persistence: Option>, + ) -> Self { + let mut is_new = false; + let session = if let Ok(session) = StoredSession::load(session_id, &dir) { + session } else { + is_new = true; let mut session = StoredSession::new(model_spec, cwd); session.id = session_id; - let mut store = Self { dir, session }; + session + }; + let identity = SessionIdentity::root(SessionRef::from(session_id)); + let mut store = Self { + dir, + session, + state_persistence, + identity, + }; + store.hydrate_plugin_state(); + if is_new { if let Err(error) = store.update_turn_metadata(mode, None) { warn!(error, "session metadata was not persisted"); } else { store.save(); } - store + } + store + } + + fn hydrate_plugin_state(&self) { + let Some(state_persistence) = &self.state_persistence else { + return; + }; + if let Err(error) = + state_persistence.hydrate(&self.identity, self.session.meta.state_snapshot.clone()) + { + warn!(session_id = %self.session.id, %error, "failed to restore plugin session state"); + } + } + + fn capture_plugin_state(&mut self) { + let Some(state_persistence) = &self.state_persistence else { + return; + }; + let persisted_revision = + state_revision_or_initial(self.session.meta.state_snapshot.as_ref()); + let revision = self + .session + .meta + .revision + .max(persisted_revision.saturating_add(1)); + match state_persistence.capture(&self.identity, revision) { + Ok(snapshot) => self.session.meta.state_snapshot = Some(snapshot), + Err(error) => { + warn!(session_id = %self.session.id, %error, "failed to capture plugin session state"); + } } } fn save(&mut self) { + self.capture_plugin_state(); if let Err(e) = self.session.save(&self.dir) { warn!(error = %e, session_id = %self.session.id, "failed to persist session"); } @@ -94,6 +196,7 @@ impl SessionStore { .transpose()?; self.session.meta.mode = Some(stored_mode); self.session.meta.plan_path = stored_plan_path; + self.session.meta.revision = self.session.meta.revision.saturating_add(1); Ok(()) } @@ -123,6 +226,17 @@ impl SessionStore { } } +impl Drop for SessionStore { + fn drop(&mut self) { + let Some(state_persistence) = &self.state_persistence else { + return; + }; + if let Err(error) = state_persistence.drop_owner(self.session.id) { + warn!(session_id = %self.session.id, %error, "failed to drop plugin session state"); + } + } +} + pub struct HeadlessParams { pub model: Model, pub config: Arc, @@ -132,6 +246,7 @@ pub struct HeadlessParams { pub prompt: String, pub images: Vec, pub prompt_slots: ResolvedSlots, + pub state_persistence: Option>, pub excluded_tools: Vec<&'static str>, pub mcp_handle: Option, pub initial_wd: PathBuf, @@ -253,8 +368,13 @@ pub fn spawn(params: HeadlessParams) -> HeadlessHandle { let error_tx = event_tx.clone(); let mut history = History::new(Vec::new()); let model_spec = model.spec(); - let mut session_store = - SessionStore::open(session_ref_clone.id(), &session_cwd, &model_spec, &mode); + let mut session_store = SessionStore::open( + session_ref_clone.id(), + &session_cwd, + &model_spec, + &mode, + params.state_persistence, + ); let mut agent = Agent::new( AgentParams { provider, @@ -349,6 +469,7 @@ pub struct InteractiveParams { pub timeouts: Timeouts, pub openai_options: OpenAiOptions, pub prompt_slots: Arc, + pub state_persistence: Option>, pub excluded_tools: Vec<&'static str>, pub mcp_handle: Option, pub initial_wd: PathBuf, @@ -403,7 +524,13 @@ pub fn spawn_interactive(params: InteractiveParams) -> InteractiveHandle { }; let working_dir = params.initial_wd.to_string_lossy().into_owned(); - let store = SessionStore::open(session_id, &working_dir, ¶ms.model.spec(), ¶ms.mode); + let store = SessionStore::open( + session_id, + &working_dir, + ¶ms.model.spec(), + ¶ms.mode, + params.state_persistence.clone(), + ); let permissions = Arc::new(PermissionManager::new( params.permissions_config.clone(), params.initial_wd.clone(), @@ -604,6 +731,57 @@ mod tests { const CWD: &str = "/project"; const MODEL_SPEC: &str = "anthropic/claude-test"; + const PLUGIN: &str = "todo_write"; + + #[derive(Default)] + struct StatePersistenceProbe { + hydrated_revisions: std::sync::Mutex>>, + captured_revisions: std::sync::Mutex>, + dropped_owners: std::sync::Mutex>, + fail_capture: std::sync::atomic::AtomicBool, + } + + impl SessionStatePersistence for StatePersistenceProbe { + fn hydrate( + &self, + _identity: &SessionIdentity, + snapshot: Option, + ) -> Result<(), String> { + self.hydrated_revisions.lock().unwrap().push( + snapshot + .as_ref() + .and_then(StoredSessionStateSnapshot::state_revision), + ); + Ok(()) + } + + fn capture( + &self, + _identity: &SessionIdentity, + revision: u64, + ) -> Result { + self.captured_revisions.lock().unwrap().push(revision); + if self.fail_capture.load(std::sync::atomic::Ordering::Relaxed) { + return Err("capture failed".into()); + } + let mut snapshot = StoredSessionStateSnapshot::new(revision); + snapshot + .set_plugin_state( + PLUGIN, + 1, + n00n_storage::sessions::StoredStateScope::Root, + serde_json::json!({"todos": []}), + ) + .unwrap(); + Ok(snapshot) + } + + fn drop_owner(&self, owner: n00nId) -> Result<(), String> { + self.dropped_owners.lock().unwrap().push(owner); + Ok(()) + } + } + fn session_id() -> n00nId { SESSION_ID.parse().unwrap() } @@ -730,10 +908,95 @@ mod tests { store.update_turn_metadata(&AgentMode::Plan(path), None), Err(NON_UTF8_PLAN_PATH_ERR) ); + assert_eq!(store.session.meta.mode, original.mode); assert_eq!(store.session.meta.plan_path, original.plan_path); } + #[test] + fn session_store_hydrates_and_captures_plugin_state() { + let tmp = TempDir::new().unwrap(); + let dir = StateDir::from_path(tmp.path().to_path_buf()); + let mut persisted = StoredSession::new(MODEL_SPEC, CWD); + persisted.id = session_id(); + let mut snapshot = StoredSessionStateSnapshot::new(7); + snapshot + .set_plugin_state( + PLUGIN, + 1, + n00n_storage::sessions::StoredStateScope::Root, + serde_json::json!({"todos": [{"content": "resume", "status": "pending"}]}), + ) + .unwrap(); + persisted.meta.state_snapshot = Some(snapshot); + persisted.save(&dir).unwrap(); + + let probe = Arc::new(StatePersistenceProbe::default()); + let state_persistence: Arc = Arc::clone(&probe) as Arc<_>; + let mut store = SessionStore::open_in_with_state( + dir.clone(), + session_id(), + CWD, + MODEL_SPEC, + &AgentMode::Build, + Some(state_persistence), + ); + assert_eq!(*probe.hydrated_revisions.lock().unwrap(), vec![Some(7)]); + + store + .record_turn(&[], MODEL_SPEC.into(), &AgentMode::Build, None) + .unwrap(); + assert_eq!(*probe.captured_revisions.lock().unwrap(), vec![8]); + assert_eq!( + StoredSession::load(session_id(), &dir) + .unwrap() + .meta + .state_snapshot + .as_ref() + .and_then(StoredSessionStateSnapshot::state_revision), + Some(8) + ); + drop(store); + assert_eq!(*probe.dropped_owners.lock().unwrap(), vec![session_id()]); + } + + #[test] + fn failed_plugin_state_capture_preserves_persisted_snapshot() { + let tmp = TempDir::new().unwrap(); + let dir = StateDir::from_path(tmp.path().to_path_buf()); + let mut persisted = StoredSession::new(MODEL_SPEC, CWD); + persisted.id = session_id(); + persisted.meta.state_snapshot = Some(StoredSessionStateSnapshot::new(11)); + persisted.save(&dir).unwrap(); + + let probe = Arc::new(StatePersistenceProbe::default()); + probe + .fail_capture + .store(true, std::sync::atomic::Ordering::Relaxed); + let state_persistence: Arc = Arc::clone(&probe) as Arc<_>; + let mut store = SessionStore::open_in_with_state( + dir.clone(), + session_id(), + CWD, + MODEL_SPEC, + &AgentMode::Build, + Some(state_persistence), + ); + store + .record_turn(&[], MODEL_SPEC.into(), &AgentMode::Build, None) + .unwrap(); + + assert_eq!( + StoredSession::load(session_id(), &dir) + .unwrap() + .meta + .state_snapshot + .as_ref() + .and_then(StoredSessionStateSnapshot::state_revision), + Some(11) + ); + } + #[test] fn extract_tool_names_filters_valid_entries() { let tools = serde_json::json!([{"name": "read"}, {"type": "function"}, {"name": "bash"}]); diff --git a/n00n-lua/src/loader.rs b/n00n-lua/src/loader.rs index 211f96ca7..f34e73491 100644 --- a/n00n-lua/src/loader.rs +++ b/n00n-lua/src/loader.rs @@ -6,6 +6,7 @@ use std::sync::{Arc, LazyLock}; use std::time::Duration; use include_dir::{Dir, include_dir}; +use n00n_agent::headless::SessionStatePersistence; use n00n_agent::tools::{SessionIdentity, ToolRegistry}; use n00n_config::{PluginsConfig, RawConfig}; use n00n_storage::id::n00nId; @@ -561,6 +562,33 @@ pub struct EventHandle { prio_tx: flume::Sender, } +impl SessionStatePersistence for EventHandle { + fn hydrate( + &self, + identity: &SessionIdentity, + snapshot: Option, + ) -> Result<(), String> { + self.drop_state_owner(identity.session_id().id()) + .map_err(|error| error.to_string())?; + self.hydrate_state(identity, snapshot) + .map_err(|error| error.to_string()) + } + + fn capture( + &self, + identity: &SessionIdentity, + revision: u64, + ) -> Result { + self.capture_state(identity, revision) + .map_err(|error| error.to_string()) + } + + fn drop_owner(&self, owner: n00nId) -> Result<(), String> { + self.drop_state_owner(owner) + .map_err(|error| error.to_string()) + } +} + impl EventHandle { pub(crate) fn from_tx(tx: flume::Sender) -> Self { Self { diff --git a/n00n-ui/src/app/session.rs b/n00n-ui/src/app/session.rs index def62b944..318ff8db0 100644 --- a/n00n-ui/src/app/session.rs +++ b/n00n-ui/src/app/session.rs @@ -12,7 +12,7 @@ use n00n_providers::{Model, TokenUsage}; use n00n_storage::id::{SessionRef, n00nId}; use n00n_storage::sessions::{ StoredDelivery, StoredImageMediaType, StoredImageSource, StoredMcpPrompt, StoredMode, - StoredQueuedMessage, StoredSubagent, StoredThinking, + StoredQueuedMessage, StoredSessionStateSnapshot, StoredSubagent, StoredThinking, }; use crate::AppSession; @@ -21,6 +21,18 @@ use super::session_state::{SessionState, stored_to_rules}; use super::{App, Mode, PendingInput, PlanState}; use crate::agent::{Delivery, QueuedMessage}; +const INITIAL_STATE_REVISION: u64 = 0; + +fn state_revision_or_initial(snapshot: Option<&StoredSessionStateSnapshot>) -> u64 { + let Some(snapshot) = snapshot else { + return INITIAL_STATE_REVISION; + }; + let Some(revision) = snapshot.state_revision() else { + return INITIAL_STATE_REVISION; + }; + revision +} + /// The single content predicate: `App::save_session` persists a session /// iff this holds, and the shutdown path reuses it to tell which tabs were /// saved, so the report and the disk can never disagree. Sync the session @@ -172,7 +184,10 @@ impl App { fn session_snapshot_with_plugin_state(&mut self) -> AppSession { let mut snapshot = self.session_snapshot(); self.capture_plugin_state(); - snapshot.meta.state_snapshot = self.state.session.meta.state_snapshot.clone(); + snapshot + .meta + .state_snapshot + .clone_from(&self.state.session.meta.state_snapshot); snapshot } @@ -206,17 +221,8 @@ impl App { }; let session_id = self.state.session.id; let identity = SessionIdentity::root(SessionRef::from_id(session_id)); - let persisted_revision = match self - .state - .session - .meta - .state_snapshot - .as_ref() - .and_then(|snapshot| snapshot.state_revision()) - { - Some(revision) => revision, - None => 0, - }; + let persisted_revision = + state_revision_or_initial(self.state.session.meta.state_snapshot.as_ref()); let revision = self .state .session diff --git a/src/cmd/acp.rs b/src/cmd/acp.rs index 584728026..ee14a4247 100644 --- a/src/cmd/acp.rs +++ b/src/cmd/acp.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use color_eyre::Result; use color_eyre::eyre::Context; +use n00n_agent::headless::SessionStatePersistence; use n00n_agent::tools::ToolRegistry; use n00n_config::{load_env_files, load_permissions}; use n00n_lua::PluginHost; @@ -56,9 +57,13 @@ pub fn run(model_arg: Option<&str>, yolo: bool, no_jit: bool) -> Result<()> { config.agent.mcp_tool_desc_max_chars, )); - let prompt_slots = plugin_host - .event_handle() - .map_or_else(Default::default, |h| h.collect_prompt_slots()); + let event_handle = plugin_host.event_handle(); + let prompt_slots = event_handle.as_ref().map_or_else( + Default::default, + n00n_lua::EventHandle::collect_prompt_slots, + ); + let state_persistence = + event_handle.map(|handle| Arc::new(handle) as Arc); n00n_acp::run(n00n_acp::AcpParams { model, @@ -69,6 +74,7 @@ pub fn run(model_arg: Option<&str>, yolo: bool, no_jit: bool) -> Result<()> { initial_wd: cwd, mcp_handle, prompt_slots: Arc::new(prompt_slots), + state_persistence, yolo, session_daemon_register: Some(|state_dir, handle, model| { crate::cmd::session_daemon::register_acp_session(state_dir, handle, model) diff --git a/src/cmd/agent.rs b/src/cmd/agent.rs index 4bf615715..8147a1d51 100644 --- a/src/cmd/agent.rs +++ b/src/cmd/agent.rs @@ -363,6 +363,8 @@ struct PreparedEnv { openai_options: OpenAiOptions, mcp_handle: Option, prompt_slots: ResolvedSlots, + state_persistence: Arc, + _plugin_host: PluginHost, } fn prepare_agent_env( @@ -418,9 +420,11 @@ fn prepare_agent_env( config.agent.mcp_tool_desc_max_chars, )); - let prompt_slots = plugin_host + let event_handle = plugin_host .event_handle() - .map_or_else(Default::default, |h| h.collect_prompt_slots()); + .ok_or_else(|| eyre!("lua plugin host is unavailable"))?; + let prompt_slots = event_handle.collect_prompt_slots(); + let state_persistence = Arc::new(event_handle) as Arc; Ok(PreparedEnv { storage, @@ -433,6 +437,8 @@ fn prepare_agent_env( openai_options: OpenAiOptions::from(&config.provider), mcp_handle, prompt_slots, + state_persistence, + _plugin_host: plugin_host, }) } @@ -478,6 +484,7 @@ pub fn run(opts: &AgentRunOptions<'_>, json: bool) -> Result<()> { prompt: message, images: Vec::new(), prompt_slots: env.prompt_slots, + state_persistence: Some(Arc::clone(&env.state_persistence)), excluded_tools, mcp_handle: env.mcp_handle, initial_wd: env.cwd, @@ -729,6 +736,7 @@ fn server_unix(opts: &AgentRunOptions<'_>, agent_id: Option) -> Result<( timeouts: env.timeouts, openai_options: env.openai_options, prompt_slots: Arc::new(env.prompt_slots), + state_persistence: Some(Arc::clone(&env.state_persistence)), excluded_tools, mcp_handle: env.mcp_handle, initial_wd: env.cwd, diff --git a/src/cmd/tui.rs b/src/cmd/tui.rs index a26cb4c0d..c478e4073 100644 --- a/src/cmd/tui.rs +++ b/src/cmd/tui.rs @@ -269,10 +269,13 @@ fn run_sdk_mode( openai_options: n00n_providers::OpenAiOptions, ) -> Result<()> { let fast = stack.config.always_fast && stack.model.supports_fast(); - let prompt_slots = stack - .plugin_host - .event_handle() - .map_or_else(Default::default, |h| h.collect_prompt_slots()); + let event_handle = stack.plugin_host.event_handle(); + let prompt_slots = event_handle.as_ref().map_or_else( + Default::default, + n00n_lua::EventHandle::collect_prompt_slots, + ); + let state_persistence = event_handle + .map(|handle| Arc::new(handle) as Arc); let timeouts = stack.timeouts(); crate::sdk_mode::run(crate::sdk_mode::SdkParams { cli, @@ -282,6 +285,7 @@ fn run_sdk_mode( timeouts, openai_options, prompt_slots, + state_persistence, fast, workflow: stack.config.always_workflow, }) diff --git a/src/print.rs b/src/print.rs index 2a356714d..37106a3f7 100644 --- a/src/print.rs +++ b/src/print.rs @@ -254,6 +254,9 @@ pub fn run(model: &Model, args: PrintArgs<'_>) -> Result<()> { prompt, images, prompt_slots, + state_persistence: lua_handle.cloned().map(|handle| { + Arc::new(handle) as Arc + }), excluded_tools: vec![QUESTION_TOOL_NAME], mcp_handle, initial_wd: cwd, diff --git a/src/sdk_mode.rs b/src/sdk_mode.rs index fbf9444b2..0dd4fda63 100644 --- a/src/sdk_mode.rs +++ b/src/sdk_mode.rs @@ -17,7 +17,7 @@ use std::time::Instant; use color_eyre::Result; use color_eyre::eyre::{Context, eyre}; use flume::{Receiver, Sender}; -use n00n_agent::headless::{self, InteractiveHandle, InteractiveParams}; +use n00n_agent::headless::{self, InteractiveHandle, InteractiveParams, SessionStatePersistence}; use n00n_agent::mcp; use n00n_agent::permissions::PermissionAnswer; use n00n_agent::prompt::ResolvedSlots; @@ -445,6 +445,7 @@ pub struct SdkParams { pub timeouts: Timeouts, pub openai_options: OpenAiOptions, pub prompt_slots: ResolvedSlots, + pub state_persistence: Option>, pub fast: bool, pub workflow: bool, } @@ -465,6 +466,7 @@ pub fn run(params: SdkParams) -> Result<()> { timeouts, openai_options, prompt_slots, + state_persistence, fast, workflow, } = params; @@ -495,6 +497,7 @@ pub fn run(params: SdkParams) -> Result<()> { timeouts, openai_options, prompt_slots: Arc::new(prompt_slots), + state_persistence, excluded_tools: vec![QUESTION_TOOL_NAME], mcp_handle, initial_wd: cwd.clone(), From 8a8daaf29277b59b7af3b6e3c8e386d1fc029030 Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Mon, 3 Aug 2026 21:56:23 -0400 Subject: [PATCH 3/5] chore: add session persistence changelog --- changelog.d/332.added.md | 1 + 1 file changed, 1 insertion(+) create mode 100644 changelog.d/332.added.md diff --git a/changelog.d/332.added.md b/changelog.d/332.added.md new file mode 100644 index 000000000..78109c08a --- /dev/null +++ b/changelog.d/332.added.md @@ -0,0 +1 @@ +Sessions now restore plugin state and todo lists across restarts in UI, SDK, ACP, print, and agent modes. From 80f2b4ba8caed7bddaf01b8f15dc4429a7873bda Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Tue, 4 Aug 2026 00:41:26 -0400 Subject: [PATCH 4/5] fix: isolate persisted session state owners --- n00n-acp/src/server.rs | 117 ++++++++++++++++++--------- n00n-agent/src/headless.rs | 40 ++++++---- n00n-lua/src/api/autocmd.rs | 3 +- n00n-lua/src/api/util/ctx.rs | 13 +++ n00n-lua/src/loader.rs | 61 +++++++++++--- n00n-lua/src/state.rs | 2 +- n00n-lua/tests/plugin_host.rs | 131 ++++++++++++++++++++++++++++++- n00n-ui/src/app/session.rs | 20 +++++ n00n-ui/src/event_loop.rs | 2 + plugins/todo_write/init.lua | 144 +++++++++++++++++++++++++--------- 10 files changed, 435 insertions(+), 98 deletions(-) diff --git a/n00n-acp/src/server.rs b/n00n-acp/src/server.rs index a8b1d3591..c688a165f 100644 --- a/n00n-acp/src/server.rs +++ b/n00n-acp/src/server.rs @@ -47,6 +47,7 @@ struct SessionState { plan_path: Option, current_model: String, pending_prompt: PendingPrompt, + event_pump: smol::Task<()>, _daemon: Option, } @@ -116,7 +117,7 @@ pub async fn serve(params: AcpParams) -> color_eyre::Result<()> { handle_incoming_response(&server, &raw); } else if let Some(method) = raw.get("method").and_then(Value::as_str) { match id { - Some(id) => handle_request(&mut server, method, id, &raw, ¶ms), + Some(id) => handle_request(&mut server, method, id, &raw, ¶ms).await, None => handle_notification(&server, method), } } else if let Some(id) = id { @@ -124,6 +125,7 @@ pub async fn serve(params: AcpParams) -> color_eyre::Result<()> { } } + retire_session(&mut server).await; drop(server); writer_task.await; @@ -134,41 +136,19 @@ fn request_id(v: &Value) -> RequestId { serde_json::from_value(v.clone()).map_or(RequestId::Null, std::convert::identity) } -fn handle_request(srv: &mut Server, method: &str, id: RequestId, raw: &Value, params: &AcpParams) { +async fn handle_request( + srv: &mut Server, + method: &str, + id: RequestId, + raw: &Value, + params: &AcpParams, +) { let result = match method { "initialize" => Ok(AgentResponse::InitializeResponse( methods::initialize_response(), )), - "session/new" => parse_params::(raw).map(|req| { - let handle = spawn_session(params, req.cwd, None, Vec::new()); - let spec = params.model.spec(); - let resp = methods::new_session_response(handle.session_id.as_str()) - .config_options(vec![methods::model_config_option(&spec, &srv.model_specs)]); - install_session(srv, handle, spec, AgentMode::Build, None, params); - AgentResponse::NewSessionResponse(resp) - }), - "session/load" => parse_params::(raw).and_then(|req| { - let session_ref: SessionRef = - req.session_id.0.parse().map_err(|_| { - AcpError::resource_not_found(Some(req.session_id.0.to_string())) - })?; - let storage = n00n_storage::StateDir::resolve() - .map_err(|e| AcpError::internal_error().data(json_str(&e)))?; - let stored = load_session_from(&storage, session_ref.id())?; - let (current_mode, plan_path) = mode_and_plan_from_stored(&storage, &stored.meta) - .map_err(|e| AcpError::internal_error().data(json_str(&e)))?; - let history = stored.messages; - let sid = SessionId::from(session_ref.to_string()); - for update in translate::replay_history(&history) { - session_update(&srv.out_tx, &sid, update); - } - let handle = spawn_session(params, req.cwd, Some(session_ref), history); - let spec = params.model.spec(); - let resp = methods::load_session_response() - .config_options(vec![methods::model_config_option(&spec, &srv.model_specs)]); - install_session(srv, handle, spec, current_mode, plan_path, params); - Ok(AgentResponse::LoadSessionResponse(resp)) - }), + "session/new" => handle_new_session(srv, raw, params).await, + "session/load" => handle_load_session(srv, raw, params).await, "session/prompt" => match handle_prompt(srv, raw, &id) { Ok(()) => return, Err(e) => Err(e), @@ -180,6 +160,73 @@ fn handle_request(srv: &mut Server, method: &str, id: RequestId, raw: &Value, pa srv.respond(id, result); } +async fn handle_new_session( + srv: &mut Server, + raw: &Value, + params: &AcpParams, +) -> Result { + let req = parse_params::(raw)?; + retire_session(srv).await; + let handle = spawn_session(params, req.cwd, None, Vec::new()); + let spec = params.model.spec(); + let resp = methods::new_session_response(handle.session_id.as_str()) + .config_options(vec![methods::model_config_option(&spec, &srv.model_specs)]); + install_session(srv, handle, spec, AgentMode::Build, None, params); + Ok(AgentResponse::NewSessionResponse(resp)) +} + +async fn handle_load_session( + srv: &mut Server, + raw: &Value, + params: &AcpParams, +) -> Result { + let req = parse_params::(raw)?; + let session_ref: SessionRef = req + .session_id + .0 + .parse() + .map_err(|_| AcpError::resource_not_found(Some(req.session_id.0.to_string())))?; + let storage = n00n_storage::StateDir::resolve() + .map_err(|error| AcpError::internal_error().data(json_str(&error)))?; + let stored = load_session_from(&storage, session_ref.id())?; + let (current_mode, plan_path) = mode_and_plan_from_stored(&storage, &stored.meta) + .map_err(|error| AcpError::internal_error().data(json_str(&error)))?; + let history = stored.messages; + let sid = SessionId::from(session_ref.to_string()); + for update in translate::replay_history(&history) { + session_update(&srv.out_tx, &sid, update); + } + retire_session(srv).await; + let handle = spawn_session(params, req.cwd, Some(session_ref), history); + let spec = params.model.spec(); + let resp = methods::load_session_response() + .config_options(vec![methods::model_config_option(&spec, &srv.model_specs)]); + install_session(srv, handle, spec, current_mode, plan_path, params); + Ok(AgentResponse::LoadSessionResponse(resp)) +} + +async fn retire_session(srv: &mut Server) { + let Some(session) = srv.session.take() else { + return; + }; + let SessionState { + handle, + event_pump, + pending_prompt, + .. + } = session; + if let Some((id, _)) = take_pending(&pending_prompt) { + let response = PromptResponse::new(StopReason::Cancelled); + send( + &srv.out_tx, + Response::new(id, Ok(AgentResponse::PromptResponse(response))), + ); + } + let _ = handle.cancel_tx.try_send(()); + event_pump.cancel().await; + handle.task.cancel().await; +} + fn spawn_session( params: &AcpParams, cwd: PathBuf, @@ -216,7 +263,7 @@ fn install_session( params: &AcpParams, ) { let pending = Arc::new(Mutex::new(PendingPromptState::default())); - start_event_pump( + let event_pump = start_event_pump( handle.event_rx.clone(), handle.session_id.clone(), srv.out_tx.clone(), @@ -233,6 +280,7 @@ fn install_session( plan_path, current_model, pending_prompt: pending, + event_pump, _daemon: daemon, }); } @@ -415,7 +463,7 @@ fn start_event_pump( session_id: SessionRef, out_tx: Sender, pending: PendingPrompt, -) { +) -> smol::Task<()> { smol::spawn(async move { let sid = SessionId::from(session_id.to_string()); let mut next_request_id = FIRST_OUTGOING_REQUEST_ID; @@ -485,7 +533,6 @@ fn start_event_pump( session_update(&out_tx, &sid, update); } }) - .detach(); } fn take_pending(pending: &PendingPrompt) -> Option<(RequestId, bool)> { diff --git a/n00n-agent/src/headless.rs b/n00n-agent/src/headless.rs index 45c350abc..9dca2ebee 100644 --- a/n00n-agent/src/headless.rs +++ b/n00n-agent/src/headless.rs @@ -42,7 +42,7 @@ pub trait SessionStatePersistence: Send + Sync { &self, identity: &SessionIdentity, snapshot: Option, - ) -> Result<(), String>; + ) -> Result; /// # Errors /// Returns an error when the runtime cannot capture its current state. @@ -54,7 +54,7 @@ pub trait SessionStatePersistence: Send + Sync { /// # Errors /// Returns an error when the runtime cannot remove the owner state. - fn drop_owner(&self, owner: n00nId) -> Result<(), String>; + fn drop_owner(&self, owner: n00nId, lease: u64) -> Result<(), String>; } fn state_revision_or_initial(snapshot: Option<&StoredSessionStateSnapshot>) -> u64 { let Some(snapshot) = snapshot else { @@ -71,6 +71,7 @@ struct SessionStore { session: StoredSession, state_persistence: Option>, identity: SessionIdentity, + state_lease: Option, } impl SessionStore { @@ -128,6 +129,7 @@ impl SessionStore { session, state_persistence, identity, + state_lease: None, }; store.hydrate_plugin_state(); if is_new { @@ -140,14 +142,15 @@ impl SessionStore { store } - fn hydrate_plugin_state(&self) { + fn hydrate_plugin_state(&mut self) { let Some(state_persistence) = &self.state_persistence else { return; }; - if let Err(error) = - state_persistence.hydrate(&self.identity, self.session.meta.state_snapshot.clone()) - { - warn!(session_id = %self.session.id, %error, "failed to restore plugin session state"); + match state_persistence.hydrate(&self.identity, self.session.meta.state_snapshot.clone()) { + Ok(lease) => self.state_lease = Some(lease), + Err(error) => { + warn!(session_id = %self.session.id, %error, "failed to restore plugin session state"); + } } } @@ -228,10 +231,11 @@ impl SessionStore { impl Drop for SessionStore { fn drop(&mut self) { - let Some(state_persistence) = &self.state_persistence else { + let (Some(state_persistence), Some(lease)) = (&self.state_persistence, self.state_lease) + else { return; }; - if let Err(error) = state_persistence.drop_owner(self.session.id) { + if let Err(error) = state_persistence.drop_owner(self.session.id, lease) { warn!(session_id = %self.session.id, %error, "failed to drop plugin session state"); } } @@ -737,7 +741,8 @@ mod tests { struct StatePersistenceProbe { hydrated_revisions: std::sync::Mutex>>, captured_revisions: std::sync::Mutex>, - dropped_owners: std::sync::Mutex>, + dropped_owners: std::sync::Mutex>, + next_lease: std::sync::atomic::AtomicU64, fail_capture: std::sync::atomic::AtomicBool, } @@ -746,13 +751,15 @@ mod tests { &self, _identity: &SessionIdentity, snapshot: Option, - ) -> Result<(), String> { + ) -> Result { self.hydrated_revisions.lock().unwrap().push( snapshot .as_ref() .and_then(StoredSessionStateSnapshot::state_revision), ); - Ok(()) + Ok(self + .next_lease + .fetch_add(1, std::sync::atomic::Ordering::Relaxed)) } fn capture( @@ -776,8 +783,8 @@ mod tests { Ok(snapshot) } - fn drop_owner(&self, owner: n00nId) -> Result<(), String> { - self.dropped_owners.lock().unwrap().push(owner); + fn drop_owner(&self, owner: n00nId, lease: u64) -> Result<(), String> { + self.dropped_owners.lock().unwrap().push((owner, lease)); Ok(()) } } @@ -957,7 +964,10 @@ mod tests { Some(8) ); drop(store); - assert_eq!(*probe.dropped_owners.lock().unwrap(), vec![session_id()]); + assert_eq!( + *probe.dropped_owners.lock().unwrap(), + vec![(session_id(), 0)] + ); } #[test] diff --git a/n00n-lua/src/api/autocmd.rs b/n00n-lua/src/api/autocmd.rs index bbf7f250f..1ac3f5f6a 100644 --- a/n00n-lua/src/api/autocmd.rs +++ b/n00n-lua/src/api/autocmd.rs @@ -129,7 +129,8 @@ fn parse_string_or_seq(value: Value, what: &str) -> LuaResult> { /// `del_autocmd` later to remove the listener. /// /// Built-in events fired by the host: `"TurnStart"`, `"TurnEnd"`, -/// `"TurnError"`, `"ToolStart"`, `"ToolDone"`, `"SessionReset"`. +/// `"TurnError"`, `"ToolStart"`, `"ToolDone"`, `"SessionReset"`, +/// `"SessionFocus"`. /// Plugins can also fire their own /// events with `exec_autocmds`. /// diff --git a/n00n-lua/src/api/util/ctx.rs b/n00n-lua/src/api/util/ctx.rs index 6ee21a63b..38115453f 100644 --- a/n00n-lua/src/api/util/ctx.rs +++ b/n00n-lua/src/api/util/ctx.rs @@ -395,6 +395,19 @@ impl UserData for LuaCtx { } }); + methods.add_method("state_owner", |lua, this, scope: String| { + let scope = match parse_state_scope(&scope) { + Ok(scope) => scope, + Err(error) => return Ok((LuaValue::Nil, Some(error.to_owned()))), + }; + let access = match this.plugin_state("state_owner") { + Ok(access) => access, + Err(error) => return Ok((LuaValue::Nil, Some(error))), + }; + let owner = access.identity.owner(scope).to_string(); + Ok((LuaValue::String(lua.create_string(owner)?), None)) + }); + methods.add_method( "state_replace", |lua, this, (scope, value): (String, LuaValue)| { diff --git a/n00n-lua/src/loader.rs b/n00n-lua/src/loader.rs index f34e73491..1f731a475 100644 --- a/n00n-lua/src/loader.rs +++ b/n00n-lua/src/loader.rs @@ -1,8 +1,8 @@ use std::collections::HashMap; use std::fs; use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::{Arc, LazyLock}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, LazyLock, Mutex}; use std::time::Duration; use include_dir::{Dir, include_dir}; @@ -160,6 +160,13 @@ static BUNDLED_DIRS: LazyLock<&'static [&'static Dir<'static>]> = LazyLock::new( pub struct PluginHost { inner: Option, + state_leases: Arc, +} + +#[derive(Default)] +struct StateLeases { + next: AtomicU64, + current: Mutex>, } impl Drop for PluginHost { @@ -201,12 +208,18 @@ impl PluginHost { /// Returns an error if the Lua runtime cannot be spawned. pub fn with_jit(registry: Arc, jit: bool) -> Result { let lua = runtime::spawn(registry, *BUNDLED_DIRS, jit)?; - Ok(Self { inner: Some(lua) }) + Ok(Self { + inner: Some(lua), + state_leases: Arc::new(StateLeases::default()), + }) } #[must_use] pub fn disabled() -> Self { - Self { inner: None } + Self { + inner: None, + state_leases: Arc::new(StateLeases::default()), + } } /// Stop the Lua thread from taking new work without joining it, so the @@ -521,6 +534,7 @@ impl PluginHost { self.inner.as_ref().map(|t| EventHandle { tx: t.tx.clone(), prio_tx: t.prio_tx.clone(), + state_leases: Arc::clone(&self.state_leases), }) } @@ -560,6 +574,7 @@ pub struct EventHandle { tx: flume::Sender, /// User-initiated requests bypass queued bulk work (session restores). prio_tx: flume::Sender, + state_leases: Arc, } impl SessionStatePersistence for EventHandle { @@ -567,11 +582,22 @@ impl SessionStatePersistence for EventHandle { &self, identity: &SessionIdentity, snapshot: Option, - ) -> Result<(), String> { - self.drop_state_owner(identity.session_id().id()) + ) -> Result { + let owner = identity.session_id().id(); + let mut current = self + .state_leases + .current + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + self.drop_state_owner(owner) .map_err(|error| error.to_string())?; - self.hydrate_state(identity, snapshot) - .map_err(|error| error.to_string()) + let lease = self.state_leases.next.fetch_add(1, Ordering::Relaxed); + current.insert(owner, lease); + if let Err(error) = self.hydrate_state(identity, snapshot) { + current.remove(&owner); + return Err(error.to_string()); + } + Ok(lease) } fn capture( @@ -583,7 +609,16 @@ impl SessionStatePersistence for EventHandle { .map_err(|error| error.to_string()) } - fn drop_owner(&self, owner: n00nId) -> Result<(), String> { + fn drop_owner(&self, owner: n00nId, lease: u64) -> Result<(), String> { + let mut current = self + .state_leases + .current + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if current.get(&owner) != Some(&lease) { + return Ok(()); + } + current.remove(&owner); self.drop_state_owner(owner) .map_err(|error| error.to_string()) } @@ -594,6 +629,7 @@ impl EventHandle { Self { tx, prio_tx: flume::unbounded().0, + state_leases: Arc::new(StateLeases::default()), } } @@ -611,6 +647,7 @@ impl EventHandle { Self { tx: shared.clone(), prio_tx: shared, + state_leases: Arc::new(StateLeases::default()), } } @@ -925,7 +962,11 @@ mod tests { fn run_command_sends_correct_request() { let (prio_tx, prio_rx) = flume::bounded(8); let (tx, _rx) = flume::bounded(8); - let handle = EventHandle { tx, prio_tx }; + let handle = EventHandle { + tx, + prio_tx, + state_leases: Arc::new(StateLeases::default()), + }; handle.run_command(Arc::from("myplugin"), Arc::from("/greet"), "world".into()); let req = prio_rx.try_recv().unwrap(); match req { diff --git a/n00n-lua/src/state.rs b/n00n-lua/src/state.rs index 172635805..a0eee9c8a 100644 --- a/n00n-lua/src/state.rs +++ b/n00n-lua/src/state.rs @@ -52,7 +52,7 @@ pub(crate) struct PluginStateIdentity { } impl PluginStateIdentity { - fn owner(&self, scope: PluginStateScope) -> n00nId { + pub(crate) fn owner(&self, scope: PluginStateScope) -> n00nId { match scope { PluginStateScope::Session => self.session_id, PluginStateScope::Root => self.root_session_id, diff --git a/n00n-lua/tests/plugin_host.rs b/n00n-lua/tests/plugin_host.rs index 8e6ddf18a..9fdde0e09 100644 --- a/n00n-lua/tests/plugin_host.rs +++ b/n00n-lua/tests/plugin_host.rs @@ -10,6 +10,7 @@ use std::path::Path; use std::sync::{Arc, Mutex}; use std::time::Duration; +use n00n_agent::headless::SessionStatePersistence; use n00n_agent::template::env_vars; use n00n_agent::tools::{ ActiveTools, DescriptionContext, SessionIdentity, ToolAudience, ToolFilter, ToolRegistry, @@ -23,7 +24,7 @@ use n00n_providers::{ StreamResponse, System, TokenUsage, }; use n00n_storage::id::SessionRef; -use n00n_storage::sessions::StoredStateScope; +use n00n_storage::sessions::{StoredSessionStateSnapshot, StoredStateScope}; const TOOL_DEFINITIONS_BYTE_BUDGET: usize = 46_000; @@ -5349,6 +5350,134 @@ fn bundled_todo_panel_keeps_current_todo_stable_in_hint() { assert!(text.contains("Run tests"), "turn end cleared todos: {text}"); } +#[test] +fn stale_session_state_lease_cannot_drop_replacement_state() { + let (_registry, host) = builtins_host(); + let handle = host.event_handle().unwrap(); + let identity = SessionIdentity::root(SessionRef::generate()); + let mut first_snapshot = StoredSessionStateSnapshot::new(1); + first_snapshot + .set_plugin_state( + "todo_write", + 1, + StoredStateScope::Root, + serde_json::json!({ "todos": [{ "content": "old", "status": "pending" }] }), + ) + .unwrap(); + let first_lease = + SessionStatePersistence::hydrate(&handle, &identity, Some(first_snapshot)).unwrap(); + let mut replacement_snapshot = StoredSessionStateSnapshot::new(2); + replacement_snapshot + .set_plugin_state( + "todo_write", + 1, + StoredStateScope::Root, + serde_json::json!({ "todos": [{ "content": "replacement", "status": "pending" }] }), + ) + .unwrap(); + let replacement_lease = + SessionStatePersistence::hydrate(&handle, &identity, Some(replacement_snapshot)).unwrap(); + + SessionStatePersistence::drop_owner(&handle, identity.session_id().id(), first_lease).unwrap(); + + let captured = handle.capture_state(&identity, 3).unwrap(); + assert_eq!( + captured + .plugin_payload_for_apply("todo_write", 1, StoredStateScope::Root) + .unwrap(), + Some(&serde_json::json!({ + "todos": [{ "content": "replacement", "status": "pending" }] + })) + ); + SessionStatePersistence::drop_owner(&handle, identity.session_id().id(), replacement_lease) + .unwrap(); +} + +#[test] +fn bundled_todo_focus_uses_persisted_session_state() { + let (reg, host) = builtins_host(); + let handle = host.event_handle().unwrap(); + let ui_rx = host.ui_action_rx().unwrap(); + let focused = SessionIdentity::root(SessionRef::generate()); + let background = SessionIdentity::root(SessionRef::generate()); + handle.fire_autocmd( + "SessionFocus", + serde_json::json!({ + "session_id": focused.session_id().to_string(), + "state_snapshot": null, + }), + ); + barrier(&host); + + let entry = reg.get("todo_write").unwrap(); + let invocation = entry + .tool + .parse(&serde_json::json!({ + "todos": [{ "content": "Background work", "status": "in_progress" }] + })) + .unwrap(); + let mut ctx = n00n_agent::tools::test_support::stub_ctx(&n00n_agent::AgentMode::Build); + ctx.identity = Some(background.clone()); + smol::block_on(invocation.execute(&ctx)).output.unwrap(); + assert!(host.hint_reader().load().entries.is_empty()); + + let snapshot = handle.capture_state(&background, 1).unwrap(); + handle.fire_autocmd( + "SessionFocus", + serde_json::json!({ + "session_id": background.session_id().to_string(), + "state_snapshot": serde_json::to_value(snapshot).unwrap(), + }), + ); + barrier(&host); + + let hint = host + .hint_reader() + .load() + .entries + .iter() + .flat_map(|(_, spans)| spans.iter().map(|(text, _)| text.as_str())) + .collect::(); + assert!( + hint.contains("Background work"), + "focused todo missing: {hint}" + ); + + handle.fire_autocmd( + "ToolStart", + serde_json::json!({ + "id": "other-session-tool", + "tool": "bash", + "summary": "must stay hidden", + "session_id": focused.session_id().to_string(), + }), + ); + barrier(&host); + let toggle_id = host + .keymap_reader() + .load() + .entries + .iter() + .find(|entry| entry.desc == "Toggle todo panel") + .unwrap() + .id; + assert!(handle.run_keybind_callback(toggle_id)); + let n00n_lua::UiAction::OpenWin { buf, .. } = + ui_rx.recv_timeout(Duration::from_secs(2)).unwrap() + else { + panic!("todo panel did not open"); + }; + let panel = buf + .read() + .iter() + .flat_map(|line| line.spans.iter().map(|span| span.text.as_str())) + .collect::(); + assert!( + !panel.contains("must stay hidden"), + "panel leaked activity: {panel}" + ); +} + #[test] fn bundled_todo_persists_root_scoped_state() { let (reg, host) = builtins_host(); diff --git a/n00n-ui/src/app/session.rs b/n00n-ui/src/app/session.rs index 318ff8db0..10c4805b2 100644 --- a/n00n-ui/src/app/session.rs +++ b/n00n-ui/src/app/session.rs @@ -191,6 +191,24 @@ impl App { snapshot } + pub(crate) fn fire_session_focus_autocmd(&mut self) { + self.capture_plugin_state(); + let state_snapshot = match self.state.session.meta.state_snapshot.as_ref() { + Some(snapshot) => match serde_json::to_value(snapshot) { + Ok(value) => value, + Err(error) => { + tracing::warn!(%error, "failed to encode focused plugin session state"); + serde_json::Value::Null + } + }, + None => serde_json::Value::Null, + }; + self.fire_session_autocmd( + "SessionFocus", + serde_json::json!({ "state_snapshot": state_snapshot }), + ); + } + pub(crate) fn checkpoint_session(&mut self) { let snapshot = self.session_snapshot(); if session_has_content(&snapshot) { @@ -421,6 +439,7 @@ impl App { self.state.session = AppSession::new(&self.state.session.model, &self.state.session.cwd); self.hydrate_plugin_state(); self.fire_session_autocmd("SessionReset", serde_json::json!({})); + self.fire_session_focus_autocmd(); vec![Action::NewSession] } @@ -486,6 +505,7 @@ impl App { } self.reset_ui_chrome(); self.restore_display(); + self.fire_session_focus_autocmd(); self.enqueue_save(); self.loaded_session_snapshot() diff --git a/n00n-ui/src/event_loop.rs b/n00n-ui/src/event_loop.rs index fdaafd92a..5f7caaa82 100644 --- a/n00n-ui/src/event_loop.rs +++ b/n00n-ui/src/event_loop.rs @@ -529,6 +529,7 @@ impl<'t> EventLoop<'t> { for w in startup_warnings { app.flash(w); } + app.fire_session_focus_autocmd(); let (submission_persist_tx, submission_persist_rx) = flume::unbounded(); Ok(Self { @@ -1069,6 +1070,7 @@ impl<'t> EventLoop<'t> { } self.sessions[self.focused].app.save_session(); self.focused = idx; + self.sessions[self.focused].app.fire_session_focus_autocmd(); } /// Focus a live session, or bring a stored one up: in place when the diff --git a/plugins/todo_write/init.lua b/plugins/todo_write/init.lua index 4c1962c0c..a24fb4f40 100644 --- a/plugins/todo_write/init.lua +++ b/plugins/todo_write/init.lua @@ -10,6 +10,7 @@ local seen_first = false local running = {} local running_order = {} local activity_expanded = false +local focused_session = nil local render_panel local STATUS_MARKERS = { @@ -45,10 +46,14 @@ local function current_todo() return nil end +local function session_is_focused(session_id) + return focused_session == nil or session_id == nil or session_id == focused_session +end + local function running_count() local count = 0 for _, activity in pairs(running) do - if activity.tool ~= "todo_write" then + if activity.tool ~= "todo_write" and session_is_focused(activity.session_id) then count = count + 1 end end @@ -58,7 +63,7 @@ end local function current_activity() for i = #running_order, 1, -1 do local activity = running[running_order[i]] - if activity and activity.tool ~= "todo_write" then + if activity and activity.tool ~= "todo_write" and session_is_focused(activity.session_id) then local label = compact_text(activity.summary ~= "" and activity.summary or activity.tool) if activity.subagent and activity.subagent ~= "" then label = activity.subagent .. ": " .. label @@ -68,6 +73,10 @@ local function current_activity() end end +local function activity_key(data) + return (data.session_id or "") .. "\0" .. data.id +end + local function prune_running_order() local compact = {} for _, id in ipairs(running_order) do @@ -126,9 +135,10 @@ local function ensure_win(visible) }) end -local function build_lines() +local function build_lines(todo_items, include_activity) + todo_items = todo_items or items local lines = {} - local activity = current_activity() + local activity = include_activity ~= false and current_activity() or nil if activity then local width = (win and win.width or n00n.ui.terminal_size().cols) - 16 local truncated = n00n.ui.truncate_text(activity, math.max(width, 8)) @@ -141,7 +151,7 @@ local function build_lines() if activity_expanded then for _, id in ipairs(running_order) do local detail = running[id] - if detail and detail.tool ~= "todo_write" then + if detail and detail.tool ~= "todo_write" and session_is_focused(detail.session_id) then local summary = compact_text(detail.summary ~= "" and detail.summary or detail.tool) local owner = detail.subagent and detail.subagent ~= "" and (detail.subagent .. " ยท ") or "" lines[#lines + 1] = { @@ -153,7 +163,7 @@ local function build_lines() end end end - for _, item in ipairs(items) do + for _, item in ipairs(todo_items) do local marker = STATUS_MARKERS[item.status] or STATUS_MARKERS.pending lines[#lines + 1] = { { marker[1] .. " " .. item.content, marker[2] }, @@ -214,38 +224,56 @@ n00n.api.register_tool({ end, restore = function(input) - items = input.todos or {} - if #items == 0 then + local restored_items = input.todos or {} + if #restored_items == 0 then return nil end - update_hint() - return ToolView.restore_lines(build_lines(), { max_lines = DEFAULT_PREVIEW_LINES, keep = "head" }) + return ToolView.restore_lines(build_lines(restored_items, false), { + max_lines = DEFAULT_PREVIEW_LINES, + keep = "head", + }) end, handler = function(input, ctx) - items = input.todos or {} + local owner, owner_err = ctx:state_owner("root") + if owner_err then + error(owner_err) + end + local next_items = input.todos or {} local _, state_err - if #items == 0 then + if #next_items == 0 then _, state_err = ctx:state_remove("root") else - _, state_err = ctx:state_replace("root", { todos = items }) + _, state_err = ctx:state_replace("root", { todos = next_items }) end if state_err then error(state_err) end - if #items == 0 then - if win and win:is_open() then - win:hide() + local is_focused = focused_session == nil or owner == focused_session + if is_focused then + focused_session = owner + items = next_items + end + if #next_items == 0 then + if is_focused then + if win and win:is_open() then + win:hide() + end + n00n.ui.set_status_hint(nil) end - n00n.ui.set_status_hint(nil) return "Todos cleared" end - local first = not seen_first - seen_first = true - render_panel(first) + if is_focused then + local first = not seen_first + seen_first = true + render_panel(first) + end return { llm_output = "", - body = ToolView.restore_lines(build_lines(), { max_lines = DEFAULT_PREVIEW_LINES, keep = "head" }), + body = ToolView.restore_lines(build_lines(next_items, false), { + max_lines = DEFAULT_PREVIEW_LINES, + keep = "head", + }), } end, }) @@ -289,17 +317,21 @@ n00n.api.create_autocmd("ToolStart", { if not data.id or data.tool == "todo_write" then return end - if running[data.id] then - running[data.id] = nil + local key = activity_key(data) + if running[key] then + running[key] = nil prune_running_order() end - running[data.id] = { + running[key] = { tool = data.tool or "tool", summary = data.summary or "", subagent = data.subagent, + session_id = data.session_id, } - running_order[#running_order + 1] = data.id - refresh_activity() + running_order[#running_order + 1] = key + if session_is_focused(data.session_id) then + refresh_activity() + end end, }) @@ -307,27 +339,69 @@ n00n.api.create_autocmd("ToolDone", { callback = function(ev) local data = ev.data or {} if data.id then - running[data.id] = nil + running[activity_key(data)] = nil prune_running_order() if running_count() == 0 then activity_expanded = false end - refresh_activity() + if session_is_focused(data.session_id) then + refresh_activity() + end + end + end, +}) +local function focused_items(data) + local snapshot = data.state_snapshot + local plugins = type(snapshot) == "table" and snapshot.plugins or nil + local plugin = type(plugins) == "table" and plugins.todo_write or nil + local root = type(plugin) == "table" and plugin.root or nil + local payload = type(root) == "table" and root.payload or nil + return type(payload) == "table" and type(payload.todos) == "table" and payload.todos or {} +end + +n00n.api.create_autocmd("SessionFocus", { + callback = function(ev) + local data = ev.data or {} + focused_session = data.session_id + items = focused_items(data) + seen_first = #items > 0 + activity_expanded = false + if #items == 0 then + if win and win:is_open() then + win:hide() + end + n00n.ui.set_status_hint(nil) + elseif win and win:is_open() and win:is_visible() then + render_panel(true) + else + update_hint() end end, }) -local function clear_activity() - running = {} - running_order = {} - activity_expanded = false - refresh_activity() +local function clear_activity(ev) + local data = ev and ev.data or {} + for id, activity in pairs(running) do + if data.session_id == nil or activity.session_id == data.session_id then + running[id] = nil + end + end + prune_running_order() + if session_is_focused(data.session_id) then + activity_expanded = false + refresh_activity() + end end -local function clear_todos() +local function clear_todos(ev) + local data = ev and ev.data or {} + if not session_is_focused(data.session_id) then + clear_activity(ev) + return + end items = {} seen_first = false - clear_activity() + clear_activity(ev) if win and win:is_open() then win:hide() end From 483c5479725fd38925343eac89470f45a4d81c81 Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Tue, 4 Aug 2026 10:30:36 -0400 Subject: [PATCH 5/5] docs: regenerate lua-api docs --- site/docs/content/lua-api/_index.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/site/docs/content/lua-api/_index.md b/site/docs/content/lua-api/_index.md index e31191420..7a70872d0 100644 --- a/site/docs/content/lua-api/_index.md +++ b/site/docs/content/lua-api/_index.md @@ -451,7 +451,8 @@ Listen for one or more events. Returns an id you can pass to `del_autocmd` later to remove the listener. Built-in events fired by the host: `"TurnStart"`, `"TurnEnd"`, -`"TurnError"`, `"ToolStart"`, `"ToolDone"`, `"SessionReset"`. +`"TurnError"`, `"ToolStart"`, `"ToolDone"`, `"SessionReset"`, +`"SessionFocus"`. Plugins can also fire their own events with `exec_autocmds`.