use std::sync::Arc;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::types::AgentMessage;
use super::super::types::{SessionError, SessionErrorCode};
use theway_contract::dag::PersistedRun;
pub use theway_contract::session::{JsonlSessionMetadata, SessionImportOrigin, SessionMetadata};
pub const SESSION_GRAPH_STATE_CUSTOM_TYPE: &str = "session_graph_state";
#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum SessionTreeEntry {
Message {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
message: AgentMessage,
},
ThinkingLevelChange {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "thinkingLevel")]
thinking_level: String,
},
ModelChange {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
provider: String,
#[serde(rename = "modelId")]
model_id: String,
},
Compaction {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
summary: String,
#[serde(rename = "firstKeptEntryId")]
first_kept_entry_id: String,
#[serde(rename = "tokensBefore")]
tokens_before: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
details: Option<Value>,
#[serde(default, rename = "fromHook", skip_serializing_if = "Option::is_none")]
from_hook: Option<bool>,
},
BranchSummary {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "fromId")]
from_id: String,
summary: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
details: Option<Value>,
#[serde(default, rename = "fromHook", skip_serializing_if = "Option::is_none")]
from_hook: Option<bool>,
},
Custom {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "customType")]
custom_type: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
data: Option<Value>,
},
CustomMessage {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "customType")]
custom_type: String,
content: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
details: Option<Value>,
display: bool,
},
Label {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "targetId")]
target_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
label: Option<String>,
},
SessionInfo {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
name: Option<String>,
},
Leaf {
id: String,
#[serde(rename = "parentId")]
parent_id: Option<String>,
timestamp: String,
#[serde(rename = "targetId")]
target_id: Option<String>,
},
}
impl SessionTreeEntry {
pub fn id(&self) -> &str {
match self {
Self::Message { id, .. }
| Self::ThinkingLevelChange { id, .. }
| Self::ModelChange { id, .. }
| Self::Compaction { id, .. }
| Self::BranchSummary { id, .. }
| Self::Custom { id, .. }
| Self::CustomMessage { id, .. }
| Self::Label { id, .. }
| Self::SessionInfo { id, .. }
| Self::Leaf { id, .. } => id,
}
}
pub fn parent_id(&self) -> Option<&str> {
match self {
Self::Message { parent_id, .. }
| Self::ThinkingLevelChange { parent_id, .. }
| Self::ModelChange { parent_id, .. }
| Self::Compaction { parent_id, .. }
| Self::BranchSummary { parent_id, .. }
| Self::Custom { parent_id, .. }
| Self::CustomMessage { parent_id, .. }
| Self::Label { parent_id, .. }
| Self::SessionInfo { parent_id, .. }
| Self::Leaf { parent_id, .. } => parent_id.as_deref(),
}
}
pub fn type_str(&self) -> &'static str {
match self {
Self::Message { .. } => "message",
Self::ThinkingLevelChange { .. } => "thinking_level_change",
Self::ModelChange { .. } => "model_change",
Self::Compaction { .. } => "compaction",
Self::BranchSummary { .. } => "branch_summary",
Self::Custom { .. } => "custom",
Self::CustomMessage { .. } => "custom_message",
Self::Label { .. } => "label",
Self::SessionInfo { .. } => "session_info",
Self::Leaf { .. } => "leaf",
}
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct SessionContext {
pub messages: Vec<AgentMessage>,
pub thinking_level: String,
pub model: Option<SessionContextModel>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SessionContextModel {
pub provider: String,
#[serde(rename = "modelId")]
pub model_id: String,
}
#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SubagentJobSnapshot {
pub id: String,
pub agent: String,
pub source: String,
pub run_id: Option<String>,
pub node_id: Option<String>,
pub session_id: Option<String>,
pub status: String,
pub started_at: Option<i64>,
pub completed_at: Option<i64>,
pub attempt: u32,
pub total_attempts: u32,
pub input_tokens: u64,
pub output_tokens: u64,
pub chars: u64,
pub tools_called: u64,
pub turn: u32,
pub error: Option<String>,
pub output_tail: String,
pub truncated: bool,
pub live_preview: Option<String>,
pub tps: Option<f64>,
pub cps: Option<f64>,
}
#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SessionGraphState {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub dags: Vec<PersistedRun>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub subagents: Vec<SubagentJobSnapshot>,
}
pub fn latest_session_graph_state(entries: &[SessionTreeEntry]) -> Option<SessionGraphState> {
entries.iter().rev().find_map(|entry| {
let SessionTreeEntry::Custom {
custom_type, data, ..
} = entry
else {
return None;
};
if custom_type != SESSION_GRAPH_STATE_CUSTOM_TYPE {
return None;
}
serde_json::from_value(data.clone()?).ok()
})
}
#[async_trait]
pub trait SessionStorage: Send + Sync {
async fn get_metadata_json(&self) -> Result<Value, SessionError>;
async fn get_leaf_id(&self) -> Result<Option<String>, SessionError>;
async fn set_leaf_id(&self, id: Option<String>) -> Result<(), SessionError>;
async fn create_entry_id(&self) -> Result<String, SessionError>;
async fn append_entry(&self, entry: SessionTreeEntry) -> Result<(), SessionError>;
async fn append_entries(&self, entries: Vec<SessionTreeEntry>) -> Result<(), SessionError> {
for entry in entries {
self.append_entry(entry).await?;
}
Ok(())
}
async fn get_entry(&self, id: &str) -> Result<Option<SessionTreeEntry>, SessionError>;
async fn get_entries(&self) -> Result<Vec<SessionTreeEntry>, SessionError>;
async fn get_path_to_root(
&self,
leaf_id: Option<&str>,
) -> Result<Vec<SessionTreeEntry>, SessionError>;
async fn find_entries(&self, entry_type: &str) -> Result<Vec<SessionTreeEntry>, SessionError>;
async fn get_label(&self, id: &str) -> Result<Option<String>, SessionError>;
}
pub use crate::agent::context::assembly::build_session_context;
pub use crate::agent::context::collapse::{COMPACT_CONTEXT_CUSTOM_TYPE, CompactContext};
#[derive(Clone)]
pub struct Session {
storage: Arc<dyn SessionStorage>,
}
impl Session {
pub fn new(storage: Arc<dyn SessionStorage>) -> Self {
Self { storage }
}
pub fn storage(&self) -> &Arc<dyn SessionStorage> {
&self.storage
}
fn not_found(msg: impl Into<String>) -> SessionError {
SessionError {
code: SessionErrorCode::NotFound,
message: msg.into(),
}
}
fn now_rfc3339() -> String {
chrono::Utc::now().to_rfc3339()
}
pub async fn leaf_id(&self) -> Result<Option<String>, SessionError> {
self.storage.get_leaf_id().await
}
pub async fn get_entry(&self, id: &str) -> Result<Option<SessionTreeEntry>, SessionError> {
self.storage.get_entry(id).await
}
pub async fn entries(&self) -> Result<Vec<SessionTreeEntry>, SessionError> {
self.storage.get_entries().await
}
pub async fn branch(
&self,
from_id: Option<&str>,
) -> Result<Vec<SessionTreeEntry>, SessionError> {
let leaf = match from_id {
Some(id) => Some(id.to_string()),
None => self.storage.get_leaf_id().await?,
};
self.storage.get_path_to_root(leaf.as_deref()).await
}
pub async fn build_context(&self) -> Result<SessionContext, SessionError> {
let branch = self.branch(None).await?;
Ok(build_session_context(&branch))
}
pub async fn label(&self, id: &str) -> Result<Option<String>, SessionError> {
self.storage.get_label(id).await
}
pub async fn session_name(&self) -> Result<Option<String>, SessionError> {
let entries = self.storage.find_entries("session_info").await?;
for entry in entries.into_iter().rev() {
if let SessionTreeEntry::SessionInfo {
name: Some(name), ..
} = entry
{
let trimmed = name.trim();
if !trimmed.is_empty() {
return Ok(Some(trimmed.to_string()));
}
}
}
Ok(None)
}
async fn append_typed(&self, entry: SessionTreeEntry) -> Result<String, SessionError> {
let id = entry.id().to_string();
self.storage.append_entry(entry).await?;
Ok(id)
}
pub async fn append_message(&self, message: AgentMessage) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
self.append_typed(SessionTreeEntry::Message {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
message,
})
.await
}
pub async fn append_messages(
&self,
messages: Vec<AgentMessage>,
) -> Result<Vec<String>, SessionError> {
let mut parent = self.storage.get_leaf_id().await?;
let mut ids = Vec::with_capacity(messages.len());
let mut entries = Vec::with_capacity(messages.len());
for message in messages {
let id = self.storage.create_entry_id().await?;
entries.push(SessionTreeEntry::Message {
id: id.clone(),
parent_id: parent,
timestamp: Self::now_rfc3339(),
message,
});
parent = Some(id.clone());
ids.push(id);
}
self.storage.append_entries(entries).await?;
Ok(ids)
}
pub async fn append_thinking_level_change(
&self,
thinking_level: impl Into<String>,
) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
self.append_typed(SessionTreeEntry::ThinkingLevelChange {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
thinking_level: thinking_level.into(),
})
.await
}
pub async fn append_model_change(
&self,
provider: impl Into<String>,
model_id: impl Into<String>,
) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
self.append_typed(SessionTreeEntry::ModelChange {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
provider: provider.into(),
model_id: model_id.into(),
})
.await
}
pub async fn append_compaction(
&self,
summary: impl Into<String>,
first_kept_entry_id: impl Into<String>,
tokens_before: u64,
details: Option<Value>,
from_hook: bool,
) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
self.append_typed(SessionTreeEntry::Compaction {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
summary: summary.into(),
first_kept_entry_id: first_kept_entry_id.into(),
tokens_before,
details,
from_hook: if from_hook { Some(true) } else { None },
})
.await
}
pub async fn append_custom(
&self,
custom_type: impl Into<String>,
data: Option<Value>,
) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
self.append_typed(SessionTreeEntry::Custom {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
custom_type: custom_type.into(),
data,
})
.await
}
pub async fn session_id(&self) -> Result<Option<String>, SessionError> {
let metadata = self.storage.get_metadata_json().await?;
Ok(metadata
.get("id")
.and_then(|value| value.as_str())
.map(str::to_string))
}
pub async fn append_session_graph_state(
&self,
state: &SessionGraphState,
) -> Result<String, SessionError> {
let mut data = serde_json::to_value(state).map_err(|error| SessionError {
code: SessionErrorCode::StorageFailure,
message: format!("serialize session graph state: {error}"),
})?;
if let Some(object) = data.as_object_mut() {
object.insert(
"updatedAt".to_string(),
serde_json::Value::String(Self::now_rfc3339()),
);
}
self.append_custom(SESSION_GRAPH_STATE_CUSTOM_TYPE, Some(data))
.await
}
pub async fn session_graph_state(&self) -> Result<Option<SessionGraphState>, SessionError> {
let entries = self.entries().await?;
Ok(latest_session_graph_state(&entries))
}
pub async fn latest_compaction_summary(&self) -> Result<Option<String>, SessionError> {
let entries = self.entries().await?;
Ok(entries.iter().rev().find_map(|entry| {
let SessionTreeEntry::Compaction { summary, .. } = entry else {
return None;
};
(!summary.is_empty()).then(|| summary.clone())
}))
}
pub async fn latest_collapse_summary(&self) -> Result<Option<String>, SessionError> {
if let Some(summary) = self.latest_compaction_summary().await? {
return Ok(Some(summary));
}
let entries = self.entries().await?;
Ok(entries.iter().rev().find_map(|entry| {
crate::agent::context::collapse::compact_context_from_entry(entry)
.map(|context| context.compact_text)
.filter(|text| !text.trim().is_empty())
}))
}
pub async fn compact_context(&self) -> Result<Option<CompactContext>, SessionError> {
let entries = self.entries().await?;
Ok(crate::agent::context::collapse::latest_compact_context(
&entries,
))
}
pub async fn collapse_node_id(&self) -> Result<Option<String>, SessionError> {
let metadata = self.storage.get_metadata_json().await?;
Ok(metadata
.get("collapseNodeId")
.and_then(Value::as_str)
.map(str::to_string))
}
pub async fn append_session_name(
&self,
name: impl Into<String>,
) -> Result<String, SessionError> {
let id = self.storage.create_entry_id().await?;
let parent = self.storage.get_leaf_id().await?;
let n = name.into().trim().to_string();
self.append_typed(SessionTreeEntry::SessionInfo {
id,
parent_id: parent,
timestamp: Self::now_rfc3339(),
name: Some(n),
})
.await
}
pub async fn move_to(
&self,
entry_id: Option<&str>,
summary: Option<BranchSummaryInput>,
) -> Result<Option<String>, SessionError> {
if let Some(id) = entry_id {
if self.storage.get_entry(id).await?.is_none() {
return Err(Self::not_found(format!("Entry {id} not found")));
}
}
self.storage.set_leaf_id(entry_id.map(String::from)).await?;
let Some(summary) = summary else {
return Ok(None);
};
let id = self.storage.create_entry_id().await?;
let from_id = entry_id.map(String::from).unwrap_or_else(|| "root".into());
let entry = SessionTreeEntry::BranchSummary {
id,
parent_id: entry_id.map(String::from),
timestamp: Self::now_rfc3339(),
from_id,
summary: summary.summary,
details: summary.details,
from_hook: if summary.from_hook { Some(true) } else { None },
};
Ok(Some(self.append_typed(entry).await?))
}
}
#[derive(Clone, Debug, Default)]
pub struct BranchSummaryInput {
pub summary: String,
pub details: Option<Value>,
pub from_hook: bool,
}
#[cfg(test)]
tests_bridge_macro::tests_bridge!("agent/session/session");
#[cfg(test)]
mod session_linecov_tests {
tests_bridge_macro::tests_bridge!("agent/session/session/linecov");
}