use super::{Session, SessionEvent, SessionEventKind, SessionManager, TurnStatus};
use crate::{output::redact_sensitive_text, persistence::CrossProcessFileLock};
use serde::Serialize;
use sha2::{Digest, Sha256};
use std::{collections::BTreeSet, fs};
mod ownership;
pub(crate) const FRONTEND_PAGE_LIMIT: usize = 32;
pub(crate) const PAGE_BYTES: usize = 16 * 1024;
const TEXT_BYTES: usize = 2048;
pub(crate) const SCAN_ENTRIES: usize = 20_000;
pub(crate) const REPLAY_BYTES: usize = 2 * 1024 * 1024;
pub(crate) type SessionWriterLease = CrossProcessFileLock;
#[derive(Debug)]
pub(crate) struct FrontendSnapshotBusy;
impl std::fmt::Display for FrontendSnapshotBusy {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("session snapshot is busy")
}
}
impl std::error::Error for FrontendSnapshotBusy {}
impl Session {
pub(crate) fn try_frontend_writer(&self) -> anyhow::Result<Option<SessionWriterLease>> {
self.acquire_frontend_writer(false)
}
pub(crate) fn try_daemon_writer(&self) -> anyhow::Result<Option<SessionWriterLease>> {
self.acquire_frontend_writer(true)
}
}
#[derive(Serialize)]
pub(crate) struct FrontendSessionSummary {
session_id: String,
title: Option<String>,
updated_at: String,
ownership_hint: ownership::SessionOwnershipHint,
}
#[derive(Serialize)]
pub(crate) struct FrontendSessionList {
sessions: Vec<FrontendSessionSummary>,
next_after: Option<String>,
truncated: bool,
incomplete: bool,
}
impl SessionManager {
pub(crate) fn frontend_list(
&self,
after: Option<&str>,
limit: usize,
) -> anyhow::Result<FrontendSessionList> {
anyhow::ensure!(
(1..=FRONTEND_PAGE_LIMIT).contains(&limit),
"invalid page limit"
);
if let Some(after) = after {
super::validate_session_id(after.to_owned())?;
}
let mut ids = BTreeSet::new();
let mut incomplete = false;
if self.root.exists() {
super::store::validate_session_root(&self.root)?;
for (index, entry) in fs::read_dir(&self.root)?.enumerate() {
if index >= SCAN_ENTRIES {
return Err(anyhow::Error::new(super::BoundedReadError::BudgetExceeded(
"session discovery limit exceeded".to_owned(),
)));
}
let entry = entry?;
let Some((id, _)) = super::metadata::session_from_entry(entry.path()) else {
continue;
};
if id.len() > 128 || super::validate_session_id(id.clone()).is_err() {
incomplete = true;
continue;
}
if after.is_some_and(|after| id.as_str() <= after) {
continue;
}
ids.insert(id);
if ids.len() > limit + 1 {
ids.pop_last();
}
}
}
let truncated = ids.len() > limit;
let mut sessions = Vec::new();
let mut last = None;
for id in ids.into_iter().take(limit) {
last = Some(id.clone());
let summary = self
.open_existing(id)
.and_then(super::metadata::frontend_session_metadata_summary);
match summary {
Ok(summary) => sessions.push(FrontendSessionSummary {
ownership_hint: summary.session.frontend_ownership_hint(),
session_id: summary.session.id().to_owned(),
title: summary.latest_title.and_then(|title| {
super::sanitize_session_title(&redact_sensitive_text(&title))
}),
updated_at: chrono::DateTime::<chrono::Utc>::from(summary.latest_activity_time)
.to_rfc3339(),
}),
Err(_) => incomplete = true,
}
}
Ok(FrontendSessionList {
sessions,
next_after: if truncated { last } else { None },
truncated,
incomplete,
})
}
}
#[derive(Debug, Clone, Serialize)]
pub(crate) struct FrontendReplayEvent {
pub(crate) cursor: String,
#[serde(flatten)]
pub(crate) item: FrontendReplayItem,
}
#[derive(Debug, Clone, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub(crate) enum FrontendReplayItem {
Message {
role: String,
text: String,
part: usize,
last_part: bool,
},
Terminal {
status: &'static str,
},
Summary {
text: String,
part: usize,
last_part: bool,
},
}
#[derive(Debug, Serialize)]
pub(crate) struct FrontendReplayPage {
pub(crate) events: Vec<FrontendReplayEvent>,
pub(crate) next_cursor: String,
pub(crate) snapshot_cursor: String,
pub(crate) truncated: bool,
pub(crate) gap: bool,
pub(crate) resync_required: bool,
}
impl Session {
pub(crate) fn frontend_replay(
&self,
after: Option<&str>,
limit: usize,
) -> anyhow::Result<FrontendReplayPage> {
anyhow::ensure!(
(1..=FRONTEND_PAGE_LIMIT).contains(&limit),
"invalid page limit"
);
let root = self
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("missing root"))?;
super::store::validate_session_root(root)?;
let snapshot =
CrossProcessFileLock::try_acquire(&self.path)?.ok_or(FrontendSnapshotBusy)?;
let read = self.read_recent_events_tolerant_tail(100_000, REPLAY_BYTES, REPLAY_BYTES)?;
drop(snapshot);
let mut gap = !read.diagnostics.is_empty();
let (checkpoint, diagnostics) =
super::latest_valid_compaction_checkpoint_for_replay(self.id(), &read.events);
gap |= !diagnostics.is_empty()
|| read
.events
.iter()
.any(|event| event.session_id != self.id());
let start = checkpoint
.as_ref()
.map_or(0, |checkpoint| checkpoint.cutoff_event_count);
let mut items = Vec::new();
if let Some(checkpoint) = &checkpoint {
push_text(&mut items, None, &checkpoint.summary);
}
let mut turn = Vec::new();
for event in read
.events
.into_iter()
.skip(start)
.filter(|event| event.session_id == self.id())
{
match event.kind() {
Some(
SessionEventKind::UserInput
| SessionEventKind::AssistantChunk
| SessionEventKind::AssistantOutput,
) => {
if event
.payload
.get("text")
.and_then(serde_json::Value::as_str)
.is_some()
{
turn.push(event);
} else {
gap = true;
}
}
Some(SessionEventKind::ToolCall | SessionEventKind::ToolResult) => turn.push(event),
Some(SessionEventKind::TurnStatus) => {
let Some(payload) = event.turn_status_payload() else {
gap = true;
continue;
};
let status = match payload.status {
TurnStatus::Complete => Some("completed"),
TurnStatus::Cancelled => Some("cancelled"),
TurnStatus::Failed => Some("failed"),
TurnStatus::CompactionRequired => Some("compaction_required"),
TurnStatus::Incomplete => None,
};
turn.push(event);
if let Some(status) = status {
project_turn(self.id(), &mut turn, &mut items);
items.push(FrontendReplayItem::Terminal { status });
}
}
_ => {}
}
}
project_turn(self.id(), &mut turn, &mut items);
let mut hash = Sha256::new();
hash.update(b"frontend-replay-v1\0");
hash.update(self.id().as_bytes());
let mut cursors = vec![cursor(0, &hash)];
for (index, item) in items.iter().enumerate() {
let encoded = serde_json::to_vec(item)?;
hash.update((encoded.len() as u64).to_le_bytes());
hash.update(encoded);
cursors.push(cursor(index + 1, &hash));
}
let snapshot_cursor = cursors.last().cloned().unwrap_or_default();
let offset = match after {
None => Some(0),
Some(value) if value.len() <= 128 => cursors.iter().position(|cursor| cursor == value),
Some(_) => None,
};
let Some(offset) = offset else {
return Ok(FrontendReplayPage {
events: Vec::new(),
next_cursor: cursors[0].clone(),
snapshot_cursor,
truncated: !items.is_empty(),
gap,
resync_required: true,
});
};
let mut page = FrontendReplayPage {
events: Vec::new(),
next_cursor: cursors[offset].clone(),
snapshot_cursor,
truncated: offset + 1 < cursors.len(),
gap,
resync_required: false,
};
for (index, item) in items.into_iter().enumerate().skip(offset).take(limit) {
page.events.push(FrontendReplayEvent {
cursor: cursors[index + 1].clone(),
item,
});
page.next_cursor = cursors[index + 1].clone();
page.truncated = index + 2 < cursors.len();
if serde_json::to_vec(&page)?.len() > PAGE_BYTES {
page.events.pop();
page.next_cursor = cursors[index].clone();
page.truncated = true;
break;
}
}
Ok(page)
}
}
fn cursor(index: usize, hash: &Sha256) -> String {
format!(
"r1_{index}_{}",
crate::hex::lower_hex(hash.clone().finalize())
)
}
fn project_turn(id: &str, events: &mut Vec<SessionEvent>, items: &mut Vec<FrontendReplayItem>) {
for message in crate::context::build_frontend_conversation_from_events(id, events) {
match message.role {
crate::providers::MessageRole::User | crate::providers::MessageRole::Assistant => {
push_text(items, Some(message.role.as_api_str()), &message.content);
}
_ => {}
}
}
events.clear();
}
fn push_text(items: &mut Vec<FrontendReplayItem>, role: Option<&str>, text: &str) {
let text = redact_sensitive_text(text);
let mut rest = text.as_str();
let mut part = 0;
while !rest.is_empty() {
let mut end = rest.len().min(TEXT_BYTES);
while !rest.is_char_boundary(end) {
end -= 1;
}
let text = rest[..end].to_owned();
rest = &rest[end..];
let last_part = rest.is_empty();
items.push(match role {
Some(role) => FrontendReplayItem::Message {
role: role.to_owned(),
text,
part,
last_part,
},
None => FrontendReplayItem::Summary {
text,
part,
last_part,
},
});
part += 1;
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use tempfile::TempDir;
fn append(session: &Session, kind: SessionEventKind, payload: serde_json::Value) {
session
.append_with_outcome(&SessionEvent::new_kind(
kind,
session.id().to_owned(),
"/private/server-owned".into(),
payload,
))
.unwrap();
}
#[test]
fn replay_preserves_response_boundaries_without_tool_content() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
for status in ["complete", "cancelled", "failed"] {
for boundary in [SessionEventKind::ToolCall, SessionEventKind::ToolResult] {
let session = manager.create().unwrap();
append(
&session,
SessionEventKind::UserInput,
json!({"text": "question"}),
);
append(
&session,
SessionEventKind::AssistantChunk,
json!({"text": "Checking."}),
);
append(
&session,
boundary,
json!({
"id": "private_call", "name": "read",
"arguments": {"path": "private argument"},
"call_id": "private_call",
"result": {"tool_name": "read", "success": true, "content": "private result"}
}),
);
append(
&session,
SessionEventKind::AssistantChunk,
json!({"text": "Partial answer."}),
);
append(
&session,
SessionEventKind::TurnStatus,
json!({
"status": status, "assistant_text": "Checking.\n\nPartial answer."
}),
);
let page = session.frontend_replay(None, 32).unwrap();
let responses: Vec<_> = page
.events
.iter()
.filter_map(|event| match &event.item {
FrontendReplayItem::Message { role, text, .. } if role == "assistant" => {
Some(text.as_str())
}
_ => None,
})
.collect();
assert_eq!(
responses,
["Checking.", "Partial answer."],
"{status} {boundary:?}"
);
assert_eq!(page.events.len(), 4);
assert!(!serde_json::to_string(&page).unwrap().contains("private"));
}
}
}
#[test]
fn replay_skips_malformed_messages_without_recovery_diagnostics() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
for status in ["cancelled", "failed"] {
for kind in [
SessionEventKind::UserInput,
SessionEventKind::AssistantChunk,
SessionEventKind::AssistantOutput,
] {
for payload in [
json!({"text": {"private": "malformed payload"}}),
json!({"private": "missing text"}),
] {
let session = manager.create().unwrap();
let events = [
(SessionEventKind::UserInput, json!({"text": "question"})),
(kind, payload),
(
SessionEventKind::TurnStatus,
json!({"status": status, "assistant_text": "Partial answer."}),
),
]
.into_iter()
.map(|(kind, payload)| {
SessionEvent::new_kind(
kind,
session.id().to_owned(),
"/private/server-owned".into(),
payload,
)
})
.collect::<Vec<_>>();
for event in &events {
session.append_with_outcome(event).unwrap();
}
let messages = crate::context::build_frontend_conversation_from_events(
session.id(),
&events,
);
assert_eq!(
messages
.iter()
.map(|message| message.content.as_str())
.collect::<Vec<_>>(),
["question", "Partial answer."],
"{status} {kind:?}"
);
let page = session.frontend_replay(None, 32).unwrap();
assert!(page.gap, "{status} {kind:?}");
assert_eq!(page.events.len(), 3);
let items = page
.events
.iter()
.map(|event| &event.item)
.collect::<Vec<_>>();
assert_eq!(
serde_json::to_value(items).unwrap(),
json!([
{"kind": "message", "role": "user", "text": "question", "part": 0, "last_part": true},
{"kind": "message", "role": "assistant", "text": "Partial answer.", "part": 0, "last_part": true},
{"kind": "terminal", "status": status}
])
);
let encoded = serde_json::to_string(&page).unwrap();
assert!(!encoded.contains("private"));
assert!(!encoded.contains("malformed payload"));
}
}
}
}
#[test]
fn replay_byte_budget_includes_page_envelope() {
let temp = TempDir::new().unwrap();
let session = SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let text = "x".repeat(7 * TEXT_BYTES + 700);
append(
&session,
SessionEventKind::AssistantOutput,
json!({"text": text}),
);
let mut page = session.frontend_replay(None, 32).unwrap();
assert!(page.truncated);
let mut replayed = String::new();
loop {
assert!(serde_json::to_vec(&page).unwrap().len() <= PAGE_BYTES);
assert!(!page.events.is_empty());
for event in &page.events {
if let FrontendReplayItem::Message { text, .. } = &event.item {
replayed.push_str(text);
}
}
if !page.truncated {
break;
}
let previous = page.next_cursor.clone();
page = session.frontend_replay(Some(&previous), 32).unwrap();
assert_ne!(page.next_cursor, previous);
}
assert_eq!(replayed, text);
assert_eq!(page.next_cursor, page.snapshot_cursor);
}
#[test]
fn replay_pages_preserve_partial_terminals_and_continue_after_restart() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for status in ["cancelled", "failed"] {
append(
&session,
SessionEventKind::UserInput,
json!({"text": format!("question {status}")}),
);
append(
&session,
SessionEventKind::ToolResult,
json!({"output": "private tool data"}),
);
append(
&session,
SessionEventKind::Diagnostic,
json!({"error": "private diagnostic"}),
);
append(
&session,
SessionEventKind::TurnStatus,
json!({"status": status, "assistant_text": format!("partial {status}")}),
);
}
let first = session.frontend_replay(None, 2).unwrap();
assert!(first.truncated);
assert!(
matches!(&first.events[0].item, FrontendReplayItem::Message { role, text, .. } if role == "user" && text == "question cancelled")
);
assert!(
matches!(&first.events[1].item, FrontendReplayItem::Message { role, text, .. } if role == "assistant" && text == "partial cancelled")
);
let reopened = SessionManager::new(temp.path().join("sessions"))
.open_existing(session.id())
.unwrap();
let second = reopened
.frontend_replay(Some(&first.next_cursor), 32)
.unwrap();
assert_eq!(second.events.len(), 4);
assert!(matches!(
second.events[0].item,
FrontendReplayItem::Terminal {
status: "cancelled"
}
));
assert!(matches!(
second.events[3].item,
FrontendReplayItem::Terminal { status: "failed" }
));
assert!(!second.truncated);
let encoded = serde_json::to_string(&second).unwrap();
for private in [
"private tool data",
"private diagnostic",
"server-owned",
"jsonl",
"session_path",
] {
assert!(!encoded.contains(private));
}
append(
&session,
SessionEventKind::UserInput,
json!({"text": "next"}),
);
append(
&session,
SessionEventKind::AssistantOutput,
json!({"text": "done"}),
);
append(
&session,
SessionEventKind::TurnStatus,
json!({"status": "complete"}),
);
let next = reopened
.frontend_replay(Some(&second.next_cursor), 32)
.unwrap();
assert!(!next.resync_required);
assert_eq!(next.events.len(), 3);
assert_eq!(next.next_cursor, next.snapshot_cursor);
let other = manager.create().unwrap();
let stale = other.frontend_replay(Some(&next.next_cursor), 32).unwrap();
assert!(stale.resync_required);
assert!(stale.events.is_empty());
}
#[test]
fn replay_bounds_escaped_text_and_reports_corruption_and_changed_prefix() {
use std::io::Write;
let temp = TempDir::new().unwrap();
let session = SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let append_lock = CrossProcessFileLock::acquire(session.path()).unwrap();
let busy = session.frontend_replay(None, 32).unwrap_err();
assert!(busy.downcast_ref::<FrontendSnapshotBusy>().is_some());
drop(append_lock);
append(
&session,
SessionEventKind::UserInput,
json!({"text": "question"}),
);
let text = "\u{0001}🦀".repeat(6000);
append(
&session,
SessionEventKind::AssistantOutput,
json!({"text": text}),
);
append(
&session,
SessionEventKind::TurnStatus,
json!({"status": "complete"}),
);
let first = session.frontend_replay(None, 32).unwrap();
assert!(first.truncated);
assert!(serde_json::to_vec(&first).unwrap().len() <= PAGE_BYTES);
let mut output = String::new();
let mut page = first;
loop {
assert!(serde_json::to_vec(&page).unwrap().len() <= PAGE_BYTES);
for event in &page.events {
if let FrontendReplayItem::Message { role, text, .. } = &event.item
&& role == "assistant"
{
output.push_str(text);
}
}
if !page.truncated {
break;
}
page = session
.frontend_replay(Some(&page.next_cursor), 32)
.unwrap();
}
assert_eq!(output, text);
super::super::record_session_compaction(
&session,
temp.path(),
"bounded summary",
"test",
"test",
3,
)
.unwrap();
let compacted = session
.frontend_replay(Some(&page.next_cursor), 32)
.unwrap();
assert!(compacted.resync_required);
let resync = session.frontend_replay(None, 32).unwrap();
assert!(
matches!(&resync.events[0].item, FrontendReplayItem::Summary { text, .. } if text == "bounded summary")
);
fs::OpenOptions::new()
.append(true)
.open(session.path())
.unwrap()
.write_all(b"private malformed line\n")
.unwrap();
assert!(session.frontend_replay(None, 32).unwrap().gap);
fs::write(session.path(), b"").unwrap();
let stale = session
.frontend_replay(Some(&page.next_cursor), 32)
.unwrap();
assert!(stale.resync_required);
assert!(stale.events.is_empty());
assert!(!session.frontend_replay(None, 32).unwrap().resync_required);
}
#[test]
fn bounded_tail_resync_and_list_continuation_do_not_expose_storage() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let mut ids = Vec::new();
for _ in 0..3 {
let session = manager.create().unwrap();
append(
&session,
SessionEventKind::Diagnostic,
json!({"text": "created"}),
);
ids.push(session.id().to_owned());
}
ids.sort();
let first = manager.frontend_list(None, 1).unwrap();
assert_eq!(first.sessions[0].session_id, ids[0]);
assert!(first.truncated);
let next = manager
.frontend_list(first.next_after.as_deref(), 2)
.unwrap();
assert_eq!(
next.sessions
.iter()
.map(|item| &item.session_id)
.collect::<Vec<_>>(),
vec![&ids[1], &ids[2]]
);
assert!(!next.truncated);
assert!(manager.frontend_list(None, 33).is_err());
let session = manager.open_existing(&ids[0]).unwrap();
append(
&session,
SessionEventKind::Diagnostic,
json!({"text": "x".repeat(REPLAY_BYTES + 10)}),
);
append(
&session,
SessionEventKind::UserInput,
json!({"text": "recent"}),
);
append(
&session,
SessionEventKind::TurnStatus,
json!({"status": "cancelled", "assistant_text": "tail partial"}),
);
let page = session.frontend_replay(None, 32).unwrap();
assert!(page.gap);
assert_eq!(page.events.len(), 3);
assert!(!page.resync_required);
}
#[test]
fn writer_conflict_and_dead_owner_recovery_cross_process() {
const ROOT: &str = "MAGI_TEST_FRONTEND_WRITER_ROOT";
const MODE: &str = "MAGI_TEST_FRONTEND_WRITER_MODE";
if let Some(root) = std::env::var_os(ROOT) {
let session = SessionManager::new(root.into())
.open("writer-test")
.unwrap();
if std::env::var(MODE).unwrap() == "conflict" {
assert_eq!(session.frontend_ownership_hint().label(), "daemon-owned");
}
let lease = session.try_frontend_writer().unwrap();
if std::env::var(MODE).unwrap() == "conflict" {
assert!(lease.is_none());
return;
}
assert!(lease.is_some());
std::process::exit(0);
}
let temp = TempDir::new().unwrap();
let root = temp.path().join("sessions");
let session = SessionManager::new(root.clone())
.open("writer-test")
.unwrap();
let owner = session.try_daemon_writer().unwrap().unwrap();
let child = |mode: &str| {
let output = std::process::Command::new(std::env::current_exe().unwrap())
.args(["--exact", "sessions::frontend::tests::writer_conflict_and_dead_owner_recovery_cross_process", "--quiet"])
.env(ROOT, &root).env(MODE, mode).output().unwrap();
assert!(
output.status.success(),
"{}",
String::from_utf8_lossy(&output.stderr)
);
};
child("conflict");
drop(owner);
child("abandon");
assert!(session.try_frontend_writer().unwrap().is_some());
}
}