use serde::{Deserialize, Serialize};
use serde_json::json;
use thiserror::Error;
use uuid::Uuid;
use crate::event::{
Event, EventError, EventKind, EventStore, SessionReplacementRecord, history_sha256,
};
use crate::workspace::{Workspace, WorkspaceError, WorkspaceState};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionState {
Initializing,
Working,
WaitingForTool,
CandidateReady,
Verifying,
Repairing,
Paused,
AcceptedAwaitingAuthority,
Completed,
Failed,
BudgetExhausted,
}
impl SessionState {
#[must_use]
pub const fn is_terminal(self) -> bool {
matches!(
self,
Self::AcceptedAwaitingAuthority
| Self::Completed
| Self::Failed
| Self::BudgetExhausted
)
}
}
#[derive(Debug, Error)]
pub enum SessionError {
#[error(transparent)]
Event(#[from] EventError),
#[error("session {0} does not exist")]
NotFound(String),
#[error("session history is missing its user goal")]
MissingGoal,
#[error("invalid persisted session state: {0}")]
InvalidState(String),
#[error("session history contains events after terminal state")]
EventsAfterTerminal,
#[error("only infrastructure-failed sessions may be replaced")]
IneligiblePredecessor,
#[error("failed predecessor has no eligible candidate provenance")]
MissingCandidateProvenance,
#[error("candidate SHA-256 has an invalid format")]
InvalidCandidateIdentity,
#[error("candidate identity mismatch (expected {expected}, found {actual})")]
CandidateMismatch { expected: String, actual: String },
#[error("candidate workspace provenance is malformed: {0}")]
InvalidCandidateProvenance(String),
#[error("replacement session chains are not supported")]
ReplacementChainUnsupported,
#[error("replacement session is not in its verification-only candidate state")]
InvalidReplacementExecutionState,
#[error("replacement workspace identity is missing or malformed")]
InvalidReplacementWorkspaceProvenance,
#[error("replacement workspace identity differs from its persisted candidate source")]
ReplacementWorkspaceMismatch,
#[error(transparent)]
Workspace(#[from] WorkspaceError),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Session {
pub id: String,
pub goal: String,
pub state: SessionState,
pub model_turns: u32,
pub tool_calls: u32,
pub repair_cycles: u32,
pub started_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CreatedReplacementSession {
pub session: Session,
pub replacement: SessionReplacementRecord,
}
impl Session {
pub fn create(store: &mut EventStore, goal: impl Into<String>) -> Result<Self, SessionError> {
let goal = goal.into();
let id = format!("session_{}", Uuid::new_v4().simple());
let created = store.append(
&id,
EventKind::SessionCreated,
&json!({
"state": SessionState::Initializing,
"harness_name": env!("CARGO_PKG_NAME"),
"harness_version": env!("CARGO_PKG_VERSION")
}),
)?;
store.append(&id, EventKind::UserGoal, &json!({"goal": goal}))?;
let mut session = Self {
id,
goal,
state: SessionState::Initializing,
model_turns: 0,
tool_calls: 0,
repair_cycles: 0,
started_at_ms: created.created_at_ms,
};
session.transition(store, SessionState::Working, "initialized")?;
Ok(session)
}
pub fn reconstruct(store: &EventStore, id: &str) -> Result<Self, SessionError> {
let events = store.events(id)?;
if events.is_empty() {
return Err(SessionError::NotFound(id.to_owned()));
}
let goal = events
.iter()
.find(|event| event.kind == EventKind::UserGoal)
.and_then(|event| event.payload.get("goal"))
.and_then(|value| value.as_str())
.ok_or(SessionError::MissingGoal)?
.to_owned();
let started_at_ms = events[0].created_at_ms;
let mut session = Self {
id: id.to_owned(),
goal,
state: SessionState::Initializing,
model_turns: 0,
tool_calls: 0,
repair_cycles: 0,
started_at_ms,
};
let mut terminal_seen = false;
for event in &events {
if event.kind.is_composition() {
continue;
}
if terminal_seen {
return Err(SessionError::EventsAfterTerminal);
}
match event.kind {
EventKind::ModelRequest => session.model_turns += 1,
EventKind::ToolRequest => session.tool_calls += 1,
EventKind::RepairStarted => session.repair_cycles += 1,
EventKind::TerminalState => {
session.state = state_from_event(event)?;
terminal_seen = true;
}
_ => {
if let Some(state) = event.payload.get("state") {
session.state = serde_json::from_value(state.clone())
.map_err(|_| SessionError::InvalidState(state.to_string()))?;
}
}
}
}
let replacement = store.replacement_for_session(id)?;
if replacement.is_none()
&& events
.iter()
.any(|event| event.kind == EventKind::SessionReplaced)
{
return Err(EventError::MalformedReplacement(
"replacement event has no durable relation".to_owned(),
)
.into());
}
Ok(session)
}
pub fn create_replacement(
store: &mut EventStore,
workspace: &Workspace,
predecessor_session_id: &str,
expected_candidate_sha256: &str,
) -> Result<CreatedReplacementSession, SessionError> {
if !valid_sha256(expected_candidate_sha256) {
return Err(SessionError::InvalidCandidateIdentity);
}
if store
.replacement_for_session(predecessor_session_id)?
.is_some()
{
return Err(SessionError::ReplacementChainUnsupported);
}
if let Some(existing) = store.replacement_for_predecessor(predecessor_session_id)? {
return Err(EventError::ReplacementAlreadyExists {
predecessor: existing.predecessor_session_id,
replacement: existing.replacement_session_id,
}
.into());
}
let predecessor = Self::reconstruct(store, predecessor_session_id)?;
let predecessor_events = store.events(predecessor_session_id)?;
let core_events: Vec<&Event> = predecessor_events
.iter()
.filter(|event| !event.kind.is_composition())
.collect();
let Some(terminal) = core_events.last() else {
return Err(SessionError::IneligiblePredecessor);
};
if predecessor.state != SessionState::Failed
|| terminal.kind != EventKind::TerminalState
|| terminal.payload["state"] != json!(SessionState::Failed)
|| terminal.payload["reason"] != "falsegreen_infrastructure_failure"
{
return Err(SessionError::IneligiblePredecessor);
}
let failure = core_events
.get(core_events.len().saturating_sub(2))
.filter(|event| event.kind == EventKind::FalsegreenResult)
.ok_or(SessionError::IneligiblePredecessor)?;
if failure.payload["verdict"] != "insufficient_evidence"
|| failure.payload["error"].as_str().is_none_or(str::is_empty)
{
return Err(SessionError::IneligiblePredecessor);
}
let (candidate_event, git_state_event, candidate_summary) =
eligible_candidate_provenance(&predecessor_events)?;
let persisted_candidate = candidate_event.payload["candidate_sha256"]
.as_str()
.ok_or(SessionError::MissingCandidateProvenance)?;
if persisted_candidate != expected_candidate_sha256 {
return Err(SessionError::CandidateMismatch {
expected: expected_candidate_sha256.to_owned(),
actual: persisted_candidate.to_owned(),
});
}
let actual_candidate = workspace.candidate_sha256()?;
if actual_candidate != expected_candidate_sha256 {
return Err(SessionError::CandidateMismatch {
expected: expected_candidate_sha256.to_owned(),
actual: actual_candidate,
});
}
let persisted_workspace: WorkspaceState =
serde_json::from_value(git_state_event.payload.clone())
.map_err(|error| SessionError::InvalidCandidateProvenance(error.to_string()))?;
let current_workspace = workspace.state()?;
if persisted_workspace != current_workspace {
return Err(SessionError::InvalidCandidateProvenance(
"current Git/workspace identity differs from the candidate source".to_owned(),
));
}
let authority_events: Vec<&Event> = predecessor_events
.iter()
.filter(|event| {
event.kind == EventKind::Checkpoint
&& event.payload["checkpoint_kind"] == "acceptance_authority"
})
.collect();
if authority_events.len() != 1 {
return Err(SessionError::InvalidCandidateProvenance(
"expected exactly one acceptance-authority binding".to_owned(),
));
}
let falsegreen_task_id = authority_events[0].payload["falsegreen_task_id"]
.as_str()
.filter(|value| !value.is_empty())
.ok_or_else(|| {
SessionError::InvalidCandidateProvenance(
"acceptance-authority task binding is missing".to_owned(),
)
})?
.to_owned();
let predecessor_history_sha256 = history_sha256(&predecessor_events)?;
let replacement_session_id = format!("session_{}", Uuid::new_v4().simple());
let record = SessionReplacementRecord {
predecessor_session_id: predecessor_session_id.to_owned(),
replacement_session_id: replacement_session_id.clone(),
predecessor_state: "failed".to_owned(),
predecessor_history_sha256: predecessor_history_sha256.clone(),
candidate_sha256: expected_candidate_sha256.to_owned(),
candidate_event_sequence: candidate_event.sequence,
source_git_state_sequence: git_state_event.sequence,
falsegreen_task_id: falsegreen_task_id.clone(),
created_at_ms: 0,
};
let replacement_payload = json!({
"predecessor_session_id": predecessor_session_id,
"predecessor_state": "failed",
"predecessor_history_sha256": predecessor_history_sha256,
"candidate_sha256": expected_candidate_sha256,
"candidate_event_sequence": candidate_event.sequence,
"source_git_state_sequence": git_state_event.sequence,
"falsegreen_task_id": falsegreen_task_id,
"topology": "single_direct_replacement_no_chains"
});
let mut authority_payload = authority_events[0].payload.clone();
authority_payload["replacement_predecessor_session_id"] = json!(predecessor_session_id);
let event_payloads = vec![
(
EventKind::SessionCreated,
json!({
"state": SessionState::Initializing,
"harness_name": env!("CARGO_PKG_NAME"),
"harness_version": env!("CARGO_PKG_VERSION")
}),
),
(EventKind::UserGoal, json!({"goal": predecessor.goal})),
(
EventKind::Checkpoint,
json!({
"state": SessionState::Working,
"previous_state": SessionState::Initializing,
"reason": "replacement_initialized"
}),
),
(EventKind::SessionReplaced, replacement_payload),
(EventKind::Checkpoint, authority_payload),
(
EventKind::CandidateReady,
json!({
"summary": candidate_summary,
"candidate_sha256": expected_candidate_sha256,
"carried_forward": true,
"predecessor_session_id": predecessor_session_id,
"predecessor_candidate_event_sequence": candidate_event.sequence
}),
),
(EventKind::GitState, git_state_event.payload.clone()),
(
EventKind::Checkpoint,
json!({
"checkpoint_kind": "replacement_candidate_validation",
"state": SessionState::CandidateReady,
"previous_state": SessionState::Working,
"reason": "authorized_candidate_carried_forward",
"expected_candidate_sha256": expected_candidate_sha256,
"actual_candidate_sha256": expected_candidate_sha256,
"matched": true,
"predecessor_session_id": predecessor_session_id
}),
),
];
let replacement = store.create_session_replacement(record, &event_payloads)?;
let session = Self::reconstruct(store, &replacement_session_id)?;
Ok(CreatedReplacementSession {
session,
replacement,
})
}
pub fn validate_replacement_resume(
&self,
store: &EventStore,
workspace: &Workspace,
) -> Result<Option<SessionReplacementRecord>, SessionError> {
let Some(replacement) = store.replacement_for_session(&self.id)? else {
return Ok(None);
};
if self.state.is_terminal() {
return Ok(Some(replacement));
}
if self.state != SessionState::CandidateReady
|| self.model_turns != 0
|| self.tool_calls != 0
|| self.repair_cycles != 0
{
return Err(SessionError::InvalidReplacementExecutionState);
}
let events = store.events(&self.id)?;
let candidate_index = events
.iter()
.rposition(|event| {
event.kind == EventKind::CandidateReady
&& event.payload["carried_forward"].as_bool() == Some(true)
&& event.payload["candidate_sha256"].as_str()
== Some(&replacement.candidate_sha256)
})
.ok_or(SessionError::InvalidReplacementExecutionState)?;
if events[candidate_index + 1..].iter().any(|event| {
matches!(
event.kind,
EventKind::ModelRequest
| EventKind::ModelResponse
| EventKind::ToolRequest
| EventKind::ToolResult
| EventKind::FileMutation
| EventKind::CommandExecution
| EventKind::FalsegreenResult
| EventKind::RepairStarted
| EventKind::GenUiActionPresented
| EventKind::GenUiActionConfirmationRequested
| EventKind::GenUiActionConfirmed
| EventKind::GenUiActionExecutionStarted
| EventKind::GenUiActionExecutionCompleted
| EventKind::GenUiActionExecutionFailed
| EventKind::GenUiActionExecutionUnknown
| EventKind::TerminalState
)
}) {
return Err(SessionError::InvalidReplacementExecutionState);
}
let git_states: Vec<&Event> = events[candidate_index + 1..]
.iter()
.filter(|event| event.kind == EventKind::GitState)
.collect();
if git_states.len() != 1 {
return Err(SessionError::InvalidReplacementWorkspaceProvenance);
}
let persisted_workspace: WorkspaceState =
serde_json::from_value(git_states[0].payload.clone())
.map_err(|_| SessionError::InvalidReplacementWorkspaceProvenance)?;
let actual_candidate = workspace.candidate_sha256()?;
if actual_candidate != replacement.candidate_sha256 {
return Err(SessionError::CandidateMismatch {
expected: replacement.candidate_sha256,
actual: actual_candidate,
});
}
if persisted_workspace != workspace.state()? {
return Err(SessionError::ReplacementWorkspaceMismatch);
}
Ok(Some(replacement))
}
pub fn transition(
&mut self,
store: &mut EventStore,
next: SessionState,
reason: &str,
) -> Result<(), SessionError> {
let previous = self.state;
let kind = if next.is_terminal() {
EventKind::TerminalState
} else {
EventKind::Checkpoint
};
store.append(
&self.id,
kind,
&json!({"state": next, "previous_state": previous, "reason": reason}),
)?;
self.state = next;
Ok(())
}
}
fn eligible_candidate_provenance(
events: &[Event],
) -> Result<(&Event, &Event, String), SessionError> {
let candidate_index = events
.iter()
.rposition(|event| event.kind == EventKind::CandidateReady)
.ok_or(SessionError::MissingCandidateProvenance)?;
let candidate = &events[candidate_index];
let digest = candidate.payload["candidate_sha256"]
.as_str()
.filter(|value| valid_sha256(value))
.ok_or(SessionError::MissingCandidateProvenance)?;
let summary = candidate.payload["summary"]
.as_str()
.filter(|value| !value.trim().is_empty())
.ok_or(SessionError::MissingCandidateProvenance)?
.to_owned();
let boundary = events[..candidate_index]
.iter()
.rev()
.find(|event| {
event.kind == EventKind::Checkpoint
&& event.payload["checkpoint_kind"] == "turn_boundary"
})
.filter(|event| event.payload["candidate_sha256"].as_str() == Some(digest))
.ok_or(SessionError::MissingCandidateProvenance)?;
let git_states: Vec<&Event> = events[candidate_index + 1..]
.iter()
.filter(|event| event.kind == EventKind::GitState)
.collect();
let source_activity_after_candidate = events[candidate_index + 1..].iter().any(|event| {
matches!(
event.kind,
EventKind::ModelRequest
| EventKind::ModelResponse
| EventKind::ToolRequest
| EventKind::ToolResult
| EventKind::FileMutation
| EventKind::CommandExecution
| EventKind::CandidateReady
| EventKind::GenUiActionPresented
| EventKind::GenUiActionConfirmationRequested
| EventKind::GenUiActionConfirmed
| EventKind::GenUiActionExecutionStarted
| EventKind::GenUiActionExecutionCompleted
| EventKind::GenUiActionExecutionFailed
| EventKind::GenUiActionExecutionUnknown
)
});
if git_states.len() != 1
|| boundary.sequence >= candidate.sequence
|| source_activity_after_candidate
{
return Err(SessionError::MissingCandidateProvenance);
}
Ok((candidate, git_states[0], summary))
}
fn valid_sha256(value: &str) -> bool {
value.len() == 64
&& value
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
}
fn state_from_event(event: &Event) -> Result<SessionState, SessionError> {
let value = event
.payload
.get("state")
.ok_or_else(|| SessionError::InvalidState(event.payload.to_string()))?;
serde_json::from_value(value.clone()).map_err(|_| SessionError::InvalidState(value.to_string()))
}
#[cfg(test)]
mod tests {
use serde_json::json;
use crate::event::{EventKind, EventStore};
use super::{Session, SessionState};
#[test]
fn reconstructs_counts_and_terminal_state() {
let mut store = EventStore::open_memory().expect("store");
let mut session = Session::create(&mut store, "fix it").expect("create");
store
.append(&session.id, EventKind::ModelRequest, &json!({}))
.expect("model request");
store
.append(&session.id, EventKind::ToolRequest, &json!({}))
.expect("tool request");
store
.append(&session.id, EventKind::ToolResult, &json!({"ok": true}))
.expect("tool result");
session
.transition(&mut store, SessionState::Failed, "test")
.expect("transition");
let loaded = Session::reconstruct(&store, &session.id).expect("reconstruct");
assert_eq!(loaded.goal, "fix it");
assert_eq!(loaded.model_turns, 1);
assert_eq!(loaded.tool_calls, 1);
assert_eq!(loaded.state, SessionState::Failed);
}
}