use std::{
collections::BTreeMap,
fs::File,
path::{Path, PathBuf},
};
use anyhow::Context as _;
use kcode_session_control_journal::{Journal, Record};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
const LIFECYCLE_SIDEBAND: &str = "session_lifecycle";
const COMMAND_SIDEBAND: &str = "session_command";
const STOP_SIDEBAND: &str = "session_stop";
const CONTROL_EXTENSION: &str = "session-control";
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum OpenMode {
CreateNew,
OpenOrCreate,
ExistingOnly,
}
#[derive(Clone, Debug, Default)]
pub struct ControlProjection {
pub lifecycle: Option<SessionRecord>,
pub commands: BTreeMap<String, SessionCommand>,
pub stop_requests: BTreeMap<String, SessionStopRequest>,
}
#[derive(Clone, Debug)]
pub enum ControlUpdate {
Lifecycle(SessionRecord),
Command(SessionCommand),
StopRequest(SessionStopRequest),
}
impl ControlUpdate {
pub fn projected(mut self) -> Self {
if let Self::Lifecycle(record) = &mut self {
record.state = control_state(&record.state);
}
self
}
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SessionRecord {
pub id: String,
pub phase: String,
pub started_at: String,
pub updated_at: String,
pub state: Value,
pub provenance_id: Option<String>,
pub version: i64,
pub last_user_message_at: Option<String>,
pub ended_at: Option<String>,
pub ingress_failure_count: i64,
pub ingress_failures: Value,
pub ingress_next_attempt_at: Option<String>,
#[serde(default, skip_serializing_if = "is_false")]
pub summary: bool,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SessionCommand {
pub id: String,
pub conversation_id: String,
pub sequence: i64,
pub kind: String,
pub payload: Value,
pub status: String,
pub cancel_requested: bool,
pub outcome: Option<Value>,
pub created_at: String,
pub processing_started_at: Option<String>,
pub completed_at: Option<String>,
pub idempotency_id: String,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SessionStopRequest {
pub id: String,
pub session_id: String,
pub scope: String,
pub status: String,
pub outcome: Option<Value>,
pub requested_at: String,
pub completed_at: Option<String>,
pub idempotency_id: String,
}
pub struct SessionControl {
directory: PathBuf,
path: PathBuf,
journal: Journal,
}
impl SessionControl {
pub fn open(
directory: impl AsRef<Path>,
session_id: &str,
mode: OpenMode,
) -> anyhow::Result<Option<Self>> {
let directory = directory.as_ref().to_path_buf();
let path = control_path(&directory, session_id);
let journal = match mode {
OpenMode::CreateNew => Journal::create(path.clone())?,
OpenMode::OpenOrCreate => match Journal::open(path.clone())? {
Some(journal) => journal,
None => Journal::create(path.clone())?,
},
OpenMode::ExistingOnly => {
let Some(journal) = Journal::open(path.clone())? else {
return Ok(None);
};
journal
}
};
Ok(Some(Self {
directory,
path,
journal,
}))
}
pub fn projection(&self) -> ControlProjection {
project_records(self.journal.records())
}
pub fn append(
&mut self,
recorded_at: impl Into<String>,
update: ControlUpdate,
) -> anyhow::Result<ControlUpdate> {
let update = update.projected();
let (kind, value) = match &update {
ControlUpdate::Lifecycle(record) => (
LIFECYCLE_SIDEBAND,
serde_json::to_value(record).context("encoding session lifecycle record")?,
),
ControlUpdate::Command(command) => (
COMMAND_SIDEBAND,
serde_json::to_value(command).context("encoding session command record")?,
),
ControlUpdate::StopRequest(request) => (
STOP_SIDEBAND,
serde_json::to_value(request).context("encoding session stop record")?,
),
};
self.journal.append(kind, recorded_at, value)?;
Ok(update)
}
pub fn delete(self) -> anyhow::Result<()> {
let Self {
directory,
path,
journal,
} = self;
drop(journal);
if path.exists() {
std::fs::remove_file(&path).with_context(|| format!("removing {}", path.display()))?;
sync_directory(&directory)?;
}
Ok(())
}
pub fn compact_directory(directory: impl AsRef<Path>) -> anyhow::Result<()> {
let directory = directory.as_ref();
let mut paths = std::fs::read_dir(directory)?
.filter_map(Result::ok)
.map(|entry| entry.path())
.filter(|path| {
path.extension().and_then(|value| value.to_str()) == Some(CONTROL_EXTENSION)
})
.collect::<Vec<_>>();
paths.sort();
for path in paths {
compact_journal(&path)?;
}
Ok(())
}
}
fn is_false(value: &bool) -> bool {
!*value
}
fn control_path(directory: &Path, session_id: &str) -> PathBuf {
directory.join(format!("{session_id}.{CONTROL_EXTENSION}"))
}
fn project_records(records: &[Record]) -> ControlProjection {
let latest_lifecycle = records
.iter()
.rev()
.find(|record| record.kind == LIFECYCLE_SIDEBAND)
.and_then(|record| serde_json::from_value(record.value.clone()).ok());
let mut commands = BTreeMap::new();
let mut stop_requests = BTreeMap::new();
for record in records {
match record.kind.as_str() {
COMMAND_SIDEBAND => {
if let Ok(command) = serde_json::from_value::<SessionCommand>(record.value.clone())
{
commands.insert(command.id.clone(), command);
}
}
STOP_SIDEBAND => {
if let Ok(request) =
serde_json::from_value::<SessionStopRequest>(record.value.clone())
{
stop_requests.insert(request.id.clone(), request);
}
}
_ => {}
}
}
ControlProjection {
lifecycle: latest_lifecycle,
commands,
stop_requests,
}
}
fn control_state(value: &Value) -> Value {
const KEYS: &[&str] = &[
"format",
"version",
"stateVersion",
"sessionId",
"sessionType",
"sourceSessionType",
"channel",
"freeTime",
"selfTimeIntent",
"orchestration",
"provenanceId",
"rustLibSessionId",
"rootNodeIds",
"referenceRootNodeIds",
"startedAt",
"pendingTurn",
"pendingExternalEventId",
"roundsUsed",
"completed",
"sessionObjectId",
"commitReceipt",
"commitAuthor",
"providerModel",
"kwebPlan",
"startIdempotencyId",
"ingressSource",
"chatendMetadata",
"sessionStatus",
"historyIngress",
];
let mut output = Map::new();
for key in KEYS {
if let Some(item) = value.get(*key) {
if *key == "commitReceipt" && item.is_null() {
continue;
}
let item = if *key == "historyIngress" {
control_state(item)
} else {
item.clone()
};
output.insert((*key).into(), item);
}
}
Value::Object(output)
}
fn compact_journal(path: &Path) -> anyhow::Result<()> {
let original_bytes = std::fs::metadata(path)
.with_context(|| format!("reading metadata for {}", path.display()))?
.len();
if original_bytes >= 16 * 1024 * 1024 {
tracing::info!(
path = %path.display(),
original_bytes,
"Compacting legacy Session History control journal"
);
}
let mut journal = Journal::open(path.to_path_buf())?
.with_context(|| format!("session-control journal {} disappeared", path.display()))?;
let repaired_bytes = std::fs::metadata(path)?.len();
let tail_repaired = repaired_bytes != original_bytes;
let mut latest_lifecycle = None;
let mut latest_commands = BTreeMap::<String, (u64, Record)>::new();
let mut latest_stop_requests = BTreeMap::<String, (u64, Record)>::new();
let mut retained_other = Vec::new();
let mut needs_rewrite = false;
for (sequence, mut record) in journal.records().iter().cloned().enumerate() {
let sequence = sequence as u64;
match record.kind.as_str() {
LIFECYCLE_SIDEBAND => {
if let Some(state) = record.value.get_mut("state") {
let projected = control_state(state);
if *state != projected {
*state = projected;
needs_rewrite = true;
}
}
if latest_lifecycle.replace((sequence, record)).is_some() {
needs_rewrite = true;
}
}
COMMAND_SIDEBAND => {
let id = record
.value
.get("id")
.and_then(Value::as_str)
.context("session command record has no ID")?
.to_owned();
if latest_commands.insert(id, (sequence, record)).is_some() {
needs_rewrite = true;
}
}
STOP_SIDEBAND => {
let id = record
.value
.get("id")
.and_then(Value::as_str)
.context("session stop record has no ID")?
.to_owned();
if latest_stop_requests
.insert(id, (sequence, record))
.is_some()
{
needs_rewrite = true;
}
}
_ => retained_other.push((sequence, record)),
}
}
if needs_rewrite {
let mut retained = retained_other;
retained.extend(latest_lifecycle);
retained.extend(latest_commands.into_values());
retained.extend(latest_stop_requests.into_values());
retained.sort_by_key(|(sequence, _)| *sequence);
journal.replace(retained.into_iter().map(|(_, record)| record))?;
}
if needs_rewrite || tail_repaired {
tracing::info!(
path = %path.display(),
original_bytes,
compacted_bytes = std::fs::metadata(path)?.len(),
"Compacted Session History control journal"
);
}
Ok(())
}
fn sync_directory(path: &Path) -> anyhow::Result<()> {
File::open(path)
.with_context(|| format!("opening directory {} for sync", path.display()))?
.sync_all()
.with_context(|| format!("syncing directory {}", path.display()))
}
#[cfg(test)]
mod tests {
use std::{
fs::{self, OpenOptions},
io::Write as _,
sync::atomic::{AtomicU64, Ordering},
time::{SystemTime, UNIX_EPOCH},
};
use serde_json::json;
use super::*;
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
fn root(label: &str) -> PathBuf {
let path = std::env::temp_dir().join(format!(
"kcode-session-control-state-{label}-{}-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos(),
NEXT_ROOT.fetch_add(1, Ordering::Relaxed),
));
fs::create_dir(&path).unwrap();
path
}
fn lifecycle(id: &str, version: i64, state: Value) -> SessionRecord {
SessionRecord {
id: id.into(),
phase: "active".into(),
started_at: "2026-08-02T00:00:00Z".into(),
updated_at: format!("2026-08-02T00:00:0{version}Z"),
state,
provenance_id: None,
version,
last_user_message_at: None,
ended_at: None,
ingress_failure_count: 0,
ingress_failures: json!([]),
ingress_next_attempt_at: None,
summary: false,
}
}
fn command(id: &str, status: &str, sequence: i64) -> SessionCommand {
SessionCommand {
id: id.into(),
conversation_id: "session-1".into(),
sequence,
kind: "message".into(),
payload: json!({"text":"hello"}),
status: status.into(),
cancel_requested: false,
outcome: None,
created_at: "2026-08-02T00:00:00Z".into(),
processing_started_at: None,
completed_at: None,
idempotency_id: format!("command-{id}"),
}
}
fn stop(id: &str, status: &str) -> SessionStopRequest {
SessionStopRequest {
id: id.into(),
session_id: "session-1".into(),
scope: "turn".into(),
status: status.into(),
outcome: None,
requested_at: "2026-08-02T00:00:00Z".into(),
completed_at: None,
idempotency_id: format!("stop-{id}"),
}
}
#[test]
fn opening_modes_preserve_create_and_absence_distinctions() {
let root = root("open-modes");
assert!(
SessionControl::open(&root, "missing", OpenMode::ExistingOnly)
.unwrap()
.is_none()
);
let created = SessionControl::open(&root, "new", OpenMode::CreateNew)
.unwrap()
.unwrap();
assert!(control_path(&root, "new").is_file());
assert!(SessionControl::open(&root, "new", OpenMode::CreateNew).is_err());
drop(created);
let opened = SessionControl::open(&root, "new", OpenMode::OpenOrCreate)
.unwrap()
.unwrap();
drop(opened);
let created_on_absence = SessionControl::open(&root, "other", OpenMode::OpenOrCreate)
.unwrap()
.unwrap();
assert!(control_path(&root, "other").is_file());
drop(created_on_absence);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn typed_append_projects_recursive_lifecycle_state_and_latest_values() {
let root = root("typed-projection");
let mut control = SessionControl::open(&root, "session-1", OpenMode::CreateNew)
.unwrap()
.unwrap();
let persisted = control
.append(
"t1",
ControlUpdate::Lifecycle(lifecycle(
"session-1",
1,
json!({
"sessionType":"conversation",
"chatendText":"discard",
"commitReceipt":null,
"historyIngress":{
"completed":true,
"commitReceipt":{"sessionObjectId":"A1234567"},
"boxes":{"1":{"text":"discard"}},
"chatendText":"discard"
}
}),
)),
)
.unwrap();
let ControlUpdate::Lifecycle(persisted) = persisted else {
panic!("append changed update kind");
};
assert_eq!(persisted.state["sessionType"], "conversation");
assert!(persisted.state.get("chatendText").is_none());
assert!(persisted.state.get("commitReceipt").is_none());
assert_eq!(
persisted.state["historyIngress"]["commitReceipt"]["sessionObjectId"],
"A1234567"
);
assert!(
persisted.state["historyIngress"]
.get("chatendText")
.is_none()
);
assert!(persisted.state["historyIngress"].get("boxes").is_none());
control
.append(
"t2",
ControlUpdate::Lifecycle(lifecycle(
"session-1",
2,
json!({"sessionType":"conversation","pendingTurn":true}),
)),
)
.unwrap();
control
.append("t3", ControlUpdate::Command(command("a", "pending", 1)))
.unwrap();
control
.append("t4", ControlUpdate::Command(command("a", "complete", 1)))
.unwrap();
control
.append("t5", ControlUpdate::Command(command("b", "pending", 2)))
.unwrap();
control
.append("t6", ControlUpdate::StopRequest(stop("s", "pending")))
.unwrap();
control
.append("t7", ControlUpdate::StopRequest(stop("s", "complete")))
.unwrap();
let projection = control.projection();
assert_eq!(projection.lifecycle.unwrap().version, 2);
assert_eq!(projection.commands.len(), 2);
assert_eq!(projection.commands["a"].status, "complete");
assert_eq!(projection.commands["b"].status, "pending");
assert_eq!(projection.stop_requests.len(), 1);
assert_eq!(projection.stop_requests["s"].status, "complete");
drop(control);
let reopened = SessionControl::open(&root, "session-1", OpenMode::ExistingOnly)
.unwrap()
.unwrap();
assert_eq!(reopened.projection().lifecycle.unwrap().version, 2);
drop(reopened);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn startup_compaction_retains_unknown_kinds_and_survivor_order() {
let root = root("compaction-order");
let path = control_path(&root, "session-1");
let mut journal = Journal::create(path.clone()).unwrap();
journal
.append("unknown-first", "t0", json!({"value":0}))
.unwrap();
journal
.append(
LIFECYCLE_SIDEBAND,
"t1",
serde_json::to_value(lifecycle(
"session-1",
1,
json!({"sessionType":"conversation","chatendText":"old"}),
))
.unwrap(),
)
.unwrap();
journal
.append(
COMMAND_SIDEBAND,
"t2",
serde_json::to_value(command("a", "pending", 1)).unwrap(),
)
.unwrap();
journal
.append("unknown-middle", "t3", json!({"value":3}))
.unwrap();
journal
.append(
LIFECYCLE_SIDEBAND,
"t4",
serde_json::to_value(lifecycle(
"session-1",
2,
json!({
"sessionType":"conversation",
"pendingTurn":true,
"historyIngress":{"chatendText":"discard","completed":true}
}),
))
.unwrap(),
)
.unwrap();
journal
.append(
STOP_SIDEBAND,
"t5",
serde_json::to_value(stop("s", "pending")).unwrap(),
)
.unwrap();
journal
.append(
COMMAND_SIDEBAND,
"t6",
serde_json::to_value(command("a", "complete", 1)).unwrap(),
)
.unwrap();
journal
.append("unknown-last", "t7", json!({"value":7}))
.unwrap();
journal
.append(
STOP_SIDEBAND,
"t8",
serde_json::to_value(stop("s", "complete")).unwrap(),
)
.unwrap();
drop(journal);
SessionControl::compact_directory(&root).unwrap();
let compacted = Journal::open(path).unwrap().unwrap();
let kinds = compacted
.records()
.iter()
.map(|record| record.kind.as_str())
.collect::<Vec<_>>();
assert_eq!(
kinds,
[
"unknown-first",
"unknown-middle",
LIFECYCLE_SIDEBAND,
COMMAND_SIDEBAND,
"unknown-last",
STOP_SIDEBAND,
]
);
let projected = project_records(compacted.records());
assert_eq!(projected.lifecycle.as_ref().unwrap().version, 2);
assert_eq!(
projected.lifecycle.unwrap().state["historyIngress"]["completed"],
true
);
assert_eq!(projected.commands["a"].status, "complete");
assert_eq!(projected.stop_requests["s"].status, "complete");
drop(compacted);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn malformed_typed_records_keep_the_existing_tolerance_and_errors() {
let root = root("malformed");
let path = control_path(&root, "session-1");
let mut journal = Journal::create(path.clone()).unwrap();
journal
.append(
LIFECYCLE_SIDEBAND,
"t1",
serde_json::to_value(lifecycle("session-1", 1, json!({}))).unwrap(),
)
.unwrap();
journal
.append(LIFECYCLE_SIDEBAND, "t2", json!({"not":"a lifecycle"}))
.unwrap();
journal
.append(COMMAND_SIDEBAND, "t3", json!({"id":"partial"}))
.unwrap();
drop(journal);
let control = SessionControl::open(&root, "session-1", OpenMode::ExistingOnly)
.unwrap()
.unwrap();
let projection = control.projection();
assert!(projection.lifecycle.is_none());
assert!(projection.commands.is_empty());
drop(control);
assert!(SessionControl::compact_directory(&root).is_ok());
let malformed_path = control_path(&root, "missing-id");
let mut malformed = Journal::create(malformed_path).unwrap();
malformed
.append(COMMAND_SIDEBAND, "t1", json!({"status":"pending"}))
.unwrap();
drop(malformed);
assert!(
SessionControl::compact_directory(&root)
.unwrap_err()
.to_string()
.contains("session command record has no ID")
);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn opening_repairs_incomplete_tail_but_rejects_complete_corruption() {
let root = root("integrity");
let path = control_path(&root, "tail");
let mut control = SessionControl::open(&root, "tail", OpenMode::CreateNew)
.unwrap()
.unwrap();
control
.append(
"t1",
ControlUpdate::Lifecycle(lifecycle("tail", 1, json!({}))),
)
.unwrap();
drop(control);
let complete = fs::read(&path).unwrap();
let mut file = OpenOptions::new().append(true).open(&path).unwrap();
file.write_all(b"incomplete tail").unwrap();
file.sync_all().unwrap();
drop(file);
let repaired = SessionControl::open(&root, "tail", OpenMode::ExistingOnly)
.unwrap()
.unwrap();
assert_eq!(repaired.projection().lifecycle.unwrap().version, 1);
drop(repaired);
assert_eq!(fs::read(&path).unwrap(), complete);
let corrupt_path = control_path(&root, "corrupt");
let mut corrupt = SessionControl::open(&root, "corrupt", OpenMode::CreateNew)
.unwrap()
.unwrap();
corrupt
.append(
"t1",
ControlUpdate::Lifecycle(lifecycle("corrupt", 1, json!({}))),
)
.unwrap();
drop(corrupt);
let mut bytes = fs::read(&corrupt_path).unwrap();
bytes[0] = if bytes[0] == b'0' { b'1' } else { b'0' };
fs::write(&corrupt_path, bytes).unwrap();
assert!(SessionControl::open(&root, "corrupt", OpenMode::ExistingOnly).is_err());
fs::remove_dir_all(root).unwrap();
}
#[test]
fn delete_removes_only_the_control_file() {
let root = root("delete");
let unrelated = root.join("keep.session-log");
fs::write(&unrelated, b"log").unwrap();
let control = SessionControl::open(&root, "session-1", OpenMode::CreateNew)
.unwrap()
.unwrap();
let path = control_path(&root, "session-1");
control.delete().unwrap();
assert!(!path.exists());
assert_eq!(fs::read(unrelated).unwrap(), b"log");
fs::remove_dir_all(root).unwrap();
}
}