Skip to main content

lattice_ai/acp/
conversation.rs

1//! Structured conversation model for the ACP agent UI (AU‑1).
2//!
3//! The ACP supervisor receives structured `SessionUpdate`s (message chunks,
4//! thought chunks, tool calls, tool-call updates). AU‑1 stops flattening the
5//! *conversation* ones to text `AiLogger` records and instead folds them into a
6//! [`Conversation`] — a turn/block tree the `ai-conversation` mode (AU‑2)
7//! projects into the `*ai:opencode*` buffer. Trace sources (`Client` /
8//! `Lifecycle`) still flow to `AiLogger`; this completes the conversation/trace
9//! split.
10//!
11//! [`Conversation::apply`] is pure and unit-testable (no I/O, no locks).
12//! [`ConversationStore`] wraps it with a shared mutex + a [`ConversationUpdated`]
13//! bus publish so the mode can live-tail.
14
15use std::collections::HashMap;
16use std::sync::{Arc, Mutex};
17
18use agent_client_protocol::schema::v1::{
19    ContentBlock, PermissionOption, PermissionOptionId, PermissionOptionKind, SessionUpdate,
20    ToolCallStatus, ToolKind,
21};
22use std::fmt;
23use tokio::sync::oneshot;
24
25use lattice_agent::SessionKey;
26
27/// Who produced a turn.
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum Role {
30    User,
31    Assistant,
32}
33
34/// Execution state of a tool call, mapped from ACP's `ToolCallStatus`.
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub enum ToolStatus {
37    Pending,
38    Running,
39    Ok,
40    Err,
41}
42
43impl ToolStatus {
44    fn from_acp(status: ToolCallStatus) -> Self {
45        match status {
46            ToolCallStatus::Pending => ToolStatus::Pending,
47            ToolCallStatus::InProgress => ToolStatus::Running,
48            ToolCallStatus::Completed => ToolStatus::Ok,
49            ToolCallStatus::Failed => ToolStatus::Err,
50            // `ToolCallStatus` is `#[non_exhaustive]`.
51            _ => ToolStatus::Pending,
52        }
53    }
54}
55
56/// Review state of an agent-proposed file edit. AU‑4 drives the transitions and
57/// attaches the `review_diff` session; AU‑1 only defines the shape.
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59pub enum EditStatus {
60    Proposed,
61    Accepted,
62    Rejected,
63}
64
65/// AUX‑1: status of an inline permission request shown in the conversation.
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum PermissionStatus {
68    Pending,
69    Allowed,
70    Denied,
71}
72
73/// PU-B.2b: a pending permission request projected for the popup menu.
74/// `ai-permission-mode`'s `on_activate` reads [`ConversationStore::
75/// oldest_pending_permission`] and renders `title`, optional `description`, and
76/// one line per `options` entry (wire order); the select handler maps the chosen
77/// line back to `options[i].option_id`.
78#[derive(Debug, Clone)]
79pub struct PendingPermissionView {
80    pub id: String,
81    pub title: String,
82    pub description: Option<String>,
83    pub options: Vec<PermissionOption>,
84}
85
86/// AUX‑3: global processing status of the current agent session — shown in the
87/// conversation buffer's headerline.
88#[derive(Debug, Clone, PartialEq, Eq, Default)]
89pub enum SessionStatus {
90    /// No active turn: the agent is not processing anything.
91    #[default]
92    Idle,
93    /// The agent is streaming a text/thought response.
94    Thinking,
95    /// The agent is executing a tool call.
96    Executing {
97        /// Human-readable tool name (e.g. "edit parse.rs").
98        tool: String,
99    },
100    /// The agent is awaiting the user's decision on a permission request.
101    AwaitingPermission,
102}
103
104impl fmt::Display for SessionStatus {
105    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
106        match self {
107            SessionStatus::Idle => write!(f, "Ready"),
108            SessionStatus::Thinking => write!(f, "Thinking\u{2026}"),
109            SessionStatus::Executing { tool } => write!(f, "Working: {tool}"),
110            SessionStatus::AwaitingPermission => write!(f, "Awaiting your approval\u{2026}"),
111        }
112    }
113}
114
115/// AUX‑2: token usage snapshot from a `UsageUpdate` notification.
116#[derive(Debug, Clone, PartialEq)]
117pub struct Cost {
118    pub amount: f64,
119    pub currency: String,
120}
121
122/// AUX‑2: latest token and cost snapshot from the ACP `usage_update`.
123#[derive(Debug, Clone, PartialEq)]
124pub struct UsageSnapshot {
125    pub used: u64,
126    pub size: u64,
127    pub cost: Option<Cost>,
128}
129
130/// One renderable unit within a turn.
131#[derive(Debug, Clone, PartialEq, Eq)]
132pub enum Block {
133    /// Streamed message text.
134    Text(String),
135    /// Streamed reasoning / thinking. TCF: the projection prefixes each line
136    /// with `│`, and a multi-line block is folded **closed by default** via
137    /// [`ReasoningFoldSource`](crate::acp::tool_fold::ReasoningFoldSource) — the
138    /// fold keeps the head line visible and hides the rest until `za`. (Before
139    /// TCF this comment claimed folding the projection did not do.)
140    Reasoning(String),
141    /// A tool invocation and its live status.
142    ///
143    /// TCF: `kind`, `input` and `output` are the detail an expanded tool call
144    /// shows. `input`/`output` are the agent's raw JSON, pretty-printed at
145    /// ingest (stored as `String`, not `serde_json::Value`, so `Block` stays
146    /// `Eq`). Detail commonly arrives on a later `ToolCallUpdate`, not the
147    /// initial call, so it is merged in as it lands.
148    ToolCall {
149        id: String,
150        title: String,
151        status: ToolStatus,
152        kind: ToolKind,
153        input: Option<String>,
154        output: Option<String>,
155    },
156    /// An agent-proposed file edit (AU‑4 wires the diff review).
157    Edit { path: String, status: EditStatus },
158    /// AUX‑1: an inline permission request awaiting or reflecting user action.
159    Permission {
160        id: String,
161        title: String,
162        description: Option<String>,
163        options: Vec<PermissionOption>,
164        status: PermissionStatus,
165    },
166}
167
168/// One turn: a role and its ordered blocks.
169#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct Turn {
171    pub role: Role,
172    pub blocks: Vec<Block>,
173}
174
175/// The full conversation as a turn/block tree. Streaming *extends* the last
176/// block; earlier turns are never rewritten.
177#[derive(Debug, Clone, Default, PartialEq)]
178pub struct Conversation {
179    pub turns: Vec<Turn>,
180    /// AUX‑2: latest token usage snapshot from `usage_update` notifications.
181    pub usage: Option<UsageSnapshot>,
182    /// AUX‑3: global processing status derived by the supervisor from the active
183    /// turn's state. Set via [`ConversationStore::set_status`].
184    pub status: SessionStatus,
185}
186
187impl Conversation {
188    /// Fold one ACP `SessionUpdate` into the model. Pure: no I/O, no locks.
189    /// Non-conversation updates (plans, mode changes, ...) are ignored;
190    /// `SessionUpdate` and `ContentBlock` are `#[non_exhaustive]`, so a
191    /// catch-all closes each match.
192    pub fn apply(&mut self, update: &SessionUpdate) {
193        match update {
194            SessionUpdate::UserMessageChunk(chunk) => {
195                if let ContentBlock::Text(t) = &chunk.content {
196                    self.extend_text(Role::User, &t.text);
197                }
198            }
199            SessionUpdate::AgentMessageChunk(chunk) => {
200                if let ContentBlock::Text(t) = &chunk.content {
201                    self.extend_text(Role::Assistant, &t.text);
202                }
203            }
204            SessionUpdate::AgentThoughtChunk(chunk) => {
205                if let ContentBlock::Text(t) = &chunk.content {
206                    self.extend_reasoning(&t.text);
207                }
208            }
209            SessionUpdate::ToolCall(tc) => {
210                self.push_tool_call(
211                    tc.tool_call_id.0.to_string(),
212                    tc.title.clone(),
213                    ToolStatus::from_acp(tc.status),
214                    tc.kind,
215                    tc.raw_input.as_ref().map(pretty_json),
216                    tc.raw_output.as_ref().map(pretty_json),
217                );
218            }
219            SessionUpdate::ToolCallUpdate(u) => {
220                self.merge_tool_update(
221                    u.tool_call_id.0.as_ref(),
222                    u.fields.status.map(ToolStatus::from_acp),
223                    u.fields.kind,
224                    u.fields.raw_input.as_ref().map(pretty_json),
225                    u.fields.raw_output.as_ref().map(pretty_json),
226                );
227            }
228            // AUX‑2: accumulate the latest usage snapshot.
229            SessionUpdate::UsageUpdate(u) => {
230                self.usage = Some(UsageSnapshot {
231                    used: u.used,
232                    size: u.size,
233                    cost: u.cost.as_ref().map(|c| Cost {
234                        amount: c.amount,
235                        currency: c.currency.clone(),
236                    }),
237                });
238            }
239            _ => {}
240        }
241    }
242
243    /// AU‑3: append a complete user prompt as a new `User` turn. Each Enter is
244    /// a distinct turn (the terminal-REPL model), so — unlike chunk-streamed
245    /// agent output — this always opens a fresh turn rather than extending.
246    pub fn push_user_text(&mut self, text: &str) {
247        self.turns.push(Turn {
248            role: Role::User,
249            blocks: vec![Block::Text(text.to_string())],
250        });
251    }
252
253    /// Append `text` to the last block if it is `Text` in a turn of `role`;
254    /// otherwise open a new block / turn as needed.
255    fn extend_text(&mut self, role: Role, text: &str) {
256        match self.turns.last_mut() {
257            Some(turn) if turn.role == role => match turn.blocks.last_mut() {
258                Some(Block::Text(s)) => s.push_str(text),
259                _ => turn.blocks.push(Block::Text(text.to_string())),
260            },
261            _ => self.turns.push(Turn {
262                role,
263                blocks: vec![Block::Text(text.to_string())],
264            }),
265        }
266    }
267
268    /// Append `text` to the last `Reasoning` block (assistant turn), opening one
269    /// as needed.
270    fn extend_reasoning(&mut self, text: &str) {
271        match self.turns.last_mut() {
272            Some(turn) if turn.role == Role::Assistant => match turn.blocks.last_mut() {
273                Some(Block::Reasoning(s)) => s.push_str(text),
274                _ => turn.blocks.push(Block::Reasoning(text.to_string())),
275            },
276            _ => self.turns.push(Turn {
277                role: Role::Assistant,
278                blocks: vec![Block::Reasoning(text.to_string())],
279            }),
280        }
281    }
282
283    /// Push a new `ToolCall` block onto the current (or a fresh) assistant turn.
284    #[allow(clippy::too_many_arguments)]
285    fn push_tool_call(
286        &mut self,
287        id: String,
288        title: String,
289        status: ToolStatus,
290        kind: ToolKind,
291        input: Option<String>,
292        output: Option<String>,
293    ) {
294        let block = Block::ToolCall {
295            id,
296            title,
297            status,
298            kind,
299            input,
300            output,
301        };
302        match self.turns.last_mut() {
303            Some(turn) if turn.role == Role::Assistant => turn.blocks.push(block),
304            _ => self.turns.push(Turn {
305                role: Role::Assistant,
306                blocks: vec![block],
307            }),
308        }
309    }
310
311    /// AUX‑1: push a `Permission` block onto the current (or a fresh) assistant turn.
312    fn push_permission_block(
313        &mut self,
314        id: String,
315        title: String,
316        description: Option<String>,
317        options: Vec<PermissionOption>,
318    ) {
319        let block = Block::Permission {
320            id,
321            title,
322            description,
323            options,
324            status: PermissionStatus::Pending,
325        };
326        match self.turns.last_mut() {
327            Some(turn) if turn.role == Role::Assistant => turn.blocks.push(block),
328            _ => self.turns.push(Turn {
329                role: Role::Assistant,
330                blocks: vec![block],
331            }),
332        }
333    }
334
335    /// PU-B.2: kind of the option `option_id` within the `Permission` block
336    /// `id` (searched newest-first), or `None` if the block or option is gone.
337    /// Used by `resolve_permission` to derive the inline block status from the
338    /// agent's actual option rather than a fixed 4-way bucket.
339    fn permission_option_kind(
340        &self,
341        id: &str,
342        option_id: &PermissionOptionId,
343    ) -> Option<PermissionOptionKind> {
344        for turn in self.turns.iter().rev() {
345            for block in turn.blocks.iter().rev() {
346                if let Block::Permission {
347                    id: bid, options, ..
348                } = block
349                    && bid == id
350                {
351                    return options
352                        .iter()
353                        .find(|o| &o.option_id == option_id)
354                        .map(|o| o.kind);
355                }
356            }
357        }
358        None
359    }
360
361    /// AUX‑1: update the status of the `Permission` block with `id`.
362    fn update_permission_status(&mut self, id: &str, new_status: PermissionStatus) {
363        for turn in self.turns.iter_mut().rev() {
364            for block in turn.blocks.iter_mut().rev() {
365                if let Block::Permission {
366                    id: bid, status, ..
367                } = block
368                    && bid == id
369                {
370                    *status = new_status;
371                    return;
372                }
373            }
374        }
375    }
376
377    /// TCF: merge a `ToolCallUpdate` into the tool-call block with `id`
378    /// (searched newest-first). Each `Some` field overwrites; `None` leaves the
379    /// existing value — detail (input/output/kind) commonly arrives on an
380    /// update after the initial call, so this accumulates rather than replaces.
381    fn merge_tool_update(
382        &mut self,
383        id: &str,
384        new_status: Option<ToolStatus>,
385        new_kind: Option<ToolKind>,
386        new_input: Option<String>,
387        new_output: Option<String>,
388    ) {
389        for turn in self.turns.iter_mut().rev() {
390            for block in turn.blocks.iter_mut().rev() {
391                if let Block::ToolCall {
392                    id: bid,
393                    status,
394                    kind,
395                    input,
396                    output,
397                    ..
398                } = block
399                    && bid == id
400                {
401                    if let Some(s) = new_status {
402                        *status = s;
403                    }
404                    if let Some(k) = new_kind {
405                        *kind = k;
406                    }
407                    if new_input.is_some() {
408                        *input = new_input;
409                    }
410                    if new_output.is_some() {
411                        *output = new_output;
412                    }
413                    return;
414                }
415            }
416        }
417    }
418}
419
420/// TCF: pretty-print a raw tool JSON payload for the expanded view. Falls back
421/// to the compact `Display` form if pretty-printing somehow fails.
422fn pretty_json(value: &serde_json::Value) -> String {
423    serde_json::to_string_pretty(value).unwrap_or_else(|_| value.to_string())
424}
425
426/// Fired after the [`ConversationStore`] mutates. The `ai-conversation` mode
427/// subscribes to it and re-projects (mirrors `AiLogPushed`).
428#[derive(Debug, Clone)]
429pub struct ConversationUpdated {
430    pub session: SessionKey,
431}
432
433lattice_protocol::register_event!(
434    ConversationUpdated,
435    "ai.conversation-updated",
436    "Fired after the ACP agent conversation model changes; drives live refresh \
437     of the *ai:<provider>* conversation buffer.",
438    "lattice-ai",
439);
440
441/// Fired by the `ai-conversation` mode's drain AFTER a re-projection edit has
442/// LANDED in the buffer. Boot wakes the editor actor on this (via
443/// `wake_on_event`) so a streamed agent response repaints WITHOUT a keystroke,
444/// and the per-tick prompt-focus callback runs.
445///
446/// Distinct from [`ConversationUpdated`], which the supervisor fires when the
447/// *model* changes — that is BEFORE the drain re-projects the buffer, so waking
448/// on it would repaint stale content. Sequencing the wake after the owner-write
449/// edit (this event) is what makes the last streamed chunk paint reliably.
450#[derive(Debug, Clone, Default)]
451pub struct ConversationProjected;
452
453lattice_protocol::register_event!(
454    ConversationProjected,
455    "ai.conversation-projected",
456    "Fired after the *ai:<provider>* conversation buffer is re-projected; wakes \
457     the render loop so streamed agent responses repaint without a keystroke.",
458    "lattice-ai",
459);
460
461/// Shared, mutable conversation state plus a bus publisher. The supervisor holds
462/// one clone and calls [`ConversationStore::apply`]; the `ai-conversation` mode
463/// reads [`ConversationStore::snapshot`] on each `ConversationUpdated`.
464#[derive(Clone)]
465pub struct ConversationStore {
466    inner: Arc<Mutex<Conversation>>,
467    publish: Arc<dyn Fn(ConversationUpdated) + Send + Sync>,
468    /// AUX‑1: pending permission request responders keyed by tool-call id.
469    /// PU-B.2: the oneshot carries the agent's chosen `PermissionOptionId`
470    /// (wire order, any arity), not the fixed 4-way `PermissionOutcome` — the
471    /// menu resolves by the option the agent actually offered.
472    pending_permissions:
473        Arc<Mutex<HashMap<String, (SessionKey, oneshot::Sender<PermissionOptionId>)>>>,
474}
475
476impl ConversationStore {
477    /// Build a store whose mutations publish `ConversationUpdated` via `publish`.
478    pub fn new(publish: Arc<dyn Fn(ConversationUpdated) + Send + Sync>) -> Self {
479        Self {
480            inner: Arc::new(Mutex::new(Conversation::default())),
481            publish,
482            pending_permissions: Arc::new(Mutex::new(HashMap::new())),
483        }
484    }
485
486    /// Fold `update` into the conversation for `session`, then publish.
487    pub fn apply(&self, session: &SessionKey, update: &SessionUpdate) {
488        {
489            let mut conv = self.inner.lock().expect("conversation mutex poisoned");
490            conv.apply(update);
491        }
492        (self.publish)(ConversationUpdated {
493            session: session.clone(),
494        });
495    }
496
497    /// AU‑3: fold a locally-composed user prompt into the conversation as a
498    /// `User` turn, then publish. ACP agents don't echo the user's prompt back,
499    /// so the supervisor calls this when it sends a prompt to make the user's
500    /// turn appear in the transcript immediately.
501    pub fn push_user_text(&self, session: &SessionKey, text: &str) {
502        {
503            let mut conv = self.inner.lock().expect("conversation mutex poisoned");
504            conv.push_user_text(text);
505        }
506        (self.publish)(ConversationUpdated {
507            session: session.clone(),
508        });
509    }
510
511    /// AUX‑1: push a permission block and register its oneshot responder. The
512    /// supervisor calls this when `classify_permission` returns `AskUser`; the
513    /// receiver is awaited in [`handle_permission`](supervisor::handle_permission).
514    pub fn push_permission_request(
515        &self,
516        session: &SessionKey,
517        id: String,
518        title: String,
519        description: Option<String>,
520        options: Vec<PermissionOption>,
521        responder: oneshot::Sender<PermissionOptionId>,
522    ) {
523        {
524            let mut conv = self.inner.lock().expect("conversation mutex poisoned");
525            conv.push_permission_block(id.clone(), title, description, options);
526            self.pending_permissions
527                .lock()
528                .expect("pending_permissions mutex poisoned")
529                .insert(id, (session.clone(), responder));
530        }
531        (self.publish)(ConversationUpdated {
532            session: session.clone(),
533        });
534    }
535
536    /// PU-B.2: resolve a pending permission request by the agent's chosen
537    /// `option_id`. Derives the inline block status from that option's `kind`,
538    /// updates the block, and sends the `option_id` through the oneshot the
539    /// supervisor is parked on (which answers ACP with `Selected(option_id)`),
540    /// then publishes. No-op when `id` is unknown (already resolved, deferred
541    /// and re-resolved, or never registered).
542    pub fn resolve_permission(&self, id: &str, option_id: PermissionOptionId) {
543        let session = {
544            let mut conv = self.inner.lock().expect("conversation mutex poisoned");
545            // Fail closed: an option whose kind we can't read (missing block, or
546            // a future `#[non_exhaustive]` kind) marks the request Denied rather
547            // than leaving it Pending or optimistically Allowed.
548            let status = match conv.permission_option_kind(id, &option_id) {
549                Some(PermissionOptionKind::AllowOnce | PermissionOptionKind::AllowAlways) => {
550                    PermissionStatus::Allowed
551                }
552                _ => PermissionStatus::Denied,
553            };
554            conv.update_permission_status(id, status);
555            let mut pending = self
556                .pending_permissions
557                .lock()
558                .expect("pending_permissions mutex poisoned");
559            let (session, sender) = match pending.remove(id) {
560                Some(entry) => entry,
561                None => return, // already resolved — no-op
562            };
563            let _ = sender.send(option_id);
564            session
565        };
566        (self.publish)(ConversationUpdated { session });
567    }
568
569    /// PU-B.2b: the oldest still-`Pending` permission request, projected for
570    /// the popup menu (`ai-permission-mode` reads it in `on_activate`). Turns
571    /// and blocks are in wire order, so the first `Pending` block is the oldest.
572    pub fn oldest_pending_permission(&self) -> Option<PendingPermissionView> {
573        self.oldest_pending_permission_where(|_| true)
574    }
575
576    /// PU-B.3: the oldest `Pending` request whose id satisfies `keep`. The
577    /// auto-open tick callback passes `|id| !deferred(id)` so an `Esc`-deferred
578    /// request is skipped (the queue advances past it) without re-opening it.
579    pub fn oldest_pending_permission_where(
580        &self,
581        keep: impl Fn(&str) -> bool,
582    ) -> Option<PendingPermissionView> {
583        let conv = self.inner.lock().expect("conversation mutex poisoned");
584        conv.turns
585            .iter()
586            .flat_map(|t| t.blocks.iter())
587            .find_map(|b| match b {
588                Block::Permission {
589                    id,
590                    title,
591                    description,
592                    options,
593                    status: PermissionStatus::Pending,
594                } if keep(id) => Some(PendingPermissionView {
595                    id: id.clone(),
596                    title: title.clone(),
597                    description: description.clone(),
598                    options: options.clone(),
599                }),
600                _ => None,
601            })
602    }
603
604    /// AUX‑3: set the global processing status and publish a
605    /// `ConversationUpdated` so the headerline re-renders.
606    pub fn set_status(&self, session: &SessionKey, status: SessionStatus) {
607        {
608            let mut conv = self.inner.lock().expect("conversation mutex poisoned");
609            conv.status = status;
610        }
611        (self.publish)(ConversationUpdated {
612            session: session.clone(),
613        });
614    }
615
616    /// Cheap-ish clone of the current conversation for projection.
617    pub fn snapshot(&self) -> Conversation {
618        self.inner
619            .lock()
620            .expect("conversation mutex poisoned")
621            .clone()
622    }
623}
624
625#[cfg(test)]
626mod tests {
627    #![allow(clippy::unwrap_used, clippy::panic)]
628    use super::*;
629    use agent_client_protocol::schema::v1::{
630        ContentChunk, PermissionOptionKind, TextContent, ToolCall as AcpToolCall, ToolCallUpdate,
631        ToolCallUpdateFields,
632    };
633
634    fn text_chunk(text: &str) -> ContentChunk {
635        ContentChunk::new(ContentBlock::Text(TextContent::new(text)))
636    }
637
638    #[test]
639    fn agent_text_chunks_extend_one_block() {
640        let mut c = Conversation::default();
641        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk("hi")));
642        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk(" there")));
643        assert_eq!(c.turns.len(), 1);
644        assert_eq!(c.turns[0].role, Role::Assistant);
645        assert_eq!(c.turns[0].blocks, vec![Block::Text("hi there".to_string())]);
646    }
647
648    /// AU‑3: a user prompt lands as its own `User` turn (the REPL model:
649    /// each Enter is distinct, never merged into agent output).
650    #[test]
651    fn push_user_text_appends_a_user_turn() {
652        let mut c = Conversation::default();
653        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk("hello")));
654        c.push_user_text("refactor parse_args");
655        assert_eq!(c.turns.len(), 2);
656        assert_eq!(c.turns[1].role, Role::User);
657        assert_eq!(
658            c.turns[1].blocks,
659            vec![Block::Text("refactor parse_args".to_string())]
660        );
661        // Two consecutive prompts are two distinct turns, not merged.
662        c.push_user_text("again");
663        assert_eq!(c.turns.len(), 3);
664        assert_eq!(c.turns[2].role, Role::User);
665    }
666
667    /// AU‑3: `ConversationStore::push_user_text` mutates the shared store and
668    /// publishes a `ConversationUpdated` so the mode's drain re-projects.
669    #[test]
670    fn store_push_user_text_mutates_and_publishes() {
671        use std::sync::atomic::{AtomicUsize, Ordering};
672        let published = Arc::new(AtomicUsize::new(0));
673        let p = published.clone();
674        let store = ConversationStore::new(Arc::new(move |_ev| {
675            p.fetch_add(1, Ordering::SeqCst);
676        }));
677        store.push_user_text(&SessionKey::new("opencode", 1), "hi");
678        assert_eq!(published.load(Ordering::SeqCst), 1, "one publish");
679        let snap = store.snapshot();
680        assert_eq!(snap.turns.len(), 1);
681        assert_eq!(snap.turns[0].role, Role::User);
682    }
683
684    #[test]
685    fn thought_chunks_extend_a_reasoning_block() {
686        let mut c = Conversation::default();
687        c.apply(&SessionUpdate::AgentThoughtChunk(text_chunk("think")));
688        c.apply(&SessionUpdate::AgentThoughtChunk(text_chunk("ing")));
689        assert_eq!(
690            c.turns[0].blocks,
691            vec![Block::Reasoning("thinking".to_string())]
692        );
693    }
694
695    #[test]
696    fn tool_call_then_update_mutates_in_place_by_id() {
697        let mut c = Conversation::default();
698        let mut tc = AcpToolCall::new("tc-1", "edit parse.rs");
699        tc.status = ToolCallStatus::InProgress;
700        c.apply(&SessionUpdate::ToolCall(tc));
701        assert_eq!(c.turns.len(), 1);
702        match &c.turns[0].blocks[0] {
703            Block::ToolCall { id, status, .. } => {
704                assert_eq!(id, "tc-1");
705                assert_eq!(*status, ToolStatus::Running);
706            }
707            other => panic!("expected ToolCall, got {other:?}"),
708        }
709
710        let update = ToolCallUpdate::new(
711            "tc-1",
712            ToolCallUpdateFields::new().status(ToolCallStatus::Completed),
713        );
714        c.apply(&SessionUpdate::ToolCallUpdate(update));
715        // Same block, updated status -- not a second block.
716        assert_eq!(c.turns[0].blocks.len(), 1);
717        match &c.turns[0].blocks[0] {
718            Block::ToolCall { status, .. } => assert_eq!(*status, ToolStatus::Ok),
719            other => panic!("expected ToolCall, got {other:?}"),
720        }
721    }
722
723    /// TCF: the tool-call detail (args on the initial call, output on a later
724    /// update) is captured, pretty-printed, and merged in place — not
725    /// discarded. This is what the expanded view will show.
726    #[test]
727    fn tool_call_captures_and_merges_input_then_output() {
728        let mut c = Conversation::default();
729        let mut tc = AcpToolCall::new("tc-1", "bash");
730        tc.raw_input = Some(serde_json::json!({ "cmd": "echo hello" }));
731        c.apply(&SessionUpdate::ToolCall(tc));
732        match &c.turns[0].blocks[0] {
733            Block::ToolCall { input, output, .. } => {
734                let input = input.as_ref().expect("input captured on the initial call");
735                assert!(input.contains("\"cmd\""), "pretty JSON input: {input}");
736                assert!(input.contains("echo hello"), "input value: {input}");
737                assert!(output.is_none(), "no output yet");
738            }
739            other => panic!("expected ToolCall, got {other:?}"),
740        }
741
742        // Output arrives on a later update — must merge, not clobber the input.
743        let update = ToolCallUpdate::new(
744            "tc-1",
745            ToolCallUpdateFields::new().raw_output(serde_json::json!("hello\n")),
746        );
747        c.apply(&SessionUpdate::ToolCallUpdate(update));
748        match &c.turns[0].blocks[0] {
749            Block::ToolCall { input, output, .. } => {
750                assert!(input.is_some(), "input survives the update merge");
751                let output = output.as_ref().expect("output captured on the update");
752                assert!(output.contains("hello"), "output value: {output}");
753            }
754            other => panic!("expected ToolCall, got {other:?}"),
755        }
756    }
757
758    #[test]
759    fn text_after_tool_call_opens_a_new_text_block_same_turn() {
760        let mut c = Conversation::default();
761        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk("before")));
762        c.apply(&SessionUpdate::ToolCall(AcpToolCall::new("t", "run")));
763        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk("after")));
764        assert_eq!(c.turns.len(), 1, "all assistant activity is one turn");
765        assert_eq!(c.turns[0].blocks.len(), 3);
766        assert_eq!(c.turns[0].blocks[2], Block::Text("after".to_string()));
767    }
768
769    // ── AUX‑1: permission block tests ──
770
771    fn test_permission_option(
772        id: &'static str,
773        name: &'static str,
774        kind: PermissionOptionKind,
775    ) -> PermissionOption {
776        PermissionOption::new(id, name, kind)
777    }
778
779    #[test]
780    fn permission_block_created_in_assistant_turn() {
781        let mut c = Conversation::default();
782        c.apply(&SessionUpdate::AgentMessageChunk(text_chunk("working")));
783        c.push_permission_block(
784            "perm-1".to_string(),
785            "Allow agent to run cargo test?".to_string(),
786            None,
787            vec![test_permission_option(
788                "a1",
789                "Allow once",
790                PermissionOptionKind::AllowOnce,
791            )],
792        );
793        assert_eq!(c.turns.len(), 1);
794        assert_eq!(c.turns[0].role, Role::Assistant);
795        assert_eq!(c.turns[0].blocks.len(), 2);
796        match &c.turns[0].blocks[1] {
797            Block::Permission {
798                id,
799                title,
800                status,
801                options,
802                ..
803            } => {
804                assert_eq!(id, "perm-1");
805                assert_eq!(title, "Allow agent to run cargo test?");
806                assert_eq!(*status, PermissionStatus::Pending);
807                assert_eq!(options.len(), 1);
808            }
809            other => panic!("expected Permission block, got {other:?}"),
810        }
811    }
812
813    #[test]
814    fn permission_block_opens_fresh_assistant_turn_when_no_previous() {
815        let mut c = Conversation::default();
816        c.push_permission_block("perm-1".to_string(), "Allow?".to_string(), None, vec![]);
817        assert_eq!(c.turns.len(), 1);
818        assert_eq!(c.turns[0].role, Role::Assistant);
819    }
820
821    #[test]
822    fn permission_block_updated_by_id() {
823        let mut c = Conversation::default();
824        c.push_permission_block("perm-1".to_string(), "Allow?".to_string(), None, vec![]);
825        c.update_permission_status("perm-1", PermissionStatus::Allowed);
826        match &c.turns[0].blocks[0] {
827            Block::Permission { status, .. } => assert_eq!(*status, PermissionStatus::Allowed),
828            other => panic!("expected Permission, got {other:?}"),
829        }
830    }
831
832    #[test]
833    fn store_push_permission_request_creates_block_and_registers_responder() {
834        let published = Arc::new(std::sync::atomic::AtomicUsize::new(0));
835        let p = published.clone();
836        let store = ConversationStore::new(Arc::new(move |_ev| {
837            p.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
838        }));
839        let session = SessionKey::new("opencode", 1);
840        let (tx, rx) = tokio::sync::oneshot::channel();
841
842        store.push_permission_request(
843            &session,
844            "perm-1".to_string(),
845            "Allow?".to_string(),
846            None,
847            vec![],
848            tx,
849        );
850
851        // One publish on push
852        assert_eq!(published.load(std::sync::atomic::Ordering::SeqCst), 1);
853        // Block exists in snapshot
854        let snap = store.snapshot();
855        assert_eq!(snap.turns.len(), 1);
856        assert!(matches!(
857            &snap.turns[0].blocks[0],
858            Block::Permission { id, status: PermissionStatus::Pending, .. } if id == "perm-1"
859        ));
860        // Responder is registered — dropping rx won't hang the test
861        drop(rx);
862    }
863
864    #[test]
865    fn store_resolve_permission_sends_outcome_and_publishes() {
866        let published = Arc::new(std::sync::atomic::AtomicUsize::new(0));
867        let p = published.clone();
868        let store = ConversationStore::new(Arc::new(move |_ev| {
869            p.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
870        }));
871        let session = SessionKey::new("opencode", 1);
872        let (tx, rx) = tokio::sync::oneshot::channel();
873
874        // PU-B.2: resolve by the agent's actual option; the block status derives
875        // from that option's kind (AllowOnce → Allowed).
876        let opt = test_permission_option("a1", "Allow once", PermissionOptionKind::AllowOnce);
877        let chosen = opt.option_id.clone();
878        store.push_permission_request(
879            &session,
880            "perm-1".to_string(),
881            "Allow?".to_string(),
882            None,
883            vec![opt],
884            tx,
885        );
886
887        // Reset publish count — the push already fired once
888        published.store(0, std::sync::atomic::Ordering::SeqCst);
889
890        store.resolve_permission("perm-1", chosen.clone());
891
892        // Block updated
893        let snap = store.snapshot();
894        match &snap.turns[0].blocks[0] {
895            Block::Permission { status, .. } => assert_eq!(*status, PermissionStatus::Allowed),
896            other => panic!("expected Permission, got {other:?}"),
897        }
898        // Oneshot delivered the chosen option id
899        assert_eq!(
900            rx.blocking_recv(),
901            Ok(chosen),
902            "responder must receive the chosen option id",
903        );
904        // Publish fired
905        assert_eq!(
906            published.load(std::sync::atomic::Ordering::SeqCst),
907            1,
908            "resolve must publish",
909        );
910    }
911
912    #[test]
913    fn store_resolve_permission_noop_for_unknown_id() {
914        let published = Arc::new(std::sync::atomic::AtomicUsize::new(0));
915        let p = published.clone();
916        let store = ConversationStore::new(Arc::new(move |_ev| {
917            p.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
918        }));
919
920        // No pending permission with this id → no-op, no publish
921        let opt = test_permission_option("a1", "Allow once", PermissionOptionKind::AllowOnce);
922        store.resolve_permission("nonexistent", opt.option_id);
923        assert_eq!(published.load(std::sync::atomic::Ordering::SeqCst), 0);
924    }
925
926    #[test]
927    fn oldest_pending_permission_returns_the_first_in_wire_order() {
928        let store = ConversationStore::new(Arc::new(|_| {}));
929        let session = SessionKey::new("opencode", 1);
930        let (tx1, _rx1) = tokio::sync::oneshot::channel();
931        let (tx2, _rx2) = tokio::sync::oneshot::channel();
932        store.push_permission_request(
933            &session,
934            "perm-1".to_string(),
935            "First?".to_string(),
936            None,
937            vec![test_permission_option(
938                "a",
939                "Allow",
940                PermissionOptionKind::AllowOnce,
941            )],
942            tx1,
943        );
944        store.push_permission_request(
945            &session,
946            "perm-2".to_string(),
947            "Second?".to_string(),
948            None,
949            vec![test_permission_option(
950                "b",
951                "Allow",
952                PermissionOptionKind::AllowOnce,
953            )],
954            tx2,
955        );
956
957        let oldest = store
958            .oldest_pending_permission()
959            .expect("a pending request exists");
960        assert_eq!(oldest.id, "perm-1", "oldest = first in wire order");
961        assert_eq!(oldest.title, "First?");
962        assert_eq!(oldest.options.len(), 1);
963
964        // Resolving the first surfaces the second as the new oldest.
965        store.resolve_permission("perm-1", oldest.options[0].option_id.clone());
966        assert_eq!(
967            store.oldest_pending_permission().map(|v| v.id),
968            Some("perm-2".to_string()),
969        );
970    }
971
972    // ── AUX‑2: usage update tests ──
973
974    fn update(used: u64, size: u64, cost: Option<(f64, &str)>) -> SessionUpdate {
975        let mut u = agent_client_protocol::schema::v1::UsageUpdate::new(used, size);
976        if let Some((amt, cur)) = cost {
977            u = u.cost(agent_client_protocol::schema::v1::Cost::new(amt, cur));
978        }
979        SessionUpdate::UsageUpdate(u)
980    }
981
982    #[test]
983    fn usage_update_stored() {
984        let mut c = Conversation::default();
985        c.apply(&update(53000, 200000, Some((0.045, "USD"))));
986        assert_eq!(
987            c.usage,
988            Some(UsageSnapshot {
989                used: 53000,
990                size: 200000,
991                cost: Some(Cost {
992                    amount: 0.045,
993                    currency: "USD".to_string(),
994                }),
995            }),
996        );
997    }
998
999    #[test]
1000    fn usage_update_overwrites() {
1001        let mut c = Conversation::default();
1002        c.apply(&update(1000, 200000, None));
1003        assert!(c.usage.as_ref().unwrap().cost.is_none());
1004        // Second update overwrites
1005        c.apply(&update(53000, 200000, Some((0.045, "USD"))));
1006        assert_eq!(c.usage.unwrap().used, 53000);
1007    }
1008
1009    /// The wire payload opencode 1.17.18 actually emits. The other usage tests
1010    /// build `SessionUpdate` in Rust, so they never exercise deserialization —
1011    /// the seam where a schema mismatch would silently drop the update.
1012    #[test]
1013    fn usage_update_from_opencode_wire_json() {
1014        let raw = serde_json::json!({
1015            "sessionUpdate": "usage_update",
1016            "used": 27744,
1017            "size": 200000,
1018            "cost": { "amount": 0, "currency": "USD" }
1019        });
1020        let update: SessionUpdate =
1021            serde_json::from_value(raw).expect("opencode usage_update should deserialize");
1022        let mut c = Conversation::default();
1023        c.apply(&update);
1024        let usage = c
1025            .usage
1026            .expect("usage should be populated from the wire payload");
1027        assert_eq!(usage.used, 27744);
1028        assert_eq!(usage.size, 200000);
1029        assert_eq!(usage.cost.map(|c| c.currency), Some("USD".to_string()));
1030    }
1031
1032    #[test]
1033    fn usage_update_no_cost() {
1034        let mut c = Conversation::default();
1035        c.apply(&update(53000, 200000, None));
1036        assert_eq!(c.usage.as_ref().unwrap().used, 53000);
1037        assert_eq!(c.usage.as_ref().unwrap().size, 200000);
1038        assert!(c.usage.as_ref().unwrap().cost.is_none());
1039    }
1040
1041    // ── AUX‑3: status tests ──
1042
1043    #[test]
1044    fn status_idle_default() {
1045        let c = Conversation::default();
1046        assert_eq!(c.status, SessionStatus::Idle);
1047    }
1048
1049    #[test]
1050    fn set_status_updates_and_publishes() {
1051        use std::sync::atomic::{AtomicUsize, Ordering};
1052        let published = Arc::new(AtomicUsize::new(0));
1053        let p = published.clone();
1054        let store = ConversationStore::new(Arc::new(move |_ev| {
1055            p.fetch_add(1, Ordering::SeqCst);
1056        }));
1057        let session = SessionKey::new("opencode", 1);
1058
1059        store.set_status(&session, SessionStatus::Thinking);
1060        assert_eq!(published.load(Ordering::SeqCst), 1);
1061        assert_eq!(store.snapshot().status, SessionStatus::Thinking);
1062
1063        store.set_status(&session, SessionStatus::Idle);
1064        assert_eq!(store.snapshot().status, SessionStatus::Idle);
1065    }
1066
1067    #[test]
1068    fn status_display_formats() {
1069        assert_eq!(SessionStatus::Idle.to_string(), "Ready");
1070        assert!(SessionStatus::Thinking.to_string().contains("Thinking"));
1071        let exec = SessionStatus::Executing {
1072            tool: "edit".into(),
1073        };
1074        assert_eq!(exec.to_string(), "Working: edit");
1075        assert!(
1076            SessionStatus::AwaitingPermission
1077                .to_string()
1078                .contains("Awaiting")
1079        );
1080    }
1081}