use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use harn_session_store::{
EventId, ListFilter, ReadRange, SessionMeta, SessionStatus, SessionStore, StoreError,
StoredEvent, MAX_READ_BATCH,
};
use rust_decimal::prelude::ToPrimitive;
use rust_decimal::Decimal;
use serde_json::json;
use super::types::{LlmUsageRecord, RunChildRecord, RunRecord, RunTraceSpanRecord, ToolCallRecord};
use crate::agent_sessions::event_facts as facts;
use crate::value::VmError;
pub const PROJECTION_SOURCE: &str = "harn.session_store.v1";
pub const AGENT_SESSION_WORKFLOW_ID: &str = "agent-session";
pub const UNRECOVERABLE_FIELDS: [&str; 4] = [
"usage.total_duration_ms",
"trace_spans[].duration_ms",
"policy",
"replay_fixture",
];
pub async fn project_run_record_from_session(
store: &dyn SessionStore,
session_id: &str,
) -> Result<RunRecord, VmError> {
let meta = store
.describe(session_id)
.await
.map_err(|error| match error {
StoreError::NotFound(_) => VmError::Runtime(format!(
"runs: no session '{session_id}' in this store. `harn session list` shows the \
sessions this workspace has persisted."
)),
other => VmError::Runtime(format!("runs: failed to describe session: {other}")),
})?;
let events = drain_events(store, session_id).await?;
let children = child_records(store, session_id).await?;
let root = root_session_id(store, &meta).await?;
Ok(assemble(meta, events, children, root))
}
async fn root_session_id(store: &dyn SessionStore, meta: &SessionMeta) -> Result<String, VmError> {
let mut visited = std::collections::HashSet::from([meta.id.clone()]);
let mut current = meta.parent_session_id.clone();
let mut root = meta.id.clone();
while let Some(parent) = current {
if !visited.insert(parent.clone()) {
break;
}
match store.describe(&parent).await {
Ok(parent_meta) => {
root = parent_meta.id.clone();
current = parent_meta.parent_session_id;
}
Err(StoreError::NotFound(_)) => {
root = parent;
break;
}
Err(error) => {
return Err(VmError::Runtime(format!(
"runs: failed to walk session lineage: {error}"
)))
}
}
}
Ok(root)
}
pub async fn materialize_session_run_record(
root: &Path,
session_id: &str,
out: Option<&Path>,
) -> Result<String, VmError> {
let store =
crate::stdlib::session_store::open_existing_canonical_store(root)?.ok_or_else(|| {
VmError::Runtime(format!(
"runs: no session store under {}. A projected run record needs \
`.harn/session-store.sqlite`; pass the workspace root that holds it.",
root.display()
))
})?;
let run = project_run_record_from_session(&store, session_id).await?;
let path = out
.map(Path::to_path_buf)
.unwrap_or_else(|| default_projection_path(root, session_id));
super::persistence::save_run_record(&run, Some(&path.to_string_lossy()))
}
pub fn default_projection_path(root: &Path, session_id: &str) -> PathBuf {
crate::runtime_paths::run_root(root).join(format!("{session_id}.json"))
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
pub struct SessionRunSummary {
pub session_id: String,
pub session_status: String,
pub title: Option<String>,
pub parent_session_id: Option<String>,
pub created_at: String,
pub updated_at: String,
pub event_count: usize,
pub input_tokens: u64,
pub output_tokens: u64,
pub cost_usd_micros: u64,
pub run_record_path: Option<String>,
}
pub async fn list_session_runs(
root: &Path,
limit: Option<usize>,
) -> Result<Vec<SessionRunSummary>, VmError> {
let Some(store) = crate::stdlib::session_store::open_existing_canonical_store(root)? else {
return Ok(Vec::new());
};
let sessions = store
.list(ListFilter {
limit,
sort_by: harn_session_store::ListSortKey::CreatedAt,
order: harn_session_store::ListOrder::Descending,
..ListFilter::default()
})
.await
.map_err(|error| VmError::Runtime(format!("runs: failed to list sessions: {error}")))?;
Ok(sessions
.into_iter()
.map(|meta| {
let record = default_projection_path(root, &meta.id);
SessionRunSummary {
session_status: status_discriminator(&meta.status).to_string(),
title: meta.title.clone(),
parent_session_id: meta.parent_session_id.clone(),
created_at: meta.created_at.clone(),
updated_at: meta.updated_at.clone(),
event_count: meta.event_count,
input_tokens: meta.usage_input,
output_tokens: meta.usage_output,
cost_usd_micros: meta.usage_cost_usd_micros,
run_record_path: record
.is_file()
.then(|| record.to_string_lossy().into_owned()),
session_id: meta.id,
}
})
.collect())
}
async fn drain_events(
store: &dyn SessionStore,
session_id: &str,
) -> Result<Vec<StoredEvent>, VmError> {
let mut all = Vec::new();
let mut cursor: Option<EventId> = None;
loop {
let page = store
.read(
session_id,
ReadRange {
from_event_id: cursor,
to_event_id: None,
limit: Some(MAX_READ_BATCH),
},
)
.await
.map_err(|error| {
VmError::Runtime(format!("runs: failed to read session events: {error}"))
})?;
let next = page.next_cursor;
all.extend(page.events);
match next {
Some(next_cursor) => cursor = Some(next_cursor),
None => break,
}
}
Ok(all)
}
async fn child_records(
store: &dyn SessionStore,
session_id: &str,
) -> Result<Vec<RunChildRecord>, VmError> {
let children = store
.list(ListFilter {
parent_session_id: Some(session_id.to_string()),
..ListFilter::default()
})
.await
.map_err(|error| {
VmError::Runtime(format!("runs: failed to list child sessions: {error}"))
})?;
Ok(children
.into_iter()
.map(|child| RunChildRecord {
worker_id: child.id.clone(),
worker_name: child.persona.clone().unwrap_or_default(),
session_id: Some(child.id.clone()),
parent_session_id: Some(session_id.to_string()),
task: child.title.clone().unwrap_or_default(),
status: run_status_for(&child.status, None).to_string(),
started_at: child.created_at.clone(),
finished_at: child.closed_at.clone(),
run_id: Some(child.id.clone()),
..RunChildRecord::default()
})
.collect())
}
#[derive(Default)]
struct SessionFold {
task: Option<String>,
usage: LlmUsageRecord,
models: Vec<String>,
providers: Vec<String>,
cache_read_tokens: i64,
cache_write_tokens: i64,
provider_attempts: i64,
rate_limited_attempts: i64,
empty_completion_attempts: i64,
other_retry_attempts: i64,
total_cost: Decimal,
tools: Vec<ToolCallRecord>,
tool_index: BTreeMap<String, usize>,
iteration: usize,
max_iteration: usize,
terminal: Option<TerminalFacts>,
llm_calls: Vec<LlmCallFacts>,
}
struct LlmCallFacts {
at_ms: i64,
model: Option<String>,
provider: Option<String>,
input_tokens: i64,
output_tokens: i64,
cache_read_tokens: i64,
cache_write_tokens: i64,
cost_usd: Option<f64>,
}
struct TerminalFacts {
final_status: Option<String>,
stop_reason: Option<String>,
error: Option<String>,
class: Option<String>,
at: String,
}
fn assemble(
meta: SessionMeta,
events: Vec<StoredEvent>,
children: Vec<RunChildRecord>,
root_run_id: String,
) -> RunRecord {
let mut fold = SessionFold::default();
for event in &events {
fold.absorb(event);
}
let status = run_status_for(
&meta.status,
fold.terminal
.as_ref()
.and_then(|t| t.final_status.as_deref()),
)
.to_string();
let finished_at = meta
.closed_at
.clone()
.or_else(|| fold.terminal.as_ref().map(|t| t.at.clone()));
let mut metadata = BTreeMap::new();
metadata.insert(
"projected_from".to_string(),
json!({
"source": PROJECTION_SOURCE,
"session_id": meta.id,
"session_status": status_discriminator(&meta.status),
"session_event_count": meta.event_count,
"not_recoverable_from_session": UNRECOVERABLE_FIELDS,
}),
);
metadata.insert(
"wall_clock_ms".to_string(),
json!(meta.updated_at_ms.saturating_sub(meta.created_at_ms)),
);
if fold.max_iteration > 0 {
metadata.insert("iterations".to_string(), json!(fold.max_iteration));
}
if fold.cache_read_tokens > 0 || fold.cache_write_tokens > 0 {
metadata.insert(
"cache_tokens".to_string(),
json!({"read": fold.cache_read_tokens, "write": fold.cache_write_tokens}),
);
}
if !fold.providers.is_empty() {
metadata.insert("providers".to_string(), json!(fold.providers));
}
if fold.provider_attempts > fold.usage.call_count {
metadata.insert(
"provider_attempts".to_string(),
json!({
"total": fold.provider_attempts,
"retries": fold.provider_attempts - fold.usage.call_count,
"rate_limited": fold.rate_limited_attempts,
"empty_completion": fold.empty_completion_attempts,
"other": fold.other_retry_attempts,
}),
);
}
if let Some(terminal) = &fold.terminal {
if let Some(stop_reason) = &terminal.stop_reason {
metadata.insert("stop_reason".to_string(), json!(stop_reason));
}
if let Some(class) = &terminal.class {
metadata.insert("terminal_class".to_string(), json!(class));
}
if let Some(error) = &terminal.error {
metadata.insert("terminal_error".to_string(), json!(error));
}
}
let usage = LlmUsageRecord {
models: fold.models.clone(),
total_cost: fold.total_cost.to_f64().unwrap_or_default(),
..fold.usage
};
RunRecord {
type_name: "run".to_string(),
id: meta.id.clone(),
workflow_id: AGENT_SESSION_WORKFLOW_ID.to_string(),
workflow_name: meta.persona.clone(),
task: meta.title.clone().or(fold.task).unwrap_or_default(),
status,
started_at: meta.created_at.clone(),
finished_at,
parent_run_id: meta.parent_session_id.clone(),
root_run_id: Some(root_run_id),
child_runs: children,
usage: (usage.call_count > 0).then_some(usage),
trace_spans: llm_call_spans(&meta, &fold.llm_calls),
tool_recordings: fold.tools,
execution: None,
metadata,
..RunRecord::default()
}
}
impl SessionFold {
fn absorb(&mut self, event: &StoredEvent) {
match event.kind.discriminator() {
"message" => self.absorb_message(event),
"tool_call" => self.absorb_tool_call(event),
"tool_call_update" => self.absorb_tool_update(event),
"tool_result" => self.absorb_tool_result(event),
"llm_call" => self.absorb_llm_call(event),
"loop_checkpoint" => self.absorb_checkpoint(event),
"agent_run_terminal" => self.absorb_terminal(event),
_ => {}
}
}
fn absorb_message(&mut self, event: &StoredEvent) {
if self.task.is_some() {
return;
}
let is_user = event.actor.as_deref() == Some("user")
|| facts::semantic_string(&event.payload, &facts::ROLE).as_deref() == Some("user");
if is_user {
self.task = facts::semantic_string(&event.payload, &facts::TEXT);
}
}
fn absorb_llm_call(&mut self, event: &StoredEvent) {
let payload = &event.payload;
self.llm_calls.push(LlmCallFacts {
at_ms: event.ts_ms,
model: facts::string_at(payload, facts::MODEL),
provider: facts::string_at(payload, facts::PROVIDER),
input_tokens: facts::i64_at(payload, facts::INPUT_TOKENS).unwrap_or(0),
output_tokens: facts::i64_at(payload, facts::OUTPUT_TOKENS).unwrap_or(0),
cache_read_tokens: facts::i64_at(payload, facts::CACHE_READ_TOKENS).unwrap_or(0),
cache_write_tokens: facts::i64_at(payload, facts::CACHE_WRITE_TOKENS).unwrap_or(0),
cost_usd: facts::f64_at(payload, facts::COST_USD),
});
self.usage.call_count += 1;
self.usage.input_tokens += facts::i64_at(payload, facts::INPUT_TOKENS).unwrap_or(0);
self.usage.output_tokens += facts::i64_at(payload, facts::OUTPUT_TOKENS).unwrap_or(0);
if let Some(cost) = facts::f64_at(payload, facts::COST_USD) {
self.total_cost += Decimal::from_f64_retain(cost).unwrap_or_default();
}
self.cache_read_tokens += facts::i64_at(payload, facts::CACHE_READ_TOKENS).unwrap_or(0);
self.cache_write_tokens += facts::i64_at(payload, facts::CACHE_WRITE_TOKENS).unwrap_or(0);
self.provider_attempts += facts::i64_at(payload, facts::PROVIDER_ATTEMPTS_TOTAL)
.filter(|total| *total > 0)
.unwrap_or(1);
self.rate_limited_attempts +=
facts::i64_at(payload, facts::PROVIDER_ATTEMPTS_RATE_LIMITED).unwrap_or(0);
self.empty_completion_attempts +=
facts::i64_at(payload, facts::PROVIDER_ATTEMPTS_EMPTY).unwrap_or(0);
self.other_retry_attempts +=
facts::i64_at(payload, facts::PROVIDER_ATTEMPTS_OTHER).unwrap_or(0);
if let Some(model) = facts::string_at(payload, facts::MODEL) {
push_distinct(&mut self.models, model);
}
if let Some(provider) = facts::string_at(payload, facts::PROVIDER) {
push_distinct(&mut self.providers, provider);
}
}
fn absorb_checkpoint(&mut self, event: &StoredEvent) {
if facts::string_at(&event.payload, facts::CHECKPOINT_KIND).as_deref()
!= Some("iteration_start")
{
return;
}
if let Some(iteration) = facts::i64_at(&event.payload, facts::ITERATION) {
self.iteration = usize::try_from(iteration).unwrap_or(0);
self.max_iteration = self.max_iteration.max(self.iteration);
}
}
fn absorb_tool_call(&mut self, event: &StoredEvent) {
let payload = &event.payload;
let Some(tool_call_id) = facts::string_at(payload, facts::TOOL_CALL_ID) else {
return;
};
if self.tool_index.contains_key(&tool_call_id) {
return;
}
let args = facts::semantic_value(payload, &[facts::TOOL_RAW_INPUT])
.unwrap_or(serde_json::Value::Null);
let tool_name = facts::semantic_string(payload, &facts::TOOL_NAME_ANY).unwrap_or_default();
self.tool_index
.insert(tool_call_id.clone(), self.tools.len());
self.tools.push(ToolCallRecord {
args_hash: super::types::tool_fixture_hash(&tool_name, &args),
tool_name,
tool_use_id: tool_call_id,
iteration: self.iteration,
timestamp: event.ts.clone(),
..ToolCallRecord::default()
});
}
fn absorb_tool_update(&mut self, event: &StoredEvent) {
let payload = &event.payload;
let Some(record) = self.tool_for(payload) else {
return;
};
let status = facts::string_at(payload, facts::TOOL_STATUS);
match status.as_deref() {
Some("completed") | Some("failed") | Some("rejected") => {
record.is_rejected = status.as_deref() == Some("rejected");
if let Some(duration) = facts::i64_at(payload, facts::TOOL_DURATION_MS) {
record.duration_ms = u64::try_from(duration).unwrap_or(0);
}
}
_ => {}
}
}
fn absorb_tool_result(&mut self, event: &StoredEvent) {
let payload = &event.payload;
let text = facts::semantic_string(payload, &facts::TEXT).unwrap_or_default();
let Some(record) = self.tool_for(payload) else {
return;
};
record.result = text;
}
fn tool_for(&mut self, payload: &serde_json::Value) -> Option<&mut ToolCallRecord> {
let tool_call_id = facts::string_at(payload, facts::TOOL_CALL_ID)?;
let index = *self.tool_index.get(&tool_call_id)?;
self.tools.get_mut(index)
}
fn absorb_terminal(&mut self, event: &StoredEvent) {
self.terminal = Some(TerminalFacts {
final_status: facts::string_at(&event.payload, facts::FINAL_STATUS),
stop_reason: facts::string_at(&event.payload, facts::STOP_REASON),
error: facts::string_at(&event.payload, facts::TERMINAL_ERROR),
class: facts::string_at(&event.payload, facts::TERMINAL_CLASS),
at: event.ts.clone(),
});
}
}
fn run_status_for(session_status: &SessionStatus, final_status: Option<&str>) -> &'static str {
if let Some(final_status) = final_status.filter(|value| !value.is_empty()) {
return if crate::llm::session_status_indicates_error(final_status) {
"failed"
} else {
"completed"
};
}
match session_status {
SessionStatus::Open => "running",
SessionStatus::Closed => "completed",
SessionStatus::SoftDeleted | SessionStatus::HardDeleted => "deleted",
}
}
fn status_discriminator(status: &SessionStatus) -> &'static str {
match status {
SessionStatus::Open => "open",
SessionStatus::Closed => "closed",
SessionStatus::SoftDeleted => "soft_deleted",
SessionStatus::HardDeleted => "hard_deleted",
}
}
fn llm_call_spans(meta: &SessionMeta, calls: &[LlmCallFacts]) -> Vec<RunTraceSpanRecord> {
calls
.iter()
.enumerate()
.map(|(index, call)| {
let mut metadata = BTreeMap::from([
(
crate::tracing::meta::INPUT_TOKENS.to_string(),
json!(call.input_tokens),
),
(
crate::tracing::meta::OUTPUT_TOKENS.to_string(),
json!(call.output_tokens),
),
(
crate::tracing::meta::CACHE_READ_TOKENS.to_string(),
json!(call.cache_read_tokens),
),
(
crate::tracing::meta::CACHE_WRITE_TOKENS.to_string(),
json!(call.cache_write_tokens),
),
("duration_available".to_string(), json!(false)),
]);
if let Some(model) = &call.model {
metadata.insert(crate::tracing::meta::MODEL.to_string(), json!(model));
}
if let Some(provider) = &call.provider {
metadata.insert(crate::tracing::meta::PROVIDER.to_string(), json!(provider));
}
RunTraceSpanRecord {
trace_id: meta.id.clone(),
span_id: index as u64 + 1,
parent_id: None,
kind: "llm_call".to_string(),
name: call.model.clone().unwrap_or_else(|| "llm_call".to_string()),
start_ms: u64::try_from(call.at_ms.saturating_sub(meta.created_at_ms)).unwrap_or(0),
duration_ms: 0,
ttft_ms: None,
metadata,
links: Vec::new(),
cost_usd: call.cost_usd,
}
})
.collect()
}
fn push_distinct(values: &mut Vec<String>, value: String) {
if !values.contains(&value) {
values.push(value);
}
}
#[cfg(test)]
mod tests;