use std::fs::{File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use crate::error::{Error, Result};
use crate::sidecar::{now_rfc3339, NativeTurn};
use crate::ChatMessage;
pub const JOURNAL_RECORD_VERSION: u8 = 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum QueueKind {
Steer,
FollowUp,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlanEntry {
pub step: String,
pub status: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "op", rename_all = "snake_case")]
pub enum JournalOp {
Message {
message: Box<NativeTurn>,
},
Rewind {
to: usize,
},
Unrewind,
Enqueue {
queue: QueueKind,
text: String,
},
Dequeue {
queue: QueueKind,
count: usize,
},
Plan {
steps: Vec<PlanEntry>,
},
Rename {
from: String,
to: String,
},
Checkpoint {
messages: usize,
},
ModelChange {
record: crate::model_change::ModelChangeRecord,
},
Usage {
record: crate::usage_log::UsageRecord,
},
Upgrade {
from_version: u32,
to_version: u32,
original: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JournalRecord {
pub supercode_journal: u8,
pub ts: String,
#[serde(flatten)]
pub op: JournalOp,
}
pub struct SessionJournal {
file: Option<File>,
path: PathBuf,
fixed_timestamp: Option<String>,
}
impl SessionJournal {
pub fn open_append(path: &Path) -> Result<Self> {
Ok(SessionJournal {
file: None,
path: path.to_path_buf(),
fixed_timestamp: None,
})
}
pub fn with_fixed_timestamp(mut self, ts: impl Into<String>) -> Self {
self.fixed_timestamp = Some(ts.into());
self
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn append(&mut self, op: JournalOp) -> Result<()> {
let record = JournalRecord {
supercode_journal: JOURNAL_RECORD_VERSION,
ts: self
.fixed_timestamp
.clone()
.unwrap_or_else(crate::sidecar::now_rfc3339),
op,
};
let mut line = serde_json::to_string(&record).map_err(Error::Decode)?;
line.push('\n');
let file = match &mut self.file {
Some(file) => file,
none => {
if let Some(parent) = self.path.parent() {
std::fs::create_dir_all(parent)?;
}
none.insert(
OpenOptions::new()
.create(true)
.append(true)
.open(&self.path)?,
)
}
};
file.write_all(line.as_bytes())?;
file.flush()?;
Ok(())
}
pub fn append_message(&mut self, msg: &ChatMessage) -> Result<()> {
let mut turn = NativeTurn::from(msg);
if let Some(ts) = &self.fixed_timestamp {
turn.ts.clone_from(ts);
turn.metadata.insert("timestamp".to_string(), ts.clone());
}
self.append(JournalOp::Message {
message: Box::new(turn),
})
}
}
#[derive(Debug, Clone, Default)]
pub struct JournalState {
pub messages: Vec<ChatMessage>,
pub unpersisted: Vec<ChatMessage>,
pub checkpoint_messages: Option<usize>,
pub undo_stack: Vec<Vec<ChatMessage>>,
pub steer_queue: Vec<String>,
pub follow_up_queue: Vec<String>,
pub plan: Vec<PlanEntry>,
pub renames: Vec<(String, String)>,
pub upgrades: Vec<(u32, u32, String)>,
pub model_changes: Vec<crate::model_change::ModelChangeRecord>,
pub usage: Vec<crate::usage_log::UsageRecord>,
pub records: usize,
pub skipped: usize,
}
impl JournalState {
pub fn queue(&self, kind: QueueKind) -> &[String] {
match kind {
QueueKind::Steer => &self.steer_queue,
QueueKind::FollowUp => &self.follow_up_queue,
}
}
pub fn plan_pairs(&self) -> Vec<(String, String)> {
self.plan
.iter()
.map(|s| (s.step.clone(), s.status.clone()))
.collect()
}
}
pub fn replay_str(text: &str) -> JournalState {
let mut state = JournalState::default();
for line in text.lines() {
if line.trim().is_empty() {
continue;
}
let Ok(record) = serde_json::from_str::<JournalRecord>(line) else {
state.skipped += 1;
continue;
};
if record.supercode_journal != JOURNAL_RECORD_VERSION {
state.skipped += 1;
continue;
}
state.records += 1;
match record.op {
JournalOp::Message { message } => {
let message = message.into_message();
state.unpersisted.push(message.clone());
state.messages.push(message);
}
JournalOp::Rewind { to } => {
let to = to.min(state.messages.len());
let tail = state.messages.split_off(to);
state.undo_stack.push(tail);
state.unpersisted.clear();
state.checkpoint_messages = None;
}
JournalOp::Unrewind => {
if let Some(mut tail) = state.undo_stack.pop() {
state.messages.append(&mut tail);
}
state.unpersisted.clear();
state.checkpoint_messages = None;
}
JournalOp::Enqueue { queue, text } => match queue {
QueueKind::Steer => state.steer_queue.push(text),
QueueKind::FollowUp => state.follow_up_queue.push(text),
},
JournalOp::Dequeue { queue, count } => {
let q = match queue {
QueueKind::Steer => &mut state.steer_queue,
QueueKind::FollowUp => &mut state.follow_up_queue,
};
let count = count.min(q.len());
q.drain(..count);
}
JournalOp::Plan { steps } => state.plan = steps,
JournalOp::Checkpoint { messages } => {
state.unpersisted.clear();
state.checkpoint_messages = Some(messages);
}
JournalOp::Rename { from, to } => state.renames.push((from, to)),
JournalOp::Upgrade {
from_version,
to_version,
original,
} => state.upgrades.push((from_version, to_version, original)),
JournalOp::ModelChange { record } => state.model_changes.push(record),
JournalOp::Usage { record } => state.usage.push(record),
}
}
state
}
pub fn replay(path: &Path) -> Result<Option<JournalState>> {
match std::fs::read_to_string(path) {
Ok(text) => Ok(Some(replay_str(&text))),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(e.into()),
}
}
pub fn timestamp() -> String {
now_rfc3339()
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RestoreReport {
pub recovered_messages: usize,
pub restored_queue: usize,
pub restored_plan: usize,
pub restored_rewinds: usize,
pub tree_loaded: bool,
pub journal_armed: bool,
}
pub fn arm(
agent: &mut crate::Agent,
store: &crate::store::SessionStore,
name: &str,
) -> RestoreReport {
let mut report = RestoreReport::default();
if !agent.config().session_persist {
return report;
}
let state = store.load_journal(name).ok().flatten();
if let Some(state) = &state {
if !state.unpersisted.is_empty()
&& state.checkpoint_messages.unwrap_or(1) == agent.history().len()
{
agent.append_recovered_messages(&state.unpersisted);
report.recovered_messages = state.unpersisted.len();
}
if agent.config().session_queue_persist {
agent.restore_queues(&state.steer_queue, &state.follow_up_queue);
report.restored_queue = state.steer_queue.len() + state.follow_up_queue.len();
}
if !state.undo_stack.is_empty() {
report.restored_rewinds = state.undo_stack.len();
agent.restore_rewind_undo(state.undo_stack.clone());
}
}
if agent.config().todos_persist {
let plan = store
.load_plan(name)
.ok()
.flatten()
.filter(|p| !p.is_empty())
.or_else(|| state.as_ref().map(|s| s.plan.clone()))
.unwrap_or_default();
if !plan.is_empty() {
report.restored_plan = plan.len();
agent.set_plan(plan);
}
}
if agent.config().session_tree_enabled {
match store.load_tree(name) {
Ok(Some(tree)) => {
report.tree_loaded = true;
agent.set_session_tree(tree);
}
_ => agent.rebuild_session_tree_from_history(),
}
}
if agent.config().session_append_only {
if let Ok(journal) = store.open_journal(name) {
agent.set_journal(journal);
if agent.history().len() > 1 {
agent.journal_checkpoint(agent.history().len());
}
report.journal_armed = true;
}
}
report
}
pub fn checkpoint(
agent: &crate::Agent,
store: &crate::store::SessionStore,
name: &str,
messages: usize,
) {
agent.journal_checkpoint(messages);
if agent.config().todos_persist {
let _ = store.save_plan(name, &agent.plan());
}
if let Some(tree) = agent.session_tree() {
let _ = store.save_tree(name, tree);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Role;
fn msg(role: Role, text: &str) -> ChatMessage {
match role {
Role::Assistant => ChatMessage::assistant(text),
_ => ChatMessage::user(text),
}
}
#[test]
fn every_append_is_a_flushed_line_readable_by_another_handle() {
let dir = std::env::temp_dir().join(format!("sc-journal-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("a.journal.jsonl");
let _ = std::fs::remove_file(&path);
let mut j = SessionJournal::open_append(&path).unwrap();
j.append_message(&msg(Role::User, "one")).unwrap();
let seen = replay(&path).unwrap().unwrap();
assert_eq!(seen.messages.len(), 1);
j.append_message(&msg(Role::Assistant, "two")).unwrap();
let seen = replay(&path).unwrap().unwrap();
assert_eq!(seen.messages.len(), 2);
let _ = std::fs::remove_file(&path);
}
#[test]
fn a_rewind_is_recorded_and_invertible_without_losing_bytes() {
let text = [
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"user","content":"a"}}"#,
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"assistant","content":"b"}}"#,
r#"{"supercode_journal":1,"ts":"t","op":"rewind","to":1}"#,
]
.join("\n");
let state = replay_str(&text);
assert_eq!(state.messages.len(), 1);
assert_eq!(state.undo_stack.len(), 1);
assert_eq!(state.undo_stack[0][0].content.as_deref(), Some("b"));
assert!(text.contains(r#""content":"b""#));
let undone = format!(
"{text}\n{}",
r#"{"supercode_journal":1,"ts":"t","op":"unrewind"}"#
);
let state = replay_str(&undone);
assert_eq!(state.messages.len(), 2);
assert_eq!(state.messages[1].content.as_deref(), Some("b"));
assert!(state.undo_stack.is_empty());
}
#[test]
fn queue_records_fold_into_the_still_pending_inputs() {
let text = [
r#"{"supercode_journal":1,"ts":"t","op":"enqueue","queue":"steer","text":"s1"}"#,
r#"{"supercode_journal":1,"ts":"t","op":"enqueue","queue":"follow_up","text":"f1"}"#,
r#"{"supercode_journal":1,"ts":"t","op":"enqueue","queue":"follow_up","text":"f2"}"#,
r#"{"supercode_journal":1,"ts":"t","op":"dequeue","queue":"follow_up","count":1}"#,
]
.join("\n");
let state = replay_str(&text);
assert_eq!(state.queue(QueueKind::Steer), ["s1"]);
assert_eq!(state.queue(QueueKind::FollowUp), ["f2"]);
}
#[test]
fn a_torn_trailing_line_is_skipped_not_fatal() {
let text = format!(
"{}\n{}",
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"user","content":"a"}}"#,
r#"{"supercode_journal":1,"ts":"t","op":"messa"#
);
let state = replay_str(&text);
assert_eq!(state.messages.len(), 1);
assert_eq!(state.skipped, 1);
}
#[test]
fn messages_after_the_last_checkpoint_are_the_ones_a_crash_would_lose() {
let text = [
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"user","content":"a"}}"#,
r#"{"supercode_journal":1,"ts":"t","op":"checkpoint","messages":2}"#,
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"user","content":"b"}}"#,
r#"{"supercode_journal":1,"ts":"t","op":"message","message":{"supercode_turn":1,"ts":"t","role":"assistant","content":"c"}}"#,
]
.join("\n");
let state = replay_str(&text);
assert_eq!(state.messages.len(), 3);
assert_eq!(state.checkpoint_messages, Some(2));
let lost: Vec<_> = state
.unpersisted
.iter()
.map(|m| m.content.clone().unwrap_or_default())
.collect();
assert_eq!(lost, ["b", "c"]);
}
#[test]
fn plan_records_replace_rather_than_merge() {
let text = [
r#"{"supercode_journal":1,"ts":"t","op":"plan","steps":[{"step":"one","status":"pending"}]}"#,
r#"{"supercode_journal":1,"ts":"t","op":"plan","steps":[{"step":"one","status":"completed"},{"step":"two","status":"pending"}]}"#,
]
.join("\n");
let state = replay_str(&text);
assert_eq!(
state.plan_pairs(),
vec![
("one".to_string(), "completed".to_string()),
("two".to_string(), "pending".to_string())
]
);
}
}