1mod activity;
2mod items;
3mod tool_calls;
4mod turn;
5
6pub use activity::Activity;
7pub use items::{ConversationContent, ConversationId, ConversationItem, ConversationItemId, ItemState, Revision};
8pub use tool_calls::{SubAgentState, SubAgentToolCall, ToolCall, ToolStatus};
9pub use turn::{TurnFinished, TurnInProgress, TurnPhase};
10
11use crate::client::AcpEvent;
12use crate::notifications::SubAgentProgressParams;
13use agent_client_protocol::schema::{MaybeUndefined, v2 as acp};
14use items::MessageRole;
15use std::collections::{HashMap, HashSet};
16use std::fmt::Display;
17
18#[derive(Debug)]
20pub struct Conversation {
21 id: ConversationId,
22 revision: Revision,
23 items: Vec<ConversationItem>,
24 tool_index: HashMap<String, usize>,
25 message_index: HashMap<acp::MessageId, usize>,
26 pending_user: Option<usize>,
27 next_item_id: u64,
28 turn: TurnPhase,
29 turn_ended: bool,
30 activity: Activity,
31 compactions: HashSet<acp::CompactionId>,
32 context_usage: Option<acp::UsageUpdate>,
33 plan: Option<acp::PlanItems>,
34}
35
36impl Default for Conversation {
37 fn default() -> Self {
38 Self::new()
39 }
40}
41
42impl Conversation {
43 pub fn new() -> Self {
44 Self {
45 id: ConversationId::next(),
46 revision: Revision::default(),
47 items: Vec::new(),
48 tool_index: HashMap::new(),
49 message_index: HashMap::new(),
50 pending_user: None,
51 next_item_id: 0,
52 turn: TurnPhase::Idle,
53 turn_ended: false,
54 activity: Activity::Idle,
55 compactions: HashSet::new(),
56 context_usage: None,
57 plan: None,
58 }
59 }
60
61 pub fn apply_event(&mut self, event: &AcpEvent) -> Option<TurnFinished> {
62 match event {
63 AcpEvent::SessionUpdate(notification) => return self.apply_update(¬ification.update),
64 AcpEvent::SubAgentProgress(progress) => self.apply_sub_agent_progress(progress),
65 AcpEvent::ContextCleared(_) => self.clear(),
66 AcpEvent::ConnectionClosed => self.connection_closed(),
67 AcpEvent::AuthMethodsUpdated(_)
68 | AcpEvent::McpNotification(_)
69 | AcpEvent::GitDiffEvent(_)
70 | AcpEvent::ElicitationRequest { .. } => {}
71 }
72 None
73 }
74
75 pub fn start_prompt(&mut self, echo: Option<Vec<acp::ContentBlock>>) -> Result<(), TurnInProgress> {
76 if !self.turn.is_idle() {
77 return Err(TurnInProgress);
78 }
79 if let Some(content) = echo {
80 self.push(ItemState::Open, ConversationContent::User(content));
81 let index = self.items.len() - 1;
82 self.items[index].preserve_user_display = true;
83 self.pending_user = Some(index);
84 }
85 self.begin_turn(TurnPhase::Submitting);
86 Ok(())
87 }
88
89 pub fn accept_prompt(&mut self) {
91 self.set_turn(self.turn.accepted());
92 }
93
94 pub fn fail_prompt(&mut self, reason: impl Display) {
95 if self.turn.waiting_for_response() {
96 self.end_turn(ToolStatus::Failed, Some(&reason.to_string()));
97 }
98 self.set_turn(TurnPhase::Idle);
99 }
100
101 pub fn clear(&mut self) {
102 *self = Self { turn: self.turn.finished(), revision: self.revision, ..Self::new() };
103 self.advance();
104 }
105
106 pub fn append_user_content(&mut self, content: Vec<acp::ContentBlock>) -> ConversationItemId {
107 self.push(ItemState::Sealed, ConversationContent::User(content))
108 }
109
110 pub fn append_notice(&mut self, text: impl Into<String>) -> ConversationItemId {
111 self.push(ItemState::Sealed, ConversationContent::Notice(text.into()))
112 }
113
114 pub fn id(&self) -> ConversationId {
115 self.id
116 }
117
118 pub fn revision(&self) -> Revision {
119 self.revision
120 }
121
122 pub fn items(&self) -> &[ConversationItem] {
123 &self.items
124 }
125
126 pub fn turn(&self) -> TurnPhase {
127 self.turn
128 }
129
130 pub fn activity(&self) -> Activity {
131 self.activity
132 }
133
134 pub fn plan(&self) -> Option<&acp::PlanItems> {
135 self.plan.as_ref()
136 }
137
138 pub fn context_usage(&self) -> Option<&acp::UsageUpdate> {
139 self.context_usage.as_ref()
140 }
141
142 pub fn is_compacting(&self) -> bool {
143 !self.compactions.is_empty()
144 }
145
146 pub fn any_running(&self) -> bool {
147 self.items.iter().any(|item| match item.content() {
148 ConversationContent::Tool(tool_call) => tool_call.is_running(),
149 _ => false,
150 })
151 }
152
153 fn apply_update(&mut self, update: &acp::SessionUpdate) -> Option<TurnFinished> {
154 if matches!(update, acp::SessionUpdate::StateUpdate(acp::StateUpdate::Running(_))) && self.turn.is_idle() {
155 self.begin_turn(TurnPhase::Running);
156 }
157 if self.turn.waiting_for_response()
158 && let Some(activity) = Activity::after(update)
159 {
160 self.set_activity(activity);
161 }
162 match update {
163 acp::SessionUpdate::CompactionUpdate(update) => {
164 if !self.turn_ended {
165 self.apply_compaction(update);
166 }
167 }
168 acp::SessionUpdate::StateUpdate(acp::StateUpdate::Idle(idle)) if self.turn.waiting_for_response() => {
169 return Some(self.finish_turn(idle.stop_reason.clone()));
170 }
171 acp::SessionUpdate::UserMessage(message) => {
172 self.upsert_message(MessageRole::User, &message.message_id, &message.content);
173 }
174 acp::SessionUpdate::AgentMessage(message) => {
175 self.upsert_message(MessageRole::Assistant, &message.message_id, &message.content);
176 }
177 acp::SessionUpdate::AgentThought(message) => {
178 self.upsert_message(MessageRole::Thought, &message.message_id, &message.content);
179 }
180 acp::SessionUpdate::UserMessageChunk(chunk) => self.append_message_chunk(MessageRole::User, chunk),
181 acp::SessionUpdate::AgentMessageChunk(chunk) => self.append_message_chunk(MessageRole::Assistant, chunk),
182 acp::SessionUpdate::AgentThoughtChunk(chunk) => self.append_message_chunk(MessageRole::Thought, chunk),
183 acp::SessionUpdate::ToolCallUpdate(update) => {
184 let index = self.tool_slot(&update.tool_call_id);
185 self.update_tool(index, |tool_call| tool_call.apply_update(update));
186 }
187 acp::SessionUpdate::ToolCallContentChunk(chunk) => {
188 let index = self.tool_slot(&chunk.tool_call_id);
189 self.update_tool(index, |tool_call| tool_call.append_content(chunk.content.clone()));
190 }
191 acp::SessionUpdate::PlanUpdate(update) => {
192 if let acp::PlanUpdateContent::Items(items) = &update.plan
193 && self.plan.as_ref() != Some(items)
194 {
195 self.plan = Some(items.clone());
196 self.advance();
197 }
198 }
199 acp::SessionUpdate::UsageUpdate(usage) if self.context_usage.as_ref() != Some(usage) => {
200 self.context_usage = Some(usage.clone());
201 self.advance();
202 }
203 _ => {}
204 }
205 None
206 }
207
208 fn apply_sub_agent_progress(&mut self, progress: &SubAgentProgressParams) {
209 if self.turn_ended {
210 return;
211 }
212 let Some(&index) = self.tool_index.get(&progress.parent_tool_id) else {
213 return;
214 };
215 self.update_tool(index, |tool_call| tool_call.apply_sub_agent_progress(progress));
216 }
217
218 fn connection_closed(&mut self) {
219 self.set_turn(TurnPhase::Idle);
220 self.turn_ended = true;
221 self.set_activity(Activity::Idle);
222 }
223
224 fn begin_turn(&mut self, phase: TurnPhase) {
225 self.turn = phase;
226 self.turn_ended = false;
227 self.activity = Activity::Thinking;
228 self.advance();
229 }
230
231 fn set_turn(&mut self, turn: TurnPhase) {
232 if self.turn != turn {
233 self.turn = turn;
234 self.advance();
235 }
236 }
237
238 fn set_activity(&mut self, activity: Activity) {
239 if self.activity != activity {
240 self.activity = activity;
241 self.advance();
242 }
243 }
244
245 fn finish_turn(&mut self, stop_reason: Option<acp::StopReason>) -> TurnFinished {
246 let status = match stop_reason {
247 Some(acp::StopReason::Cancelled) => ToolStatus::Cancelled,
248 _ => ToolStatus::Success,
249 };
250 self.end_turn(status, None);
251 TurnFinished { stop_reason }
252 }
253
254 fn end_turn(&mut self, status: ToolStatus, error: Option<&str>) {
255 self.turn = self.turn.finished();
256 self.turn_ended = true;
257 self.compactions.clear();
258 self.activity = Activity::Idle;
259 self.pending_user = None;
260 let revision = self.advance();
261 for item in self.items.iter_mut().filter(|item| item.is_open()) {
262 if let ConversationContent::Tool(tool_call) = &mut item.content {
263 tool_call.finalize(status, error);
264 }
265 item.state = ItemState::Sealed;
266 item.touch(revision, false);
267 }
268 }
269
270 fn apply_compaction(&mut self, update: &acp::CompactionUpdate) {
271 let changed = match update.status {
272 acp::CompactionStatus::InProgress => self.compactions.insert(update.compaction_id.clone()),
273 acp::CompactionStatus::Completed | acp::CompactionStatus::Failed | acp::CompactionStatus::Cancelled => {
274 self.compactions.remove(&update.compaction_id)
275 }
276 _ => false,
277 };
278 if changed {
279 self.advance();
280 }
281 }
282
283 fn upsert_message(
284 &mut self,
285 role: MessageRole,
286 message_id: &acp::MessageId,
287 content: &MaybeUndefined<Vec<acp::ContentBlock>>,
288 ) {
289 let blocks = match content {
290 MaybeUndefined::Value(blocks) => blocks.clone(),
291 MaybeUndefined::Null if self.message_index.contains_key(message_id) => Vec::new(),
292 MaybeUndefined::Null | MaybeUndefined::Undefined => return,
293 };
294 let index = self.message_slot(role, message_id.clone());
295 let item = &mut self.items[index];
296 if item.preserve_user_display {
297 return;
298 }
299 let content = role.content(blocks);
300 if item.content != content {
301 item.content = content;
302 self.touch(index, true);
303 }
304 }
305
306 fn append_message_chunk(&mut self, role: MessageRole, chunk: &acp::ContentChunk) {
307 let index = self.message_slot(role, chunk.message_id.clone());
308 let item = &mut self.items[index];
309 if item.preserve_user_display {
310 return;
311 }
312 let rewrites = !item.is_open() || role == MessageRole::User;
313 if role.blocks_mut(&mut item.content).is_some_and(|blocks| items::append_block(blocks, &chunk.content)) {
314 self.touch(index, rewrites);
315 }
316 }
317
318 fn message_slot(&mut self, role: MessageRole, message_id: acp::MessageId) -> usize {
319 if let Some(&index) = self.message_index.get(&message_id) {
320 return index;
321 }
322 let index = if role == MessageRole::User { self.pending_user.take() } else { None }.unwrap_or_else(|| {
323 let index = self.items.len();
324 self.push(ItemState::Open, role.content(Vec::new()));
325 index
326 });
327 self.items[index].message_id = Some(message_id.clone());
328 self.message_index.insert(message_id, index);
329 index
330 }
331
332 fn tool_slot(&mut self, id: &acp::ToolCallId) -> usize {
333 if let Some(&index) = self.tool_index.get(id.0.as_ref()) {
334 return index;
335 }
336 let index = self.items.len();
337 self.push(
338 ItemState::Open,
339 ConversationContent::Tool(ToolCall::from_update(&acp::ToolCallUpdate::new(id.clone()))),
340 );
341 self.tool_index.insert(id.to_string(), index);
342 index
343 }
344
345 fn push(&mut self, state: ItemState, content: ConversationContent) -> ConversationItemId {
346 let id = ConversationItemId(self.next_item_id);
347 self.next_item_id = self.next_item_id.saturating_add(1);
348 let revision = self.advance();
349 self.items.push(ConversationItem::new(id, revision, state, content));
350 id
351 }
352
353 fn update_tool(&mut self, index: usize, apply: impl FnOnce(&mut ToolCall)) {
354 let item = &mut self.items[index];
355 let rewrites = !item.is_open();
356 let ConversationContent::Tool(tool_call) = &mut item.content else {
357 return;
358 };
359 let previous = tool_call.clone();
360 apply(tool_call);
361 if *tool_call == previous {
362 return;
363 }
364 item.state = if tool_call.rendering_final() { ItemState::Sealed } else { ItemState::Open };
365 self.touch(index, rewrites);
366 }
367
368 fn touch(&mut self, index: usize, rewrites: bool) {
369 let revision = self.advance();
370 self.items[index].touch(revision, rewrites);
371 }
372
373 fn advance(&mut self) -> Revision {
374 self.revision.advance()
375 }
376}