use aion_core::{
AssistantCommand, AssistantConfigOption, AssistantSessionId, AssistantSessionState,
AssistantSessionSummary, Payload,
};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::StoreError;
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct AssistantSessionRecord {
pub session_id: AssistantSessionId,
pub subject: String,
pub harness: String,
pub account: Option<String>,
pub title: Option<String>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
pub turns: u64,
pub mcp_token_digest: Option<String>,
#[serde(default)]
pub commands: Vec<AssistantCommand>,
#[serde(default)]
pub config_options: Vec<AssistantConfigOption>,
}
impl AssistantSessionRecord {
pub fn encode(&self) -> Result<Vec<u8>, StoreError> {
serde_json::to_vec(self).map_err(|error| StoreError::Serialization(error.to_string()))
}
pub fn decode(bytes: &[u8]) -> Result<Self, StoreError> {
serde_json::from_slice(bytes).map_err(|error| StoreError::Serialization(error.to_string()))
}
#[must_use]
pub fn summary(
&self,
state: AssistantSessionState,
reason: Option<String>,
) -> AssistantSessionSummary {
AssistantSessionSummary {
session_id: self.session_id,
harness: self.harness.clone(),
account: self.account.clone(),
state,
reason,
created_at: self.created_at,
updated_at: self.updated_at,
turns: self.turns,
title: self.title.clone(),
commands: self.commands.clone(),
config_options: self.config_options.clone(),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct AssistantTranscriptEvent {
pub index: u64,
pub recorded_at: DateTime<Utc>,
pub payload: Payload,
}
impl AssistantTranscriptEvent {
pub fn encode(&self) -> Result<Vec<u8>, StoreError> {
serde_json::to_vec(self).map_err(|error| StoreError::Serialization(error.to_string()))
}
pub fn decode(bytes: &[u8]) -> Result<Self, StoreError> {
serde_json::from_slice(bytes).map_err(|error| StoreError::Serialization(error.to_string()))
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct UndecodableAssistantSession {
pub session_id: String,
pub error: String,
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct AssistantSessionListing {
pub sessions: Vec<AssistantSessionRecord>,
pub undecodable: Vec<UndecodableAssistantSession>,
}
impl AssistantSessionListing {
pub fn sort(&mut self) {
self.sessions
.sort_by(|left, right| match left.created_at.cmp(&right.created_at) {
std::cmp::Ordering::Equal => left.session_id.cmp(&right.session_id),
ordering => ordering,
});
self.undecodable
.sort_by(|left, right| left.session_id.cmp(&right.session_id));
}
}
#[async_trait]
pub trait AssistantSessionStore: Send + Sync + 'static {
async fn put_assistant_session(&self, record: AssistantSessionRecord)
-> Result<(), StoreError>;
async fn get_assistant_session(
&self,
session_id: &AssistantSessionId,
) -> Result<Option<AssistantSessionRecord>, StoreError>;
async fn list_assistant_sessions(&self) -> Result<AssistantSessionListing, StoreError>;
async fn append_assistant_transcript_event(
&self,
session_id: &AssistantSessionId,
recorded_at: DateTime<Utc>,
payload: Payload,
) -> Result<u64, StoreError>;
async fn assistant_transcript_head(
&self,
session_id: &AssistantSessionId,
) -> Result<u64, StoreError>;
async fn assistant_transcript(
&self,
session_id: &AssistantSessionId,
after: Option<u64>,
) -> Result<Vec<AssistantTranscriptEvent>, StoreError>;
async fn put_assistant_default_harness(
&self,
subject: &str,
harness: &str,
) -> Result<(), StoreError>;
async fn assistant_default_harness(&self, subject: &str) -> Result<Option<String>, StoreError>;
}
#[cfg(test)]
mod tests {
use aion_core::ContentType;
use chrono::TimeZone;
use super::*;
fn instant(offset: i64) -> Result<DateTime<Utc>, StoreError> {
Utc.with_ymd_and_hms(2026, 8, 29, 6, 0, 0)
.single()
.map(|base| base + chrono::Duration::seconds(offset))
.ok_or_else(|| StoreError::Backend("test instant must be valid".to_owned()))
}
fn record(offset: i64) -> Result<AssistantSessionRecord, StoreError> {
Ok(AssistantSessionRecord {
session_id: AssistantSessionId::new(uuid::Uuid::from_u128(7)),
subject: String::from("operator"),
harness: String::from("claude"),
account: Some(String::from("work")),
title: Some(String::from("fix the check")),
created_at: instant(offset)?,
updated_at: instant(offset + 10)?,
turns: 2,
mcp_token_digest: Some("0".repeat(64)),
commands: vec![AssistantCommand {
name: String::from("compact"),
description: String::from("compact the conversation"),
input_hint: None,
}],
config_options: Vec::new(),
})
}
#[test]
fn a_session_record_round_trips() -> Result<(), StoreError> {
let expected = record(0)?;
assert_eq!(
AssistantSessionRecord::decode(&expected.encode()?)?,
expected
);
Ok(())
}
#[test]
fn a_record_carries_no_status_of_its_own() -> Result<(), StoreError> {
let summary = record(0)?.summary(
AssistantSessionState::Dormant,
Some(String::from("process_exited")),
);
assert_eq!(summary.state, AssistantSessionState::Dormant);
assert_eq!(summary.reason.as_deref(), Some("process_exited"));
Ok(())
}
#[test]
fn a_transcript_event_round_trips() -> Result<(), StoreError> {
let event = AssistantTranscriptEvent {
index: 4,
recorded_at: instant(0)?,
payload: Payload::new(ContentType::Json, b"{\"type\":\"delta\"}".to_vec()),
};
assert_eq!(AssistantTranscriptEvent::decode(&event.encode()?)?, event);
Ok(())
}
#[test]
fn the_listing_orders_by_created_at_then_id() -> Result<(), StoreError> {
let mut early = record(0)?;
early.session_id = AssistantSessionId::new(uuid::Uuid::from_u128(2));
let mut tied = record(0)?;
tied.session_id = AssistantSessionId::new(uuid::Uuid::from_u128(1));
let late = record(100)?;
let mut listing = AssistantSessionListing {
sessions: vec![late.clone(), early.clone(), tied.clone()],
undecodable: vec![
UndecodableAssistantSession {
session_id: String::from("b"),
error: String::from("bad"),
},
UndecodableAssistantSession {
session_id: String::from("a"),
error: String::from("bad"),
},
],
};
listing.sort();
assert_eq!(listing.sessions, vec![tied, early, late]);
assert_eq!(
listing
.undecodable
.iter()
.map(|row| row.session_id.as_str())
.collect::<Vec<_>>(),
vec!["a", "b"]
);
Ok(())
}
}