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, 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/// One session's conversation items and state of the current turn.
18#[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(&notification.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}