Skip to main content

aether_sessions/
transcript.rs

1use 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}