use std::collections::BTreeMap;
use agent_client_protocol::schema::ProtocolVersion as AcpProtocolVersion;
use agent_client_protocol::schema::v1::{
AgentCapabilities, AvailableCommand, ContentBlock, Implementation, SessionConfigOption,
SessionModeState, SessionUpdate,
};
use anyhow::{Context, Result, anyhow, bail};
use serde::{Deserialize, Serialize};
use crate::config::HarnessKind;
use crate::elicitation::ElicitationRequest;
use serde_json::Value;
use sha2::{Digest, Sha256};
use super::capacity::{CAPACITY_STOP_REASON, CapacityRetry};
use super::{
RELAY_EVENT_DIGEST_DOMAIN, RELAY_EVENT_DIGEST_DOMAIN_V2, RELAY_EVENT_GENESIS_DIGEST,
RELAY_STATE_VERSION, RELAY_TRUNCATION_FLOOR,
};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(
tag = "type",
content = "data",
rename_all = "snake_case",
deny_unknown_fields
)]
pub enum RelayCommand {
Prompt {
prompt: Vec<ContentBlock>,
},
RunUserShell {
command: String,
},
CancelUserShell {
shell_command_id: String,
},
RemoveQueuedPrompt {
queued_command_id: String,
},
ClearQueuedPrompts,
SetConfig {
key: String,
value: String,
},
GoalControl {
action: crate::goal::GoalControlAction,
},
SetSessionMode {
mode_id: String,
},
CancelTurn,
Cancel,
Close {
barrier_command_id: String,
expected: RelayCursor,
},
BeginCheckpoint {
reason: Option<String>,
},
CompleteCheckpoint {
barrier_command_id: String,
},
ReleaseCheckpoint {
barrier_command_id: String,
},
AdvanceRecoveryFloor {
through: RelayCursor,
},
RecordNotice {
text: String,
},
}
impl RelayCommand {
pub fn minimum_protocol(&self) -> u32 {
match self {
Self::RunUserShell { .. } | Self::CancelUserShell { .. } => 5,
Self::GoalControl { .. } => 11,
Self::CancelTurn => 7,
Self::Prompt { prompt } if crate::attachment::has_references(prompt) => 8,
_ => super::RELAY_MIN_PROTOCOL_VERSION,
}
}
pub fn is_queue_entry(&self) -> bool {
matches!(self, Self::Prompt { .. } | Self::SetConfig { .. })
}
pub fn is_relay_local(&self) -> bool {
matches!(
self,
Self::RemoveQueuedPrompt { .. }
| Self::ClearQueuedPrompts
| Self::CompleteCheckpoint { .. }
| Self::ReleaseCheckpoint { .. }
| Self::AdvanceRecoveryFloor { .. }
| Self::RecordNotice { .. }
)
}
pub fn is_effectful_acp(&self) -> bool {
matches!(
self,
Self::Prompt { .. }
| Self::SetConfig { .. }
| Self::GoalControl { .. }
| Self::SetSessionMode { .. }
| Self::CancelTurn
| Self::Cancel
| Self::Close { .. }
)
}
pub fn is_effectful_user_shell(&self) -> bool {
matches!(
self,
Self::RunUserShell { .. } | Self::CancelUserShell { .. }
)
}
pub const fn kind(&self) -> RelayCommandKind {
match self {
Self::Prompt { .. } => RelayCommandKind::Prompt,
Self::RunUserShell { .. } => RelayCommandKind::RunUserShell,
Self::CancelUserShell { .. } => RelayCommandKind::CancelUserShell,
Self::RemoveQueuedPrompt { .. } => RelayCommandKind::RemoveQueuedPrompt,
Self::ClearQueuedPrompts => RelayCommandKind::ClearQueuedPrompts,
Self::SetConfig { .. } => RelayCommandKind::SetConfig,
Self::GoalControl { .. } => RelayCommandKind::GoalControl,
Self::SetSessionMode { .. } => RelayCommandKind::SetSessionMode,
Self::CancelTurn => RelayCommandKind::CancelTurn,
Self::Cancel => RelayCommandKind::Cancel,
Self::Close { .. } => RelayCommandKind::Close,
Self::BeginCheckpoint { .. } => RelayCommandKind::BeginCheckpoint,
Self::CompleteCheckpoint { .. } => RelayCommandKind::CompleteCheckpoint,
Self::ReleaseCheckpoint { .. } => RelayCommandKind::ReleaseCheckpoint,
Self::AdvanceRecoveryFloor { .. } => RelayCommandKind::AdvanceRecoveryFloor,
Self::RecordNotice { .. } => RelayCommandKind::RecordNotice,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RelayCommandKind {
Prompt,
RunUserShell,
CancelUserShell,
RemoveQueuedPrompt,
ClearQueuedPrompts,
SetConfig,
GoalControl,
SetSessionMode,
CancelTurn,
Cancel,
Close,
BeginCheckpoint,
CompleteCheckpoint,
ReleaseCheckpoint,
AdvanceRecoveryFloor,
RecordNotice,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueuedRelayPrompt {
pub command_id: String,
pub created_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ActiveRelayPrompt {
pub command_id: String,
pub created_at_ms: i64,
pub started_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ActiveUserShell {
pub command_id: String,
pub command: String,
pub created_at_ms: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub started_at_ms: Option<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ActiveAgentTerminal {
pub terminal_id: String,
pub command: String,
pub started_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BackgroundCommand {
#[serde(default)]
pub id: String,
pub started_at_ms: i64,
pub command: String,
#[serde(default)]
pub can_stop: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BackgroundTaskStopTarget {
HostedTerminal { terminal_id: String },
ClaudeAsyncTask { task_id: String },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct HarnessTurn {
pub started_at_ms: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UserShellStatus {
Exited,
Signaled,
TimedOut,
Cancelled,
Interrupted,
Failed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct UserShellResult {
pub command: String,
pub stdout: String,
pub stderr: String,
#[serde(default)]
pub stdout_truncated: bool,
#[serde(default)]
pub stderr_truncated: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub signal: Option<String>,
pub duration_ms: u64,
pub status: UserShellStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
impl UserShellResult {
pub fn prompt_context(&self) -> String {
fn escaped(text: &str) -> String {
text.replace('&', "&")
.replace('<', "<")
.replace('>', ">")
}
let status = match self.status {
UserShellStatus::Exited => "exited",
UserShellStatus::Signaled => "signaled",
UserShellStatus::TimedOut => "timed_out",
UserShellStatus::Cancelled => "cancelled",
UserShellStatus::Interrupted => "interrupted",
UserShellStatus::Failed => "failed",
};
let mut result = format!("status: {status}\nduration_ms: {}", self.duration_ms);
if let Some(exit_code) = self.exit_code {
result.push_str(&format!("\nexit_code: {exit_code}"));
}
if let Some(signal) = &self.signal {
result.push_str(&format!("\nsignal: {}", escaped(signal)));
}
if let Some(error) = &self.error {
result.push_str(&format!("\nerror: {}", escaped(error)));
}
if !self.stdout.is_empty() {
result.push_str(&format!("\nstdout:\n{}", escaped(&self.stdout)));
}
if !self.stderr.is_empty() {
result.push_str(&format!("\nstderr:\n{}", escaped(&self.stderr)));
}
format!(
"<user_shell_command>\n<command>{}</command>\n<result>{result}</result>\n</user_shell_command>",
escaped(&self.command)
)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StoredQueuedRelayCommand {
pub command_id: String,
#[serde(flatten)]
pub payload: StoredQueuedRelayPayload,
pub created_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum StoredQueuedRelayPayload {
Prompt { prompt: Vec<ContentBlock> },
SetConfig { key: String, value: String },
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StoredActiveRelayPrompt {
pub command_id: String,
pub prompt: Vec<ContentBlock>,
pub created_at_ms: i64,
pub started_at_ms: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RelayExecutionState {
Idle,
Running,
Closing,
Closed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RelayCursor {
pub ordinal: u64,
pub digest: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RelayOperationalState {
#[serde(default)]
pub goal: crate::goal::GoalState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub capacity_retry: Option<CapacityRetry>,
pub session_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub activity_turn_started_at_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub store_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idle_since_ms: Option<i64>,
pub execution: RelayExecutionState,
pub latest_ordinal: u64,
pub latest_digest: String,
pub acknowledged_through: u64,
pub acknowledged_digest: String,
pub recovery_floor_ordinal: u64,
pub recovery_floor_digest: String,
pub native_session_id: Option<String>,
#[serde(default)]
pub native_continuity_lost: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub acp_ready: Option<bool>,
#[serde(default)]
pub checkpoint_only: bool,
pub agent_capabilities: Option<Box<AgentCapabilities>>,
pub agent_info: Option<Implementation>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub steering_supported: Option<bool>,
pub config_options: Vec<SessionConfigOption>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub modes: Option<SessionModeState>,
pub available_commands: Vec<AvailableCommand>,
pub config: BTreeMap<String, String>,
pub active_prompt: Option<ActiveRelayPrompt>,
pub queued_prompts: Vec<QueuedRelayPrompt>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub active_user_shells: Vec<ActiveUserShell>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub active_agent_terminals: Vec<ActiveAgentTerminal>,
pub checkpoint_barrier: Option<String>,
pub checkpoint_ready: Option<RelayCursor>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_acp_activity_at_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_step_started_at_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub foreground_tool_started_at_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub harness_turn: Option<HarnessTurn>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_harness_turn_started_ordinal: Option<u64>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub background_commands: Vec<BackgroundCommand>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub background_work_known: Option<bool>,
}
impl RelayOperationalState {
#[must_use]
pub fn native_session_is_ready(&self) -> bool {
!self.checkpoint_only
&& self.execution != RelayExecutionState::Closed
&& self.native_session_id.is_some()
&& self.acp_ready.unwrap_or(true)
}
#[must_use]
pub fn is_quiet(&self) -> bool {
self.goal.pending_resume.is_none()
&& self.goal.decision.is_none()
&& !self.goal.active()
&& !self.goal.running()
&& self.execution == RelayExecutionState::Idle
&& self.acp_ready != Some(false)
&& self.background_work_known != Some(false)
&& self.active_prompt.is_none()
&& self.harness_turn.is_none()
&& self.queued_prompts.is_empty()
&& self.active_user_shells.is_empty()
&& self.active_agent_terminals.is_empty()
&& self.foreground_tool_started_at_ms.is_none()
&& self.background_commands.is_empty()
&& self.checkpoint_barrier.is_none()
}
#[must_use]
pub fn safe_to_replace(&self, harness: HarnessKind) -> bool {
self.is_quiet()
&& (harness != HarnessKind::Codex || self.goal.synchronized())
&& (harness != HarnessKind::Kimi || self.background_work_known == Some(true))
}
#[must_use]
pub fn safe_for_checkpoint(&self, harness: HarnessKind) -> bool {
self.checkpoint_background_blocker(harness).is_none()
}
pub fn checkpoint_background_blocker(&self, harness: HarnessKind) -> Option<&'static str> {
if self.checkpoint_only || self.execution == RelayExecutionState::Closed {
None
} else if harness == HarnessKind::Codex && !self.goal.synchronized() {
Some("Codex goal and execution state is not synchronized; checkpoint deferred")
} else if self.goal.active() || self.goal.running() || self.goal.decision.is_some() {
Some("an active goal owns this session; pause the goal before checkpointing")
} else if harness != HarnessKind::Kimi {
None
} else if self.background_work_known.is_none() {
Some(
"Kimi worker has not reported background-agent synchronization support; checkpoint requires a worker reporting synchronized task state",
)
} else if self.background_work_known == Some(false) {
Some(
"Kimi background-agent state is not synchronized; checkpoint requires a synchronized empty task list",
)
} else if !self.background_commands.is_empty() {
Some(
"Kimi background agents are still active; checkpoint requires their completion and a synchronized empty task list",
)
} else {
None
}
}
}
pub const RELAY_EVENT_FORMAT_V1: u8 = 1;
pub const RELAY_EVENT_FORMAT_V2: u8 = 2;
fn default_relay_event_format() -> u8 {
RELAY_EVENT_FORMAT_V1
}
fn is_relay_event_format_v1(format: &u8) -> bool {
*format == RELAY_EVENT_FORMAT_V1
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RelayEvent {
#[serde(
default = "default_relay_event_format",
skip_serializing_if = "is_relay_event_format_v1"
)]
pub format: u8,
pub ordinal: u64,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub previous_digest: String,
pub digest: String,
pub recorded_at_ms: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub command_id: Option<String>,
pub observation: RelayObservation,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", content = "data", rename_all = "snake_case")]
pub enum RelayObservation {
AgentInitialized {
protocol_version: AcpProtocolVersion,
capabilities: Box<AgentCapabilities>,
agent_info: Option<Implementation>,
},
SessionOpened {
native_session_id: String,
resumed: bool,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
native_continuity_lost: bool,
},
SessionConfigured {
config_options: Vec<SessionConfigOption>,
},
SessionModesConfigured {
modes: Option<SessionModeState>,
},
SessionUpdate {
update: Box<SessionUpdate>,
},
PermissionAutoApproved {
option_id: String,
option_name: String,
},
ElicitationRequested {
request: ElicitationRequest,
},
ElicitationResolved {
elicitation_id: String,
action: String,
},
ElicitationsCleared,
CommandQueued {
command_id: String,
command: RelayCommand,
created_at_ms: i64,
},
CommandStarted {
command_id: String,
started_at_ms: i64,
},
CommandCompleted {
command_id: String,
outcome: RelayCommandOutcome,
},
CommandRejected {
command_id: String,
command: RelayCommandKind,
message: String,
},
CommandInterrupted {
command_id: String,
command: RelayCommandKind,
message: String,
},
UserShellOutput {
command_id: String,
command: String,
stdout: String,
stderr: String,
stdout_truncated: bool,
stderr_truncated: bool,
},
ConfigurationUpdated {
key: String,
value: String,
},
CheckpointReady {
command_id: String,
through: u64,
},
Warning {
message: String,
},
SessionRestarted,
TerminalOutput {
terminal_id: String,
output: String,
truncated: bool,
exit_code: Option<u32>,
signal: Option<String>,
},
Notice {
message: String,
},
HarnessTurnStarted {
started_at_ms: i64,
},
HarnessTurnSettled {
origin: Option<String>,
#[serde(default)]
prompt_in_flight: bool,
},
Closing,
Closed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "type", content = "data", rename_all = "snake_case")]
pub enum RelayCommandOutcome {
Prompt {
stop_reason: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
#[serde(default, skip_serializing_if = "Option::is_none")]
usage: Option<crate::usage::TokenUsage>,
},
UserShell {
result: UserShellResult,
},
UserShellCancelled,
Configured,
GoalControlled,
SessionModeSet,
Cancelled,
Steered {
queued_command_id: String,
},
Closed,
QueueChanged {
removed_command_ids: Vec<String>,
},
CheckpointCompleted,
CheckpointReleased,
RecoveryFloorAdvanced,
NoticeRecorded,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ClaimedSteeringPrompt {
#[serde(skip)]
pub attachment_root: Option<std::path::PathBuf>,
pub queued_command_id: String,
pub prompt: Vec<ContentBlock>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ClaimedRelayCommand {
pub command_id: String,
pub accepted_ordinal: u64,
pub command: RelayCommand,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hidden_prompt_context: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub steering_prompt: Option<ClaimedSteeringPrompt>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PendingPromptContext {
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attached_command_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PendingUserShellContext {
pub shell_command_id: String,
pub accepted_ordinal: u64,
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attached_command_id: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RelayDispatchState {
Queued,
Pending,
InFlight,
Completed,
Rejected,
Interrupted,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RelayDispatchRecord {
pub command: RelayCommand,
pub state: RelayDispatchState,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StoredHarnessTurn {
pub started_at_ms: i64,
pub first_ordinal: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct HandledRelayCommand {
pub command: RelayCommand,
pub accepted_ordinal: u64,
pub terminal_ordinal: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RelaySnapshot {
#[serde(default)]
pub goal: crate::goal::GoalState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub capacity_retry: Option<CapacityRetry>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub activity_turn_started_at_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub store_id: Option<String>,
pub format_version: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idle_since_ms: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub activity_was_idle: Option<bool>,
pub session_id: String,
pub execution: RelayExecutionState,
pub latest_ordinal: u64,
pub latest_digest: String,
pub acknowledged_through: u64,
pub acknowledged_digest: String,
pub recovery_floor_ordinal: u64,
pub recovery_floor_digest: String,
pub native_session_id: Option<String>,
#[serde(default)]
pub native_continuity_lost: bool,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub native_session_used: bool,
pub agent_capabilities: Option<Box<AgentCapabilities>>,
pub agent_info: Option<Implementation>,
pub config_options: Vec<SessionConfigOption>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub modes: Option<SessionModeState>,
pub available_commands: Vec<AvailableCommand>,
pub config: BTreeMap<String, String>,
pub active_prompt: Option<StoredActiveRelayPrompt>,
pub queued_prompts: Vec<StoredQueuedRelayCommand>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub pending_prompt_context: Option<PendingPromptContext>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub pending_user_shell_contexts: Vec<PendingUserShellContext>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub active_user_shells: BTreeMap<String, ActiveUserShell>,
pub checkpoint_barrier: Option<String>,
pub checkpoint_ready_through: Option<u64>,
pub checkpoint_ready_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub harness_turn: Option<StoredHarnessTurn>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_harness_turn_started_ordinal: Option<u64>,
pub handled_commands: BTreeMap<String, HandledRelayCommand>,
pub dispatches: BTreeMap<String, RelayDispatchRecord>,
}
impl RelaySnapshot {
pub fn new(session_id: String) -> Self {
Self {
goal: Default::default(),
capacity_retry: None,
activity_turn_started_at_ms: None,
store_id: None,
format_version: RELAY_STATE_VERSION,
idle_since_ms: None,
activity_was_idle: None,
session_id,
execution: RelayExecutionState::Idle,
latest_ordinal: 0,
latest_digest: RELAY_EVENT_GENESIS_DIGEST.to_owned(),
acknowledged_through: 0,
acknowledged_digest: RELAY_EVENT_GENESIS_DIGEST.to_owned(),
recovery_floor_ordinal: 0,
recovery_floor_digest: RELAY_EVENT_GENESIS_DIGEST.to_owned(),
native_session_id: None,
native_continuity_lost: false,
native_session_used: false,
agent_capabilities: None,
agent_info: None,
config_options: Vec::new(),
modes: None,
available_commands: Vec::new(),
config: BTreeMap::new(),
active_prompt: None,
queued_prompts: Vec::new(),
pending_prompt_context: None,
pending_user_shell_contexts: Vec::new(),
active_user_shells: BTreeMap::new(),
checkpoint_barrier: None,
checkpoint_ready_through: None,
checkpoint_ready_digest: None,
harness_turn: None,
last_harness_turn_started_ordinal: None,
handled_commands: BTreeMap::new(),
dispatches: BTreeMap::new(),
}
}
pub fn operational_state(&self) -> RelayOperationalState {
RelayOperationalState {
goal: self.goal.clone(),
capacity_retry: self.capacity_retry.clone().filter(|r| !r.submitted),
activity_turn_started_at_ms: self.activity_turn_started_at_ms,
store_id: self.store_id.clone(),
session_id: self.session_id.clone(),
idle_since_ms: self.idle_since_ms,
execution: self.execution,
latest_ordinal: self.latest_ordinal,
latest_digest: self.latest_digest.clone(),
acknowledged_through: self.acknowledged_through,
acknowledged_digest: self.acknowledged_digest.clone(),
recovery_floor_ordinal: self.recovery_floor_ordinal,
recovery_floor_digest: self.recovery_floor_digest.clone(),
native_session_id: self.native_session_id.clone(),
native_continuity_lost: self.native_continuity_lost,
checkpoint_only: false,
acp_ready: None,
agent_capabilities: self.agent_capabilities.clone(),
agent_info: self.agent_info.clone(),
steering_supported: None,
config_options: self.config_options.clone(),
modes: self.modes.clone(),
available_commands: self.available_commands.clone(),
config: self.config.clone(),
active_prompt: self.active_prompt.as_ref().map(|prompt| ActiveRelayPrompt {
command_id: prompt.command_id.clone(),
created_at_ms: prompt.created_at_ms,
started_at_ms: prompt.started_at_ms,
}),
queued_prompts: self
.queued_prompts
.iter()
.map(|prompt| QueuedRelayPrompt {
command_id: prompt.command_id.clone(),
created_at_ms: prompt.created_at_ms,
})
.collect(),
active_user_shells: self.active_user_shells.values().cloned().collect(),
active_agent_terminals: Vec::new(),
checkpoint_barrier: self.checkpoint_barrier.clone(),
checkpoint_ready: self
.checkpoint_ready_through
.zip(self.checkpoint_ready_digest.as_ref())
.map(|(ordinal, digest)| RelayCursor {
ordinal,
digest: digest.clone(),
}),
last_acp_activity_at_ms: None,
current_step_started_at_ms: None,
foreground_tool_started_at_ms: None,
harness_turn: self.harness_turn.map(|turn| HarnessTurn {
started_at_ms: turn.started_at_ms,
}),
last_harness_turn_started_ordinal: self.last_harness_turn_started_ordinal,
background_commands: Vec::new(),
background_work_known: None,
}
}
pub fn retained_through(&self) -> u64 {
self.acknowledged_through.min(self.recovery_floor_ordinal)
}
pub fn retained_digest(&self) -> &str {
if self.acknowledged_through <= self.recovery_floor_ordinal {
&self.acknowledged_digest
} else {
&self.recovery_floor_digest
}
}
}
pub fn ensure_serialized_budget(
value: &impl Serialize,
budget: usize,
description: &str,
) -> Result<()> {
let size = serde_json::to_vec(value)
.with_context(|| format!("serialize {description} for size validation"))?
.len();
ensure_byte_budget(size, budget, description)
}
pub fn ensure_byte_budget(size: usize, budget: usize, description: &str) -> Result<()> {
if size > budget {
bail!("{description} is too large ({size} bytes; maximum {budget})");
}
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum JsonSegment {
Key(String),
Index(usize),
}
fn longest_string_path(value: &Value) -> Option<(Vec<JsonSegment>, usize)> {
fn walk(
value: &Value,
path: &mut Vec<JsonSegment>,
best: &mut Option<(Vec<JsonSegment>, usize)>,
) {
match value {
Value::String(text) => {
if best.as_ref().is_none_or(|(_, length)| text.len() > *length) {
*best = Some((path.clone(), text.len()));
}
}
Value::Array(items) => {
for (index, item) in items.iter().enumerate() {
path.push(JsonSegment::Index(index));
walk(item, path, best);
path.pop();
}
}
Value::Object(entries) => {
for (key, entry) in entries {
path.push(JsonSegment::Key(key.clone()));
walk(entry, path, best);
path.pop();
}
}
_ => {}
}
}
let mut best = None;
walk(value, &mut Vec::new(), &mut best);
best
}
fn string_at_path<'a>(value: &'a mut Value, path: &[JsonSegment]) -> Option<&'a mut String> {
let mut cursor = value;
for segment in path {
cursor = match (segment, cursor) {
(JsonSegment::Key(key), Value::Object(entries)) => entries.get_mut(key)?,
(JsonSegment::Index(index), Value::Array(items)) => items.get_mut(*index)?,
_ => return None,
};
}
match cursor {
Value::String(text) => Some(text),
_ => None,
}
}
fn truncate_with_marker(text: &mut String, keep: usize) {
let mut end = keep.min(text.len());
while end > 0 && !text.is_char_boundary(end) {
end -= 1;
}
let dropped = text.len() - end;
text.truncate(end);
text.push_str(&format!("… [mj truncated {dropped} bytes]"));
}
#[cfg_attr(not(unix), allow(dead_code))]
pub fn truncate_start_with_marker(text: &mut String, keep: usize) -> bool {
if text.len() <= keep {
return false;
}
let mut start = text.len() - keep;
while start < text.len() && !text.is_char_boundary(start) {
start += 1;
}
let dropped = start;
text.drain(..start);
text.insert_str(0, &format!("[mj dropped {dropped} earlier bytes]\n"));
true
}
pub fn clamp_observation(observation: RelayObservation, budget: usize) -> Result<RelayObservation> {
let mut size = serde_json::to_vec(&observation)
.context("measure relay observation")?
.len();
if size <= budget {
return Ok(observation);
}
let mut value =
serde_json::to_value(&observation).context("serialize relay observation for clamping")?;
let original = size;
while size > budget {
let Some((path, length)) = longest_string_path(&value) else {
break;
};
if length <= RELAY_TRUNCATION_FLOOR {
break;
}
let Some(text) = string_at_path(&mut value, &path) else {
break;
};
let keep = length
.saturating_sub(size - budget + 64)
.max(RELAY_TRUNCATION_FLOOR);
truncate_with_marker(text, keep);
size = serde_json::to_vec(&value)
.context("measure clamped relay observation")?
.len();
}
if size > budget {
return Ok(RelayObservation::Warning {
message: format!(
"dropped an observation that cannot be recorded: {original} bytes exceeds the {budget} byte event budget and its payload is not truncatable"
),
});
}
match serde_json::from_value(value) {
Ok(clamped) => {
tracing::warn!(
original,
clamped = size,
"truncated an oversized relay observation"
);
Ok(clamped)
}
Err(error) => Ok(RelayObservation::Warning {
message: format!(
"dropped an observation of {original} bytes: it could not be re-read after truncation: {error}"
),
}),
}
}
#[derive(Serialize)]
struct RelayEventDigestPayload<'a, O: Serialize> {
ordinal: u64,
previous_digest: &'a str,
recorded_at_ms: i64,
#[serde(skip_serializing_if = "Option::is_none")]
command_id: Option<&'a str>,
observation: &'a O,
}
#[derive(Serialize)]
struct RelayEventDigestPayloadV2<'a, O: Serialize> {
ordinal: u64,
recorded_at_ms: i64,
#[serde(skip_serializing_if = "Option::is_none")]
command_id: Option<&'a str>,
observation: &'a O,
}
#[derive(Serialize)]
#[serde(tag = "type", content = "data", rename_all = "snake_case")]
enum LegacyFlaggedObservation<'a> {
SessionOpened {
native_session_id: &'a str,
resumed: bool,
native_continuity_lost: bool,
},
}
fn digest_over(domain: &[u8], encoded: &[u8]) -> String {
let mut hasher = Sha256::new();
hasher.update(domain);
hasher.update(encoded);
format!("{:x}", hasher.finalize())
}
pub fn relay_event_digest(event: &RelayEvent) -> Result<String> {
relay_event_digest_over(event, &event.observation)
}
fn relay_event_digest_over<O: Serialize>(event: &RelayEvent, observation: &O) -> Result<String> {
match event.format {
RELAY_EVENT_FORMAT_V1 => {
validate_relay_digest(&event.previous_digest, "previous event digest")?;
let payload = RelayEventDigestPayload {
ordinal: event.ordinal,
previous_digest: &event.previous_digest,
recorded_at_ms: event.recorded_at_ms,
command_id: event.command_id.as_deref(),
observation,
};
let encoded =
serde_json::to_vec(&payload).context("serialize relay event digest payload")?;
Ok(digest_over(RELAY_EVENT_DIGEST_DOMAIN, &encoded))
}
RELAY_EVENT_FORMAT_V2 => {
if !event.previous_digest.is_empty() {
bail!(
"v2 relay event {} must not carry a previous_digest",
event.ordinal
);
}
let payload = RelayEventDigestPayloadV2 {
ordinal: event.ordinal,
recorded_at_ms: event.recorded_at_ms,
command_id: event.command_id.as_deref(),
observation,
};
let encoded =
serde_json::to_vec(&payload).context("serialize relay event digest payload")?;
Ok(digest_over(RELAY_EVENT_DIGEST_DOMAIN_V2, &encoded))
}
other => bail!(
"unknown relay event format {other} at event {}",
event.ordinal
),
}
}
pub fn validate_relay_event(
previous_ordinal: u64,
previous_digest: &str,
event: &RelayEvent,
) -> Result<()> {
validate_relay_digest(previous_digest, "previous cursor digest")?;
let expected_ordinal = previous_ordinal
.checked_add(1)
.ok_or_else(|| anyhow!("relay event ordinal exhausted"))?;
if event.ordinal != expected_ordinal {
bail!(
"relay event gap: expected {expected_ordinal}, found {}",
event.ordinal
);
}
if event.format == RELAY_EVENT_FORMAT_V1 && event.previous_digest != previous_digest {
bail!(
"relay event {} previous digest does not match cursor",
event.ordinal
);
}
validate_relay_event_self(event)
}
pub fn validate_relay_event_self(event: &RelayEvent) -> Result<()> {
validate_relay_digest(&event.digest, "event digest")?;
let expected_digest = relay_event_digest(event)?;
if event.digest == expected_digest {
return Ok(());
}
if let RelayObservation::SessionOpened {
native_session_id,
resumed,
native_continuity_lost: false,
} = &event.observation
{
let legacy = LegacyFlaggedObservation::SessionOpened {
native_session_id,
resumed: *resumed,
native_continuity_lost: false,
};
if event.digest == relay_event_digest_over(event, &legacy)? {
return Ok(());
}
}
bail!("relay event {} digest is invalid", event.ordinal);
}
pub fn validate_relay_digest(digest: &str, name: &str) -> Result<()> {
if digest.len() != 64
|| !digest
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
{
bail!("{name} must be 64 lowercase hexadecimal characters");
}
Ok(())
}
pub fn observation_changes_state(observation: &RelayObservation) -> bool {
match observation {
RelayObservation::AgentInitialized { .. }
| RelayObservation::SessionOpened { .. }
| RelayObservation::SessionConfigured { .. }
| RelayObservation::SessionModesConfigured { .. }
| RelayObservation::CommandQueued { .. }
| RelayObservation::CommandStarted { .. }
| RelayObservation::CommandCompleted { .. }
| RelayObservation::CommandRejected { .. }
| RelayObservation::CommandInterrupted { .. }
| RelayObservation::ConfigurationUpdated { .. }
| RelayObservation::CheckpointReady { .. }
| RelayObservation::SessionRestarted
| RelayObservation::HarnessTurnStarted { .. }
| RelayObservation::HarnessTurnSettled { .. }
| RelayObservation::Closing
| RelayObservation::Closed => true,
RelayObservation::SessionUpdate { update } => matches!(
update.as_ref(),
SessionUpdate::AvailableCommandsUpdate(_)
| SessionUpdate::ConfigOptionUpdate(_)
| SessionUpdate::CurrentModeUpdate(_)
| SessionUpdate::SessionInfoUpdate(_)
),
RelayObservation::PermissionAutoApproved { .. }
| RelayObservation::ElicitationRequested { .. }
| RelayObservation::ElicitationResolved { .. }
| RelayObservation::ElicitationsCleared
| RelayObservation::Warning { .. }
| RelayObservation::UserShellOutput { .. }
| RelayObservation::TerminalOutput { .. }
| RelayObservation::Notice { .. } => false,
}
}
pub fn apply_relay_event(snapshot: &mut RelaySnapshot, event: &RelayEvent) -> Result<()> {
validate_relay_event(snapshot.latest_ordinal, &snapshot.latest_digest, event)?;
match &event.observation {
RelayObservation::AgentInitialized {
capabilities,
agent_info,
..
} => {
snapshot.agent_capabilities = Some(capabilities.clone());
snapshot.agent_info = agent_info.clone();
}
RelayObservation::SessionOpened {
native_session_id,
native_continuity_lost,
resumed,
} => {
snapshot.native_session_id = Some(native_session_id.clone());
snapshot.native_continuity_lost = *native_continuity_lost;
if *resumed {
snapshot.native_session_used = true;
}
}
RelayObservation::SessionConfigured { config_options } => {
snapshot.config_options = config_options.clone();
}
RelayObservation::SessionModesConfigured { modes } => {
snapshot.modes = modes.clone();
}
RelayObservation::CommandQueued {
command_id,
command,
created_at_ms,
} => {
if cancels_capacity_retry(command) {
if let Some(retry) = snapshot.capacity_retry.as_mut()
&& retry.command_id == *command_id
&& matches!(command, RelayCommand::Prompt { .. })
{
retry.submitted = true;
} else {
snapshot.capacity_retry = None;
}
}
snapshot.handled_commands.insert(
command_id.clone(),
HandledRelayCommand {
command: command.clone(),
accepted_ordinal: event.ordinal,
terminal_ordinal: None,
},
);
snapshot.dispatches.insert(
command_id.clone(),
RelayDispatchRecord {
command: command.clone(),
state: RelayDispatchState::Queued,
},
);
let payload = match command {
RelayCommand::Prompt { prompt } => Some(StoredQueuedRelayPayload::Prompt {
prompt: prompt.clone(),
}),
RelayCommand::SetConfig { key, value } => {
Some(StoredQueuedRelayPayload::SetConfig {
key: key.clone(),
value: value.clone(),
})
}
_ => None,
};
if let Some(payload) = payload {
snapshot.queued_prompts.push(StoredQueuedRelayCommand {
command_id: command_id.clone(),
payload,
created_at_ms: *created_at_ms,
});
}
if let RelayCommand::RunUserShell { command } = command {
snapshot.active_user_shells.insert(
command_id.clone(),
ActiveUserShell {
command_id: command_id.clone(),
command: command.clone(),
created_at_ms: *created_at_ms,
started_at_ms: None,
},
);
}
if matches!(command, RelayCommand::Close { .. }) {
snapshot.execution = RelayExecutionState::Closing;
}
}
RelayObservation::CommandStarted {
command_id,
started_at_ms,
} => {
let dispatch = snapshot
.dispatches
.get_mut(command_id)
.ok_or_else(|| anyhow!("started unknown relay command {command_id}"))?;
dispatch.state = RelayDispatchState::Pending;
match &dispatch.command {
RelayCommand::Prompt { .. } => {
let index = snapshot
.queued_prompts
.iter()
.position(|queued| queued.command_id == *command_id)
.ok_or_else(|| anyhow!("started prompt {command_id} was not queued"))?;
let queued = snapshot.queued_prompts.remove(index);
let StoredQueuedRelayPayload::Prompt { prompt } = queued.payload else {
bail!("queued command {command_id} is not a prompt");
};
snapshot.execution = RelayExecutionState::Running;
snapshot.activity_turn_started_at_ms = Some(*started_at_ms);
snapshot.active_prompt = Some(StoredActiveRelayPrompt {
command_id: queued.command_id,
prompt,
created_at_ms: queued.created_at_ms,
started_at_ms: *started_at_ms,
});
}
RelayCommand::SetConfig { .. } => {
let index = snapshot
.queued_prompts
.iter()
.position(|queued| queued.command_id == *command_id)
.ok_or_else(|| {
anyhow!("started configuration change {command_id} was not queued")
})?;
snapshot.queued_prompts.remove(index);
}
RelayCommand::Close { .. } => snapshot.execution = RelayExecutionState::Closing,
RelayCommand::RunUserShell { .. } => {
let shell = snapshot
.active_user_shells
.get_mut(command_id)
.ok_or_else(|| anyhow!("started unknown user shell {command_id}"))?;
shell.started_at_ms = Some(*started_at_ms);
}
RelayCommand::BeginCheckpoint { .. } => {
if snapshot.checkpoint_barrier.is_some() {
bail!("checkpoint barrier started while another barrier was active");
}
snapshot.checkpoint_barrier = Some(command_id.clone());
snapshot.checkpoint_ready_through = None;
snapshot.checkpoint_ready_digest = None;
}
_ => {}
}
}
RelayObservation::CommandCompleted {
command_id,
outcome,
} => {
let command = snapshot
.dispatches
.get(command_id)
.ok_or_else(|| anyhow!("completed unknown relay command {command_id}"))?
.command
.clone();
snapshot
.dispatches
.get_mut(command_id)
.expect("dispatch disappeared")
.state = RelayDispatchState::Completed;
snapshot
.handled_commands
.get_mut(command_id)
.ok_or_else(|| anyhow!("completed command {command_id} is not in the ledger"))?
.terminal_ordinal = Some(event.ordinal);
if let RelayCommandOutcome::Prompt { stop_reason, .. } = outcome {
let accepted = snapshot.handled_commands[command_id].accepted_ordinal;
let superseded = snapshot.handled_commands.values().any(|handled| {
handled.accepted_ordinal > accepted && cancels_capacity_retry(&handled.command)
});
let attempt = snapshot
.capacity_retry
.as_ref()
.filter(|retry| retry.command_id == *command_id)
.map_or(1, |retry| retry.attempt.saturating_add(1));
snapshot.capacity_retry = if stop_reason == CAPACITY_STOP_REASON && !superseded {
Some(CapacityRetry::new(
attempt,
event.ordinal,
event.recorded_at_ms,
))
} else {
None
};
}
match (command, outcome) {
(RelayCommand::Prompt { .. }, RelayCommandOutcome::Prompt { .. }) => {
if snapshot
.active_prompt
.as_ref()
.map(|active| &active.command_id)
== Some(command_id)
{
snapshot.active_prompt = None;
}
if !snapshot.goal.running() {
snapshot.harness_turn = None;
}
if snapshot.execution == RelayExecutionState::Running
&& !snapshot.goal.running()
{
snapshot.execution = RelayExecutionState::Idle;
}
if snapshot
.pending_prompt_context
.as_ref()
.and_then(|context| context.attached_command_id.as_deref())
== Some(command_id.as_str())
{
snapshot.pending_prompt_context = None;
}
snapshot.pending_user_shell_contexts.retain(|context| {
context.attached_command_id.as_deref() != Some(command_id.as_str())
});
}
(RelayCommand::RunUserShell { .. }, RelayCommandOutcome::UserShell { result }) => {
snapshot.active_user_shells.remove(command_id);
let accepted_ordinal = snapshot
.handled_commands
.get(command_id)
.ok_or_else(|| anyhow!("completed user shell is not in the ledger"))?
.accepted_ordinal;
snapshot
.pending_user_shell_contexts
.push(PendingUserShellContext {
shell_command_id: command_id.clone(),
accepted_ordinal,
text: result.prompt_context(),
attached_command_id: None,
});
}
(RelayCommand::CancelUserShell { .. }, RelayCommandOutcome::UserShellCancelled) => {
}
(
RelayCommand::RemoveQueuedPrompt { queued_command_id },
RelayCommandOutcome::QueueChanged {
removed_command_ids,
},
) => {
let expected = snapshot
.queued_prompts
.iter()
.any(|queued| queued.command_id == queued_command_id)
.then_some(vec![queued_command_id]);
if expected.as_deref() != Some(removed_command_ids.as_slice()) {
bail!("removed queue outcome does not match the durable queue");
}
terminalize_removed_prompts(snapshot, removed_command_ids, event.ordinal)?;
}
(
RelayCommand::ClearQueuedPrompts,
RelayCommandOutcome::QueueChanged {
removed_command_ids,
},
) => {
let expected: Vec<String> = snapshot
.queued_prompts
.iter()
.map(|queued| queued.command_id.clone())
.collect();
if expected != *removed_command_ids {
bail!("cleared queue outcome does not match the durable queue");
}
terminalize_removed_prompts(snapshot, removed_command_ids, event.ordinal)?;
}
(RelayCommand::SetConfig { key, value }, RelayCommandOutcome::Configured) => {
snapshot.config.insert(key.clone(), value.clone());
crate::acp::AcceptedSessionConfig::record_completed(
&mut snapshot.config,
&key,
&value,
&snapshot.config_options,
);
}
(RelayCommand::SetSessionMode { mode_id }, RelayCommandOutcome::SessionModeSet) => {
snapshot.config.insert("mode".to_owned(), mode_id);
}
(RelayCommand::GoalControl { .. }, RelayCommandOutcome::GoalControlled) => {}
(RelayCommand::Cancel, RelayCommandOutcome::Cancelled)
| (RelayCommand::CancelTurn, RelayCommandOutcome::Cancelled) => {}
(RelayCommand::Cancel, RelayCommandOutcome::Steered { queued_command_id }) => {
let queued = snapshot
.queued_prompts
.first()
.ok_or_else(|| anyhow!("steered prompt is no longer queued"))?;
if queued.command_id != *queued_command_id
|| !matches!(queued.payload, StoredQueuedRelayPayload::Prompt { .. })
{
bail!("steered prompt is not the queued prompt head");
}
let target = snapshot
.dispatches
.get_mut(queued_command_id)
.ok_or_else(|| anyhow!("steered unknown queued prompt"))?;
if target.state != RelayDispatchState::Queued
|| !matches!(target.command, RelayCommand::Prompt { .. })
{
bail!("steered target is not a queued prompt");
}
target.state = RelayDispatchState::Completed;
snapshot
.handled_commands
.get_mut(queued_command_id)
.ok_or_else(|| anyhow!("steered prompt is not in the ledger"))?
.terminal_ordinal = Some(event.ordinal);
snapshot.queued_prompts.remove(0);
if snapshot
.pending_prompt_context
.as_ref()
.and_then(|context| context.attached_command_id.as_deref())
== Some(queued_command_id.as_str())
{
snapshot.pending_prompt_context = None;
}
snapshot.pending_user_shell_contexts.retain(|context| {
context.attached_command_id.as_deref() != Some(queued_command_id.as_str())
});
}
(RelayCommand::Close { .. }, RelayCommandOutcome::Closed) => {
snapshot.execution = RelayExecutionState::Closed;
snapshot.active_prompt = None;
}
(
RelayCommand::CompleteCheckpoint { barrier_command_id },
RelayCommandOutcome::CheckpointCompleted,
) => {
if snapshot.checkpoint_barrier.as_deref() != Some(&barrier_command_id) {
bail!("checkpoint completion does not match the active barrier");
}
let ready_through = snapshot
.checkpoint_ready_through
.ok_or_else(|| anyhow!("checkpoint barrier was not ready"))?;
let ready_digest = snapshot
.checkpoint_ready_digest
.clone()
.ok_or_else(|| anyhow!("checkpoint barrier ready digest is missing"))?;
snapshot.recovery_floor_ordinal = ready_through;
snapshot.recovery_floor_digest = ready_digest;
snapshot.checkpoint_barrier = None;
snapshot.checkpoint_ready_through = None;
snapshot.checkpoint_ready_digest = None;
if let Some(barrier) = snapshot.dispatches.get_mut(&barrier_command_id) {
barrier.state = RelayDispatchState::Completed;
}
if let Some(barrier) = snapshot.handled_commands.get_mut(&barrier_command_id) {
barrier.terminal_ordinal = Some(event.ordinal);
}
}
(
RelayCommand::ReleaseCheckpoint { barrier_command_id },
RelayCommandOutcome::CheckpointReleased,
) => {
if snapshot.checkpoint_barrier.as_deref() != Some(&barrier_command_id) {
bail!("checkpoint release does not match the active barrier");
}
if snapshot.checkpoint_ready_through.is_none() {
bail!("checkpoint barrier was not ready");
}
snapshot.checkpoint_barrier = None;
snapshot.checkpoint_ready_through = None;
snapshot.checkpoint_ready_digest = None;
if let Some(barrier) = snapshot.dispatches.get_mut(&barrier_command_id) {
barrier.state = RelayDispatchState::Completed;
}
if let Some(barrier) = snapshot.handled_commands.get_mut(&barrier_command_id) {
barrier.terminal_ordinal = Some(event.ordinal);
}
}
(
RelayCommand::AdvanceRecoveryFloor { through },
RelayCommandOutcome::RecoveryFloorAdvanced,
) => {
if through.ordinal < snapshot.recovery_floor_ordinal {
bail!("recovery floor cannot move back");
}
snapshot.recovery_floor_ordinal = through.ordinal;
snapshot.recovery_floor_digest = through.digest;
}
(RelayCommand::RecordNotice { .. }, RelayCommandOutcome::NoticeRecorded) => {}
(RelayCommand::BeginCheckpoint { .. }, _) => {
bail!("checkpoint barriers complete through checkpoint-ready")
}
(command, outcome) => {
bail!(
"relay command {:?} has incompatible completion outcome {outcome:?}",
command.kind()
)
}
}
}
RelayObservation::CommandRejected {
command_id,
command: observed_command,
message,
}
| RelayObservation::CommandInterrupted {
command_id,
command: observed_command,
message,
} => {
let state = if matches!(event.observation, RelayObservation::CommandRejected { .. }) {
RelayDispatchState::Rejected
} else {
RelayDispatchState::Interrupted
};
let command = snapshot
.dispatches
.get(command_id)
.ok_or_else(|| anyhow!("terminated unknown relay command {command_id}"))?
.command
.clone();
if command.kind() != *observed_command {
bail!("terminated command {command_id} has the wrong command identity");
}
snapshot
.dispatches
.get_mut(command_id)
.expect("dispatch disappeared")
.state = state;
snapshot
.handled_commands
.get_mut(command_id)
.ok_or_else(|| anyhow!("terminated command {command_id} is not in the ledger"))?
.terminal_ordinal = Some(event.ordinal);
snapshot
.queued_prompts
.retain(|queued| queued.command_id != *command_id);
snapshot.active_user_shells.remove(command_id);
if let RelayCommand::RunUserShell { command } = &command {
let accepted_ordinal = snapshot
.handled_commands
.get(command_id)
.expect("terminated shell command disappeared from the ledger")
.accepted_ordinal;
let result = UserShellResult {
command: command.clone(),
stdout: String::new(),
stderr: String::new(),
stdout_truncated: false,
stderr_truncated: false,
exit_code: None,
signal: None,
duration_ms: 0,
status: if state == RelayDispatchState::Rejected {
UserShellStatus::Failed
} else {
UserShellStatus::Interrupted
},
error: Some(message.clone()),
};
snapshot
.pending_user_shell_contexts
.push(PendingUserShellContext {
shell_command_id: command_id.clone(),
accepted_ordinal,
text: result.prompt_context(),
attached_command_id: None,
});
}
if snapshot
.active_prompt
.as_ref()
.map(|active| &active.command_id)
== Some(command_id)
{
snapshot.active_prompt = None;
if !snapshot.goal.running() {
snapshot.harness_turn = None;
snapshot.execution = RelayExecutionState::Idle;
}
}
if snapshot
.pending_prompt_context
.as_ref()
.and_then(|context| context.attached_command_id.as_deref())
== Some(command_id.as_str())
{
snapshot
.pending_prompt_context
.as_mut()
.expect("pending prompt context disappeared")
.attached_command_id = None;
}
for context in &mut snapshot.pending_user_shell_contexts {
if context.attached_command_id.as_deref() == Some(command_id.as_str()) {
context.attached_command_id = None;
}
}
if matches!(command, RelayCommand::BeginCheckpoint { .. })
&& snapshot.checkpoint_barrier.as_deref() == Some(command_id)
{
snapshot.checkpoint_barrier = None;
snapshot.checkpoint_ready_through = None;
snapshot.checkpoint_ready_digest = None;
}
if matches!(command, RelayCommand::Close { .. })
&& snapshot.execution == RelayExecutionState::Closing
{
snapshot.execution = RelayExecutionState::Idle;
}
}
RelayObservation::ConfigurationUpdated { key, value } => {
snapshot.config.insert(key.clone(), value.clone());
crate::acp::AcceptedSessionConfig::record_completed(
&mut snapshot.config,
key,
value,
&snapshot.config_options,
);
}
RelayObservation::CheckpointReady {
command_id,
through,
} => {
let Some(dispatch) = snapshot.dispatches.get(command_id) else {
bail!("checkpoint ready for unknown command {command_id}");
};
if !matches!(dispatch.command, RelayCommand::BeginCheckpoint { .. }) {
bail!("checkpoint ready for non-barrier command {command_id}");
}
if snapshot.checkpoint_barrier.as_deref() != Some(command_id) {
bail!("checkpoint ready does not match the active barrier");
}
if *through != event.ordinal {
bail!("checkpoint ready frontier does not match its event ordinal");
}
snapshot.checkpoint_ready_through = Some(*through);
snapshot.checkpoint_ready_digest = Some(event.digest.clone());
}
RelayObservation::HarnessTurnStarted { started_at_ms } => {
snapshot.activity_turn_started_at_ms = Some(*started_at_ms);
snapshot.harness_turn = Some(StoredHarnessTurn {
started_at_ms: *started_at_ms,
first_ordinal: event.ordinal,
});
snapshot.last_harness_turn_started_ordinal = Some(event.ordinal);
if snapshot.execution == RelayExecutionState::Idle {
snapshot.execution = RelayExecutionState::Running;
}
}
RelayObservation::HarnessTurnSettled { .. } => {
snapshot.harness_turn = None;
if snapshot.active_prompt.is_none()
&& snapshot.execution == RelayExecutionState::Running
{
snapshot.execution = RelayExecutionState::Idle;
}
}
RelayObservation::SessionRestarted => {
snapshot.goal.restart();
if snapshot.harness_turn.take().is_some()
&& snapshot.active_prompt.is_none()
&& snapshot.execution == RelayExecutionState::Running
{
snapshot.execution = RelayExecutionState::Idle;
}
}
RelayObservation::Closing => {
snapshot.capacity_retry = None;
snapshot.harness_turn = None;
snapshot.execution = RelayExecutionState::Closing;
}
RelayObservation::Closed => {
snapshot.capacity_retry = None;
snapshot.activity_turn_started_at_ms = None;
snapshot.harness_turn = None;
snapshot.execution = RelayExecutionState::Closed;
snapshot.active_prompt = None;
}
RelayObservation::SessionUpdate { update } => match update.as_ref() {
SessionUpdate::SessionInfoUpdate(_) => {
snapshot.goal.apply(update)?;
}
SessionUpdate::AvailableCommandsUpdate(update) => {
snapshot.available_commands = update.available_commands.clone();
}
SessionUpdate::ConfigOptionUpdate(update) => {
snapshot.config_options = update.config_options.clone();
}
SessionUpdate::CurrentModeUpdate(update) => {
if let Some(modes) = snapshot.modes.as_mut() {
modes.current_mode_id = update.current_mode_id.clone();
}
snapshot
.config
.insert("mode".to_owned(), update.current_mode_id.to_string());
}
_ => {}
},
RelayObservation::PermissionAutoApproved { .. }
| RelayObservation::ElicitationRequested { .. }
| RelayObservation::ElicitationResolved { .. }
| RelayObservation::ElicitationsCleared
| RelayObservation::Warning { .. }
| RelayObservation::UserShellOutput { .. }
| RelayObservation::TerminalOutput { .. }
| RelayObservation::Notice { .. } => {}
}
snapshot.latest_ordinal = event.ordinal;
snapshot.latest_digest = event.digest.clone();
Ok(())
}
pub fn releases_history(command: &RelayCommand) -> bool {
matches!(
command,
RelayCommand::CompleteCheckpoint { .. } | RelayCommand::AdvanceRecoveryFloor { .. }
)
}
fn terminalize_removed_prompts(
snapshot: &mut RelaySnapshot,
removed_command_ids: &[String],
terminal_ordinal: u64,
) -> Result<()> {
for command_id in removed_command_ids {
let dispatch = snapshot
.dispatches
.get_mut(command_id)
.ok_or_else(|| anyhow!("removed unknown queued command {command_id}"))?;
if !dispatch.command.is_queue_entry() || dispatch.state != RelayDispatchState::Queued {
bail!("removed command {command_id} is not a queued command");
}
dispatch.state = RelayDispatchState::Rejected;
snapshot
.handled_commands
.get_mut(command_id)
.ok_or_else(|| anyhow!("removed command {command_id} is not in the ledger"))?
.terminal_ordinal = Some(terminal_ordinal);
}
snapshot.queued_prompts.retain(|queued| {
!removed_command_ids
.iter()
.any(|command_id| command_id == &queued.command_id)
});
Ok(())
}
pub fn validate_relay_snapshot_frontiers(snapshot: &RelaySnapshot) -> Result<()> {
if snapshot.acknowledged_through > snapshot.latest_ordinal {
bail!("relay acknowledgement is ahead of the event frontier");
}
if snapshot.recovery_floor_ordinal > snapshot.latest_ordinal {
bail!("relay recovery floor is ahead of the event frontier");
}
validate_relay_digest(&snapshot.latest_digest, "relay latest digest")?;
validate_relay_digest(
&snapshot.acknowledged_digest,
"relay acknowledgement digest",
)?;
validate_relay_digest(
&snapshot.recovery_floor_digest,
"relay recovery floor digest",
)?;
if (snapshot.latest_ordinal == 0) != (snapshot.latest_digest == RELAY_EVENT_GENESIS_DIGEST) {
bail!("relay latest frontier and genesis digest disagree");
}
if (snapshot.acknowledged_through == 0)
!= (snapshot.acknowledged_digest == RELAY_EVENT_GENESIS_DIGEST)
{
bail!("relay acknowledgement frontier and genesis digest disagree");
}
if (snapshot.recovery_floor_ordinal == 0)
!= (snapshot.recovery_floor_digest == RELAY_EVENT_GENESIS_DIGEST)
{
bail!("relay recovery floor and genesis digest disagree");
}
Ok(())
}
fn cancels_capacity_retry(command: &RelayCommand) -> bool {
matches!(
command,
RelayCommand::Prompt { .. }
| RelayCommand::Cancel
| RelayCommand::CancelTurn
| RelayCommand::SetConfig { .. }
| RelayCommand::GoalControl { .. }
| RelayCommand::SetSessionMode { .. }
| RelayCommand::Close { .. }
)
}
#[cfg(test)]
mod native_continuity_encoding_tests {
use super::*;
const RECORDED_BEFORE_THE_FLAG: &str = r#"{"format":2,"ordinal":493425,"digest":"28c6574464b1361ad00dedd7cb04dac0c617948e290936082f1fd354334b262e","recorded_at_ms":1789421283703,"observation":{"type":"session_opened","data":{"native_session_id":"fe6031fb-9af5-49ca-a587-5a816de61f86","resumed":false}}}"#;
#[test]
fn a_session_opened_record_without_the_flag_still_verifies() {
let event: RelayEvent = serde_json::from_str(RECORDED_BEFORE_THE_FLAG).unwrap();
validate_relay_event_self(&event).expect("the stored digest must still match");
let encoded = serde_json::to_string(&event.observation).unwrap();
assert!(
!encoded.contains("native_continuity_lost"),
"false must not be written: {encoded}"
);
}
const RECORDED_BY_THE_FLAGGING_BUILD: &str = r#"{"format":2,"ordinal":188,"digest":"83a8ded900c35360e8998e6b7710902475145370bdaa88ce1fa186cb5d108636","recorded_at_ms":1789436362441,"observation":{"type":"session_opened","data":{"native_session_id":"020bb831-ec40-49ac-bbcb-a702135deb5d","resumed":true,"native_continuity_lost":false}}}"#;
#[test]
fn a_record_that_wrote_the_flag_as_false_still_verifies() {
let event: RelayEvent = serde_json::from_str(RECORDED_BY_THE_FLAGGING_BUILD).unwrap();
validate_relay_event_self(&event).expect("the legacy encoding must still verify");
let mut tampered = event;
tampered.recorded_at_ms += 1;
assert!(validate_relay_event_self(&tampered).is_err());
}
#[test]
fn a_lost_native_session_is_written_and_read_back() {
let observation = RelayObservation::SessionOpened {
native_session_id: "fresh".into(),
resumed: false,
native_continuity_lost: true,
};
let encoded = serde_json::to_string(&observation).unwrap();
assert!(encoded.contains("\"native_continuity_lost\":true"));
let decoded: RelayObservation = serde_json::from_str(&encoded).unwrap();
assert_eq!(decoded, observation);
}
}