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