From 6cef85e4aaba6721a76519b6974bde60a18fa25d Mon Sep 17 00:00:00 2001 From: Tryanks Date: Wed, 30 Sep 2026 20:20:02 +0800 Subject: [PATCH] refactor(session): let the host own the running turn and its requests A client folds a byte-bounded window of a session's log and used to decide liveness itself: mark_idle unless the status said the turn ran. That lost the pending question and the working indicator whenever a stale cached status settled the fold, and #554 patched it with a refold on every running flip plus back-to-back history pages until the running turn's start arrived, downloading the whole running turn (often several MB). The host already holds the complete fold and knows whether its provider is live, so SessionStatus now carries the running turn (its index in the whole log and start time), the pending approvals and the pending question, gated on an in-flight turn. The client folds held records purely and settles which held turn runs from the status, in place and without discarding anything, whenever either topic changes. An unready plan is hidden while its turn is not running instead of being discarded. PROTOCOL_VERSION 7. --- crates/core/src/session.rs | 58 +++- crates/protocol/src/event.rs | 14 +- crates/protocol/src/lib.rs | 5 +- crates/runtime/src/app/approvals.rs | 22 +- crates/runtime/src/app/command_validation.rs | 2 +- crates/runtime/src/app/events.rs | 11 + crates/runtime/src/app/snapshots.rs | 36 ++- crates/runtime/src/app/tests.rs | 63 +++- crates/ui/src/chat/mod.rs | 13 +- .../ui/src/composer/components/user_input.rs | 20 +- crates/ui/src/shell.rs | 5 +- crates/ui/src/store/history.rs | 20 +- crates/ui/src/store/intents.rs | 4 +- crates/ui/src/store/mod.rs | 297 ++++++++++-------- crates/ui/src/store/snapshots.rs | 13 +- 15 files changed, 387 insertions(+), 196 deletions(-) diff --git a/crates/core/src/session.rs b/crates/core/src/session.rs index a3d8d39c..1430b2be 100644 --- a/crates/core/src/session.rs +++ b/crates/core/src/session.rs @@ -485,13 +485,23 @@ pub struct ProposedPlan { /// A structured question set the agent is waiting on, or working past, from /// [`AgentEvent::UserInputRequested`]. -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct PendingUserInput { pub request_id: String, pub questions: Vec, pub delivery: UserInputDelivery, } +/// The turn a live provider is running, located in the session's whole log so +/// a fold of any window of it can find the turn among its own. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunningTurn { + /// Index of the turn among every turn the whole log folds to. + pub turn: u64, + /// When the turn began, if the record that opened it carried a time. + pub started_at: Option, +} + /// Folded view of a session's event history. #[derive(Debug, Clone, Default)] pub struct Timeline { @@ -584,6 +594,39 @@ impl Timeline { } } + /// The running turn of a fold that holds the whole log. + pub fn running_turn(&self) -> Option { + let turn = self.turns.len().checked_sub(1)?; + (self.turn_running && self.turns[turn].running).then(|| RunningTurn { + turn: turn as u64, + started_at: self.turns[turn].start_ts, + }) + } + + /// Take liveness from the host instead of the records: only the turn at + /// `running` runs, timed from `started_at` when given. A fold of a window + /// cannot tell a live turn from one whose provider stopped without a + /// record, and a window cut inside the turn misses its start. Nothing the + /// records built is discarded, so a later settle can revive the turn. + pub fn settle_running_turn(&mut self, running: Option, started_at: Option) { + let running = running.filter(|turn| *turn < self.turns.len()); + self.turn_running = running.is_some(); + for (index, turn) in self.turns.iter_mut().enumerate() { + turn.running = running == Some(index); + } + if let (Some(turn), Some(started_at)) = (running, started_at) { + self.turns[turn].start_ts = Some(started_at); + } + } + + /// The proposed plan to show: a streamed plan stays display-only while + /// its turn runs and is dropped once the turn stops without finalizing it. + pub fn shown_proposed_plan(&self) -> Option<&ProposedPlan> { + self.proposed_plan + .as_ref() + .filter(|plan| plan.ready || self.turns.get(plan.turn).is_some_and(|turn| turn.running)) + } + /// The latest finalized plan, unless that exact provider item has already /// been implemented, handed off, or dismissed. pub fn plan_ready(&self) -> Option<&ProposedPlan> { @@ -2539,6 +2582,15 @@ mod tests { interrupted = timeline.clone(); interrupted.mark_idle(); assert!(interrupted.proposed_plan.is_none()); + // A client settled idle hides the streamed plan but can revive it. + let mut settled = timeline.clone(); + settled.settle_running_turn(None, None); + assert!(settled.shown_proposed_plan().is_none()); + settled.settle_running_turn(Some(0), None); + assert_eq!( + settled.shown_proposed_plan().unwrap().markdown, + "# Plan\nstep one" + ); timeline.apply_at( None, @@ -2550,6 +2602,10 @@ mod tests { timeline.apply_at(None, &turn_completed()); timeline.mark_idle(); assert_eq!(timeline.plan_ready().unwrap().markdown, "# Final plan"); + assert_eq!( + timeline.shown_proposed_plan().unwrap().markdown, + "# Final plan" + ); for resolution in [ agent::PlanResolution::Implemented, agent::PlanResolution::Dismissed, diff --git a/crates/protocol/src/event.rs b/crates/protocol/src/event.rs index 2de12dd5..80fb782a 100644 --- a/crates/protocol/src/event.rs +++ b/crates/protocol/src/event.rs @@ -2,15 +2,15 @@ use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use agent::{ - ApprovalMode, InteractionMode, OptionDescriptor, OptionSelection, ProviderCommand, - ProviderKind, RewindMode, + ApprovalMode, ApprovalRequest, InteractionMode, OptionDescriptor, OptionSelection, + ProviderCommand, ProviderKind, RewindMode, }; use serde::{Deserialize, Serialize}; use tcode_core::{ git::{GitAction, GitStatus}, project::{Project, SessionMeta, WorktreeInfo}, provider_status::ProviderSnapshot, - session::{ReviewComment, StoredEvent}, + session::{PendingUserInput, ReviewComment, RunningTurn, StoredEvent}, settings::Settings, ui::{TerminalSplitDirection, WorkspaceMode}, }; @@ -268,8 +268,12 @@ pub struct SessionStatus { #[serde(default)] pub stopping: bool, pub working: bool, - pub pending_approval: bool, - pub pending_user_input: bool, + /// The live turn and the requests it waits on. A client's fold of a + /// history window can place them but cannot decide them: the window may + /// begin after they opened, and a provider can stop without a record. + pub running_turn: Option, + pub pending_approvals: Vec, + pub pending_user_input: Option, #[serde(rename = "supports_steering")] pub steering_supported: bool, pub provider_option_descriptors: Vec, diff --git a/crates/protocol/src/lib.rs b/crates/protocol/src/lib.rs index d63298f2..3deab44d 100644 --- a/crates/protocol/src/lib.rs +++ b/crates/protocol/src/lib.rs @@ -37,8 +37,9 @@ pub use wire::{ // Version 4 adds client-generated command deduplication keys; version 5 moves // authentication into the transport, so hello carries no token; version 6 // sends index and history changes instead of whole replacements and -// compresses the native transport. -pub const PROTOCOL_VERSION: u32 = 6; +// compresses the native transport; version 7 carries the running turn and +// the requests it waits on in the session status. +pub const PROTOCOL_VERSION: u32 = 7; #[cfg(test)] mod tests; diff --git a/crates/runtime/src/app/approvals.rs b/crates/runtime/src/app/approvals.rs index 600c4d77..736477b7 100644 --- a/crates/runtime/src/app/approvals.rs +++ b/crates/runtime/src/app/approvals.rs @@ -31,10 +31,6 @@ impl AppState { self.approval_requests(session_id).first() } - pub(super) fn has_approval(&self, session_id: &str) -> bool { - !self.approval_requests(session_id).is_empty() - } - pub(super) fn clear_approvals(&mut self, session_id: &str) { self.approvals.remove(session_id); } @@ -115,6 +111,7 @@ mod tests { let (commands, receiver) = smol::channel::unbounded(); let mut child = live_session("child", commands); child.meta.parent_session_id = Some("parent".into()); + child.turn_in_flight = true; state.sessions.push(child.meta.clone()); if parked { state.residents.parked.insert("child".into(), child); @@ -134,11 +131,13 @@ mod tests { let status = state.child_status_json(&state.sessions[0], &Timeline::default()); assert_eq!(status["approval_request_id"], "first"); assert_eq!(status["waiting_approval"], "command `cargo test`"); - assert!( + assert_eq!( state .session_status_snapshot("child") .unwrap() - .pending_approval + .pending_approvals + .len(), + 2 ); state @@ -148,7 +147,7 @@ mod tests { matches!(receiver.try_recv(), Ok(SessionCommand::RespondApproval { request_id, decision: ApprovalDecision::Deny }) if request_id == "first") ); assert_eq!(state.first_approval("child").unwrap().id, "second"); - assert!(state.has_approval("parent")); + assert!(!state.approval_requests("parent").is_empty()); drop(receiver); assert!( @@ -169,16 +168,17 @@ mod tests { }, ); assert!( - !state + state .session_status_snapshot("child") .unwrap() - .pending_approval + .pending_approvals + .is_empty() ); assert_eq!( state.child_status_json(&state.sessions[0], &Timeline::default())["approval_request_id"], serde_json::Value::Null ); - assert!(state.has_approval("parent")); + assert!(!state.approval_requests("parent").is_empty()); state.record_approval_event( "parent", &AgentEvent::TurnCompleted { @@ -187,7 +187,7 @@ mod tests { usage: None, }, ); - assert!(!state.has_approval("parent")); + assert!(state.approval_requests("parent").is_empty()); } } } diff --git a/crates/runtime/src/app/command_validation.rs b/crates/runtime/src/app/command_validation.rs index 19e3c78f..3cbb2afb 100644 --- a/crates/runtime/src/app/command_validation.rs +++ b/crates/runtime/src/app/command_validation.rs @@ -461,7 +461,7 @@ mod tests { ))); } assert!( - state.has_approval(&id), + !state.approval_requests(&id).is_empty(), "failed delivery must leave the approval pending" ); } diff --git a/crates/runtime/src/app/events.rs b/crates/runtime/src/app/events.rs index 3f2d0784..31abff37 100644 --- a/crates/runtime/src/app/events.rs +++ b/crates/runtime/src/app/events.rs @@ -544,6 +544,17 @@ impl AppState { self.record_event_at(session_id, ts, event, cx); } + /// Deliver an event as a session's provider would. + #[cfg(any(test, feature = "test-support"))] + pub fn provider_event_for_test( + &mut self, + session_id: &str, + event: AgentEvent, + cx: &mut HostCx, + ) { + self.on_event(session_id, event, cx); + } + /// Give a new session an immediate first-message fallback, then ask a fresh /// background provider session for a concise title. The hidden request has /// no resume cursor or MCP servers, so it never enters the conversation or diff --git a/crates/runtime/src/app/snapshots.rs b/crates/runtime/src/app/snapshots.rs index 57580383..975e02a3 100644 --- a/crates/runtime/src/app/snapshots.rs +++ b/crates/runtime/src/app/snapshots.rs @@ -170,12 +170,13 @@ impl AppState { .ids() .filter_map(|id| { let session = self.resident(id)?; + let (approvals, user_input) = self.open_requests(id, session); Some(( id.to_string(), ( session.has_work(), - self.has_approval(id), - session.timeline.pending_user_input.is_some(), + !approvals.is_empty(), + user_input.is_some(), session.background_task_count > 0 && !session.turn_in_flight && session.queue.is_empty(), @@ -301,6 +302,25 @@ impl AppState { }) } + /// The requests a session waits on. A provider shut down without a closing + /// record leaves them open in the timeline; only an in-flight turn waits. + fn open_requests<'a>( + &'a self, + session_id: &str, + session: &'a ActiveSession, + ) -> ( + &'a [agent::ApprovalRequest], + Option<&'a tcode_core::session::PendingUserInput>, + ) { + if !session.turn_in_flight { + return (&[], None); + } + ( + self.approval_requests(session_id), + session.timeline.pending_user_input.as_ref(), + ) + } + /// Build the complete non-event-stream status projection for one resident /// session. This is the sole constructor for the replicated status domain. pub fn session_status_snapshot(&self, session_id: &str) -> Option { @@ -330,8 +350,8 @@ impl AppState { ) }) }); - let pending_approval = self.has_approval(session_id); let terminal_preferences = self.terminal_preferences_for(session); + let (approvals, user_input) = self.open_requests(session_id, session); Some(SessionStatus { session_id: session_id.to_string(), title: meta.title.clone(), @@ -384,8 +404,14 @@ impl AppState { turn_running: session.turn_in_flight, stopping: session.interrupt_requested, working: session.has_work(), - pending_approval, - pending_user_input: session.timeline.pending_user_input.is_some(), + // A provider shut down without a closing record leaves its turn + // running in the timeline; only an in-flight turn is live. + running_turn: session + .turn_in_flight + .then(|| session.timeline.running_turn()) + .flatten(), + pending_approvals: approvals.to_vec(), + pending_user_input: user_input.cloned(), steering_supported: session.can_steer(), provider_option_descriptors, provider_option_selections: meta.option_selections.clone(), diff --git a/crates/runtime/src/app/tests.rs b/crates/runtime/src/app/tests.rs index 563769e6..975848fa 100644 --- a/crates/runtime/src/app/tests.rs +++ b/crates/runtime/src/app/tests.rs @@ -3882,7 +3882,7 @@ fn child_approval_policy_distinguishes_report_tools_and_deduplicates_notices() { "{provider:?} {mode:?} {tool:?}" ); assert!(parent_receiver.try_recv().is_err()); - assert!(!state.has_approval("child")); + assert!(state.approval_requests("child").is_empty()); } else { assert!( child_receiver.try_recv().is_err(), @@ -7738,6 +7738,67 @@ fn interrupt_reports_stopping_until_the_turn_completes() { assert!(!status.turn_running); } +#[test] +fn status_carries_the_running_turn_and_its_question_only_while_in_flight() { + let store = TestStore::new("status-running-turn"); + let mut state = AppState::new((*store).clone()); + let context = TestAppContext::default(); + let mut cx = context.host_cx(); + let id = state.start_draft("fixture".into(), std::env::temp_dir(), &mut cx); + let (commands, _receiver) = smol::channel::unbounded(); + state.resident_mut(&id).unwrap().runtime = Runtime::Live(commands); + for turn in ["turn-1", "turn-2"] { + state.on_event( + &id, + AgentEvent::TurnStarted { + turn_id: turn.into(), + }, + &mut cx, + ); + if turn == "turn-1" { + state.on_event( + &id, + AgentEvent::TurnCompleted { + turn_id: turn.into(), + status: TurnStatus::Completed, + usage: None, + }, + &mut cx, + ); + } + } + state.on_event( + &id, + AgentEvent::UserInputRequested { + request_id: "ask".into(), + questions: Vec::new(), + delivery: agent::UserInputDelivery::Blocking, + }, + &mut cx, + ); + let status = state.session_status_snapshot(&id).unwrap(); + let started_at = state.resident(&id).unwrap().timeline.turns[1].start_ts; + assert!(started_at.is_some()); + assert_eq!( + status.running_turn, + Some(tcode_core::session::RunningTurn { + turn: 1, + started_at + }) + ); + assert_eq!( + status.pending_user_input.map(|pending| pending.request_id), + Some("ask".into()) + ); + + // Shut down without a closing record: the timeline still holds both. + state.resident_mut(&id).unwrap().shutdown_to_idle(); + assert!(state.resident(&id).unwrap().timeline.turn_running); + let status = state.session_status_snapshot(&id).unwrap(); + assert_eq!(status.running_turn, None); + assert_eq!(status.pending_user_input, None); +} + /// Times the initial snapshot plus the pages a client fetches while scrolling /// to the top of a real long thread. Run with /// `TCODE_HISTORY_BENCH_LOG=/path/to/thread.jsonl cargo test -p tcode-runtime diff --git a/crates/ui/src/chat/mod.rs b/crates/ui/src/chat/mod.rs index 37c90883..e9074245 100644 --- a/crates/ui/src/chat/mod.rs +++ b/crates/ui/src/chat/mod.rs @@ -685,8 +685,7 @@ impl ChatView { &timeline.turns, &timeline.entries, timeline - .proposed_plan - .as_ref() + .shown_proposed_plan() .map(|plan| (plan.turn, plan.item_id.as_str(), plan.markdown.as_str())), &self.expanded, continuity, @@ -943,7 +942,7 @@ impl ChatView { } } } - if let Some(plan) = &timeline.proposed_plan { + if let Some(plan) = timeline.shown_proposed_plan() { let id = format!("plan:{}", plan.item_id); if decisions.build.contains(&id) { texts.push(( @@ -1012,8 +1011,7 @@ impl ChatView { .with_active_timeline(|timeline| { if let Some(item_id) = id.strip_prefix("plan:") { return timeline - .proposed_plan - .as_ref() + .shown_proposed_plan() .filter(|plan| plan.item_id == item_id) .and_then(|plan| rows_of_turn(&self.rows, plan.turn).last()); } @@ -1428,8 +1426,7 @@ impl ChatView { .read(cx) .with_active_timeline(|timeline| { timeline - .proposed_plan - .as_ref() + .shown_proposed_plan() .filter(|plan| plan.turn == index) .map(|plan| (plan.item_id.clone(), plan.markdown.clone())) }) @@ -2828,7 +2825,7 @@ fn markdown_entries_for_residency( }); } } - if let Some(plan) = &timeline.proposed_plan + if let Some(plan) = timeline.shown_proposed_plan() && let Some(row) = rows_of_turn(rows, plan.turn).last() && scope.includes(row) { diff --git a/crates/ui/src/composer/components/user_input.rs b/crates/ui/src/composer/components/user_input.rs index 57213113..599a1c22 100644 --- a/crates/ui/src/composer/components/user_input.rs +++ b/crates/ui/src/composer/components/user_input.rs @@ -819,13 +819,19 @@ mod tests { .unwrap(); let (session_id, timeline) = smol::block_on(host.update_state_for_test(|state, cx| { let id = state.start_draft("async-question".into(), std::env::temp_dir(), cx); - let active = state.residents.live.get_mut(&id).unwrap(); - active.timeline.pending_user_input = Some(PendingUserInput { - request_id: "ask".into(), - questions: vec![question("first"), question("second")], - delivery: agent::UserInputDelivery::Async, - }); - (id, active.timeline.clone()) + for event in [ + agent::AgentEvent::TurnStarted { + turn_id: "turn".into(), + }, + agent::AgentEvent::UserInputRequested { + request_id: "ask".into(), + questions: vec![question("first"), question("second")], + delivery: agent::UserInputDelivery::Async, + }, + ] { + state.provider_event_for_test(&id, event, cx); + } + (id.clone(), state.residents.live[&id].timeline.clone()) })) .unwrap(); let store = cx.new(|cx| WorkspaceStore::new(host.link(), cx)); diff --git a/crates/ui/src/shell.rs b/crates/ui/src/shell.rs index fb7e56af..9baaf0bf 100644 --- a/crates/ui/src/shell.rs +++ b/crates/ui/src/shell.rs @@ -4906,8 +4906,9 @@ mod tests { turn_running: false, stopping: false, working: false, - pending_approval: false, - pending_user_input: false, + running_turn: None, + pending_approvals: Vec::new(), + pending_user_input: None, steering_supported: false, provider_option_descriptors: Vec::new(), provider_option_selections: Vec::new(), diff --git a/crates/ui/src/store/history.rs b/crates/ui/src/store/history.rs index d2627700..d060dbdc 100644 --- a/crates/ui/src/store/history.rs +++ b/crates/ui/src/store/history.rs @@ -27,27 +27,14 @@ impl WorkspaceStore { pub(super) fn load_pending_chat_history(&mut self, cx: &mut Context) { if self.history_error.is_none() - && (self.pending_chat_turn.as_ref().is_some_and(|(id, turn)| { + && self.pending_chat_turn.as_ref().is_some_and(|(id, turn)| { self.selected_session_id.as_ref() == Some(id) && *turn < self.session_turn_offset - }) || self.running_turn_opening_unheld()) + }) { self.load_earlier_messages(cx); } } - /// A window cut inside the running turn holds none of the records that - /// opened it, so its fold can neither show the turn as live nor time it. - /// Such a window holds only that turn; any earlier turn in it would mean - /// the running one is held from its start. - fn running_turn_opening_unheld(&self) -> bool { - self.active_turn_running() - && self.session_replica.as_ref().is_some_and(|(id, timeline)| { - self.selected_session_id.as_ref() == Some(id) - && !timeline.turn_running - && timeline.turns.len() == 1 - }) - } - /// Geometry is reported after layout, including the first frame on restore. /// Event counts cannot predict the height of folded turns. pub(crate) fn update_history_window(&mut self, screens: f32, cx: &mut Context) { @@ -114,10 +101,11 @@ impl WorkspaceStore { let held = store.session_records.entry(session_id.clone()).or_default(); held.splice(0..0, records); store.session_from.insert(session_id.clone(), from); - let timeline = store.fold_held_records(&session_id); + let mut timeline = store.fold_held_records(&session_id); store.session_turn_offset = store .session_turn_offset .saturating_sub(timeline.turns.len().saturating_sub(previous_turns)); + store.settle_running_turn(&mut timeline); store.session_replica = Some((session_id, timeline)); } Err(error) => { diff --git a/crates/ui/src/store/intents.rs b/crates/ui/src/store/intents.rs index 1043cca6..53f18cdc 100644 --- a/crates/ui/src/store/intents.rs +++ b/crates/ui/src/store/intents.rs @@ -236,8 +236,8 @@ impl WorkspaceStore { status.session_id.clone(), ( status.working, - status.pending_approval, - status.pending_user_input, + !status.pending_approvals.is_empty(), + status.pending_user_input.is_some(), Self::status_background_only(status), ), ); diff --git a/crates/ui/src/store/mod.rs b/crates/ui/src/store/mod.rs index db962ca0..accfcc61 100644 --- a/crates/ui/src/store/mod.rs +++ b/crates/ui/src/store/mod.rs @@ -961,12 +961,16 @@ impl WorkspaceStore { self.session_statuses .insert(session_id.clone(), status.clone()); if self.selected_session_id.as_ref() == Some(session_id) { - let was_running = self.active_turn_running(); let mut status = status.clone(); status.native_rewind_prefill_available = self.native_rewind_prefills.contains_key(session_id); self.session_status_replica = Some(status); - self.reconcile_replica_liveness(session_id, was_running); + if let Some((replica_id, mut timeline)) = self.session_replica.take() { + if replica_id == *session_id { + self.settle_running_turn(&mut timeline); + } + self.session_replica = Some((replica_id, timeline)); + } self.sync_terminal_topics(); self.sync_active_conversation_ui(); self.background_session_flags.remove(session_id); @@ -975,8 +979,8 @@ impl WorkspaceStore { session_id.clone(), ( status.working, - status.pending_approval, - status.pending_user_input, + !status.pending_approvals.is_empty(), + status.pending_user_input.is_some(), Self::status_background_only(status), ), ); @@ -1035,11 +1039,12 @@ impl WorkspaceStore { if self.session_catching_up { return; } - let timeline = self.fold_held_records(session_id); + let mut timeline = self.fold_held_records(session_id); self.baseline_topics.insert(envelope.topic.clone()); self.hydrated_sessions.insert(session_id.clone()); self.session_turn_offset = (*total_turns as usize).saturating_sub(timeline.turns.len()); + self.settle_running_turn(&mut timeline); self.session_replica = Some((session_id.clone(), timeline)); } (Topic::SessionEvents { session_id }, ServerEvent::SessionEvent(record)) => { @@ -1332,37 +1337,27 @@ impl WorkspaceStore { .is_some_and(|status| status.turn_running) } - /// The host's status, not the held records, decides whether a turn is - /// live: records folded after their provider stopped still end running. fn fold_held_records(&self, session_id: &str) -> Timeline { - let mut timeline = - Timeline::fold_stored(self.session_records.get(session_id).into_iter().flatten()); - if !self.active_turn_running() { - timeline.mark_idle(); - } - timeline + Timeline::fold_stored(self.session_records.get(session_id).into_iter().flatten()) } - /// Status and events are separate topics, so the status that settles a - /// fold can arrive after it. Idling is cheap to apply in place; reviving - /// needs a refold because `mark_idle` discarded the live turn's state. - fn reconcile_replica_liveness(&mut self, session_id: &str, was_running: bool) { - let running = self.active_turn_running(); - if running == was_running { - return; - } - let Some((replica_id, timeline)) = self.session_replica.as_mut() else { - return; - }; - if replica_id != session_id { - return; - } - if !running { - timeline.mark_idle(); - } else if !timeline.turn_running { - let timeline = self.fold_held_records(session_id); - self.session_replica = Some((session_id.to_owned(), timeline)); - } + /// Records folded after their provider stopped still end running, and a + /// window cut inside the running turn misses its start: the host's status + /// places the running turn among the held ones. Status and events are + /// separate topics, so each settles the replica whenever it changes. + fn settle_running_turn(&self, timeline: &mut Timeline) { + let running = self + .session_status_replica + .as_ref() + .and_then(|status| status.running_turn); + timeline.settle_running_turn( + running.and_then(|running| { + usize::try_from(running.turn) + .ok()? + .checked_sub(self.session_turn_offset) + }), + running.and_then(|running| running.started_at), + ); } fn suppress_task_auto_open_if_running(&mut self) { @@ -1908,7 +1903,7 @@ impl WorkspaceStore { self.session_status_replica .as_ref() .filter(|status| status.session_id == session_id) - .map(|status| status.pending_approval) + .map(|status| !status.pending_approvals.is_empty()) .or_else(|| { self.background_session_flags .get(session_id) @@ -1921,7 +1916,7 @@ impl WorkspaceStore { self.session_status_replica .as_ref() .filter(|status| status.session_id == session_id) - .map(|status| status.pending_user_input) + .map(|status| status.pending_user_input.is_some()) .or_else(|| { self.background_session_flags .get(session_id) @@ -2874,8 +2869,7 @@ impl WorkspaceStore { self.with_active_timeline(|timeline| { ( timeline - .proposed_plan - .as_ref() + .shown_proposed_plan() .map(|plan| plan.markdown.clone()), timeline.plan_steps.clone(), ) @@ -3720,16 +3714,34 @@ mod tests { status } - fn with_turn_running( + fn with_running_turn( status: &tcode_protocol::SessionStatus, - running: bool, + running: Option, + question: Option, ) -> tcode_protocol::SessionStatus { let mut status = status.clone(); - status.turn_running = running; - status.working = running; + status.turn_running = running.is_some(); + status.working = running.is_some(); + status.running_turn = running; + status.pending_user_input = question; status } + fn question() -> tcode_core::session::PendingUserInput { + tcode_core::session::PendingUserInput { + request_id: "ask".into(), + questions: vec![agent::UserInputQuestion { + id: "which".into(), + header: "Next".into(), + question: "Which way?".into(), + options: Vec::new(), + multi_select: false, + prefill: None, + }], + delivery: agent::UserInputDelivery::Blocking, + } + } + fn recorded(ts: u64, event: AgentEvent) -> SessionEventRecord { SessionEventRecord { ts: Some(ts), @@ -3787,66 +3799,110 @@ mod tests { } #[gpui::test] - fn status_arriving_after_the_snapshot_settles_the_running_turn(cx: &mut TestAppContext) { + fn the_status_settles_the_running_turn_whichever_topic_arrives_first(cx: &mut TestAppContext) { let status = draft_status(); let id = status.session_id.clone(); - let (to_host, _outgoing) = async_channel::unbounded(); - let (_incoming, from_host) = async_channel::unbounded(); - let link = tcode_client::HostLink::new(to_host, from_host); - let workspace = cx.new(|cx| { - WorkspaceStore::new_attached(link, WorkspaceAttachment::Local, None, None, false, cx) - }); - let status_event = |running| EventEnvelope { + let status_event = |running: Option, question| EventEnvelope { request_id: None, topic: Topic::SessionStatus { session_id: id.clone(), }, - event: ServerEvent::SessionStatusReplaced(with_turn_running(&status, running)), + event: ServerEvent::SessionStatusReplaced(with_running_turn( + &status, + running.map(|turn| tcode_core::session::RunningTurn { + turn, + started_at: Some(1_000), + }), + question, + )), }; - workspace.update(cx, |store, cx| { - store.selected_session_id = Some(id.clone()); - // The status cached from the last visit, before this turn began. - store.session_status_replica = Some(with_turn_running(&status, false)); - store.apply_domain_event( - &session_snapshot( - &id, - 0, - vec![ - recorded( - 1_000, - AgentEvent::TurnStarted { - turn_id: "turn".into(), - }, - ), - reply(2_000), - ], + let snapshot = session_snapshot( + &id, + 0, + vec![ + recorded( + 1_000, + AgentEvent::TurnStarted { + turn_id: "turn".into(), + }, ), - cx, - ); - store.apply_domain_event(&status_event(true), cx); - }); - assert_eq!(live_turn(&workspace, cx), Some(1_000)); + reply(2_000), + recorded( + 3_000, + AgentEvent::UserInputRequested { + request_id: "ask".into(), + questions: question().questions, + delivery: agent::UserInputDelivery::Blocking, + }, + ), + ], + ); + for status_first in [false, true] { + let (to_host, _outgoing) = async_channel::unbounded(); + let (_incoming, from_host) = async_channel::unbounded(); + let link = tcode_client::HostLink::new(to_host, from_host); + let workspace = cx.new(|cx| { + WorkspaceStore::new_attached( + link, + WorkspaceAttachment::Local, + None, + None, + false, + cx, + ) + }); + let open_question = |cx: &TestAppContext| { + workspace.read_with(cx, |store, _| { + store + .composer_state() + .pending_user_input + .map(|pending| pending.request_id) + }) + }; + workspace.update(cx, |store, cx| { + store.selected_session_id = Some(id.clone()); + // The status cached from the last visit, before this turn began. + store.session_status_replica = Some(status.clone()); + if status_first { + store.apply_domain_event(&status_event(Some(0), Some(question())), cx); + store.apply_domain_event(&snapshot, cx); + } else { + store.apply_domain_event(&snapshot, cx); + assert_eq!(store.with_active_timeline(|t| t.turn_running), Some(false)); + store.apply_domain_event(&status_event(Some(0), Some(question())), cx); + } + }); + assert_eq!(live_turn(&workspace, cx), Some(1_000)); + assert_eq!(open_question(cx).as_deref(), Some("ask")); - workspace.update(cx, |store, cx| { - store.apply_domain_event(&status_event(false), cx) - }); - assert_eq!(live_turn(&workspace, cx), None); + // The next turn runs before its opening record arrives. + workspace.update(cx, |store, cx| { + store.apply_domain_event(&status_event(Some(1), None), cx) + }); + assert_eq!(live_turn(&workspace, cx), None); + + workspace.update(cx, |store, cx| { + store.apply_domain_event(&status_event(None, None), cx) + }); + assert_eq!(live_turn(&workspace, cx), None); + assert_eq!(open_question(cx), None); + } } #[gpui::test] - fn a_window_cut_inside_the_running_turn_loads_back_to_its_start(cx: &mut TestAppContext) { - let status = with_turn_running(&draft_status(), true); + fn a_window_cut_inside_the_running_turn_is_timed_from_the_host(cx: &mut TestAppContext) { + let status = with_running_turn( + &draft_status(), + Some(tcode_core::session::RunningTurn { + turn: 0, + started_at: Some(1_000), + }), + None, + ); let id = status.session_id.clone(); let (to_host, outgoing) = async_channel::unbounded(); - let (incoming, from_host) = async_channel::unbounded(); + let (_incoming, from_host) = async_channel::unbounded(); let link = tcode_client::HostLink::new(to_host, from_host); - let pump_link = link.clone(); - let executor = cx.background_executor.clone(); - let _pump = cx.background_executor.spawn(async move { - pump_link - .pump_with_timer(|| executor.timer(std::time::Duration::from_millis(25))) - .await; - }); let workspace = cx.new(|cx| { WorkspaceStore::new_attached(link, WorkspaceAttachment::Local, None, None, false, cx) }); @@ -3856,46 +3912,18 @@ mod tests { store.apply_domain_event(&session_snapshot(&id, 3, vec![reply(4_000)]), cx); }); cx.run_until_parked(); - assert_eq!(live_turn(&workspace, cx), None); - - let request = std::iter::from_fn(|| outgoing.try_recv().ok()) - .map(|line| tcode_protocol::decode_client_line(&line).unwrap()) - .find(|request| { - matches!( + assert_eq!(live_turn(&workspace, cx), Some(1_000)); + assert!( + !std::iter::from_fn(|| outgoing.try_recv().ok()) + .map(|line| tcode_protocol::decode_client_line(&line).unwrap()) + .any(|request| matches!( request.payload, tcode_protocol::ClientPayload::Query( - tcode_protocol::Query::SessionHistoryPage { before: 3, .. } + tcode_protocol::Query::SessionHistoryPage { .. } ) - ) - }) - .expect("the page holding the turn's start"); - incoming - .try_send( - tcode_protocol::encode_line(&tcode_protocol::HostMessage::QueryResult { - id: request.id, - result: Ok(tcode_protocol::QueryResponse::SessionHistoryPage { - records: vec![ - recorded( - 1_000, - AgentEvent::TurnStarted { - turn_id: "turn".into(), - }, - ), - reply(2_000), - reply(3_000), - ], - from: 0, - end: 3, - truncated: false, - }), - }) - .unwrap(), - ) - .unwrap(); - wait_until(cx, &workspace, "the turn's start applied", |cx| { - live_turn(&workspace, cx).is_some() - }); - assert_eq!(live_turn(&workspace, cx), Some(1_000)); + )), + "the turn's start comes with the status, not from earlier pages" + ); } fn test_host(store: SessionStore) -> SpawnedHost { @@ -4814,8 +4842,19 @@ mod tests { .expect("first session status"); parked.turn_running = true; parked.working = true; - parked.pending_user_input = true; - parked.pending_approval = true; + parked.pending_user_input = Some(tcode_core::session::PendingUserInput { + request_id: "ask".into(), + questions: Vec::new(), + delivery: agent::UserInputDelivery::Blocking, + }); + parked.pending_approvals = vec![agent::ApprovalRequest { + id: "approve".into(), + turn_id: None, + kind: agent::ApprovalKind::FileRead { + detail: "Cargo.toml".into(), + }, + options: Vec::new(), + }]; store.apply_domain_event( &EventEnvelope { request_id: None, @@ -4832,8 +4871,8 @@ mod tests { next.cwd = second.cwd.clone(); next.turn_running = false; next.working = false; - next.pending_user_input = false; - next.pending_approval = false; + next.pending_user_input = None; + next.pending_approvals.clear(); store.apply_domain_event( &EventEnvelope { request_id: None, @@ -4870,8 +4909,8 @@ mod tests { let mut finished = store.session_statuses[&first.id].clone(); finished.turn_running = false; finished.working = false; - finished.pending_user_input = false; - finished.pending_approval = false; + finished.pending_user_input = None; + finished.pending_approvals.clear(); store.apply_domain_event( &EventEnvelope { request_id: None, diff --git a/crates/ui/src/store/snapshots.rs b/crates/ui/src/store/snapshots.rs index 3ccd8885..8552f2bc 100644 --- a/crates/ui/src/store/snapshots.rs +++ b/crates/ui/src/store/snapshots.rs @@ -176,7 +176,7 @@ pub(crate) fn composer_state( .map(|status| status.provider_commands.clone()) .unwrap_or_default(), attachments_dir: status.map(|status| status.attachments_dir.clone()), - pending_user_input: timeline.and_then(|timeline| timeline.pending_user_input.clone()), + pending_user_input: status.and_then(|status| status.pending_user_input.clone()), active_model: status.map(|status| ComposerActiveModel { provider: status.provider, model: status.requested_model.clone(), @@ -215,8 +215,8 @@ pub(crate) fn composer_state( checkout, turn_running: status.is_some_and(|status| status.turn_running), stopping: status.is_some_and(|status| status.stopping), - pending_approval: timeline.and_then(|timeline| timeline.pending_approvals.first().cloned()), - pending_approval_count: timeline.map_or(0, |timeline| timeline.pending_approvals.len()), + pending_approval: status.and_then(|status| status.pending_approvals.first().cloned()), + pending_approval_count: status.map_or(0, |status| status.pending_approvals.len()), } } @@ -241,7 +241,7 @@ pub(crate) fn panel_state( right_panel_expanded: ui.is_some_and(|ui| ui.right_panel_expanded), terminal_open: ui.is_some_and(|ui| ui.terminal_open), terminal_height: ui.map_or(240., |ui| ui.terminal_height), - plan_tab_active: timeline.is_some_and(|timeline| timeline.proposed_plan.is_some()) + plan_tab_active: timeline.is_some_and(|timeline| timeline.shown_proposed_plan().is_some()) || status.is_some_and(|status| status.interaction_mode == agent::InteractionMode::Plan), } } @@ -275,8 +275,9 @@ mod tests { turn_running: false, stopping: false, working: false, - pending_approval: false, - pending_user_input: false, + running_turn: None, + pending_approvals: Vec::new(), + pending_user_input: None, steering_supported: true, provider_option_descriptors: Vec::new(), provider_option_selections: Vec::new(),