aether_sessions/
transcript.rs1use crate::model::{SessionEvent, UserEvent};
2use aether_core::events::{AgentEvent, ContextEvent, MessageEvent, ToolEvent, TurnEvent, task_created_result};
3use llm::{AssistantReasoning, ChatMessage, Context, MessageId, ToolCallError, ToolCallResult};
4
5pub fn context_from_events(events: &[SessionEvent]) -> Context {
6 let mut context = Context::new(vec![], vec![]);
7 let mut acc = MessageAccumulator::default();
8 for event in events {
9 match event {
10 SessionEvent::User(event) => {
11 acc.flush(&mut context);
12 apply_user_event(&mut context, event);
13 }
14 SessionEvent::Agent(event) => apply_agent_event(&mut context, event, &mut acc),
15 SessionEvent::Control(_) => {}
16 }
17 }
18 acc.flush(&mut context);
19 context
20}
21
22pub fn conversation_messages_from_events(events: &[SessionEvent]) -> Vec<ChatMessage> {
23 context_from_events(events).messages().iter().filter(|message| !message.is_system()).cloned().collect()
24}
25
26#[derive(Default)]
27struct MessageAccumulator {
28 message_id: Option<MessageId>,
29 text: String,
30 reasoning: String,
31 tool_results: Vec<Result<ToolCallResult, ToolCallError>>,
32 task_messages: Vec<ChatMessage>,
33}
34
35impl MessageAccumulator {
36 fn flush(&mut self, context: &mut Context) {
37 let pending = std::mem::take(self);
38 if let Some(message_id) = pending.message_id {
39 let reasoning = AssistantReasoning::from_parts(pending.reasoning, None);
40 context.push_assistant_turn(message_id, &pending.text, reasoning, pending.tool_results);
41 }
42 for message in pending.task_messages {
43 context.add_message(message);
44 }
45 }
46}
47
48fn apply_user_event(ctx: &mut Context, event: &UserEvent) {
49 match event {
50 UserEvent::Message { message_id, content, .. } => {
51 ctx.add_message(ChatMessage::user_with_id(message_id.clone(), content.clone()));
52 }
53 UserEvent::ClearContext => ctx.clear_conversation(),
54 }
55}
56
57fn apply_agent_event(ctx: &mut Context, event: &AgentEvent, acc: &mut MessageAccumulator) {
58 match event {
59 AgentEvent::Message(MessageEvent::Text { message_id, chunk, is_complete: true }) => {
60 if acc.message_id.is_some() {
61 acc.flush(ctx);
62 }
63 acc.message_id = Some(message_id.clone());
64 acc.text.clone_from(chunk);
65 }
66 AgentEvent::Message(MessageEvent::Thought { message_id, chunk, is_complete: true }) => {
67 if acc.message_id.as_ref() == Some(message_id) {
68 acc.reasoning.clone_from(chunk);
69 }
70 }
71 AgentEvent::Tool(ToolEvent::Call { .. }) => {
72 if acc.message_id.is_some() {
73 acc.flush(ctx);
74 }
75 }
76 AgentEvent::Tool(ToolEvent::Result { result, .. }) => acc.tool_results.push(Ok(result.clone())),
77 AgentEvent::Tool(ToolEvent::TaskCreated { request, task_id, .. }) => {
78 acc.tool_results.push(Ok(task_created_result(request, task_id)));
79 }
80 AgentEvent::Tool(ToolEvent::Error { error }) => acc.tool_results.push(Err(error.clone())),
81 AgentEvent::Turn(TurnEvent::AutoContinue { message_id, content, .. }) => {
82 acc.flush(ctx);
83 ctx.add_message(ChatMessage::user_with_id(message_id.clone(), content.clone()));
84 }
85 AgentEvent::Turn(TurnEvent::Ended { .. }) => acc.flush(ctx),
86 AgentEvent::Context(ContextEvent::Cleared) => {
87 ctx.clear_conversation();
88 *acc = MessageAccumulator::default();
89 }
90 AgentEvent::Context(ContextEvent::CompactionResult { message_id, summary, .. }) => {
91 acc.flush(ctx);
92 *ctx = ctx.with_compacted_summary(message_id.clone(), summary);
93 }
94 AgentEvent::Tool(
95 event @ (ToolEvent::TaskCompleted { .. } | ToolEvent::TaskFailed { .. } | ToolEvent::TaskCancelled { .. }),
96 ) => {
97 if let Some(message) = event.task_context_message() {
98 if acc.message_id.is_none() && !acc.tool_results.is_empty() {
99 acc.task_messages.push(message);
100 } else {
101 acc.flush(ctx);
102 ctx.add_message(message);
103 }
104 }
105 }
106 _ => {}
107 }
108}