Skip to main content

acp_utils/conversation/
mod.rs

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/// One session's conversation items and state of the current turn.
19#[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(&notification.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    /// The agent accepted the prompt.
90    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}