1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum Role {
30 User,
31 Assistant,
32}
33
34#[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 _ => ToolStatus::Pending,
52 }
53 }
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59pub enum EditStatus {
60 Proposed,
61 Accepted,
62 Rejected,
63}
64
65#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum PermissionStatus {
68 Pending,
69 Allowed,
70 Denied,
71}
72
73#[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#[derive(Debug, Clone, PartialEq, Eq, Default)]
89pub enum SessionStatus {
90 #[default]
92 Idle,
93 Thinking,
95 Executing {
97 tool: String,
99 },
100 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#[derive(Debug, Clone, PartialEq)]
117pub struct Cost {
118 pub amount: f64,
119 pub currency: String,
120}
121
122#[derive(Debug, Clone, PartialEq)]
124pub struct UsageSnapshot {
125 pub used: u64,
126 pub size: u64,
127 pub cost: Option<Cost>,
128}
129
130#[derive(Debug, Clone, PartialEq, Eq)]
132pub enum Block {
133 Text(String),
135 Reasoning(String),
141 ToolCall {
149 id: String,
150 title: String,
151 status: ToolStatus,
152 kind: ToolKind,
153 input: Option<String>,
154 output: Option<String>,
155 },
156 Edit { path: String, status: EditStatus },
158 Permission {
160 id: String,
161 title: String,
162 description: Option<String>,
163 options: Vec<PermissionOption>,
164 status: PermissionStatus,
165 },
166}
167
168#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct Turn {
171 pub role: Role,
172 pub blocks: Vec<Block>,
173}
174
175#[derive(Debug, Clone, Default, PartialEq)]
178pub struct Conversation {
179 pub turns: Vec<Turn>,
180 pub usage: Option<UsageSnapshot>,
182 pub status: SessionStatus,
185}
186
187impl Conversation {
188 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 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 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 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 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 #[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 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 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 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 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
420fn pretty_json(value: &serde_json::Value) -> String {
423 serde_json::to_string_pretty(value).unwrap_or_else(|_| value.to_string())
424}
425
426#[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#[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#[derive(Clone)]
465pub struct ConversationStore {
466 inner: Arc<Mutex<Conversation>>,
467 publish: Arc<dyn Fn(ConversationUpdated) + Send + Sync>,
468 pending_permissions:
473 Arc<Mutex<HashMap<String, (SessionKey, oneshot::Sender<PermissionOptionId>)>>>,
474}
475
476impl ConversationStore {
477 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 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 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 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 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 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, };
563 let _ = sender.send(option_id);
564 session
565 };
566 (self.publish)(ConversationUpdated { session });
567 }
568
569 pub fn oldest_pending_permission(&self) -> Option<PendingPermissionView> {
573 self.oldest_pending_permission_where(|_| true)
574 }
575
576 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 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 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 #[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 c.push_user_text("again");
663 assert_eq!(c.turns.len(), 3);
664 assert_eq!(c.turns[2].role, Role::User);
665 }
666
667 #[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 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 #[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 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 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 assert_eq!(published.load(std::sync::atomic::Ordering::SeqCst), 1);
853 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 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 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 published.store(0, std::sync::atomic::Ordering::SeqCst);
889
890 store.resolve_permission("perm-1", chosen.clone());
891
892 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 assert_eq!(
900 rx.blocking_recv(),
901 Ok(chosen),
902 "responder must receive the chosen option id",
903 );
904 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 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 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 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 c.apply(&update(53000, 200000, Some((0.045, "USD"))));
1006 assert_eq!(c.usage.unwrap().used, 53000);
1007 }
1008
1009 #[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 #[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}