Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 57 additions & 1 deletion crates/core/src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<UserInputQuestion>,
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<u64>,
}

/// Folded view of a session's event history.
#[derive(Debug, Clone, Default)]
pub struct Timeline {
Expand Down Expand Up @@ -584,6 +594,39 @@ impl Timeline {
}
}

/// The running turn of a fold that holds the whole log.
pub fn running_turn(&self) -> Option<RunningTurn> {
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<usize>, started_at: Option<u64>) {
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> {
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down
14 changes: 9 additions & 5 deletions crates/protocol/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
};
Expand Down Expand Up @@ -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<RunningTurn>,
pub pending_approvals: Vec<ApprovalRequest>,
pub pending_user_input: Option<PendingUserInput>,
#[serde(rename = "supports_steering")]
pub steering_supported: bool,
pub provider_option_descriptors: Vec<OptionDescriptor>,
Expand Down
5 changes: 3 additions & 2 deletions crates/protocol/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
22 changes: 11 additions & 11 deletions crates/runtime/src/app/approvals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
Expand All @@ -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
Expand All @@ -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!(
Expand All @@ -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 {
Expand All @@ -187,7 +187,7 @@ mod tests {
usage: None,
},
);
assert!(!state.has_approval("parent"));
assert!(state.approval_requests("parent").is_empty());
}
}
}
2 changes: 1 addition & 1 deletion crates/runtime/src/app/command_validation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -461,7 +461,7 @@ mod tests {
)));
}
assert!(
state.has_approval(&id),
!state.approval_requests(&id).is_empty(),
"failed delivery must leave the approval pending"
);
}
Expand Down
11 changes: 11 additions & 0 deletions crates/runtime/src/app/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
36 changes: 31 additions & 5 deletions crates/runtime/src/app/snapshots.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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<SessionStatus> {
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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(),
Expand Down
63 changes: 62 additions & 1 deletion crates/runtime/src/app/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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
Expand Down
Loading