use std::{
collections::{BTreeMap, HashSet},
future::Future,
path::{Path, PathBuf},
pin::Pin,
sync::Arc,
time::Duration,
};
use crate::opensymphony_domain::{
ConversationId, ConversationMetadata, DurationMs, HarnessAdapter, HarnessInterruptCommand,
IssueId, IssueIdentifier, NormalizedIssue, RunAttempt, RuntimeStreamState, TimestampMs,
WorkerId, WorkerOutcomeKind, WorkerOutcomeRecord,
};
#[cfg(test)]
use crate::opensymphony_domain::{RuntimeLivenessPhase, RuntimeProgressSnapshot};
use crate::opensymphony_gateway_schema::capability::HarnessCapability;
use crate::opensymphony_memory::{
IssueEvidence, IssueLinkEvidence, MemoryConfig, MemoryContextOptions, SourceFile,
context_for_issue_with_options,
};
use crate::opensymphony_workflow::{
Environment, OpenHandsConversationToolConfig, ProcessEnvironment, ResolvedWorkflow,
};
use crate::opensymphony_workspace::{
ParentRuntimeEnvelope, RunManifest, RunStatus, TerminalRuntimeEnvelope, WorkspaceError,
WorkspaceHandle, WorkspaceManager, compose_terminal_prompt,
};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use thiserror::Error;
use tokio::time::{Instant, timeout_at};
use tracing::debug;
use uuid::Uuid;
use super::{
AgentConfig, CondenserConfig, ConfirmationPolicy, Conversation, ConversationCreateRequest,
ConversationStateMirror, EventEnvelope, KnownEvent, LlmConfig, OpenHandsClient, OpenHandsError,
RuntimeEventStream, RuntimeStreamConfig, SendMessageRequest, TerminalExecutionStatus,
ToolConfig, WorkspaceConfig,
};
pub const RUNTIME_CONTRACT_VERSION: &str = "openhands-sdk-agent-server-v1";
const OPENAI_SUBSCRIPTION_CREDENTIAL_MODE: &str = "openai_subscription";
const OPENAI_CODEX_SUBSCRIPTION_BASE_URL: &str = "https://chatgpt.com/backend-api/codex";
const DEFAULT_REUSE_POLICY: &str = "per_issue";
const FRESH_EACH_RUN_REUSE_POLICY: &str = "fresh_each_run";
const INTERRUPT_RECENT_EVENT_LIMIT: usize = 20;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IssueSessionReusePolicy {
PerIssue,
FreshEachRun,
Unsupported(String),
}
impl IssueSessionReusePolicy {
fn parse(value: &str) -> Self {
match value.trim().to_ascii_lowercase().as_str() {
DEFAULT_REUSE_POLICY => Self::PerIssue,
FRESH_EACH_RUN_REUSE_POLICY => Self::FreshEachRun,
other => Self::Unsupported(other.to_owned()),
}
}
pub fn as_str(&self) -> &str {
match self {
Self::PerIssue => DEFAULT_REUSE_POLICY,
Self::FreshEachRun => FRESH_EACH_RUN_REUSE_POLICY,
Self::Unsupported(value) => value.as_str(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IssueSessionRunnerConfig {
pub reuse_policy: IssueSessionReusePolicy,
pub runtime_stream: RuntimeStreamConfig,
pub terminal_wait_timeout: Duration,
pub total_runtime_cap_ms: Option<Duration>,
pub finished_drain_timeout: Duration,
pub memory: Option<MemoryWorkerAccess>,
pub repository_instructions: Option<String>,
pub terminal_prompt: Option<String>,
pub continuation_prompt: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MemoryWorkerAccess {
pub endpoint: String,
pub token: Option<String>,
pub project: Option<String>,
pub project_set: Option<String>,
pub execution_repo: Option<String>,
pub authorized_repositories: Vec<String>,
pub run_id: Option<String>,
pub attempt: Option<u32>,
pub requires_fresh_conversation: bool,
}
impl MemoryWorkerAccess {
pub fn mcp_config(&self) -> Option<BTreeMap<String, Value>> {
let mut server = serde_json::Map::new();
server.insert("url".to_owned(), json!(self.endpoint));
if let Some(token) = &self.token {
server.insert(
"headers".to_owned(),
json!({ "Authorization": format!("Bearer {token}") }),
);
}
Some(BTreeMap::from([(
"mcpServers".to_owned(),
json!({ "opensymphony-memory": Value::Object(server) }),
)]))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WorkpadComment {
pub id: String,
pub body: String,
pub updated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OpenHandsInterruptMethod {
Interrupt,
PauseFallback,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OpenHandsInterruptAcknowledgement {
pub method: OpenHandsInterruptMethod,
pub reconciled_events: usize,
pub execution_status: Option<String>,
pub diagnostic: Option<String>,
pub timed_out: bool,
}
#[async_trait]
pub trait WorkpadCommentSource: Send + Sync {
async fn fetch_workpad_comment(&self, issue_id: &str)
-> Result<Option<WorkpadComment>, String>;
}
pub trait IssueSessionObserver {
fn on_launch(&mut self, _conversation: &ConversationMetadata) {}
fn on_launch_with_started_at(
&mut self,
conversation: &ConversationMetadata,
_started_at: Option<TimestampMs>,
) {
self.on_launch(conversation);
}
fn on_runtime_event(
&mut self,
_observed_at: TimestampMs,
_event_id: Option<String>,
_event_kind: Option<String>,
_summary: Option<String>,
_payload: Option<Value>,
) {
}
fn on_conversation_update(&mut self, _conversation: &ConversationMetadata) {}
}
impl IssueSessionObserver for () {}
#[derive(Debug, Clone)]
pub(crate) struct LivenessTracker {
event_count: u64,
input_tokens: u64,
output_tokens: u64,
execution_status: Option<String>,
last_activity_at: Option<TimestampMs>,
idle_timeout_ms: DurationMs,
total_runtime_cap_ms: Option<DurationMs>,
started_at: Option<TimestampMs>,
}
impl LivenessTracker {
#[cfg(test)]
pub fn new(idle_timeout_ms: DurationMs) -> Self {
Self::with_runtime_cap(idle_timeout_ms, None)
}
pub fn with_runtime_cap(
idle_timeout_ms: DurationMs,
total_runtime_cap_ms: Option<DurationMs>,
) -> Self {
Self {
event_count: 0,
input_tokens: 0,
output_tokens: 0,
execution_status: None,
last_activity_at: None,
idle_timeout_ms,
total_runtime_cap_ms,
started_at: None,
}
}
pub fn mark_started(&mut self, started_at: TimestampMs) {
self.started_at = Some(started_at);
self.last_activity_at = Some(started_at);
}
pub fn record_event(&mut self, event_at: TimestampMs) -> bool {
let advanced = self.advance_activity(event_at);
self.event_count = self.event_count.saturating_add(1);
advanced
}
pub fn record_reconciled_events(&mut self, count: u64, reconciled_at: TimestampMs) -> bool {
self.event_count = self.event_count.saturating_add(count);
self.advance_activity(reconciled_at)
}
pub fn record_tokens(&mut self, input: u64, output: u64, recorded_at: TimestampMs) -> bool {
let advanced = input > self.input_tokens || output > self.output_tokens;
self.input_tokens = self.input_tokens.max(input);
self.output_tokens = self.output_tokens.max(output);
if advanced {
self.advance_activity(recorded_at)
} else {
false
}
}
pub fn record_status_change(&mut self, status: &str, recorded_at: TimestampMs) -> bool {
let changed = self.execution_status.as_deref() != Some(status);
self.execution_status = Some(status.to_string());
if changed {
self.advance_activity(recorded_at);
}
changed
}
pub fn is_stalled_at(&self, now: TimestampMs) -> bool {
let Some(started_at) = self.started_at else {
return false;
};
let Some(last_activity) = self.last_activity_at else {
return false;
};
let idle_deadline = last_activity.saturating_add(self.idle_timeout_ms);
if now >= idle_deadline {
return true;
}
if let Some(cap) = self.total_runtime_cap_ms {
let hard_cap = started_at.saturating_add(cap);
if now >= hard_cap {
return true;
}
}
false
}
pub fn stall_deadline_at(&self) -> Option<TimestampMs> {
let last_activity = self.last_activity_at?;
let started_at = self.started_at?;
let idle_deadline = last_activity.saturating_add(self.idle_timeout_ms);
match self.total_runtime_cap_ms {
Some(cap) => {
let hard_cap = started_at.saturating_add(cap);
Some(idle_deadline.min(hard_cap))
}
None => Some(idle_deadline),
}
}
#[cfg(test)]
pub fn snapshot(&self, previous: &RuntimeProgressSnapshot) -> RuntimeProgressSnapshot {
let phase = match (self.last_activity_at, self.started_at) {
(None, None) => RuntimeLivenessPhase::WaitingOnPriorTurn,
(Some(_), Some(_)) => RuntimeLivenessPhase::RunningTurn,
_ => RuntimeLivenessPhase::WaitingOnPriorTurn,
};
previous
.update_with(phase)
.with_event_count(self.event_count)
.with_input_tokens(self.input_tokens)
.with_output_tokens(self.output_tokens)
.with_execution_status(self.execution_status.clone())
.with_last_activity_at(self.last_activity_at)
.with_stall_deadline_at(self.stall_deadline_at())
.build()
}
fn advance_activity(&mut self, activity_at: TimestampMs) -> bool {
if self.last_activity_at.is_some_and(|last| activity_at < last) {
return false;
}
let advanced = self.last_activity_at.is_none_or(|last| activity_at > last);
self.last_activity_at = Some(activity_at);
advanced
}
}
impl Default for IssueSessionRunnerConfig {
fn default() -> Self {
Self {
reuse_policy: IssueSessionReusePolicy::PerIssue,
runtime_stream: RuntimeStreamConfig::default(),
terminal_wait_timeout: Duration::from_secs(300),
total_runtime_cap_ms: None,
finished_drain_timeout: Duration::from_millis(100),
memory: None,
repository_instructions: None,
terminal_prompt: None,
continuation_prompt: None,
}
}
}
impl IssueSessionRunnerConfig {
pub fn from_workflow(workflow: &ResolvedWorkflow) -> Self {
let websocket = &workflow.extensions.openhands.websocket;
Self {
reuse_policy: IssueSessionReusePolicy::parse(
&workflow.extensions.openhands.conversation.reuse_policy,
),
runtime_stream: RuntimeStreamConfig {
readiness_timeout: Duration::from_millis(websocket.ready_timeout_ms),
reconnect_initial_backoff: Duration::from_millis(websocket.reconnect_initial_ms),
reconnect_max_backoff: Duration::from_millis(websocket.reconnect_max_ms),
..RuntimeStreamConfig::default()
},
terminal_wait_timeout: Duration::from_millis(
workflow.config.agent.stall_timeout_ms.unwrap_or(300_000),
),
total_runtime_cap_ms: None,
finished_drain_timeout: Duration::from_millis(100),
memory: None,
repository_instructions: None,
terminal_prompt: None,
continuation_prompt: None,
}
}
pub fn with_memory(mut self, memory: Option<MemoryWorkerAccess>) -> Self {
self.memory = memory;
self
}
pub fn with_repository_instructions(mut self, instructions: Option<String>) -> Self {
self.repository_instructions = instructions;
self
}
pub fn with_terminal_prompt(mut self, prompt: Option<String>) -> Self {
self.terminal_prompt = prompt;
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IssueSessionPromptKind {
Full,
Continuation,
}
impl IssueSessionPromptKind {
pub fn as_str(self) -> &'static str {
match self {
Self::Full => "full",
Self::Continuation => "continuation",
}
}
fn artifact_name(self) -> &'static str {
match self {
Self::Full => "last-full-prompt.md",
Self::Continuation => "last-continuation-prompt.md",
}
}
fn artifact_path(self, workspace: &WorkspaceHandle) -> PathBuf {
workspace.prompts_dir().join(self.artifact_name())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConversationLaunchProfile {
pub workspace_kind: String,
pub confirmation_policy_kind: String,
pub agent_kind: String,
pub llm_model: String,
#[serde(default = "default_llm_credential_mode")]
pub llm_credential_mode: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_api_key_env: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_base_url_env: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_subscription: Option<ConversationLaunchSubscriptionProfile>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub condenser: Option<ConversationLaunchCondenserProfile>,
pub agent_tools: Option<Vec<ToolConfig>>,
pub agent_include_default_tools: Option<Vec<String>>,
pub max_iterations: u32,
pub stuck_detection: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_api_key_fingerprint: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConversationLaunchSubscriptionProfile {
pub vendor: String,
pub access_token_env: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub account_id_env: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_directory_env: Option<String>,
pub auth_method: String,
pub open_browser: bool,
pub force_login: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConversationLaunchCondenserProfile {
pub max_size: u64,
pub keep_first: u64,
}
impl ConversationLaunchProfile {
pub fn api_key_fingerprint(&self, env: &dyn Environment) -> Option<String> {
let api_key = if self.llm_credential_mode == OPENAI_SUBSCRIPTION_CREDENTIAL_MODE {
self.llm_subscription
.as_ref()
.and_then(|subscription| env.get(&subscription.access_token_env))
} else {
self.llm_api_key_env
.as_deref()
.and_then(|env_name| env.get(env_name))
.or_else(|| env.get("LLM_API_KEY"))
}?;
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(api_key.as_bytes());
let result = hasher.finalize();
Some(result[..8].iter().map(|b| format!("{:02x}", b)).collect())
}
pub fn from_workflow(workflow: &ResolvedWorkflow) -> Result<Self, String> {
let conversation = &workflow.extensions.openhands.conversation;
let max_iterations = u32::try_from(conversation.max_iterations).map_err(|_| {
format!(
"workflow max_iterations {} exceeds u32::MAX ({})",
conversation.max_iterations,
u32::MAX
)
})?;
let llm_model = conversation
.agent
.llm
.as_ref()
.and_then(|llm| llm.model.as_ref())
.cloned()
.ok_or_else(|| {
"workflow openhands.conversation.agent.llm.model is required".to_string()
})?;
Ok(Self {
workspace_kind: "LocalWorkspace".to_string(),
confirmation_policy_kind: conversation.confirmation_policy.kind.clone(),
agent_kind: conversation.agent.kind.clone(),
llm_model,
llm_credential_mode: conversation
.agent
.llm
.as_ref()
.map(|llm| llm.credential_mode.clone())
.unwrap_or_else(default_llm_credential_mode),
llm_api_key_env: conversation
.agent
.llm
.as_ref()
.and_then(|llm| llm.api_key_env.clone()),
llm_base_url_env: conversation
.agent
.llm
.as_ref()
.and_then(|llm| llm.base_url_env.clone()),
llm_subscription: conversation
.agent
.llm
.as_ref()
.and_then(|llm| llm.subscription.as_ref())
.map(|subscription| ConversationLaunchSubscriptionProfile {
vendor: subscription.vendor.clone(),
access_token_env: subscription.access_token_env.clone(),
account_id_env: subscription.account_id_env.clone(),
auth_directory_env: subscription.auth_directory_env.clone(),
auth_method: subscription.auth_method.clone(),
open_browser: subscription.open_browser,
force_login: subscription.force_login,
}),
condenser: conversation.agent.condenser.as_ref().map(|condenser| {
ConversationLaunchCondenserProfile {
max_size: condenser.max_size,
keep_first: condenser.keep_first,
}
}),
agent_tools: conversation
.agent
.tools
.as_ref()
.map(|tools| tools.iter().map(tool_config_from_workflow).collect()),
agent_include_default_tools: conversation.agent.include_default_tools.clone(),
max_iterations,
stuck_detection: conversation.stuck_detection,
llm_api_key_fingerprint: None, })
}
pub fn to_create_request(
&self,
env: &dyn Environment,
working_dir: &Path,
persistence_dir: &Path,
conversation_id: Option<Uuid>,
) -> Result<ConversationCreateRequest, String> {
let llm = self.to_llm_config(env)?;
Ok(ConversationCreateRequest {
conversation_id: conversation_id.unwrap_or_else(Uuid::new_v4),
workspace: WorkspaceConfig {
working_dir: working_dir.display().to_string(),
kind: self.workspace_kind.clone(),
},
persistence_dir: persistence_dir.display().to_string(),
max_iterations: self.max_iterations,
stuck_detection: self.stuck_detection,
confirmation_policy: ConfirmationPolicy {
kind: self.confirmation_policy_kind.clone(),
},
agent: AgentConfig {
kind: self.agent_kind.clone(),
llm: llm.clone(),
condenser: self.condenser.as_ref().map(|condenser| {
CondenserConfig::llm_summarizing(
llm.clone(),
condenser.max_size,
condenser.keep_first,
)
}),
tools: self.agent_tools.clone(),
mcp_config: None,
include_default_tools: self.agent_include_default_tools.clone(),
},
})
}
fn to_llm_config(&self, env: &dyn Environment) -> Result<LlmConfig, String> {
if self.llm_credential_mode == OPENAI_SUBSCRIPTION_CREDENTIAL_MODE {
return self.to_openai_subscription_llm_config(env);
}
let api_key = resolve_provider_override(
env,
"openhands.conversation.agent.llm.api_key_env",
self.llm_api_key_env.as_deref(),
)?
.or_else(|| normalize_environment_value(env.get("LLM_API_KEY")));
let base_url = resolve_provider_override(
env,
"openhands.conversation.agent.llm.base_url_env",
self.llm_base_url_env.as_deref(),
)?
.or_else(|| normalize_environment_value(env.get("LLM_BASE_URL")));
Ok(LlmConfig {
model: self.llm_model.clone(),
api_key,
base_url,
usage_id: None,
extra_headers: None,
litellm_extra_body: None,
stream: None,
})
}
fn to_openai_subscription_llm_config(
&self,
env: &dyn Environment,
) -> Result<LlmConfig, String> {
let subscription = self.llm_subscription.as_ref().ok_or_else(|| {
"openhands.conversation.agent.llm.subscription is required for openai_subscription mode"
.to_string()
})?;
if subscription.vendor != "openai" {
return Err(format!(
"unsupported OpenHands subscription vendor `{}`",
subscription.vendor
));
}
let access_token = resolve_provider_override(
env,
"openhands.conversation.agent.llm.subscription.access_token_env",
Some(&subscription.access_token_env),
)?;
let account_id = subscription
.account_id_env
.as_deref()
.and_then(|env_name| normalize_environment_value(env.get(env_name)));
let mut extra_headers = BTreeMap::from([
(
"OpenAI-Beta".to_string(),
"responses=experimental".to_string(),
),
(
"User-Agent".to_string(),
format!(
"openhands-sdk (OpenSymphony; {}; {})",
std::env::consts::OS,
std::env::consts::ARCH
),
),
("originator".to_string(), "codex_cli_rs".to_string()),
]);
if let Some(account_id) = account_id {
extra_headers.insert("chatgpt-account-id".to_string(), account_id);
}
Ok(LlmConfig {
model: normalize_openai_subscription_model(&self.llm_model)?,
api_key: access_token,
base_url: Some(OPENAI_CODEX_SUBSCRIPTION_BASE_URL.to_string()),
usage_id: None,
extra_headers: Some(extra_headers),
litellm_extra_body: Some(BTreeMap::from([("store".to_string(), json!(false))])),
stream: Some(true),
})
}
}
fn default_llm_credential_mode() -> String {
"api_key".to_string()
}
fn normalize_openai_subscription_model(model: &str) -> Result<String, String> {
let normalized = model.trim();
if normalized.is_empty() {
return Err("OpenHands SDK OpenAI subscription model must not be blank".to_string());
}
if let Some((provider, bare)) = normalized.split_once('/') {
if provider != "openai" {
return Err(format!(
"OpenHands SDK OpenAI subscription model `{normalized}` must use the openai/ provider prefix"
));
}
if bare.trim().is_empty() {
return Err("OpenHands SDK OpenAI subscription model must not be blank".to_string());
}
return Ok(normalized.to_string());
}
Ok(format!("openai/{normalized}"))
}
fn resolve_provider_override(
env: &dyn Environment,
field: &'static str,
env_name: Option<&str>,
) -> Result<Option<String>, String> {
let Some(env_name) = env_name else {
return Ok(None);
};
normalize_environment_value(env.get(env_name)).ok_or_else(|| {
format!(
"{field} references environment variable `{env_name}`, but it is not set or is blank"
)
})
.map(Some)
}
fn normalize_environment_value(value: Option<String>) -> Option<String> {
value
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
}
fn tool_config_from_workflow(tool: &OpenHandsConversationToolConfig) -> ToolConfig {
ToolConfig {
name: tool.name.clone(),
params: tool.params.clone(),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LlmConfigFingerprint {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub api_key_hash: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub base_url_hash: Option<String>,
pub model: String,
}
impl LlmConfigFingerprint {
pub fn from_llm_config(llm: &LlmConfig) -> Self {
Self {
api_key_hash: None,
base_url_hash: None,
model: llm.model.clone(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IssueConversationManifest {
pub issue_id: IssueId,
pub identifier: IssueIdentifier,
pub conversation_id: ConversationId,
#[serde(default = "default_reuse_policy")]
pub reuse_policy: String,
pub server_base_url: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transport_target: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub http_auth_mode: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub websocket_auth_mode: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub websocket_query_param_name: Option<String>,
pub persistence_dir: PathBuf,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
pub last_attached_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub launch_profile: Option<ConversationLaunchProfile>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_config_fingerprint: Option<LlmConfigFingerprint>,
pub fresh_conversation: bool,
pub workflow_prompt_seeded: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reset_reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_contract_version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub runtime_envelope: Option<TerminalRuntimeEnvelope>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_runtime_envelope: Option<ParentRuntimeEnvelope>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub codex_archive_state: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_turn_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub active_run_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prepared_run_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub trigger_pending_run_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_prompt_kind: Option<IssueSessionPromptKind>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_prompt_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_prompt_path: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_execution_status: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_kind: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_summary: Option<String>,
#[serde(default)]
pub input_tokens: u64,
#[serde(default)]
pub output_tokens: u64,
#[serde(default)]
pub cache_read_tokens: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_token_accumulation_at: Option<DateTime<Utc>>,
}
impl IssueConversationManifest {
#[allow(clippy::too_many_arguments)]
fn new(
issue_id: IssueId,
identifier: IssueIdentifier,
conversation_id: ConversationId,
reuse_policy: impl Into<String>,
persistence_dir: PathBuf,
attached_at: DateTime<Utc>,
reset_reason: Option<String>,
mut launch_profile: ConversationLaunchProfile,
env: &dyn Environment,
) -> Self {
launch_profile.llm_api_key_fingerprint = launch_profile.api_key_fingerprint(env);
Self {
issue_id,
identifier,
conversation_id,
reuse_policy: reuse_policy.into(),
server_base_url: None,
transport_target: None,
http_auth_mode: None,
websocket_auth_mode: None,
websocket_query_param_name: None,
persistence_dir,
created_at: attached_at,
updated_at: attached_at,
last_attached_at: attached_at,
launch_profile: Some(launch_profile),
llm_config_fingerprint: None,
fresh_conversation: true,
workflow_prompt_seeded: false,
reset_reason,
runtime_contract_version: Some(RUNTIME_CONTRACT_VERSION.to_string()),
runtime_envelope: None,
parent_runtime_envelope: None,
codex_archive_state: None,
last_turn_id: None,
active_run_id: None,
prepared_run_id: None,
trigger_pending_run_id: None,
last_prompt_kind: None,
last_prompt_at: None,
last_prompt_path: None,
last_execution_status: None,
last_event_id: None,
last_event_kind: None,
last_event_at: None,
last_event_summary: None,
input_tokens: 0,
output_tokens: 0,
cache_read_tokens: 0,
last_token_accumulation_at: None,
}
}
fn prompt_kind(&self) -> IssueSessionPromptKind {
if self.workflow_prompt_seeded {
IssueSessionPromptKind::Continuation
} else {
IssueSessionPromptKind::Full
}
}
fn is_reusable_for(
&self,
issue: &NormalizedIssue,
expected_persistence_dir: &Path,
expected_reuse_policy: &str,
) -> bool {
self.issue_id == issue.id
&& self.identifier == issue.identifier
&& self.reuse_policy == expected_reuse_policy
&& self.persistence_dir == expected_persistence_dir
&& self.runtime_contract_version.as_deref() == Some(RUNTIME_CONTRACT_VERSION)
}
fn record_prompt(
&mut self,
prompt_kind: IssueSessionPromptKind,
prompt_path: PathBuf,
recorded_at: DateTime<Utc>,
) {
self.last_prompt_kind = Some(prompt_kind);
self.last_prompt_at = Some(recorded_at);
self.last_prompt_path = Some(prompt_path);
self.updated_at = recorded_at;
}
fn apply_runtime_snapshot(&mut self, stream: &RuntimeEventStream) {
self.last_execution_status = stream
.state_mirror()
.execution_status()
.map(ToOwned::to_owned);
if let Some(event) = stream.event_cache().items().last() {
self.last_event_id = Some(event.id.clone());
self.last_event_kind = Some(event.kind.clone());
self.last_event_at = Some(event.timestamp);
self.last_event_summary = Some(summarize_event(event));
}
if let Some((input, output, cache_read)) = stream.state_mirror().accumulated_token_usage() {
if input > self.input_tokens {
self.input_tokens = input;
}
if output > self.output_tokens {
self.output_tokens = output;
}
if cache_read > self.cache_read_tokens {
self.cache_read_tokens = cache_read;
}
}
self.updated_at = Utc::now();
}
fn apply_transport_diagnostics(
&mut self,
diagnostics: Option<&super::TransportDiagnostics>,
server_base_url: &str,
) {
self.server_base_url = Some(server_base_url.to_string());
self.transport_target =
diagnostics.map(|diagnostics| diagnostics.target_kind.as_str().to_string());
self.http_auth_mode =
diagnostics.map(|diagnostics| diagnostics.http_auth_kind.as_str().to_string());
self.websocket_auth_mode =
diagnostics.map(|diagnostics| diagnostics.websocket_auth_kind.as_str().to_string());
self.websocket_query_param_name =
diagnostics.and_then(|diagnostics| diagnostics.websocket_query_param_name.clone());
}
fn to_domain_metadata(&self, stream_state: RuntimeStreamState) -> ConversationMetadata {
ConversationMetadata {
harness_capability: None,
conversation_id: self.conversation_id.clone(),
server_base_url: self.server_base_url.clone(),
transport_target: self.transport_target.clone(),
http_auth_mode: self.http_auth_mode.clone(),
websocket_auth_mode: self.websocket_auth_mode.clone(),
websocket_query_param_name: self.websocket_query_param_name.clone(),
fresh_conversation: self.fresh_conversation,
runtime_contract_version: self.runtime_contract_version.clone(),
stream_state,
last_event_id: self.last_event_id.clone(),
last_event_kind: self.last_event_kind.clone(),
last_event_at: self.last_event_at.map(timestamp_ms_from_datetime),
last_event_summary: self.last_event_summary.clone(),
recent_activity: Vec::new(),
input_tokens: self.input_tokens,
output_tokens: self.output_tokens,
cache_read_tokens: self.cache_read_tokens,
total_tokens: self.input_tokens + self.output_tokens,
runtime_seconds: 0, next_activity_sequence: 0,
}
}
}
fn bind_runtime_envelopes(
manifest: &mut IssueConversationManifest,
run_manifest: &mut RunManifest,
) {
manifest.runtime_envelope = run_manifest.runtime_envelope.clone();
if let Some(envelope) = manifest.runtime_envelope.as_mut() {
envelope.conversation_binding = Some(manifest.conversation_id.to_string());
}
run_manifest.runtime_envelope = manifest.runtime_envelope.clone();
manifest.parent_runtime_envelope = run_manifest.parent_runtime_envelope.clone();
if let Some(envelope) = manifest.parent_runtime_envelope.as_mut() {
envelope.conversation_binding = Some(manifest.conversation_id.to_string());
}
run_manifest.parent_runtime_envelope = manifest.parent_runtime_envelope.clone();
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IssueSessionContext {
pub run_id: String,
pub issue_id: IssueId,
pub identifier: IssueIdentifier,
pub worker_id: WorkerId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempt: Option<u32>,
pub normal_retry_count: u32,
pub turn_count: u32,
pub max_turns: u32,
pub prompt_kind: IssueSessionPromptKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt_path: Option<PathBuf>,
pub conversation_id: ConversationId,
#[serde(default = "default_reuse_policy")]
pub reuse_policy: String,
pub fresh_conversation: bool,
pub workflow_prompt_seeded: bool,
pub server_base_url: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transport_target: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub http_auth_mode: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub websocket_auth_mode: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub websocket_query_param_name: Option<String>,
pub persistence_dir: PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_execution_status: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_kind: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_event_summary: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub worker_outcome: Option<WorkerOutcomeRecord>,
pub updated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IssueSessionResult {
pub prompt_kind: IssueSessionPromptKind,
pub conversation: Option<ConversationMetadata>,
pub worker_outcome: WorkerOutcomeRecord,
pub run_status: RunStatus,
}
#[derive(Debug, Error)]
pub enum IssueSessionError {
#[error(transparent)]
Workspace(#[from] WorkspaceError),
#[error("unexpected early result: {0}")]
UnexpectedEarlyResult(String),
#[error("rehydration failed: {0}")]
RehydrationFailed(String),
#[error("conversation retirement failed: {0}")]
ConversationRetirementFailed(String),
#[error("superseded conversation evidence is invalid: {0}")]
SupersededConversationEvidence(String),
}
fn decode_superseded_conversation_manifests(
raw: &str,
) -> Result<Option<Vec<IssueConversationManifest>>, IssueSessionError> {
serde_json::from_str(raw)
.map_err(|error| IssueSessionError::SupersededConversationEvidence(error.to_string()))
}
#[derive(Debug, Clone)]
struct NormalizedOutcome {
kind: WorkerOutcomeKind,
summary: String,
error: Option<String>,
}
#[derive(Debug, Clone)]
enum StateCheckResult {
Terminal(NormalizedOutcome),
StillRunningWithProgress,
NoProgress,
}
enum ReconcileResult {
Terminal(NormalizedOutcome),
Progress,
NoProgress,
}
pub struct ActiveSession {
stream: RuntimeEventStream,
manifest: IssueConversationManifest,
prompt_kind: IssueSessionPromptKind,
prompt_path: Option<PathBuf>,
}
impl ActiveSession {
fn accumulate_tokens(&mut self) {
use super::events::KnownEvent;
use chrono::DateTime;
let cutoff = self.manifest.last_token_accumulation_at;
let mut max_event_time: Option<DateTime<Utc>> = None;
let mut events_scanned = 0;
let mut llm_events_found = 0;
let mut tokens_added = (0u64, 0u64);
for event in self.stream.event_cache().items() {
events_scanned += 1;
if let Some(cutoff_time) = cutoff
&& event.timestamp <= cutoff_time
{
continue;
}
if let KnownEvent::LlmCompletionLog(llm_event) = KnownEvent::from_envelope(event) {
llm_events_found += 1;
if let Some((input, output)) = llm_event.token_usage() {
self.manifest.input_tokens += input;
self.manifest.output_tokens += output;
tokens_added.0 += input;
tokens_added.1 += output;
tracing::debug!(
input_tokens = input,
output_tokens = output,
event_id = %event.id,
"accumulated tokens from LLM completion event"
);
} else {
tracing::debug!(
event_id = %event.id,
payload = %llm_event.payload,
"LLMCompletionLogEvent has no token usage data"
);
}
}
if max_event_time.is_none_or(|current| event.timestamp > current) {
max_event_time = Some(event.timestamp);
}
}
if let Some(latest_time) = max_event_time {
self.manifest.last_token_accumulation_at = Some(latest_time);
}
if let Some((input, output, cache_read)) =
self.stream.state_mirror().accumulated_token_usage()
{
if input > self.manifest.input_tokens {
let new_input = input - self.manifest.input_tokens;
self.manifest.input_tokens = input;
tokens_added.0 += new_input;
}
if output > self.manifest.output_tokens {
let new_output = output - self.manifest.output_tokens;
self.manifest.output_tokens = output;
tokens_added.1 += new_output;
}
if cache_read > self.manifest.cache_read_tokens {
self.manifest.cache_read_tokens = cache_read;
}
tracing::debug!(
state_input = input,
state_output = output,
state_cache_read = cache_read,
manifest_input = self.manifest.input_tokens,
manifest_output = self.manifest.output_tokens,
manifest_cache_read = self.manifest.cache_read_tokens,
"updated tokens from conversation state"
);
}
tracing::debug!(
events_scanned,
llm_events_found,
input_tokens_added = tokens_added.0,
output_tokens_added = tokens_added.1,
total_input_tokens = self.manifest.input_tokens,
total_output_tokens = self.manifest.output_tokens,
total_cache_read_tokens = self.manifest.cache_read_tokens,
"token accumulation complete"
);
}
}
enum Step<T> {
Continue(T),
EarlyResult(Box<IssueSessionResult>),
}
enum ReuseSession {
Active(Box<ActiveSession>),
Reset {
reason: String,
manifest: Box<IssueConversationManifest>,
},
}
struct PreparedTurn {
conversation_id: Uuid,
prompt: String,
baseline_event_ids: HashSet<String>,
launch_reported: bool,
waited_for_prior_turn: bool,
}
#[derive(Default)]
struct LoadedManifest {
manifest: Option<IssueConversationManifest>,
reset_reason: Option<String>,
}
pub struct IssueSessionRunner {
client: OpenHandsClient,
config: IssueSessionRunnerConfig,
environment: Arc<dyn Environment + Send + Sync>,
workpad_comment_source: Option<Arc<dyn WorkpadCommentSource>>,
}
impl IssueSessionRunner {
pub fn new(client: OpenHandsClient, config: IssueSessionRunnerConfig) -> Self {
Self::with_environment(client, config, ProcessEnvironment)
}
pub fn with_environment<E>(
client: OpenHandsClient,
config: IssueSessionRunnerConfig,
environment: E,
) -> Self
where
E: Environment + Send + Sync + 'static,
{
Self {
client,
config,
environment: Arc::new(environment),
workpad_comment_source: None,
}
}
pub fn with_workpad_comment_source(
mut self,
workpad_comment_source: Arc<dyn WorkpadCommentSource>,
) -> Self {
self.workpad_comment_source = Some(workpad_comment_source);
self
}
pub fn with_repository_instructions(mut self, instructions: Option<String>) -> Self {
self.config.repository_instructions = instructions;
self
}
pub fn with_terminal_prompt(mut self, prompt: Option<String>) -> Self {
self.config.terminal_prompt = prompt;
self
}
pub fn with_continuation_prompt(mut self, prompt: Option<String>) -> Self {
self.config.continuation_prompt = prompt;
self
}
pub fn client(&self) -> &OpenHandsClient {
&self.client
}
pub fn config(&self) -> &IssueSessionRunnerConfig {
&self.config
}
pub async fn interrupt(
&self,
command: &HarnessInterruptCommand,
) -> Result<OpenHandsInterruptAcknowledgement, OpenHandsError> {
let conversation_id =
Uuid::parse_str(command.conversation_id.as_str()).map_err(|error| {
OpenHandsError::InvalidConfiguration {
detail: format!(
"interrupt command conversation id `{}` is not a UUID: {error}",
command.conversation_id
),
}
})?;
let (method, diagnostic) = match self.client.interrupt_conversation(conversation_id).await {
Ok(_) => (OpenHandsInterruptMethod::Interrupt, None),
Err(OpenHandsError::HttpStatus {
status_code, body, ..
}) if matches!(status_code, 404 | 405 | 501) => {
self.client.pause_conversation(conversation_id).await?;
(
OpenHandsInterruptMethod::PauseFallback,
Some(format!(
"OpenHands /interrupt unavailable with HTTP {status_code}; fell back to /pause: {body}"
)),
)
}
Err(error) => return Err(error),
};
let mut stream = self
.client
.attach_runtime_stream_with_recent_events(
conversation_id,
self.config.runtime_stream.clone(),
INTERRUPT_RECENT_EVENT_LIMIT,
)
.await?;
let deadline = Instant::now() + self.config.runtime_stream.readiness_timeout;
let mut reconciled_events = stream
.reconcile_recent_events(INTERRUPT_RECENT_EVENT_LIMIT)
.await?;
loop {
if stream
.state_mirror()
.execution_status()
.is_some_and(interrupt_acknowledgement_observed)
{
return Ok(OpenHandsInterruptAcknowledgement {
method,
reconciled_events,
execution_status: stream.state_mirror().execution_status().map(str::to_owned),
diagnostic,
timed_out: false,
});
}
let wait = timeout_at(deadline, stream.next_event()).await;
let reconcile_after_event = match wait {
Ok(Ok(Some(_))) => true,
Ok(Ok(None)) | Ok(Err(_)) => {
reconciled_events += stream
.reconcile_recent_events(INTERRUPT_RECENT_EVENT_LIMIT)
.await?;
false
}
Err(_) => {
return Ok(OpenHandsInterruptAcknowledgement {
method,
reconciled_events,
execution_status: stream
.state_mirror()
.execution_status()
.map(str::to_owned),
diagnostic: Some(match diagnostic {
Some(diagnostic) => format!(
"{diagnostic}; timed out waiting for interrupt acknowledgement after {:?}",
self.config.runtime_stream.readiness_timeout
),
None => format!(
"timed out waiting for interrupt acknowledgement after {:?}",
self.config.runtime_stream.readiness_timeout
),
}),
timed_out: true,
});
}
};
if reconcile_after_event {
reconciled_events += stream
.reconcile_recent_events(INTERRUPT_RECENT_EVENT_LIMIT)
.await?;
}
}
}
pub async fn run(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
issue: &NormalizedIssue,
run: &RunAttempt,
workflow: &ResolvedWorkflow,
) -> Result<IssueSessionResult, IssueSessionError> {
self.run_with_observer(
workspace_manager,
workspace,
run_manifest,
issue,
run,
workflow,
&mut (),
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn recover_with_observer<O>(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
issue: &NormalizedIssue,
run: &RunAttempt,
observer: &mut O,
) -> Result<IssueSessionResult, IssueSessionError>
where
O: IssueSessionObserver,
{
let observed_run = observed_run_for_turn(run);
let raw_manifest = workspace_manager
.read_text_artifact(workspace, &workspace.conversation_manifest_path())
.await?
.ok_or_else(|| {
IssueSessionError::RehydrationFailed(
"recovered run has no conversation manifest to reattach".to_string(),
)
})?;
let manifest =
serde_json::from_str::<IssueConversationManifest>(&raw_manifest).map_err(|error| {
IssueSessionError::RehydrationFailed(format!(
"recovered conversation manifest is invalid: {error}"
))
})?;
let trigger_pending = manifest
.trigger_pending_run_id
.as_deref()
.is_some_and(|run_id| run_id == run_manifest.run_id);
let baseline_last_event_id = manifest.last_event_id.clone();
if manifest.issue_id != issue.id
|| manifest.identifier != issue.identifier
|| manifest.runtime_contract_version.as_deref() != Some(RUNTIME_CONTRACT_VERSION)
{
return Err(IssueSessionError::RehydrationFailed(
"recovered conversation manifest is incompatible with the active issue or runtime"
.to_string(),
));
}
let conversation_id = parse_uuid(manifest.conversation_id.as_str())
.map_err(IssueSessionError::RehydrationFailed)?;
let recovery_client = match manifest.server_base_url.as_deref() {
Some(base_url) => self
.client
.with_persisted_base_url(base_url)
.map_err(|error| {
IssueSessionError::RehydrationFailed(format!(
"recovered conversation server identity is not trusted: {error}"
))
})?,
None => self.client.clone(),
};
let active_session = match self
.try_attach_and_resume(
&recovery_client,
workspace_manager,
workspace,
manifest,
conversation_id,
)
.await?
{
ReuseSession::Active(session) => *session,
ReuseSession::Reset { reason, .. } => {
return Err(IssueSessionError::RehydrationFailed(reason));
}
};
let mut active_session = active_session;
let mut pre_trigger_baseline_event_ids =
trigger_pending.then(|| all_event_ids(active_session.stream.event_cache().items()));
if trigger_pending {
match recovery_client.run_conversation(conversation_id).await {
Ok(_) => {}
Err(OpenHandsError::HttpStatus {
status_code: 409, ..
}) => {
if let Err(error) = active_session.stream.reconcile_events().await {
debug!(
%error,
conversation_id = %active_session.manifest.conversation_id,
"pre-wait reconcile failed after recovered run retry conflict, proceeding anyway"
);
}
if let Err(error) = self
.wait_for_active_turn_to_finish(&mut active_session.stream)
.await
{
return Err(IssueSessionError::RehydrationFailed(format!(
"previous OpenHands turn did not finish before recovered run retry: {error}"
)));
}
pre_trigger_baseline_event_ids =
Some(all_event_ids(active_session.stream.event_cache().items()));
if let Err(error) = recovery_client.run_conversation(conversation_id).await {
return Err(IssueSessionError::RehydrationFailed(format!(
"failed to retry recovered OpenHands run after conflict: {error}"
)));
}
}
Err(error) => {
return Err(IssueSessionError::RehydrationFailed(format!(
"failed to trigger recovered OpenHands run: {error}"
)));
}
}
active_session.manifest.trigger_pending_run_id = None;
active_session.manifest.updated_at = Utc::now();
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&active_session.manifest,
)
.await?;
run_manifest.status = RunStatus::Running;
run_manifest.started_at.get_or_insert_with(Utc::now);
run_manifest.status_detail = Some(format!(
"recovered {} prompt trigger for conversation {}",
active_session.prompt_kind.as_str(),
active_session.manifest.conversation_id
));
workspace_manager
.write_run_manifest(workspace, run_manifest)
.await?;
}
observer.on_launch_with_started_at(
&active_session
.manifest
.to_domain_metadata(RuntimeStreamState::Ready),
run_manifest.started_at.map(timestamp_ms_from_datetime),
);
let baseline_event_ids = pre_trigger_baseline_event_ids.unwrap_or_else(|| {
recovery_baseline_event_ids(
active_session.stream.event_cache().items(),
baseline_last_event_id.as_deref(),
)
});
let outcome = self
.await_terminal_outcome(&mut active_session, &baseline_event_ids, observer)
.await;
self.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
&observed_run,
active_session,
outcome,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn run_with_observer<O>(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
issue: &NormalizedIssue,
run: &RunAttempt,
workflow: &ResolvedWorkflow,
observer: &mut O,
) -> Result<IssueSessionResult, IssueSessionError>
where
O: IssueSessionObserver,
{
const MAX_CONTEXT_OVERFLOW_RETRIES: usize = 1;
let observed_run = observed_run_for_turn(run);
let active_session = match Box::pin(self.initialize_session(
workspace_manager,
workspace,
run_manifest,
&observed_run,
issue,
workflow,
))
.await?
{
Step::Continue(session) => session,
Step::EarlyResult(result) => return Ok(*result),
};
let mut launch_reported = false;
if !active_session.manifest.fresh_conversation
&& active_session
.stream
.state_mirror()
.execution_status()
.is_some_and(turn_is_in_progress)
{
observer.on_launch_with_started_at(
&active_session
.manifest
.to_domain_metadata(RuntimeStreamState::Ready),
run_manifest.started_at.map(timestamp_ms_from_datetime),
);
launch_reported = true;
}
let mut retry_count = 0;
let mut current_session = active_session;
let (final_session, outcome) = loop {
let (mut active_session, outcome, turn_launch_reported) =
match Box::pin(self.execute_turn(
workspace_manager,
workspace,
run_manifest,
&observed_run,
workflow,
issue,
run,
current_session,
launch_reported,
observer,
))
.await?
{
Step::Continue(result) => result,
Step::EarlyResult(result) => return Ok(*result),
};
launch_reported = turn_launch_reported;
if is_context_overflow_outcome(&outcome) && retry_count < MAX_CONTEXT_OVERFLOW_RETRIES {
retry_count += 1;
tracing::warn!(
conversation_id = %active_session.manifest.conversation_id,
error = ?outcome.error,
retry_count,
max_retries = MAX_CONTEXT_OVERFLOW_RETRIES,
"context overflow detected, attempting auto-recovery via rehydration"
);
let old_manifest = active_session.manifest.clone();
let _ = active_session.stream.close().await;
match self
.recover_from_context_overflow(
workspace_manager,
workspace,
run_manifest,
&observed_run,
issue,
workflow,
&old_manifest,
)
.await
{
Ok(recovered_session) => {
tracing::info!(
old_conversation_id = %old_manifest.conversation_id,
new_conversation_id = %recovered_session.manifest.conversation_id,
"context overflow recovery: rehydration succeeded, re-running turn"
);
current_session = recovered_session;
continue;
}
Err(rehydration_error) => {
tracing::warn!(
%rehydration_error,
"context overflow recovery: rehydration failed, returning original failure"
);
break (active_session, outcome);
}
}
} else if is_condenser_tool_matching_outcome(&outcome)
&& retry_count < MAX_CONTEXT_OVERFLOW_RETRIES
{
retry_count += 1;
tracing::warn!(
conversation_id = %active_session.manifest.conversation_id,
error = ?outcome.error,
retry_count,
max_retries = MAX_CONTEXT_OVERFLOW_RETRIES,
"condenser tool-matching error detected (corrupted event history), creating fresh conversation"
);
let _ = active_session.stream.close().await;
if let Err(preserve_error) = self
.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
&active_session.manifest,
)
.await
{
tracing::warn!(
%preserve_error,
conversation_id = %active_session.manifest.conversation_id,
"condenser tool-matching recovery could not preserve the old conversation before replacement"
);
break (active_session, outcome);
}
let replacement = self
.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
&observed_run,
issue,
workflow,
Some("condenser tool-matching error auto-recovery".to_string()),
)
.await;
match replacement {
Ok(Step::Continue(fresh_session)) => {
let retired = match self
.retire_conversation(
&active_session.manifest,
"condenser tool-matching recovery",
)
.await
{
Ok(()) => true,
Err(retirement_error) => {
tracing::warn!(
%retirement_error,
"condenser tool-matching recovery could not retire the old conversation after replacement"
);
false
}
};
if retired
&& let Err(clear_error) = self
.clear_superseded_conversation_manifest(
workspace_manager,
workspace,
&active_session.manifest.conversation_id,
)
.await
{
tracing::warn!(
%clear_error,
conversation_id = %active_session.manifest.conversation_id,
"condenser tool-matching recovery could not clear retired conversation evidence"
);
}
tracing::info!(
old_conversation_id = %active_session.manifest.conversation_id,
new_conversation_id = %fresh_session.manifest.conversation_id,
"condenser tool-matching recovery: fresh conversation created, re-running turn"
);
current_session = fresh_session;
continue;
}
Ok(Step::EarlyResult(result)) => {
tracing::warn!(
"condenser tool-matching recovery: failed to create fresh session, returning early result"
);
return Ok(*result);
}
Err(fresh_error) => {
tracing::warn!(
%fresh_error,
"condenser tool-matching recovery: failed to create fresh session, returning original failure"
);
break (active_session, outcome);
}
}
} else {
break (active_session, outcome);
}
};
self.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
&observed_run,
final_session,
outcome,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn execute_turn<O>(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
workflow: &ResolvedWorkflow,
issue: &NormalizedIssue,
run: &RunAttempt,
active_session: ActiveSession,
launch_reported: bool,
observer: &mut O,
) -> Result<Step<(ActiveSession, NormalizedOutcome, bool)>, IssueSessionError>
where
O: IssueSessionObserver,
{
let (mut active_session, mut prepared_turn) = match Box::pin(self.prepare_turn(
workspace_manager,
workspace,
run_manifest,
observed_run,
workflow,
issue,
run,
active_session,
launch_reported,
observer,
))
.await?
{
Step::Continue(state) => state,
Step::EarlyResult(result) => return Ok(Step::EarlyResult(result)),
};
active_session = match Box::pin(self.start_turn(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
&mut prepared_turn,
observer,
))
.await?
{
Step::Continue(session) => session,
Step::EarlyResult(result) => return Ok(Step::EarlyResult(result)),
};
let outcome = Box::pin(self.await_terminal_outcome(
&mut active_session,
&prepared_turn.baseline_event_ids,
observer,
))
.await;
Ok(Step::Continue((
active_session,
outcome,
prepared_turn.launch_reported,
)))
}
#[allow(clippy::too_many_arguments)]
async fn recover_from_context_overflow(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
old_manifest: &IssueConversationManifest,
) -> Result<ActiveSession, IssueSessionError> {
let rehydration_result = self
.rehydrate_conversation(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
old_manifest,
RehydrationOptions {
reason: "context overflow auto-recovery".to_string(),
summarize: false,
max_summary_events: 0,
},
)
.await?;
Ok(rehydration_result.session)
}
async fn initialize_session(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
) -> Result<Step<ActiveSession>, IssueSessionError> {
match &self.config.reuse_policy {
IssueSessionReusePolicy::PerIssue => {
let loaded = self
.load_existing_conversation_manifest(workspace_manager, workspace, issue, workflow)
.await?;
match loaded.manifest {
Some(manifest)
if run_manifest.runtime_envelope.is_some()
&& manifest.runtime_envelope.as_ref().is_none_or(|envelope| {
envelope.conversation_binding.as_deref()
!= Some(manifest.conversation_id.as_str())
}) => {
self.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
&manifest,
)
.await?;
tracing::warn!(
conversation_id = %manifest.conversation_id,
"skipping retirement of conversation with mismatched runtime envelope binding"
);
let replacement = self.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
Some("conversation runtime envelope binding changed; superseding conversation".into()),
)
.await?;
Ok(replacement)
}
Some(manifest)
if run_manifest.runtime_envelope.as_ref().is_some_and(|expected| {
manifest.runtime_envelope.as_ref().is_some_and(|actual| {
actual != expected
})
}) => {
self.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
&manifest,
)
.await?;
tracing::warn!(
conversation_id = %manifest.conversation_id,
"skipping retirement of conversation with an untrusted runtime envelope"
);
let replacement = self.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
Some("terminal runtime envelope changed; superseding conversation".into()),
)
.await?;
Ok(replacement)
}
Some(manifest) => match self
.try_reuse_session(workspace_manager, workspace, issue, workflow, manifest)
.await?
{
ReuseSession::Active(session) => Ok(Step::Continue(*session)),
ReuseSession::Reset { reason, manifest } => {
self.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
&manifest,
)
.await?;
let replacement = self.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
Some(reason),
)
.await?;
self.retire_replaced_conversation(
workspace_manager,
workspace,
&manifest,
&replacement,
)
.await;
Ok(replacement)
}
},
None => {
self.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
loaded.reset_reason,
)
.await
}
}
}
IssueSessionReusePolicy::FreshEachRun => {
let loaded = self
.load_existing_conversation_manifest(workspace_manager, workspace, issue, workflow)
.await?;
let safe_to_retire = loaded.manifest.as_ref().is_none_or(|manifest| {
run_manifest.runtime_envelope.as_ref().is_none_or(|expected| {
manifest.runtime_envelope.as_ref().is_some_and(|actual| {
actual == expected
&& actual.conversation_binding.as_deref()
== Some(manifest.conversation_id.as_str())
})
})
});
let reset_reason = loaded.manifest.as_ref().map_or(loaded.reset_reason, |_| {
Some(
"workflow reuse policy `fresh_each_run` requested a new conversation for this run"
.to_string(),
)
});
if let Some(manifest) = loaded.manifest.as_ref() {
self.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
manifest,
)
.await?;
if !safe_to_retire {
tracing::warn!(
conversation_id = %manifest.conversation_id,
"skipping retirement of fresh_each_run conversation with an untrusted runtime envelope"
);
} else {
tracing::debug!(
conversation_id = %manifest.conversation_id,
"deferring fresh_each_run conversation retirement until replacement is durable"
);
}
}
let replacement = self.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
reset_reason,
)
.await?;
if safe_to_retire
&& let Some(manifest) = loaded.manifest.as_ref()
{
self.retire_replaced_conversation(
workspace_manager,
workspace,
manifest,
&replacement,
)
.await;
}
Ok(replacement)
}
IssueSessionReusePolicy::Unsupported(policy) => self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
None,
failed_outcome(
"workflow configured an unsupported OpenHands conversation reuse policy",
format!(
"unsupported OpenHands conversation reuse policy `{policy}`; supported runtime policies: `{DEFAULT_REUSE_POLICY}`, `{FRESH_EACH_RUN_REUSE_POLICY}`"
),
),
)
.await
.map(Box::new)
.map(Step::EarlyResult),
}
}
#[allow(clippy::too_many_arguments)]
async fn prepare_turn<O>(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
workflow: &ResolvedWorkflow,
issue: &NormalizedIssue,
run: &RunAttempt,
mut active_session: ActiveSession,
launch_reported: bool,
_observer: &mut O,
) -> Result<Step<(ActiveSession, PreparedTurn)>, IssueSessionError>
where
O: IssueSessionObserver,
{
let mut waited_for_prior_turn = false;
if let Some(status) = active_session.stream.state_mirror().execution_status()
&& turn_is_in_progress(status)
{
if let Err(error) = self
.wait_for_active_turn_to_finish(&mut active_session.stream)
.await
{
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome(
"previous OpenHands turn did not finish before retrying",
error.to_string(),
),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
waited_for_prior_turn = true;
active_session
.manifest
.apply_runtime_snapshot(&active_session.stream);
}
if let Err(error) = write_memory_context_artifact(
workspace_manager,
workspace,
issue,
self.config.memory.as_ref(),
)
.await
{
debug!(
%error,
issue = %issue.identifier,
"memory context artifact generation failed; continuing without it"
);
}
let prompt = match self.render_prompt(workflow, issue, run, active_session.prompt_kind) {
Ok(prompt) => prompt,
Err(detail) => {
let summary = format!(
"failed to render {} prompt",
active_session.prompt_kind.as_str()
);
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome(summary, detail),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
let prompt_path = active_session.prompt_kind.artifact_path(workspace);
workspace_manager
.write_text_artifact(workspace, &prompt_path, &prompt)
.await?;
let prompt_recorded_at = Utc::now();
active_session.manifest.record_prompt(
active_session.prompt_kind,
prompt_path.clone(),
prompt_recorded_at,
);
active_session.prompt_path = Some(prompt_path);
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&active_session.manifest,
)
.await?;
let conversation_id = match parse_uuid(active_session.manifest.conversation_id.as_str()) {
Ok(conversation_id) => conversation_id,
Err(detail) => {
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome(
"conversation manifest contained an invalid conversation ID",
detail,
),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
let baseline_event_ids = active_session
.stream
.event_cache()
.items()
.iter()
.map(|event| event.id.clone())
.collect::<HashSet<_>>();
Ok(Step::Continue((
active_session,
PreparedTurn {
conversation_id,
prompt,
baseline_event_ids,
launch_reported,
waited_for_prior_turn,
},
)))
}
#[allow(clippy::too_many_arguments)]
async fn start_turn<O>(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
mut active_session: ActiveSession,
prepared_turn: &mut PreparedTurn,
observer: &mut O,
) -> Result<Step<ActiveSession>, IssueSessionError>
where
O: IssueSessionObserver,
{
active_session.manifest.prepared_run_id = Some(run_manifest.run_id.clone());
active_session.manifest.trigger_pending_run_id = None;
active_session.manifest.updated_at = Utc::now();
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&active_session.manifest,
)
.await?;
if let Err(error) = self
.client
.send_message(
prepared_turn.conversation_id,
&SendMessageRequest::user_text(prepared_turn.prompt.clone()),
)
.await
{
active_session.manifest.prepared_run_id = None;
let summary = format!(
"failed to send {} prompt event",
active_session.prompt_kind.as_str()
);
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome(summary, error.to_string()),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
active_session.manifest.active_run_id = Some(run_manifest.run_id.clone());
active_session.manifest.prepared_run_id = None;
active_session.manifest.trigger_pending_run_id = Some(run_manifest.run_id.clone());
active_session.manifest.updated_at = Utc::now();
if active_session.prompt_kind == IssueSessionPromptKind::Full {
active_session.manifest.workflow_prompt_seeded = true;
}
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&active_session.manifest,
)
.await?;
let mut had_run_conflict = false;
loop {
match self
.client
.run_conversation(prepared_turn.conversation_id)
.await
{
Ok(_) => break,
Err(OpenHandsError::HttpStatus {
status_code: 409, ..
}) => {
had_run_conflict = true;
if !prepared_turn.launch_reported {
observer.on_launch_with_started_at(
&active_session
.manifest
.to_domain_metadata(RuntimeStreamState::Ready),
run_manifest.started_at.map(timestamp_ms_from_datetime),
);
prepared_turn.launch_reported = true;
}
if let Err(error) = active_session.stream.reconcile_events().await {
debug!(
%error,
conversation_id = %active_session.manifest.conversation_id,
"pre-wait reconcile failed after run retry conflict, proceeding anyway"
);
}
if let Err(error) = self
.wait_for_active_turn_to_finish(&mut active_session.stream)
.await
{
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome(
"previous OpenHands turn did not finish after run retry conflict",
error.to_string(),
),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
active_session
.manifest
.apply_runtime_snapshot(&active_session.stream);
active_session.accumulate_tokens();
prepared_turn.baseline_event_ids.extend(
active_session
.stream
.event_cache()
.items()
.iter()
.map(|event| event.id.clone()),
);
}
Err(error) => {
return self
.finalize_active_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
active_session,
failed_outcome("failed to trigger OpenHands run", error.to_string()),
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
}
}
active_session.manifest.trigger_pending_run_id = None;
active_session.manifest.updated_at = Utc::now();
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&active_session.manifest,
)
.await?;
if (prepared_turn.waited_for_prior_turn || had_run_conflict)
&& let Err(error) = active_session.stream.reconcile_events().await
{
debug!(
%error,
conversation_id = %active_session.manifest.conversation_id,
"post-conflict reconcile failed, proceeding anyway"
);
}
run_manifest.status = RunStatus::Running;
run_manifest.started_at.get_or_insert_with(Utc::now);
run_manifest.status_detail = Some(format!(
"{} prompt sent to conversation {}",
active_session.prompt_kind.as_str(),
active_session.manifest.conversation_id
));
workspace_manager
.write_run_manifest(workspace, run_manifest)
.await?;
workspace_manager
.write_json_artifact(
workspace,
&session_context_path(workspace),
&build_session_context(
run_manifest,
observed_run,
&active_session.manifest,
active_session.prompt_kind,
active_session.prompt_path.clone(),
None,
),
)
.await?;
active_session.accumulate_tokens();
if !prepared_turn.launch_reported {
observer.on_launch_with_started_at(
&active_session
.manifest
.to_domain_metadata(RuntimeStreamState::Ready),
run_manifest.started_at.map(timestamp_ms_from_datetime),
);
prepared_turn.launch_reported = true;
}
Ok(Step::Continue(active_session))
}
async fn load_existing_conversation_manifest(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
) -> Result<LoadedManifest, IssueSessionError> {
let Some(raw) = workspace_manager
.read_text_artifact(workspace, &workspace.conversation_manifest_path())
.await?
else {
return Ok(LoadedManifest::default());
};
let manifest = match serde_json::from_str::<IssueConversationManifest>(&raw) {
Ok(manifest) => manifest,
Err(error) => {
return Ok(LoadedManifest {
manifest: None,
reset_reason: Some(format!("invalid conversation manifest: {error}")),
});
}
};
let expected_persistence_dir = configured_persistence_dir(workflow, workspace);
let expected_reuse_policy = self.config.reuse_policy.as_str();
if !manifest.is_reusable_for(issue, &expected_persistence_dir, expected_reuse_policy) {
let reset_reason = if manifest.reuse_policy != expected_reuse_policy {
format!(
"conversation manifest reuse policy `{}` does not match expected `{}`",
manifest.reuse_policy, expected_reuse_policy
)
} else {
format!(
"conversation manifest is incompatible with issue {} or the current workspace",
issue.identifier
)
};
return Ok(LoadedManifest {
manifest: None,
reset_reason: Some(reset_reason),
});
}
Ok(LoadedManifest {
manifest: Some(manifest),
reset_reason: None,
})
}
async fn preserve_superseded_conversation_manifest(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
manifest: &IssueConversationManifest,
) -> Result<(), IssueSessionError> {
let path = superseded_conversation_manifests_path(workspace);
let mut manifests = workspace_manager
.read_text_artifact(workspace, &path)
.await?
.map(|raw| decode_superseded_conversation_manifests(&raw))
.transpose()?
.flatten()
.unwrap_or_default();
if manifests
.iter()
.all(|existing| existing.conversation_id != manifest.conversation_id)
{
manifests.push(manifest.clone());
}
workspace_manager
.write_json_artifact_atomically(workspace, &path, &Some(&manifests))
.await?;
Ok(())
}
async fn clear_superseded_conversation_manifest(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
conversation_id: &ConversationId,
) -> Result<(), IssueSessionError> {
let path = superseded_conversation_manifests_path(workspace);
let Some(raw) = workspace_manager
.read_text_artifact(workspace, &path)
.await?
else {
return Ok(());
};
let Some(mut manifests) = decode_superseded_conversation_manifests(&raw)? else {
return Ok(());
};
manifests.retain(|manifest| &manifest.conversation_id != conversation_id);
let replacement = (!manifests.is_empty()).then_some(&manifests);
workspace_manager
.write_json_artifact_atomically(workspace, &path, &replacement)
.await?;
Ok(())
}
async fn try_reuse_session(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
_issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
manifest: IssueConversationManifest,
) -> Result<ReuseSession, IssueSessionError> {
let conversation_id = match parse_uuid(manifest.conversation_id.as_str()) {
Ok(conversation_id) => conversation_id,
Err(error) => {
return Ok(ReuseSession::Reset {
reason: error,
manifest: Box::new(manifest),
});
}
};
let workflow_condenser = workflow
.extensions
.openhands
.conversation
.agent
.condenser
.as_ref();
let manifest_condenser = manifest
.launch_profile
.as_ref()
.and_then(|p| p.condenser.as_ref());
if workflow_condenser.is_some() && manifest_condenser.is_none() {
return Ok(ReuseSession::Reset {
reason: "workflow now has condenser enabled, but existing conversation was created without condenser - resetting to apply condenser".to_string(),
manifest: Box::new(manifest),
});
}
if manifest.last_execution_status.as_deref() == Some("error")
&& manifest.last_event_id.is_some()
{
return Ok(ReuseSession::Reset {
reason: "previous run ended with error status, resetting to avoid potential corrupted event history".to_string(),
manifest: Box::new(manifest),
});
}
if self
.config
.memory
.as_ref()
.is_some_and(|memory| memory.requires_fresh_conversation)
{
return Ok(ReuseSession::Reset {
reason: "existing conversation has a process-scoped memory grant that cannot be refreshed in place; creating a new conversation".to_string(),
manifest: Box::new(manifest),
});
}
self.try_attach_and_resume(
&self.client,
workspace_manager,
workspace,
manifest,
conversation_id,
)
.await
}
async fn retire_conversation(
&self,
manifest: &IssueConversationManifest,
reason: &str,
) -> Result<(), IssueSessionError> {
let conversation_id = match parse_uuid(manifest.conversation_id.as_str()) {
Ok(conversation_id) => conversation_id,
Err(error) => {
tracing::warn!(
conversation_id = %manifest.conversation_id,
%reason,
%error,
"superseding conversation has no parseable ID; continuing without remote retirement"
);
return Ok(());
}
};
match self.client.delete_conversation(conversation_id).await {
Ok(())
| Err(OpenHandsError::HttpStatus {
status_code: 404, ..
}) => {}
Err(error) => {
return Err(IssueSessionError::ConversationRetirementFailed(format!(
"failed to retire conversation {} for {reason}: {error}",
manifest.conversation_id
)));
}
}
tracing::info!(
conversation_id = %manifest.conversation_id,
%reason,
"retired superseded OpenHands conversation"
);
Ok(())
}
async fn retire_replaced_conversation(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
manifest: &IssueConversationManifest,
replacement: &Step<ActiveSession>,
) {
if !matches!(replacement, Step::Continue(_)) {
return;
}
if let Err(error) = self
.retire_conversation(
manifest,
"superseded conversation after replacement became durable",
)
.await
{
tracing::warn!(
conversation_id = %manifest.conversation_id,
%error,
"failed to retire superseded conversation after durable replacement"
);
return;
}
if let Err(error) = self
.clear_superseded_conversation_manifest(
workspace_manager,
workspace,
&manifest.conversation_id,
)
.await
{
tracing::warn!(
conversation_id = %manifest.conversation_id,
%error,
"failed to clear retired superseded conversation evidence"
);
}
}
async fn try_attach_and_resume(
&self,
client: &OpenHandsClient,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
mut manifest: IssueConversationManifest,
conversation_id: Uuid,
) -> Result<ReuseSession, IssueSessionError> {
let manifest_conversation_id = manifest.conversation_id.clone();
let mut stream = match client
.attach_runtime_stream(conversation_id, self.config.runtime_stream.clone())
.await
{
Ok(stream) => stream,
Err(OpenHandsError::HttpStatus {
status_code: 404, ..
}) => {
return Ok(ReuseSession::Reset {
reason: format!(
"existing conversation {} is no longer available",
manifest_conversation_id
),
manifest: Box::new(manifest),
});
}
Err(error) => {
return Err(IssueSessionError::RehydrationFailed(format!(
"failed to attach existing conversation {}: {error}",
manifest_conversation_id
)));
}
};
if let Err(reason) =
verify_conversation_workspace(stream.conversation(), workspace.workspace_path()).await
{
let _ = stream.close().await;
return Ok(ReuseSession::Reset {
reason: format!(
"existing conversation {} workspace is incompatible: {reason}",
manifest_conversation_id
),
manifest: Box::new(manifest),
});
}
let attached_at = Utc::now();
manifest.fresh_conversation = false;
manifest.reuse_policy = self.config.reuse_policy.as_str().to_owned();
manifest.llm_config_fingerprint.get_or_insert_with(|| {
LlmConfigFingerprint::from_llm_config(&stream.conversation().agent.llm)
});
let transport_diagnostics = client.transport_diagnostics().ok();
manifest.apply_transport_diagnostics(transport_diagnostics.as_ref(), client.base_url());
manifest.runtime_contract_version = Some(RUNTIME_CONTRACT_VERSION.to_string());
manifest.last_attached_at = attached_at;
manifest.updated_at = attached_at;
manifest.reset_reason = None;
manifest.apply_runtime_snapshot(&stream);
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&manifest,
)
.await?;
workspace_manager
.write_json_artifact(
workspace,
&last_conversation_state_path(workspace),
&conversation_snapshot(&stream),
)
.await?;
let mut session = ActiveSession {
prompt_kind: manifest.prompt_kind(),
stream,
manifest,
prompt_path: None,
};
session.accumulate_tokens();
Ok(ReuseSession::Active(Box::new(session)))
}
#[allow(clippy::too_many_arguments)]
fn create_fresh_session<'a>(
&'a self,
workspace_manager: &'a WorkspaceManager,
workspace: &'a WorkspaceHandle,
run_manifest: &'a mut RunManifest,
observed_run: &'a RunAttempt,
issue: &'a NormalizedIssue,
workflow: &'a ResolvedWorkflow,
reset_reason: Option<String>,
) -> Pin<Box<dyn Future<Output = Result<Step<ActiveSession>, IssueSessionError>> + Send + 'a>>
{
Box::pin(self.create_fresh_session_inner(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
reset_reason,
))
}
#[allow(clippy::too_many_arguments)]
async fn create_fresh_session_inner(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
reset_reason: Option<String>,
) -> Result<Step<ActiveSession>, IssueSessionError> {
let launch_profile = match ConversationLaunchProfile::from_workflow(workflow) {
Ok(launch_profile) => launch_profile,
Err(detail) => {
return self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
None,
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "failed to build conversation launch profile".to_string(),
error: Some(detail),
},
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
let request = match launch_profile.to_create_request(
self.environment.as_ref(),
workspace.workspace_path(),
&configured_persistence_dir(workflow, workspace),
None,
) {
Ok(request) => request,
Err(detail) => {
return self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
None,
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "failed to build OpenHands conversation create request"
.to_string(),
error: Some(detail),
},
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
let mut request = request;
if let Some(memory) = self.config.memory.as_ref() {
request.agent.mcp_config = memory.mcp_config();
}
let persisted_request = request.without_mcp_credentials();
workspace_manager
.write_json_artifact(
workspace,
&create_conversation_request_path(workspace),
&persisted_request,
)
.await?;
let conversation = match self.client.create_conversation(&request).await {
Ok(conversation) => conversation,
Err(error) => {
return self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
None,
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "failed to create OpenHands conversation".to_string(),
error: Some(redacted_openhands_error(&error, &request)),
},
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
let created_at = Utc::now();
let mut manifest = IssueConversationManifest::new(
issue.id.clone(),
issue.identifier.clone(),
ConversationId::new(conversation.conversation_id.to_string())
.expect("UUID-backed conversation ID should not be empty"),
self.config.reuse_policy.as_str(),
configured_persistence_dir(workflow, workspace),
created_at,
reset_reason,
launch_profile.clone(),
self.environment.as_ref(),
);
manifest.llm_config_fingerprint =
Some(LlmConfigFingerprint::from_llm_config(&request.agent.llm));
bind_runtime_envelopes(&mut manifest, run_manifest);
let pending_manifest_path = pending_conversation_manifest_path(workspace);
if let Err(error) = workspace_manager
.write_json_artifact_atomically(workspace, &pending_manifest_path, &Some(&manifest))
.await
{
let retirement = self
.retire_conversation(&manifest, "failed to persist fresh conversation ownership")
.await
.err();
let mut detail = format!("failed to persist fresh conversation ownership: {error}");
if let Some(retirement) = retirement {
detail.push_str(&format!("; failed to retire conversation: {retirement}"));
}
return Err(IssueSessionError::RehydrationFailed(detail));
}
workspace_manager
.write_run_manifest(workspace, run_manifest)
.await?;
let stream = match Box::pin(self.client.attach_runtime_stream(
conversation.conversation_id,
self.config.runtime_stream.clone(),
))
.await
{
Ok(stream) => stream,
Err(error) => {
return self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
Some(build_summary_metadata(
&conversation,
true,
RuntimeStreamState::Failed,
self.client.transport_diagnostics().ok().as_ref(),
self.client.base_url(),
)),
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "failed to attach runtime stream for a fresh conversation"
.to_string(),
error: Some(error.to_string()),
},
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
};
if let Err(detail) =
verify_conversation_workspace(stream.conversation(), workspace.workspace_path()).await
{
let (failure_detail, conversation_metadata, clear_pending_ownership) =
Box::pin(self.handle_workspace_mismatch(
workspace_manager,
workspace,
run_manifest,
issue,
workflow,
&launch_profile,
&conversation,
stream,
detail,
))
.await;
let mut failure_detail = failure_detail;
if clear_pending_ownership
&& let Err(error) = workspace_manager
.write_json_artifact_atomically(
workspace,
&pending_manifest_path,
&Option::<IssueConversationManifest>::None,
)
.await
{
failure_detail.push_str(&format!(
"; failed to clear pending conversation ownership: {error}"
));
}
return self
.persist_failure_without_stream(
workspace_manager,
workspace,
run_manifest,
observed_run,
IssueSessionPromptKind::Full,
Some(conversation_metadata),
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary:
"OpenHands runtime stream workspace did not match verified checkout"
.to_owned(),
error: Some(failure_detail),
},
)
.await
.map(Box::new)
.map(Step::EarlyResult);
}
let attached_at = Utc::now();
manifest.last_attached_at = attached_at;
manifest.updated_at = attached_at;
let transport_diagnostics = self.client.transport_diagnostics().ok();
manifest
.apply_transport_diagnostics(transport_diagnostics.as_ref(), self.client.base_url());
manifest.apply_runtime_snapshot(&stream);
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&manifest,
)
.await?;
workspace_manager
.write_json_artifact(
workspace,
&pending_manifest_path,
&Option::<IssueConversationManifest>::None,
)
.await?;
workspace_manager
.write_json_artifact(
workspace,
&last_conversation_state_path(workspace),
&conversation_snapshot(&stream),
)
.await?;
let mut session = ActiveSession {
prompt_kind: manifest.prompt_kind(),
stream,
manifest,
prompt_path: None,
};
session.accumulate_tokens();
Ok(Step::Continue(session))
}
#[allow(clippy::too_many_arguments)]
async fn handle_workspace_mismatch(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &RunManifest,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
launch_profile: &ConversationLaunchProfile,
conversation: &Conversation,
mut stream: RuntimeEventStream,
detail: String,
) -> (String, ConversationMetadata, bool) {
let mut failed_manifest = IssueConversationManifest::new(
issue.id.clone(),
issue.identifier.clone(),
ConversationId::new(conversation.conversation_id.to_string())
.expect("UUID-backed conversation ID should not be empty"),
self.config.reuse_policy.as_str(),
configured_persistence_dir(workflow, workspace),
Utc::now(),
Some("verified checkout/runtime workspace mismatch".to_owned()),
launch_profile.clone(),
self.environment.as_ref(),
);
failed_manifest.runtime_envelope = run_manifest.runtime_envelope.clone();
if let Some(envelope) = failed_manifest.runtime_envelope.as_mut() {
envelope.conversation_binding = Some(failed_manifest.conversation_id.to_string());
}
failed_manifest.parent_runtime_envelope = run_manifest.parent_runtime_envelope.clone();
if let Some(envelope) = failed_manifest.parent_runtime_envelope.as_mut() {
envelope.conversation_binding = Some(failed_manifest.conversation_id.to_string());
}
failed_manifest.apply_transport_diagnostics(
self.client.transport_diagnostics().ok().as_ref(),
self.client.base_url(),
);
failed_manifest.apply_runtime_snapshot(&stream);
let _ = stream.close().await;
let preserve_error = self
.preserve_superseded_conversation_manifest(
workspace_manager,
workspace,
&failed_manifest,
)
.await
.err();
let retire_error = self
.retire_conversation(
&failed_manifest,
"verified checkout/runtime workspace mismatch",
)
.await
.err();
if preserve_error.is_none()
&& retire_error.is_none()
&& let Err(error) = self
.clear_superseded_conversation_manifest(
workspace_manager,
workspace,
&failed_manifest.conversation_id,
)
.await
{
tracing::warn!(
conversation_id = %failed_manifest.conversation_id,
%error,
"failed to clear retired workspace-mismatch conversation evidence"
);
}
let clear_pending_ownership = preserve_error.is_none() || retire_error.is_none();
let mut failure_detail = detail;
if let Some(error) = preserve_error {
failure_detail.push_str(&format!(
"; failed to preserve conversation evidence: {error}"
));
}
if let Some(error) = retire_error {
failure_detail.push_str(&format!("; failed to retire conversation: {error}"));
}
let conversation_metadata = build_summary_metadata(
conversation,
true,
RuntimeStreamState::Failed,
self.client.transport_diagnostics().ok().as_ref(),
self.client.base_url(),
);
(
failure_detail,
conversation_metadata,
clear_pending_ownership,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn rehydrate_conversation(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
issue: &NormalizedIssue,
workflow: &ResolvedWorkflow,
old_manifest: &IssueConversationManifest,
options: RehydrationOptions,
) -> Result<RehydrationResult, IssueSessionError> {
let old_conversation_id = old_manifest.conversation_id.clone();
let old_stream = if options.summarize {
match parse_uuid(old_conversation_id.as_str()) {
Ok(conversation_id) => self
.client
.attach_runtime_stream(conversation_id, self.config.runtime_stream.clone())
.await
.ok(),
Err(_) => None,
}
} else {
None
};
let context = if let Some(stream) = old_stream {
self.build_rehydration_context(&stream, options.max_summary_events)
.await
} else {
None
};
self.preserve_superseded_conversation_manifest(workspace_manager, workspace, old_manifest)
.await?;
let step = self
.create_fresh_session(
workspace_manager,
workspace,
run_manifest,
observed_run,
issue,
workflow,
Some(format!("rehydration: {}", options.reason)),
)
.await?;
let mut session = match step {
Step::Continue(session) => session,
Step::EarlyResult(result) => {
return Err(IssueSessionError::RehydrationFailed(
result
.worker_outcome
.error
.or(result.worker_outcome.summary)
.unwrap_or_else(|| "rehydration failed".to_string()),
));
}
};
let retired = if parse_uuid(old_conversation_id.as_str()).is_ok() {
match self
.retire_conversation(
old_manifest,
"superseded conversation after replacement became durable during rehydration",
)
.await
{
Ok(()) => true,
Err(error) => {
tracing::warn!(
conversation_id = %old_conversation_id,
%error,
"failed to delete old conversation after replacement became durable during rehydration; preserving superseded evidence"
);
false
}
}
} else {
false
};
if retired
&& let Err(error) = self
.clear_superseded_conversation_manifest(
workspace_manager,
workspace,
&old_manifest.conversation_id,
)
.await
{
tracing::warn!(
conversation_id = %old_conversation_id,
%error,
"failed to clear retired superseded conversation evidence after rehydration"
);
}
session.manifest.input_tokens = old_manifest.input_tokens;
session.manifest.output_tokens = old_manifest.output_tokens;
session.manifest.cache_read_tokens = old_manifest.cache_read_tokens;
session.manifest.last_token_accumulation_at = old_manifest.last_token_accumulation_at;
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&session.manifest,
)
.await?;
Ok(RehydrationResult {
session,
context,
old_conversation_id: old_conversation_id.to_string(),
})
}
async fn build_rehydration_context(
&self,
stream: &RuntimeEventStream,
max_events: usize,
) -> Option<String> {
let conversation = stream.conversation();
let summary = self.summarize_conversation(stream, max_events).await;
let context = format!(
"## Previous Conversation Context\n\n\
This conversation was rehydrated from a previous session.\n\n\
- Original conversation ID: {}\n\
- Model: {}\n\
- Max iterations: {}\n\
- Execution status: {}\n\n\
### Summary of Previous Work\n\n\
{}",
conversation.conversation_id,
conversation.agent.llm.model,
conversation.max_iterations,
conversation.execution_status,
summary.unwrap_or_else(|| "No summary available.".to_string())
);
Some(context)
}
async fn summarize_conversation(
&self,
stream: &RuntimeEventStream,
max_events: usize,
) -> Option<String> {
let events: Vec<_> = stream
.event_cache()
.items()
.iter()
.filter(|e| {
matches!(
KnownEvent::from_envelope(e),
KnownEvent::Message(_)
| KnownEvent::Observation(_)
| KnownEvent::ConversationError(_)
)
})
.take(max_events)
.collect();
if events.is_empty() {
return None;
}
let mut summary = String::new();
for event in events {
let line = match KnownEvent::from_envelope(event) {
KnownEvent::Message(msg) => {
format!("- User: {}\n", msg.text_preview.as_deref().unwrap_or("..."))
}
KnownEvent::Observation(obs) => {
format!(
"- {}: {}\n",
obs.tool_name.as_deref().unwrap_or("Result"),
obs.text_preview.as_deref().unwrap_or("...")
)
}
KnownEvent::ConversationError(err) => {
let msg = err
.payload
.get("message")
.and_then(|v| v.as_str())
.unwrap_or("Unknown error");
format!("- Error: {}\n", msg)
}
_ => String::new(),
};
summary.push_str(&line);
}
Some(summary)
}
fn render_prompt(
&self,
workflow: &ResolvedWorkflow,
issue: &NormalizedIssue,
run: &RunAttempt,
prompt_kind: IssueSessionPromptKind,
) -> Result<String, String> {
match prompt_kind {
IssueSessionPromptKind::Full => {
if let Some(prompt) = self.config.terminal_prompt.as_deref() {
return Ok(prompt.to_owned());
}
workflow
.render_prompt(issue, run.attempt.map(|attempt| attempt.get()))
.map(|prompt| {
self.config.repository_instructions.as_deref().map_or(
prompt.clone(),
|instructions| {
compose_terminal_prompt(
&prompt,
&format!("Issue: {}\nTitle: {}", issue.identifier, issue.title),
"Verified checkout facts are persisted in the run envelope.",
Some(instructions),
"trusted_host_process_cwd",
)
},
)
})
.map_err(|error| error.to_string())
}
IssueSessionPromptKind::Continuation => {
let prompt = append_memory_scope_guidance(
build_continuation_guidance(issue, run),
self.config.memory.as_ref(),
);
Ok(append_attempt_continuation(
prompt,
self.config.continuation_prompt.as_deref(),
))
}
}
}
async fn wait_for_active_turn_to_finish(
&self,
stream: &mut RuntimeEventStream,
) -> Result<(), OpenHandsError> {
if stream
.state_mirror()
.execution_status()
.is_none_or(turn_has_stopped)
{
return Ok(());
}
let deadline = Instant::now() + self.config.terminal_wait_timeout;
loop {
if stream
.state_mirror()
.execution_status()
.is_some_and(turn_has_stopped)
{
return Ok(());
}
match timeout_at(deadline, stream.next_event()).await {
Err(_) => {
let status = stream
.state_mirror()
.execution_status()
.unwrap_or("unknown");
return Err(OpenHandsError::Protocol {
operation: "wait for active turn to finish",
detail: format!(
"execution_status `{status}` did not stop within {} ms",
self.config.terminal_wait_timeout.as_millis()
),
});
}
Ok(Ok(Some(_event))) => {}
Ok(Ok(None)) => {}
Ok(Err(error)) => {
if stream
.state_mirror()
.execution_status()
.is_some_and(turn_has_stopped)
&& finished_stream_error_is_tolerable(&error)
{
return Ok(());
}
return Err(error);
}
}
}
}
async fn await_terminal_outcome<O>(
&self,
session: &mut ActiveSession,
baseline_event_ids: &HashSet<String>,
observer: &mut O,
) -> NormalizedOutcome
where
O: IssueSessionObserver,
{
let idle_timeout_ms = DurationMs::new(self.config.terminal_wait_timeout.as_millis() as u64);
let total_runtime_cap_ms = self
.config
.total_runtime_cap_ms
.map(|d| DurationMs::new(d.as_millis() as u64));
let mut tracker = LivenessTracker::with_runtime_cap(idle_timeout_ms, total_runtime_cap_ms);
tracker.mark_started(timestamp_ms_from_datetime(Utc::now()));
let mut next_token_accumulation = Instant::now() + Duration::from_secs(15);
loop {
if Instant::now() >= next_token_accumulation {
session.accumulate_tokens();
let (input, output) = (
session.manifest.input_tokens,
session.manifest.output_tokens,
);
tracker.record_tokens(input, output, timestamp_ms_from_datetime(Utc::now()));
observer.on_conversation_update(
&session
.manifest
.to_domain_metadata(RuntimeStreamState::Ready),
);
next_token_accumulation = Instant::now() + Duration::from_secs(15);
}
match Box::pin(self.terminal_outcome_from_state(
&mut session.stream,
baseline_event_ids,
observer,
))
.await
{
StateCheckResult::Terminal(outcome) => {
session.accumulate_tokens();
let (input, output) = (
session.manifest.input_tokens,
session.manifest.output_tokens,
);
tracker.record_tokens(input, output, timestamp_ms_from_datetime(Utc::now()));
return outcome;
}
StateCheckResult::StillRunningWithProgress => {
}
StateCheckResult::NoProgress => {}
}
let now = Instant::now();
let now_ts = timestamp_ms_from_datetime(Utc::now());
let event_timeout =
compute_timeout_duration(&tracker, next_token_accumulation, now, now_ts);
let next_event =
timeout_at(now + event_timeout, Box::pin(session.stream.next_event())).await;
match next_event {
Err(_) => {
if Instant::now() >= next_token_accumulation {
session.accumulate_tokens();
next_token_accumulation = Instant::now() + Duration::from_secs(15);
}
let now_ts = timestamp_ms_from_datetime(Utc::now());
if tracker.is_stalled_at(now_ts) {
match self
.handle_reconcile_progress(
session,
baseline_event_ids,
observer,
Some(&mut tracker),
)
.await
{
ReconcileResult::Terminal(outcome) => {
session.accumulate_tokens();
return outcome;
}
ReconcileResult::Progress => {
continue;
}
ReconcileResult::NoProgress => {}
}
let now_ts = timestamp_ms_from_datetime(Utc::now());
if tracker.is_stalled_at(now_ts) {
return NormalizedOutcome {
kind: WorkerOutcomeKind::Stalled,
summary:
"runtime did not reach a terminal state before the stall timeout"
.to_string(),
error: Some(format!(
"no progress observed within {} ms idle timeout",
self.config.terminal_wait_timeout.as_millis()
)),
};
}
}
}
Ok(Ok(Some(event))) => {
observe_event(observer, &event);
tracker.record_event(timestamp_ms_from_datetime(event.timestamp));
if let Some(status) = session.stream.state_mirror().execution_status() {
tracker
.record_status_change(status, timestamp_ms_from_datetime(Utc::now()));
}
}
Ok(Ok(None)) => {
match self
.handle_reconcile_progress(
session,
baseline_event_ids,
observer,
Some(&mut tracker),
)
.await
{
ReconcileResult::Terminal(outcome) => {
session.accumulate_tokens();
return outcome;
}
ReconcileResult::Progress => {
continue;
}
ReconcileResult::NoProgress => {}
}
session.accumulate_tokens();
return NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "runtime event stream ended before terminal status".to_string(),
error: Some(
"runtime event stream closed before a terminal state was observed"
.to_string(),
),
};
}
Ok(Err(error)) => {
match self
.handle_reconcile_progress(
session,
baseline_event_ids,
observer,
Some(&mut tracker),
)
.await
{
ReconcileResult::Terminal(outcome) => {
session.accumulate_tokens();
return outcome;
}
ReconcileResult::Progress => {
continue;
}
ReconcileResult::NoProgress => {}
}
return NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "runtime event stream failed before terminal status".to_string(),
error: Some(error.to_string()),
};
}
}
}
}
async fn handle_reconcile_progress<O>(
&self,
session: &mut ActiveSession,
baseline_event_ids: &HashSet<String>,
observer: &mut O,
tracker: Option<&mut LivenessTracker>,
) -> ReconcileResult
where
O: IssueSessionObserver,
{
let previously_observed = session
.stream
.event_cache()
.items()
.iter()
.map(|event| event.id.clone())
.collect::<HashSet<_>>();
if let Ok(inserted) = session.stream.reconcile_events().await
&& inserted > 0
{
observe_reconciled_events(
observer,
session.stream.event_cache().items(),
&previously_observed,
);
if let Some(tracker) = tracker {
tracker.record_reconciled_events(
inserted as u64,
timestamp_ms_from_datetime(Utc::now()),
);
}
match self
.terminal_outcome_from_state(&mut session.stream, baseline_event_ids, observer)
.await
{
StateCheckResult::Terminal(outcome) => return ReconcileResult::Terminal(outcome),
StateCheckResult::StillRunningWithProgress => return ReconcileResult::Progress,
StateCheckResult::NoProgress => return ReconcileResult::Progress,
}
}
ReconcileResult::NoProgress
}
async fn terminal_outcome_from_state<O>(
&self,
stream: &mut RuntimeEventStream,
baseline_event_ids: &HashSet<String>,
observer: &mut O,
) -> StateCheckResult
where
O: IssueSessionObserver,
{
let has_current_turn_activity = stream
.event_cache()
.items()
.iter()
.any(|event| !baseline_event_ids.contains(&event.id));
if !has_current_turn_activity {
return StateCheckResult::NoProgress;
}
if let Some(error_detail) =
latest_current_turn_error(stream.event_cache().items(), baseline_event_ids)
{
return StateCheckResult::Terminal(NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "received ConversationErrorEvent during the current run".to_string(),
error: Some(error_detail),
});
}
match stream.state_mirror().terminal_status() {
Some(TerminalExecutionStatus::Finished) => {
if self
.confirm_finished_terminal_state(stream, baseline_event_ids, observer)
.await
{
StateCheckResult::Terminal(NormalizedOutcome {
kind: WorkerOutcomeKind::Succeeded,
summary: "OpenHands execution_status `finished`".to_string(),
error: None,
})
} else {
StateCheckResult::NoProgress
}
}
Some(TerminalExecutionStatus::Error) => {
let error_detail = extract_error_detail_from_state(stream.state_mirror())
.unwrap_or_else(|| "execution_status error".to_string());
StateCheckResult::Terminal(NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "OpenHands execution_status `error`".to_string(),
error: Some(error_detail),
})
}
Some(TerminalExecutionStatus::Stuck) => StateCheckResult::Terminal(NormalizedOutcome {
kind: WorkerOutcomeKind::Stalled,
summary: "OpenHands execution_status `stuck`".to_string(),
error: Some(
stream
.state_mirror()
.execution_status()
.unwrap_or_default()
.to_string(),
),
}),
None => StateCheckResult::StillRunningWithProgress,
}
}
async fn confirm_finished_terminal_state<O>(
&self,
stream: &mut RuntimeEventStream,
baseline_event_ids: &HashSet<String>,
observer: &mut O,
) -> bool
where
O: IssueSessionObserver,
{
let deadline = Instant::now() + self.config.finished_drain_timeout;
loop {
if latest_current_turn_error(stream.event_cache().items(), baseline_event_ids).is_some()
{
return false;
}
if !matches!(
stream.state_mirror().terminal_status(),
Some(TerminalExecutionStatus::Finished)
) {
return false;
}
match timeout_at(deadline, stream.next_event()).await {
Err(_) => return true,
Ok(Ok(Some(event))) => {
observe_event(observer, &event);
continue;
}
Ok(Ok(None)) => return true,
Ok(Err(error)) => return finished_stream_error_is_tolerable(&error),
}
}
}
async fn finalize_active_session(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
mut session: ActiveSession,
outcome: NormalizedOutcome,
) -> Result<IssueSessionResult, IssueSessionError> {
session.manifest.prepared_run_id = None;
session.manifest.trigger_pending_run_id = None;
session.manifest.apply_runtime_snapshot(&session.stream);
workspace_manager
.write_json_artifact(
workspace,
&last_conversation_state_path(workspace),
&conversation_snapshot(&session.stream),
)
.await?;
let run_status = run_status_for(outcome.kind);
run_manifest.status = run_status;
run_manifest.status_detail = Some(
outcome
.error
.clone()
.unwrap_or_else(|| outcome.summary.clone()),
);
let harness_stopped = session.stream.state_mirror().terminal_status().is_some();
run_manifest.harness_stopped = harness_stopped;
workspace_manager
.finish_run(workspace, run_manifest, run_status)
.await?;
let mut worker_outcome = WorkerOutcomeRecord::from_run(
observed_run,
outcome.kind,
timestamp_ms_from_datetime(Utc::now()),
Some(outcome.summary.clone()),
outcome.error.clone(),
);
worker_outcome.harness_stopped = harness_stopped;
workspace_manager
.write_json_artifact(
workspace,
&session_context_path(workspace),
&build_session_context(
run_manifest,
observed_run,
&session.manifest,
session.prompt_kind,
session.prompt_path.clone(),
Some(worker_outcome.clone()),
),
)
.await?;
workspace_manager
.write_json_artifact(
workspace,
&workspace.conversation_manifest_path(),
&session.manifest,
)
.await?;
let conversation = session
.manifest
.to_domain_metadata(RuntimeStreamState::Closed);
let _ = session.stream.close().await;
Ok(IssueSessionResult {
prompt_kind: session.prompt_kind,
conversation: Some(conversation),
worker_outcome,
run_status,
})
}
#[allow(clippy::too_many_arguments)]
async fn persist_failure_without_stream(
&self,
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
run_manifest: &mut RunManifest,
observed_run: &RunAttempt,
prompt_kind: IssueSessionPromptKind,
conversation: Option<ConversationMetadata>,
outcome: NormalizedOutcome,
) -> Result<IssueSessionResult, IssueSessionError> {
let run_status = run_status_for(outcome.kind);
run_manifest.status = run_status;
run_manifest.status_detail = Some(
outcome
.error
.clone()
.unwrap_or_else(|| outcome.summary.clone()),
);
workspace_manager
.finish_run(workspace, run_manifest, run_status)
.await?;
let worker_outcome = WorkerOutcomeRecord::from_run(
observed_run,
outcome.kind,
timestamp_ms_from_datetime(Utc::now()),
Some(outcome.summary),
outcome.error,
);
Ok(IssueSessionResult {
prompt_kind,
conversation,
worker_outcome,
run_status,
})
}
}
async fn verify_conversation_workspace(
conversation: &Conversation,
expected_workspace: &Path,
) -> Result<(), String> {
let expected = tokio::fs::canonicalize(expected_workspace)
.await
.map_err(|error| format!("failed to canonicalize verified checkout: {error}"))?;
let actual_path = Path::new(&conversation.workspace.working_dir);
let actual = tokio::fs::canonicalize(actual_path)
.await
.map_err(|error| format!("failed to canonicalize OpenHands workspace: {error}"))?;
if actual != expected {
return Err(format!(
"OpenHands workspace `{}` does not match verified checkout `{}`",
actual.display(),
expected.display()
));
}
Ok(())
}
fn redacted_openhands_error(error: &OpenHandsError, request: &ConversationCreateRequest) -> String {
let mut detail = error.to_string();
let Some(servers) = request
.agent
.mcp_config
.as_ref()
.and_then(|config| config.get("mcpServers"))
.and_then(Value::as_object)
else {
return detail;
};
for server in servers.values() {
let Some(headers) = server.get("headers").and_then(Value::as_object) else {
continue;
};
for value in headers.values().filter_map(Value::as_str) {
let value = value.trim();
if value.is_empty() {
continue;
}
detail = detail.replace(value, "<redacted>");
if let Some(token) = value.strip_prefix("Bearer ") {
detail = detail.replace(token, "<redacted>");
}
}
}
detail
}
impl HarnessAdapter for IssueSessionRunner {
fn harness_kind(&self) -> &'static str {
"openhands_agent_server"
}
fn capabilities(&self) -> HarnessCapability {
HarnessCapability::openhands_agent_server()
}
}
fn configured_persistence_dir(workflow: &ResolvedWorkflow, workspace: &WorkspaceHandle) -> PathBuf {
workspace.workspace_path().join(
&workflow
.extensions
.openhands
.conversation
.persistence_dir_relative,
)
}
fn turn_is_in_progress(status: &str) -> bool {
!matches!(status, "idle" | "finished" | "error" | "stuck")
}
fn turn_has_stopped(status: &str) -> bool {
!turn_is_in_progress(status)
}
fn interrupt_acknowledgement_observed(status: &str) -> bool {
status == "paused" || turn_has_stopped(status)
}
pub fn build_continuation_guidance(issue: &NormalizedIssue, run: &RunAttempt) -> String {
let attempt = run
.attempt
.map(|attempt| format!("Worker retry attempt: {}.", attempt.get()))
.unwrap_or_else(|| "Worker retry attempt: initial worker lifetime.".to_string());
format!(
"Continue working on issue {}: {}.\nThe original workflow prompt is already present in this conversation, so do not resend or restate it.\nResume from the current workspace and conversation context, inspect the latest progress, and continue from where the previous worker left off.\nCurrent issue state: {}\n{}\n",
issue.identifier, issue.title, issue.state.name, attempt,
)
}
fn append_memory_scope_guidance(
mut guidance: String,
memory: Option<&MemoryWorkerAccess>,
) -> String {
let Some(memory) = memory else {
return guidance;
};
let (Some(project), Some(repo)) = (memory.project.as_deref(), memory.execution_repo.as_deref())
else {
return guidance;
};
guidance.push_str(&format!(
"\nMemory tool scope: project={project}; repo={repo}. Pass these exact values as `project` and `repo` arguments to memory.context, memory.search, and memory.related; do not use process-global scope. Sibling repositories are persisted-memory and target-snapshot reads only; live code overlays are limited to the verified execution checkout."
));
if let Some(project_set) = memory.project_set.as_deref() {
guidance.push_str(&format!(" Project set is {project_set}."));
}
if !memory.authorized_repositories.is_empty() {
guidance.push_str(&format!(
" Authorized repositories are {}; name a sibling repository explicitly for persisted reads.",
memory.authorized_repositories.join(", ")
));
}
if let Some(run_id) = memory.run_id.as_deref() {
guidance.push_str(&format!(
" Current run={run_id}; pass it as runId only for a live execution-repository overlay."
));
}
if let Some(attempt) = memory.attempt {
guidance.push_str(&format!(" Current attempt={attempt}."));
}
guidance
}
async fn write_memory_context_artifact(
workspace_manager: &WorkspaceManager,
workspace: &WorkspaceHandle,
issue: &NormalizedIssue,
memory: Option<&MemoryWorkerAccess>,
) -> Result<(), String> {
if let Some(memory) = memory {
let context = fetch_memory_context_from_server(memory, issue).await?;
return workspace_manager
.write_text_artifact(workspace, &workspace.memory_context_path(), &context)
.await
.map_err(|error| format!("failed to write memory context artifact: {error}"));
}
let config = MemoryConfig::load(workspace.workspace_path(), None).map_err(|error| {
format!(
"failed to load memory config from {}: {error}",
workspace.workspace_path().display()
)
})?;
if !config.enabled || (!config.config_path.exists() && !config.index_path.exists()) {
return Ok(());
}
let source = SourceFile {
issues: vec![IssueEvidence {
id: Some(issue.id.as_str().to_string()),
identifier: issue.identifier.as_str().to_string(),
title: issue.title.clone(),
url: issue.url.clone(),
description: issue.description.clone(),
state: Some(issue.state.name.clone()),
labels: issue.labels.clone(),
children: issue
.sub_issues
.iter()
.map(|child| IssueLinkEvidence {
id: Some(child.id.as_str().to_string()),
identifier: child.identifier.as_str().to_string(),
state: Some(child.state.clone()),
..IssueLinkEvidence::default()
})
.collect(),
blocked_by: issue
.blocked_by
.iter()
.filter_map(|blocker| {
blocker
.identifier
.as_ref()
.map(|identifier| IssueLinkEvidence {
id: blocker.id.as_ref().map(|id| id.as_str().to_string()),
identifier: identifier.as_str().to_string(),
state: blocker.state.clone(),
..IssueLinkEvidence::default()
})
})
.collect(),
..IssueEvidence::default()
}],
..SourceFile::default()
};
let options = MemoryContextOptions::for_issue(issue.identifier.as_str(), 20);
let context = context_for_issue_with_options(&config, &source, &options)
.map_err(|error| format!("failed to build memory context: {error}"))?;
workspace_manager
.write_text_artifact(workspace, &workspace.memory_context_path(), &context)
.await
.map_err(|error| format!("failed to write memory context artifact: {error}"))
}
async fn fetch_memory_context_from_server(
memory: &MemoryWorkerAccess,
issue: &NormalizedIssue,
) -> Result<String, String> {
let mut arguments = json!({
"issue": issue.identifier.to_string(),
"limit": 20,
"currentIssue": {
"id": issue.id.to_string(),
"identifier": issue.identifier.to_string(),
"title": issue.title.clone(),
"description": issue.description.clone(),
"state": issue.state.name.clone(),
"labels": issue.labels.clone(),
"children": issue.sub_issues.iter().map(|child| json!({
"id": child.id.to_string(),
"identifier": child.identifier.to_string(),
"state": child.state.clone(),
})).collect::<Vec<_>>(),
"blockedBy": issue.blocked_by.iter().filter_map(|blocker| {
blocker.identifier.as_ref().map(|identifier| json!({
"id": blocker.id.as_ref().map(ToString::to_string),
"identifier": identifier.to_string(),
"state": blocker.state.clone(),
}))
}).collect::<Vec<_>>(),
},
});
if let Value::Object(map) = &mut arguments {
if let Some(project) = &memory.project {
map.insert("project".to_string(), json!(project));
}
if let Some(project_set) = &memory.project_set {
map.insert("projectSet".to_string(), json!(project_set));
}
if let Some(repo) = &memory.execution_repo {
map.insert("repo".to_string(), json!(repo));
}
}
let request = json!({
"jsonrpc": "2.0",
"id": "opensymphony-worker-context",
"method": "tools/call",
"params": {
"name": "memory.context",
"arguments": arguments
}
});
let client = reqwest::Client::new();
let mut builder = client.post(&memory.endpoint).json(&request);
if let Some(token) = &memory.token {
builder = builder.bearer_auth(token);
}
let response = builder
.send()
.await
.map_err(|error| format!("failed to call memory server: {error}"))?;
let status = response.status();
let payload = response
.json::<Value>()
.await
.map_err(|error| format!("memory server response was not valid JSON: {error}"))?;
if !status.is_success() {
return Err(format!("memory server returned HTTP {status}: {payload}"));
}
if let Some(error) = payload.get("error") {
return Err(format!("memory server returned an MCP error: {error}"));
}
payload
.get("result")
.and_then(|result| result.get("content"))
.and_then(Value::as_array)
.and_then(|content| content.first())
.and_then(|item| item.get("text"))
.and_then(Value::as_str)
.map(ToString::to_string)
.ok_or_else(|| "memory server response omitted text content".to_string())
}
fn build_session_context(
run_manifest: &RunManifest,
observed_run: &RunAttempt,
manifest: &IssueConversationManifest,
prompt_kind: IssueSessionPromptKind,
prompt_path: Option<PathBuf>,
worker_outcome: Option<WorkerOutcomeRecord>,
) -> IssueSessionContext {
IssueSessionContext {
run_id: run_manifest.run_id.clone(),
issue_id: manifest.issue_id.clone(),
identifier: manifest.identifier.clone(),
worker_id: observed_run.worker_id.clone(),
attempt: observed_run.attempt.map(|attempt| attempt.get()),
normal_retry_count: observed_run.normal_retry_count,
turn_count: observed_run.turn_count,
max_turns: observed_run.max_turns,
prompt_kind,
prompt_path,
conversation_id: manifest.conversation_id.clone(),
reuse_policy: manifest.reuse_policy.clone(),
fresh_conversation: manifest.fresh_conversation,
workflow_prompt_seeded: manifest.workflow_prompt_seeded,
server_base_url: manifest.server_base_url.clone(),
transport_target: manifest.transport_target.clone(),
http_auth_mode: manifest.http_auth_mode.clone(),
websocket_auth_mode: manifest.websocket_auth_mode.clone(),
websocket_query_param_name: manifest.websocket_query_param_name.clone(),
persistence_dir: manifest.persistence_dir.clone(),
last_execution_status: manifest.last_execution_status.clone(),
last_event_id: manifest.last_event_id.clone(),
last_event_kind: manifest.last_event_kind.clone(),
last_event_at: manifest.last_event_at,
last_event_summary: manifest.last_event_summary.clone(),
worker_outcome,
updated_at: Utc::now(),
}
}
fn default_reuse_policy() -> String {
DEFAULT_REUSE_POLICY.to_owned()
}
fn conversation_snapshot(stream: &RuntimeEventStream) -> Conversation {
let mut conversation = stream.conversation().clone();
if let Some(status) = stream.state_mirror().execution_status() {
conversation.execution_status = status.to_string();
}
conversation.without_mcp_credentials()
}
fn build_summary_metadata(
conversation: &Conversation,
fresh_conversation: bool,
stream_state: RuntimeStreamState,
diagnostics: Option<&super::TransportDiagnostics>,
server_base_url: &str,
) -> ConversationMetadata {
ConversationMetadata {
harness_capability: None,
conversation_id: ConversationId::new(conversation.conversation_id.to_string())
.expect("UUID-backed conversation ID should not be empty"),
server_base_url: Some(server_base_url.to_string()),
transport_target: diagnostics
.map(|diagnostics| diagnostics.target_kind.as_str().to_string()),
http_auth_mode: diagnostics
.map(|diagnostics| diagnostics.http_auth_kind.as_str().to_string()),
websocket_auth_mode: diagnostics
.map(|diagnostics| diagnostics.websocket_auth_kind.as_str().to_string()),
websocket_query_param_name: diagnostics
.and_then(|diagnostics| diagnostics.websocket_query_param_name.clone()),
fresh_conversation,
runtime_contract_version: Some(RUNTIME_CONTRACT_VERSION.to_string()),
stream_state,
last_event_id: None,
last_event_kind: None,
last_event_at: None,
last_event_summary: None,
recent_activity: Vec::new(),
input_tokens: 0,
output_tokens: 0,
cache_read_tokens: 0,
total_tokens: 0,
runtime_seconds: 0,
next_activity_sequence: 0,
}
}
fn observe_event<O>(observer: &mut O, event: &EventEnvelope)
where
O: IssueSessionObserver,
{
observer.on_runtime_event(
timestamp_ms_from_datetime(event.timestamp),
Some(event.id.clone()),
Some(event.kind.clone()),
Some(summarize_event(event)),
Some(runtime_event_payload(event)),
);
}
fn observe_reconciled_events<O>(
observer: &mut O,
events: &[EventEnvelope],
previously_observed: &HashSet<String>,
) where
O: IssueSessionObserver,
{
for event in events
.iter()
.filter(|event| !previously_observed.contains(&event.id))
{
observe_event(observer, event);
}
}
fn append_attempt_continuation(mut prompt: String, continuation: Option<&str>) -> String {
if let Some(continuation) = continuation {
prompt.push_str("\n\n");
prompt.push_str(continuation);
}
prompt
}
fn failed_outcome(summary: impl Into<String>, error: impl Into<String>) -> NormalizedOutcome {
NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: summary.into(),
error: Some(error.into()),
}
}
fn is_context_overflow_error(msg: &str) -> bool {
let lower = msg.to_ascii_lowercase();
lower.contains("prompt is too long")
|| lower.contains("maximum context length")
|| lower.contains("context window")
|| lower.contains("context_length")
|| lower.contains("token limit")
|| lower.contains("context length exceeded")
|| (lower.contains("exceeds") && lower.contains("context"))
}
fn is_context_overflow_outcome(outcome: &NormalizedOutcome) -> bool {
outcome.kind == WorkerOutcomeKind::Failed
&& outcome
.error
.as_deref()
.is_some_and(is_context_overflow_error)
}
fn is_condenser_tool_matching_error(msg: &str) -> bool {
msg.starts_with("'") && msg.ends_with("'") && msg.contains("chatcmpl-tool-")
}
fn is_condenser_tool_matching_outcome(outcome: &NormalizedOutcome) -> bool {
outcome.kind == WorkerOutcomeKind::Failed
&& outcome
.error
.as_deref()
.is_some_and(is_condenser_tool_matching_error)
}
fn extract_error_detail_from_state(state: &ConversationStateMirror) -> Option<String> {
let raw = state.raw_state();
if let Some(delta) = raw.get("state_delta")
&& let Some(error) = delta.get("last_error").and_then(|e: &Value| e.as_str())
{
return Some(error.to_string());
}
if let Some(stats) = raw.get("stats")
&& let Some(error) = stats.get("last_error").and_then(|e: &Value| e.as_str())
{
return Some(error.to_string());
}
if let Some(error) = raw.get("error").and_then(|e: &Value| e.as_str()) {
return Some(error.to_string());
}
None
}
fn summarize_event(event: &EventEnvelope) -> String {
match KnownEvent::from_envelope(event) {
KnownEvent::ConversationStateUpdate(payload) => match payload.execution_status {
Some(execution_status) => format!("status: {}", execution_status),
None => "state update".to_string(),
},
KnownEvent::ConversationError(_) => format!("error: {}", event.id),
KnownEvent::LlmCompletionLog(_) => "llm completion".to_string(),
KnownEvent::Message(msg) => {
let preview = msg.text_preview.unwrap_or_else(|| "message".to_string());
format!("{}: {}", msg.role, preview)
}
KnownEvent::Action(action) => {
let tool = action.tool_name.as_deref().unwrap_or("");
let command = action_command(&action.arguments);
let msg = action.message.as_deref().or(command).unwrap_or("action");
if tool.is_empty() {
msg.to_string()
} else {
format!("{}: {}", tool, msg)
}
}
KnownEvent::Observation(obs) => {
let tool = obs.tool_name.as_deref().unwrap_or("");
let preview = obs.text_preview.as_deref().unwrap_or("");
if tool.is_empty() {
if preview.is_empty() {
"result".to_string()
} else {
format!("→ {}", preview)
}
} else {
format!("{}: {}", tool, preview)
}
}
KnownEvent::Unknown(unknown) => unknown.kind,
}
}
fn runtime_event_payload(event: &EventEnvelope) -> Value {
match KnownEvent::from_envelope(event) {
KnownEvent::Action(action) => {
let mut body = json!({
"action_id": action.action_id,
"tool_name": action.tool_name,
"message": action.message,
});
if let (Value::Object(body), Value::Object(arguments)) = (&mut body, &action.arguments)
{
for (key, value) in arguments {
body.entry(key.clone()).or_insert_with(|| value.clone());
}
}
body
}
KnownEvent::Observation(observation) => json!({
"observation_id": observation.observation_id,
"tool_name": observation.tool_name,
"exit_code": observation.exit_code,
"preview": observation.text_preview,
"content": event
.payload
.get("observation")
.and_then(|value| value.get("content"))
.cloned()
.unwrap_or(Value::Null),
}),
KnownEvent::ConversationStateUpdate(payload) => json!({
"execution_status": payload.execution_status,
}),
KnownEvent::ConversationError(error) => error.payload,
KnownEvent::LlmCompletionLog(log) => log.payload,
KnownEvent::Message(message) => json!({
"role": message.role,
"preview": message.text_preview,
"content": event.payload.get("content").cloned().unwrap_or(Value::Null),
}),
KnownEvent::Unknown(unknown) => json!({
"kind": unknown.kind,
"payload": unknown.payload,
}),
}
}
fn action_command(arguments: &Value) -> Option<&str> {
arguments
.get("command")
.or_else(|| {
arguments
.get("arguments")
.and_then(|value| value.get("command"))
})
.or_else(|| arguments.get("args").and_then(|value| value.get("command")))
.and_then(Value::as_str)
.filter(|command| !command.trim().is_empty())
}
fn latest_current_turn_error(
events: &[EventEnvelope],
baseline_event_ids: &HashSet<String>,
) -> Option<String> {
events
.iter()
.rev()
.find(|event| {
!baseline_event_ids.contains(&event.id)
&& matches!(
KnownEvent::from_envelope(event),
KnownEvent::ConversationError(_)
)
})
.map(conversation_error_detail)
}
fn recovery_baseline_event_ids(
events: &[EventEnvelope],
last_event_id: Option<&str>,
) -> HashSet<String> {
let terminal_state_event_id = events.iter().rev().find_map(|event| {
let KnownEvent::ConversationStateUpdate(payload) = KnownEvent::from_envelope(event) else {
return None;
};
payload
.execution_status
.as_deref()
.filter(|status| matches!(*status, "finished" | "error" | "stuck"))
.map(|_| event.id.as_str())
});
if let Some(baseline_len) = last_event_id.and_then(|last_event_id| {
events
.iter()
.position(|event| event.id == last_event_id)
.map(|index| index + 1)
}) {
return events
.iter()
.take(baseline_len)
.filter(|event| Some(event.id.as_str()) != terminal_state_event_id)
.map(|event| event.id.clone())
.collect();
}
events
.iter()
.filter(|event| Some(event.id.as_str()) != terminal_state_event_id)
.map(|event| event.id.clone())
.collect()
}
fn all_event_ids(events: &[EventEnvelope]) -> HashSet<String> {
events.iter().map(|event| event.id.clone()).collect()
}
fn conversation_error_detail(event: &EventEnvelope) -> String {
let message = event
.payload
.get("message")
.and_then(Value::as_str)
.map(ToOwned::to_owned)
.unwrap_or_else(|| {
serde_json::to_string(&event.payload)
.unwrap_or_else(|_| "unable to encode ConversationErrorEvent payload".to_string())
});
format!("ConversationErrorEvent {}: {}", event.id, message)
}
fn run_status_for(outcome_kind: WorkerOutcomeKind) -> RunStatus {
match outcome_kind {
WorkerOutcomeKind::Succeeded => RunStatus::Succeeded,
WorkerOutcomeKind::Cancelled => RunStatus::Cancelled,
WorkerOutcomeKind::Failed
| WorkerOutcomeKind::TimedOut
| WorkerOutcomeKind::Stalled
| WorkerOutcomeKind::Detached
| WorkerOutcomeKind::CancelFailed => RunStatus::Failed,
}
}
pub struct RehydrationResult {
pub session: ActiveSession,
pub context: Option<String>,
pub old_conversation_id: String,
}
impl std::fmt::Debug for RehydrationResult {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RehydrationResult")
.field("session", &"<ActiveSession>")
.field("context", &self.context)
.field("old_conversation_id", &self.old_conversation_id)
.finish()
}
}
#[derive(Debug, Clone)]
pub struct RehydrationOptions {
pub reason: String,
pub summarize: bool,
pub max_summary_events: usize,
}
impl Default for RehydrationOptions {
fn default() -> Self {
Self {
reason: "explicit rehydration request".to_string(),
summarize: true,
max_summary_events: 50,
}
}
}
fn observed_run_for_turn(run: &RunAttempt) -> RunAttempt {
let mut observed_run = run.clone();
if observed_run.started_at.is_none() {
observed_run = observed_run.mark_started(timestamp_ms_from_datetime(Utc::now()));
}
observed_run.record_turn_started();
observed_run
}
fn timestamp_ms_from_datetime(value: DateTime<Utc>) -> TimestampMs {
TimestampMs::new(value.timestamp_millis().max(0) as u64)
}
fn compute_timeout_duration(
tracker: &LivenessTracker,
next_token_accumulation: Instant,
now: Instant,
now_ts: TimestampMs,
) -> Duration {
let token_remaining = next_token_accumulation.saturating_duration_since(now);
let stall_remaining = tracker
.stall_deadline_at()
.map_or(Duration::MAX, |deadline_ts| {
let remaining_ms = deadline_ts.as_u64().saturating_sub(now_ts.as_u64());
Duration::from_millis(remaining_ms)
});
token_remaining
.min(stall_remaining)
.max(Duration::from_millis(1))
}
fn parse_uuid(value: &str) -> Result<Uuid, String> {
Uuid::parse_str(value).map_err(|error| format!("invalid UUID `{value}`: {error}"))
}
fn create_conversation_request_path(workspace: &WorkspaceHandle) -> PathBuf {
workspace
.openhands_dir()
.join("create-conversation-request.json")
}
pub fn pending_conversation_manifest_path(workspace: &WorkspaceHandle) -> PathBuf {
workspace.openhands_dir().join("pending-conversation.json")
}
pub fn superseded_conversation_manifests_path(workspace: &WorkspaceHandle) -> PathBuf {
workspace
.openhands_dir()
.join("superseded-conversations.json")
}
fn last_conversation_state_path(workspace: &WorkspaceHandle) -> PathBuf {
workspace
.openhands_dir()
.join("last-conversation-state.json")
}
fn session_context_path(workspace: &WorkspaceHandle) -> PathBuf {
workspace.generated_dir().join("session-context.json")
}
fn finished_stream_error_is_tolerable(error: &OpenHandsError) -> bool {
matches!(
error,
OpenHandsError::ReconnectExhausted { .. } | OpenHandsError::WebSocketClosed
)
}
#[cfg(test)]
mod tests {
use std::{
collections::{BTreeMap, HashSet},
path::Path,
sync::Arc,
};
use crate::opensymphony_domain::{
BlockerRef, ConversationId, HarnessInterruptExpectedNextState, HarnessInterruptReason,
IssueIdentifier, IssueRef, IssueState, IssueStateCategory, WorkerOutcomeKind,
};
use crate::opensymphony_testkit::{FakeOpenHandsConfig, FakeOpenHandsServer};
use axum::{Json, Router, extract::State, routing::post};
use tokio::{net::TcpListener, sync::Mutex};
use super::super::TransportConfig;
use super::*;
fn must<T, E: std::fmt::Display>(result: Result<T, E>) -> T {
match result {
Ok(value) => value,
Err(error) => panic!("{error}"),
}
}
#[derive(Default)]
struct ReconciledEventObserver {
event_ids: Vec<String>,
}
impl IssueSessionObserver for ReconciledEventObserver {
fn on_runtime_event(
&mut self,
_observed_at: TimestampMs,
event_id: Option<String>,
_event_kind: Option<String>,
_summary: Option<String>,
_payload: Option<Value>,
) {
if let Some(event_id) = event_id {
self.event_ids.push(event_id);
}
}
}
#[test]
fn reconciliation_forwards_every_new_event_in_cache_order() {
let timestamp = Utc::now();
let events = vec![
EventEnvelope::new("prior", timestamp, "runtime", "MessageEvent", json!({})),
EventEnvelope::new(
"command-start",
timestamp + chrono::Duration::milliseconds(1),
"agent",
"ActionEvent",
json!({"command":"cargo test"}),
),
EventEnvelope::new(
"command-finish",
timestamp + chrono::Duration::milliseconds(2),
"tool",
"ObservationEvent",
json!({"exit_code":0}),
),
];
let mut observer = ReconciledEventObserver::default();
observe_reconciled_events(&mut observer, &events, &HashSet::from(["prior".to_owned()]));
assert_eq!(
observer.event_ids,
vec!["command-start".to_owned(), "command-finish".to_owned()]
);
}
#[tokio::test]
async fn prior_turn_wait_keeps_its_commands_out_of_the_current_attempt() {
let server = FakeOpenHandsServer::start()
.await
.expect("fake server should start");
let client = OpenHandsClient::new(TransportConfig::new(server.base_url()));
let conversation = client
.create_conversation(&ConversationCreateRequest::doctor_probe(
"/tmp/opensymphony-prior-turn",
"/tmp/opensymphony-prior-turn/.opensymphony/openhands",
Some("fake-model".to_string()),
None,
))
.await
.expect("conversation should be created");
let mut stream = client
.attach_runtime_stream(
conversation.conversation_id,
RuntimeStreamConfig {
readiness_timeout: Duration::from_secs(2),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
)
.await
.expect("runtime stream should attach");
server
.emit_state_update(conversation.conversation_id, "running")
.await
.expect("prior turn should run");
stream
.next_event()
.await
.expect("running event")
.expect("running event payload");
let prior_command = EventEnvelope::new(
"prior-command",
Utc::now(),
"agent",
"ActionEvent",
json!({"command": "cargo test --test prior"}),
);
server
.insert_event(conversation.conversation_id, prior_command)
.await
.expect("prior command should be delivered");
server
.emit_state_update(conversation.conversation_id, "finished")
.await
.expect("prior turn should finish");
let config = IssueSessionRunnerConfig {
terminal_wait_timeout: Duration::from_secs(2),
..IssueSessionRunnerConfig::default()
};
let runner = IssueSessionRunner::new(client, config);
runner
.wait_for_active_turn_to_finish(&mut stream)
.await
.expect("prior turn should drain");
let baseline_event_ids = stream
.event_cache()
.items()
.iter()
.map(|event| event.id.clone())
.collect::<HashSet<_>>();
assert!(baseline_event_ids.contains("prior-command"));
let current_command = EventEnvelope::new(
"current-command",
Utc::now() + chrono::Duration::milliseconds(1),
"agent",
"ActionEvent",
json!({"command": "cargo test --test current"}),
);
let mut observer = ReconciledEventObserver::default();
let mut events = stream.event_cache().items().to_vec();
events.push(current_command);
observe_reconciled_events(&mut observer, &events, &baseline_event_ids);
assert_eq!(observer.event_ids, vec!["current-command".to_owned()]);
}
#[test]
fn reused_conversation_receives_current_attempt_guidance() {
let prompt = append_attempt_continuation(
"Continue work on the existing issue conversation.".to_owned(),
Some("run_id=run-parent-2 attempt=2 receipt=evidence/final-verification.json"),
);
assert!(prompt.contains("Continue work on the existing issue conversation."));
assert!(prompt.contains("run_id=run-parent-2"));
assert!(prompt.contains("attempt=2"));
assert!(prompt.contains("evidence/final-verification.json"));
}
#[test]
fn openhands_binds_parent_scope_without_creating_a_leaf_envelope() {
let root = tempfile::tempdir().expect("temporary workspace root");
let path = root.path().join("parent-root");
std::fs::create_dir_all(&path).expect("parent root should exist");
let workspace =
WorkspaceHandle::new("parent-id", "COE-PARENT", "parent-COE-PARENT", path.clone());
let parent_envelope: ParentRuntimeEnvelope = serde_json::from_value(json!({
"parent_issue_id": "parent-id",
"parent_identifier": "COE-PARENT",
"run_id": "run-parent",
"attempt": 1,
"hierarchy_generation": 7,
"workspace_path": path,
"checkouts": {},
"harness": "openhands_agent_server",
"model_profile": "default",
"requested_execution_scope": "parent_multi_checkout",
"effective_containment": "trusted_host"
}))
.expect("parent envelope should decode");
let mut run_manifest = RunManifest::new(
&workspace,
&crate::opensymphony_workspace::RunDescriptor::new("run-parent", 1)
.with_parent_runtime_envelope(Some(parent_envelope)),
);
let profile = ConversationLaunchProfile {
workspace_kind: "LocalWorkspace".to_owned(),
confirmation_policy_kind: "NeverConfirm".to_owned(),
agent_kind: "Agent".to_owned(),
llm_model: "test-model".to_owned(),
llm_credential_mode: "api_key".to_owned(),
llm_api_key_env: None,
llm_base_url_env: None,
llm_subscription: None,
condenser: None,
agent_tools: None,
agent_include_default_tools: None,
max_iterations: 10,
stuck_detection: true,
llm_api_key_fingerprint: None,
};
let conversation_id = must(ConversationId::new("conversation-parent"));
let mut manifest = IssueConversationManifest::new(
must(IssueId::new("parent-id")),
must(IssueIdentifier::new("COE-PARENT")),
conversation_id.clone(),
"per_issue",
workspace.openhands_dir(),
Utc::now(),
None,
profile,
&BTreeMap::new(),
);
bind_runtime_envelopes(&mut manifest, &mut run_manifest);
assert!(manifest.runtime_envelope.is_none());
assert_eq!(
manifest
.parent_runtime_envelope
.as_ref()
.and_then(|envelope| envelope.conversation_binding.as_deref()),
Some(conversation_id.as_str())
);
assert_eq!(
manifest.parent_runtime_envelope,
run_manifest.parent_runtime_envelope
);
}
#[test]
fn launch_profile_constructs_openai_subscription_llm_without_persisting_tokens() {
let profile = ConversationLaunchProfile {
workspace_kind: "LocalWorkspace".to_string(),
confirmation_policy_kind: "NeverConfirm".to_string(),
agent_kind: "Agent".to_string(),
llm_model: "gpt-5.2-codex".to_string(),
llm_credential_mode: "openai_subscription".to_string(),
llm_api_key_env: None,
llm_base_url_env: None,
llm_subscription: Some(ConversationLaunchSubscriptionProfile {
vendor: "openai".to_string(),
access_token_env: "OPENHANDS_OPENAI_SUBSCRIPTION_ACCESS_TOKEN".to_string(),
account_id_env: Some("OPENHANDS_OPENAI_SUBSCRIPTION_ACCOUNT_ID".to_string()),
auth_directory_env: Some("OPENHANDS_AUTH_DIR".to_string()),
auth_method: "device_code".to_string(),
open_browser: false,
force_login: false,
}),
condenser: None,
agent_tools: None,
agent_include_default_tools: None,
max_iterations: 12,
stuck_detection: true,
llm_api_key_fingerprint: None,
};
let env = BTreeMap::from([
(
"OPENHANDS_OPENAI_SUBSCRIPTION_ACCESS_TOKEN".to_string(),
"oauth-access-token".to_string(),
),
(
"OPENHANDS_OPENAI_SUBSCRIPTION_ACCOUNT_ID".to_string(),
"account-123".to_string(),
),
(
"OPENHANDS_AUTH_DIR".to_string(),
"/Users/test/.cache/openhands/auth".to_string(),
),
]);
let request = profile
.to_create_request(
&env,
Path::new("/tmp/workspace"),
Path::new("/tmp/workspace/.opensymphony/openhands"),
Some(Uuid::nil()),
)
.expect("subscription profile should construct a request");
let profile_json = serde_json::to_string(&profile).expect("profile should serialize");
assert_eq!(request.agent.llm.model, "openai/gpt-5.2-codex");
assert_eq!(
request.agent.llm.base_url.as_deref(),
Some("https://chatgpt.com/backend-api/codex")
);
assert_eq!(
request.agent.llm.api_key.as_deref(),
Some("oauth-access-token")
);
assert_eq!(request.agent.llm.stream, Some(true));
assert_eq!(
profile
.llm_subscription
.as_ref()
.and_then(|subscription| subscription.auth_directory_env.as_deref()),
Some("OPENHANDS_AUTH_DIR")
);
assert_eq!(
env.get("OPENHANDS_AUTH_DIR").map(String::as_str),
Some("/Users/test/.cache/openhands/auth")
);
assert_eq!(
request
.agent
.llm
.extra_headers
.as_ref()
.and_then(|headers| headers.get("originator"))
.map(String::as_str),
Some("codex_cli_rs")
);
assert_eq!(
request
.agent
.llm
.litellm_extra_body
.as_ref()
.and_then(|body| body.get("store")),
Some(&json!(false))
);
assert!(!profile_json.contains("oauth-access-token"));
assert!(!profile_json.contains("account-123"));
assert!(!profile_json.contains("/Users/test/.cache/openhands/auth"));
assert!(profile_json.contains("OPENHANDS_AUTH_DIR"));
assert_eq!(
profile.api_key_fingerprint(&env).as_deref(),
Some("9e6b6d3f6b582838")
);
}
#[test]
fn malformed_superseded_conversation_evidence_is_not_treated_as_empty() {
let error = decode_superseded_conversation_manifests("{not-json")
.expect_err("malformed superseded evidence must fail closed");
assert!(matches!(
error,
IssueSessionError::SupersededConversationEvidence(_)
));
}
#[test]
fn openai_subscription_model_normalization_allows_future_openai_codex_models() {
assert_eq!(
must(normalize_openai_subscription_model("gpt-5.9-codex")),
"openai/gpt-5.9-codex"
);
assert_eq!(
must(normalize_openai_subscription_model("gpt-5.5")),
"openai/gpt-5.5"
);
assert_eq!(
must(normalize_openai_subscription_model("openai/gpt-5.9-codex")),
"openai/gpt-5.9-codex"
);
assert!(
normalize_openai_subscription_model("anthropic/claude-4").is_err(),
"subscription mode should remain scoped to OpenAI provider models"
);
}
#[test]
fn action_event_payload_surfaces_command_when_message_is_generic() {
let event = EventEnvelope::new(
"evt-action",
Utc::now(),
"agent",
"ActionEvent",
serde_json::json!({
"action": {
"tool_name": "terminal",
"arguments": { "command": "npm test -- apps/desktop" }
}
}),
);
assert_eq!(
summarize_event(&event),
"terminal: npm test -- apps/desktop"
);
let payload = runtime_event_payload(&event);
assert_eq!(
payload.get("tool_name").and_then(Value::as_str),
Some("terminal")
);
assert_eq!(
payload
.get("arguments")
.and_then(|arguments| arguments.get("command"))
.and_then(Value::as_str),
Some("npm test -- apps/desktop")
);
}
#[tokio::test]
async fn fetch_memory_context_from_server_calls_mcp_context_with_worker_scope() {
let requests = Arc::new(Mutex::new(Vec::<Value>::new()));
let listener = TcpListener::bind("127.0.0.1:0")
.await
.expect("memory test listener should bind");
let address = listener
.local_addr()
.expect("memory test listener should expose an address");
let app = Router::new()
.route("/mcp", post(memory_test_mcp))
.with_state(requests.clone());
let task = tokio::spawn(async move {
axum::serve(listener, app)
.await
.expect("memory test server should run");
});
let access = MemoryWorkerAccess {
endpoint: format!("http://{address}/mcp"),
token: Some("read-token".to_string()),
project: Some("project-alpha".to_string()),
project_set: Some("set-alpha".to_string()),
execution_repo: Some("/tmp/repo-alpha".to_string()),
authorized_repositories: vec!["repo-alpha".to_string()],
run_id: None,
attempt: None,
requires_fresh_conversation: false,
};
let issue = NormalizedIssue {
id: must(IssueId::new("issue-999")),
identifier: must(IssueIdentifier::new("COE-999")),
title: "Memory server context".to_string(),
description: Some("Use deterministic worker issue facts.".to_string()),
priority: None,
state: IssueState {
id: None,
name: "In Progress".to_string(),
category: IssueStateCategory::Active,
},
branch_name: None,
pr_url: None,
pr_urls: Vec::new(),
url: None,
labels: vec!["area:memory".to_string()],
project_id: None,
project_slug: None,
project_name: None,
parent_id: None,
repository_binding: None,
blocked_by: vec![BlockerRef {
id: Some(must(IssueId::new("issue-100"))),
identifier: Some(must(IssueIdentifier::new("COE-100"))),
state: Some("Done".to_string()),
created_at: None,
updated_at: None,
}],
sub_issues: vec![IssueRef {
id: must(IssueId::new("issue-101")),
identifier: must(IssueIdentifier::new("COE-101")),
state: "Done".to_string(),
}],
created_at: None,
updated_at: None,
};
let context = fetch_memory_context_from_server(&access, &issue)
.await
.expect("memory server context should load");
assert_eq!(context, "# Memory Context: COE-999\n");
let requests = requests.lock().await;
assert_eq!(requests.len(), 1);
let request = &requests[0];
assert_eq!(request["method"], "tools/call");
assert_eq!(request["params"]["name"], "memory.context");
assert_eq!(request["params"]["arguments"]["issue"], "COE-999");
assert_eq!(
request["params"]["arguments"]["currentIssue"]["labels"][0],
"area:memory"
);
assert_eq!(
request["params"]["arguments"]["currentIssue"]["children"][0]["identifier"],
"COE-101"
);
assert_eq!(
request["params"]["arguments"]["currentIssue"]["blockedBy"][0]["identifier"],
"COE-100"
);
assert_eq!(request["params"]["arguments"]["project"], "project-alpha");
assert_eq!(request["params"]["arguments"]["projectSet"], "set-alpha");
assert_eq!(request["params"]["arguments"]["repo"], "/tmp/repo-alpha");
task.abort();
}
#[test]
fn continuation_guidance_carries_worker_memory_scope() {
let guidance = append_memory_scope_guidance(
"continue from the existing conversation".to_owned(),
Some(&MemoryWorkerAccess {
endpoint: "http://127.0.0.1:8765/mcp".to_owned(),
token: None,
project: Some("project-alpha".to_owned()),
project_set: None,
execution_repo: Some("repo-alpha".to_owned()),
authorized_repositories: vec!["repo-alpha".to_owned()],
run_id: None,
attempt: None,
requires_fresh_conversation: false,
}),
);
assert!(guidance.contains("project=project-alpha"));
assert!(guidance.contains("repo=repo-alpha"));
assert!(guidance.contains("memory.context"));
}
#[test]
fn memory_worker_access_builds_a_scoped_mcp_server_config() {
let access = MemoryWorkerAccess {
endpoint: "http://127.0.0.1:8765/mcp".to_owned(),
token: Some("worker-secret".to_owned()),
project: Some("project-alpha".to_owned()),
project_set: None,
execution_repo: Some("repo-alpha".to_owned()),
authorized_repositories: vec!["repo-alpha".to_owned(), "repo-beta".to_owned()],
run_id: None,
attempt: None,
requires_fresh_conversation: false,
};
let config = access.mcp_config().expect("memory MCP config should exist");
assert_eq!(
config["mcpServers"]["opensymphony-memory"]["url"],
"http://127.0.0.1:8765/mcp"
);
assert_eq!(
config["mcpServers"]["opensymphony-memory"]["headers"]["Authorization"],
"Bearer worker-secret"
);
}
async fn memory_test_mcp(
State(requests): State<Arc<Mutex<Vec<Value>>>>,
Json(request): Json<Value>,
) -> Json<Value> {
requests.lock().await.push(request.clone());
Json(json!({
"jsonrpc": "2.0",
"id": request.get("id").cloned().unwrap_or(Value::Null),
"result": {
"content": [
{ "type": "text", "text": "# Memory Context: COE-999\n" }
]
}
}))
}
#[tokio::test]
async fn interrupt_acknowledges_after_reconciled_paused_state() {
let server = FakeOpenHandsServer::start_with_config(FakeOpenHandsConfig {
initial_execution_status: "running",
..FakeOpenHandsConfig::default()
})
.await
.expect("fake server should start");
let client = OpenHandsClient::new(TransportConfig::new(server.base_url()));
let conversation = client
.create_conversation(&ConversationCreateRequest::doctor_probe(
"/tmp/opensymphony-live",
"/tmp/opensymphony-live/.opensymphony/openhands",
Some("fake-model".to_string()),
None,
))
.await
.expect("conversation should be created");
let runner = IssueSessionRunner::new(
client,
IssueSessionRunnerConfig {
runtime_stream: RuntimeStreamConfig {
readiness_timeout: Duration::from_millis(500),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
..IssueSessionRunnerConfig::default()
},
);
let acknowledgement = runner
.interrupt(&interrupt_command(conversation.conversation_id))
.await
.expect("interrupt should acknowledge after reconciliation");
assert_eq!(acknowledgement.method, OpenHandsInterruptMethod::Interrupt);
assert_eq!(acknowledgement.execution_status.as_deref(), Some("paused"));
assert!(acknowledgement.diagnostic.is_none());
}
#[tokio::test]
async fn interrupt_falls_back_to_pause_when_interrupt_route_is_missing() {
let server = FakeOpenHandsServer::start_with_config(FakeOpenHandsConfig {
initial_execution_status: "running",
interrupt_supported: false,
..FakeOpenHandsConfig::default()
})
.await
.expect("fake server should start");
let client = OpenHandsClient::new(TransportConfig::new(server.base_url()));
let conversation = client
.create_conversation(&ConversationCreateRequest::doctor_probe(
"/tmp/opensymphony-live",
"/tmp/opensymphony-live/.opensymphony/openhands",
Some("fake-model".to_string()),
None,
))
.await
.expect("conversation should be created");
let runner = IssueSessionRunner::new(
client,
IssueSessionRunnerConfig {
runtime_stream: RuntimeStreamConfig {
readiness_timeout: Duration::from_millis(500),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
..IssueSessionRunnerConfig::default()
},
);
let acknowledgement = runner
.interrupt(&interrupt_command(conversation.conversation_id))
.await
.expect("pause fallback should acknowledge after reconciliation");
assert_eq!(
acknowledgement.method,
OpenHandsInterruptMethod::PauseFallback
);
assert_eq!(acknowledgement.execution_status.as_deref(), Some("paused"));
assert!(
acknowledgement
.diagnostic
.as_deref()
.is_some_and(|diagnostic| diagnostic.contains("/interrupt unavailable"))
);
}
#[tokio::test]
async fn interrupt_reports_timeout_when_reconciliation_never_stops() {
let server = FakeOpenHandsServer::start_with_config(FakeOpenHandsConfig {
initial_execution_status: "running",
interrupt_execution_status: "running",
..FakeOpenHandsConfig::default()
})
.await
.expect("fake server should start");
let client = OpenHandsClient::new(TransportConfig::new(server.base_url()));
let conversation = client
.create_conversation(&ConversationCreateRequest::doctor_probe(
"/tmp/opensymphony-live",
"/tmp/opensymphony-live/.opensymphony/openhands",
Some("fake-model".to_string()),
None,
))
.await
.expect("conversation should be created");
let runner = IssueSessionRunner::new(
client,
IssueSessionRunnerConfig {
runtime_stream: RuntimeStreamConfig {
readiness_timeout: Duration::from_millis(20),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
..IssueSessionRunnerConfig::default()
},
);
let acknowledgement = runner
.interrupt(&interrupt_command(conversation.conversation_id))
.await
.expect("timeout path should return a diagnostic acknowledgement");
assert_eq!(acknowledgement.method, OpenHandsInterruptMethod::Interrupt);
assert_eq!(acknowledgement.execution_status.as_deref(), Some("running"));
assert!(
acknowledgement
.diagnostic
.as_deref()
.is_some_and(|diagnostic| diagnostic
.contains("timed out waiting for interrupt acknowledgement"))
);
}
#[test]
fn recovery_client_uses_trusted_persisted_server_base_url_and_current_auth() {
let client = OpenHandsClient::new(TransportConfig::new("http://127.0.0.1:8000").with_auth(
super::super::AuthConfig::header_api_key("x-session-api-key", "redacted"),
));
let recovery_client = client
.with_persisted_base_url("http://127.0.0.1:8000/api")
.expect("same-origin recovery transport should remain trusted");
assert_eq!(recovery_client.base_url(), "http://127.0.0.1:8000/api");
let diagnostics = recovery_client
.transport_diagnostics()
.expect("recovery transport should remain valid");
assert_eq!(
diagnostics.http_auth_kind,
super::super::TransportAuthKind::Header
);
}
#[test]
fn recovery_client_rejects_changed_persisted_server_origin() {
let client = OpenHandsClient::new(TransportConfig::new("http://127.0.0.1:8000").with_auth(
super::super::AuthConfig::header_api_key("x-session-api-key", "redacted"),
));
let error = match client.with_persisted_base_url("http://attacker.example.test:9000") {
Ok(_) => panic!("recovery must not carry auth to a changed origin"),
Err(error) => error,
};
assert!(
error
.to_string()
.contains("does not match the configured OpenHands server")
);
}
fn interrupt_command(conversation_id: Uuid) -> HarnessInterruptCommand {
HarnessInterruptCommand {
run_id: "COE-489".to_string(),
issue_id: must(IssueId::new("issue-489")),
harness_kind: "openhands_agent_server".to_string(),
conversation_id: must(ConversationId::new(conversation_id.to_string())),
turn_id: None,
reason: HarnessInterruptReason::OperatorCancel,
expected_next_state: HarnessInterruptExpectedNextState::Paused,
}
}
#[tokio::test]
async fn await_terminal_outcome_accepts_reconciled_finished_state_after_stream_close() {
let server = FakeOpenHandsServer::start()
.await
.expect("fake server should start");
let client = OpenHandsClient::new(TransportConfig::new(server.base_url()));
let conversation = client
.create_conversation(&ConversationCreateRequest::doctor_probe(
"/tmp/opensymphony-live",
"/tmp/opensymphony-live/.opensymphony/openhands",
Some("fake-model".to_string()),
None,
))
.await
.expect("conversation should be created");
let mut stream = client
.attach_runtime_stream(
conversation.conversation_id,
RuntimeStreamConfig {
readiness_timeout: Duration::from_secs(2),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
)
.await
.expect("runtime stream should attach");
let baseline_event_ids = stream
.event_cache()
.items()
.iter()
.map(|event| event.id.clone())
.collect::<HashSet<_>>();
server
.emit_state_update(conversation.conversation_id, "running")
.await
.expect("running state should be recorded");
server
.emit_state_update(conversation.conversation_id, "finished")
.await
.expect("finished state should be recorded");
stream.close().await.expect("stream should close cleanly");
let runner = IssueSessionRunner::new(
client,
IssueSessionRunnerConfig {
reuse_policy: IssueSessionReusePolicy::PerIssue,
runtime_stream: RuntimeStreamConfig {
readiness_timeout: Duration::from_secs(2),
reconnect_initial_backoff: Duration::from_millis(25),
reconnect_max_backoff: Duration::from_millis(25),
max_reconnect_attempts: 1,
replay_existing_events_on_attach: false,
},
terminal_wait_timeout: Duration::from_millis(25),
total_runtime_cap_ms: None,
finished_drain_timeout: Duration::from_millis(25),
memory: None,
repository_instructions: None,
terminal_prompt: None,
continuation_prompt: None,
},
);
let manifest = IssueConversationManifest {
issue_id: must(IssueId::new("lin_test")),
identifier: must(IssueIdentifier::new("TEST-123")),
conversation_id: must(ConversationId::new(
conversation.conversation_id.to_string(),
)),
reuse_policy: "per_issue".to_string(),
server_base_url: None,
transport_target: None,
http_auth_mode: None,
websocket_auth_mode: None,
websocket_query_param_name: None,
persistence_dir: PathBuf::from("/tmp/test"),
created_at: Utc::now(),
updated_at: Utc::now(),
last_attached_at: Utc::now(),
launch_profile: None,
llm_config_fingerprint: None,
fresh_conversation: true,
workflow_prompt_seeded: false,
reset_reason: None,
runtime_contract_version: None,
runtime_envelope: None,
parent_runtime_envelope: None,
codex_archive_state: None,
last_turn_id: None,
active_run_id: None,
prepared_run_id: None,
trigger_pending_run_id: None,
last_prompt_kind: None,
last_prompt_at: None,
last_prompt_path: None,
last_execution_status: None,
last_event_id: None,
last_event_kind: None,
last_event_at: None,
last_event_summary: None,
input_tokens: 0,
output_tokens: 0,
cache_read_tokens: 0,
last_token_accumulation_at: None,
};
let mut session = ActiveSession {
stream,
manifest,
prompt_kind: IssueSessionPromptKind::Full,
prompt_path: None,
};
let outcome = runner
.await_terminal_outcome(&mut session, &baseline_event_ids, &mut ())
.await;
assert_eq!(outcome.kind, WorkerOutcomeKind::Succeeded);
assert_eq!(
session.stream.state_mirror().terminal_status(),
Some(TerminalExecutionStatus::Finished)
);
}
#[test]
fn recovery_baseline_includes_events_through_persisted_marker() {
let events = vec![
EventEnvelope::new("old-1", Utc::now(), "runtime", "old", Value::Null),
EventEnvelope::new("old-2", Utc::now(), "runtime", "old", Value::Null),
EventEnvelope::new("current", Utc::now(), "runtime", "current", Value::Null),
];
assert_eq!(
recovery_baseline_event_ids(&events, Some("old-2")),
HashSet::from(["old-1".to_owned(), "old-2".to_owned()])
);
assert_eq!(
recovery_baseline_event_ids(&events, Some("missing")),
HashSet::from(["old-1".to_owned(), "old-2".to_owned(), "current".to_owned()])
);
}
#[test]
fn recovery_baseline_keeps_unmarked_terminal_state_current() {
let events = vec![
EventEnvelope::new(
"old-error",
Utc::now(),
"runtime",
"ConversationErrorEvent",
serde_json::json!({"message": "old failure"}),
),
EventEnvelope::new(
"terminal",
Utc::now(),
"runtime",
"ConversationStateUpdateEvent",
serde_json::json!({"execution_status": "finished"}),
),
];
assert_eq!(
recovery_baseline_event_ids(&events, None),
HashSet::from(["old-error".to_owned()])
);
}
#[test]
fn recovery_baseline_keeps_persisted_terminal_state_current() {
let events = vec![
EventEnvelope::new("old", Utc::now(), "runtime", "old", Value::Null),
EventEnvelope::new(
"terminal",
Utc::now(),
"runtime",
"ConversationStateUpdateEvent",
serde_json::json!({"execution_status": "finished"}),
),
];
assert_eq!(
recovery_baseline_event_ids(&events, Some("terminal")),
HashSet::from(["old".to_owned()])
);
}
#[test]
fn trigger_pending_recovery_baselines_all_pre_trigger_events() {
let events = vec![
EventEnvelope::new("old", Utc::now(), "runtime", "old", Value::Null),
EventEnvelope::new(
"terminal",
Utc::now(),
"runtime",
"ConversationStateUpdateEvent",
serde_json::json!({"execution_status": "finished"}),
),
];
assert_eq!(
all_event_ids(&events),
HashSet::from(["old".to_owned(), "terminal".to_owned()])
);
}
#[test]
fn token_accumulation_extracts_tokens_from_llm_completion_events() {
use chrono::Utc;
use serde_json::json;
let event1 = EventEnvelope::new(
"evt-1",
Utc::now(),
"llm",
"LLMCompletionLogEvent",
json!({
"model": "gpt-4",
"usage": {
"prompt_tokens": 100,
"completion_tokens": 50,
"total_tokens": 150
}
}),
);
let event2 = EventEnvelope::new(
"evt-2",
Utc::now() + chrono::Duration::milliseconds(100),
"llm",
"LLMCompletionLogEvent",
json!({
"model": "gpt-4",
"input_tokens": 200,
"output_tokens": 75
}),
);
let event3 = EventEnvelope::new(
"evt-3",
Utc::now() + chrono::Duration::milliseconds(200),
"llm",
"LLMCompletionLogEvent",
json!({
"model": "gpt-4",
"tokens": 300
}),
);
let mut cache = super::super::events::EventCache::new();
cache.insert(event1);
cache.insert(event2);
cache.insert(event3);
let mut total_input = 0u64;
let mut total_output = 0u64;
for event in cache.items() {
if let KnownEvent::LlmCompletionLog(llm_event) = KnownEvent::from_envelope(event)
&& let Some((input, output)) = llm_event.token_usage()
{
total_input += input;
total_output += output;
}
}
assert_eq!(total_input, 300, "input tokens should be 100 + 200 + 0");
assert_eq!(total_output, 425, "output tokens should be 50 + 75 + 300");
}
#[test]
fn condenser_tool_matching_error_detection() {
assert!(is_condenser_tool_matching_error(
"'chatcmpl-tool-bcea2761df6a8821'"
));
assert!(is_condenser_tool_matching_error("'chatcmpl-tool-abc123'"));
assert!(!is_condenser_tool_matching_error(
"prompt is too long: 1000"
));
assert!(!is_condenser_tool_matching_error("some other error"));
assert!(!is_condenser_tool_matching_error(
"tool error without quotes"
));
assert!(!is_condenser_tool_matching_error(
"'not-a-tool-matching-error'"
));
assert!(!is_condenser_tool_matching_error("'no-tool-here'"));
assert!(!is_condenser_tool_matching_error("'other-tool-id-123'"));
let tool_matching_outcome = NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "test".to_string(),
error: Some("'chatcmpl-tool-xyz'".to_string()),
};
assert!(is_condenser_tool_matching_outcome(&tool_matching_outcome));
let other_outcome = NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "test".to_string(),
error: Some("some other error".to_string()),
};
assert!(!is_condenser_tool_matching_outcome(&other_outcome));
let context_overflow_outcome = NormalizedOutcome {
kind: WorkerOutcomeKind::Failed,
summary: "test".to_string(),
error: Some("prompt is too long: 1000".to_string()),
};
assert!(!is_condenser_tool_matching_outcome(
&context_overflow_outcome
));
assert!(is_context_overflow_outcome(&context_overflow_outcome));
}
#[test]
fn liveness_tracker_slides_deadline_on_progress_signals() {
let idle = DurationMs::new(300_000); let cap = DurationMs::new(3_600_000); let mut tracker = LivenessTracker::with_runtime_cap(idle, Some(cap));
let start = TimestampMs::new(1000);
assert!(!tracker.is_stalled_at(TimestampMs::new(1000)));
tracker.mark_started(start);
assert!(!tracker.is_stalled_at(start));
assert!(!tracker.is_stalled_at(TimestampMs::new(1000 + 240_000)));
let t3 = TimestampMs::new(1000 + 180_000);
tracker.record_event(t3);
assert!(!tracker.is_stalled_at(TimestampMs::new(1000 + 180_000 + 240_000)));
let t6 = TimestampMs::new(1000 + 360_000);
tracker.record_tokens(100, 50, t6);
assert!(!tracker.is_stalled_at(TimestampMs::new(1000 + 360_000 + 240_000)));
assert!(tracker.is_stalled_at(TimestampMs::new(1000 + 360_000 + 360_000)));
tracker.mark_started(start);
tracker.record_event(TimestampMs::new(start.as_u64() + 3_500_000));
assert!(tracker.is_stalled_at(TimestampMs::new(start.as_u64() + 3_600_000 + 1)));
}
#[test]
fn liveness_tracker_snapshot_produces_correct_deltas() {
let mut tracker = LivenessTracker::new(DurationMs::new(60_000));
tracker.mark_started(TimestampMs::new(1000));
tracker.record_event(TimestampMs::new(1100));
tracker.record_tokens(50, 25, TimestampMs::new(1200));
let initial = RuntimeProgressSnapshot::initial(RuntimeLivenessPhase::RunningTurn);
let snapshot = tracker.snapshot(&initial);
assert_eq!(snapshot.phase, RuntimeLivenessPhase::RunningTurn);
assert_eq!(snapshot.event_count, 1);
assert_eq!(snapshot.event_delta, 1);
assert_eq!(snapshot.input_tokens, 50);
assert_eq!(snapshot.input_token_delta, 50);
assert_eq!(snapshot.output_tokens, 25);
assert_eq!(snapshot.output_token_delta, 25);
let snapshot2 = tracker.snapshot(&snapshot);
assert_eq!(snapshot2.event_delta, 0);
assert_eq!(snapshot2.input_token_delta, 0);
assert_eq!(snapshot2.output_token_delta, 0);
let unchanged = tracker.record_tokens(50, 25, TimestampMs::new(1300));
assert!(!unchanged);
let snapshot3 = tracker.snapshot(&snapshot2);
assert_eq!(snapshot3.input_token_delta, 0);
assert_eq!(snapshot3.output_token_delta, 0);
tracker.record_event(TimestampMs::new(1300));
tracker.record_tokens(60, 30, TimestampMs::new(1400));
let snapshot4 = tracker.snapshot(&snapshot3);
assert_eq!(snapshot4.event_delta, 1);
assert_eq!(snapshot4.input_token_delta, 10);
assert_eq!(snapshot4.output_token_delta, 5);
}
#[test]
fn liveness_tracker_waiting_on_prior_turn_phase() {
let tracker = LivenessTracker::new(DurationMs::new(60_000));
let initial = RuntimeProgressSnapshot::initial(RuntimeLivenessPhase::WaitingOnPriorTurn);
let snapshot = tracker.snapshot(&initial);
assert_eq!(snapshot.phase, RuntimeLivenessPhase::WaitingOnPriorTurn);
assert_eq!(snapshot.event_count, 0);
assert!(snapshot.last_activity_at.is_none());
}
#[test]
fn liveness_tracker_status_change_is_liveness_signal() {
let mut tracker = LivenessTracker::new(DurationMs::new(60_000));
tracker.mark_started(TimestampMs::new(1000));
let changed = tracker.record_status_change("running", TimestampMs::new(1100));
assert!(changed);
let changed = tracker.record_status_change("running", TimestampMs::new(1200));
assert!(!changed);
let changed = tracker.record_status_change("finished", TimestampMs::new(1300));
assert!(changed);
}
#[test]
fn compute_timeout_duration_uses_tracker_deadline_with_sampled_logical_time() {
let mut tracker = LivenessTracker::new(DurationMs::new(5_000));
tracker.mark_started(TimestampMs::new(1_000));
let now = Instant::now();
let timeout = compute_timeout_duration(
&tracker,
now + Duration::from_secs(30),
now,
TimestampMs::new(3_000),
);
assert_eq!(timeout, Duration::from_millis(3_000));
}
#[test]
fn compute_timeout_duration_uses_token_poll_when_tracker_has_no_deadline() {
let tracker = LivenessTracker::new(DurationMs::new(5_000));
let now = Instant::now();
let timeout = compute_timeout_duration(
&tracker,
now + Duration::from_millis(750),
now,
TimestampMs::new(3_000),
);
assert_eq!(timeout, Duration::from_millis(750));
}
#[test]
fn create_error_diagnostic_redacts_worker_grants() {
let mut request =
ConversationCreateRequest::doctor_probe("/tmp/checkout", "/tmp/state", None, None);
request.agent.mcp_config = Some(BTreeMap::from([(
"mcpServers".to_owned(),
serde_json::json!({
"opensymphony-memory": {
"headers": { "Authorization": "Bearer worker-secret" }
}
}),
)]));
let error = OpenHandsError::HttpStatus {
operation: "create_conversation",
status_code: 422,
body: r#"{"headers":{"Authorization":"Bearer worker-secret"}}"#.to_owned(),
};
let detail = redacted_openhands_error(&error, &request);
assert!(!detail.contains("worker-secret"));
assert!(detail.contains("<redacted>"));
}
#[tokio::test]
async fn conversation_workspace_must_match_verified_checkout() {
let root = tempfile::tempdir().expect("temporary workspace root should exist");
let checkout = root.path().join("checkout");
let other = root.path().join("other");
std::fs::create_dir_all(&checkout).expect("checkout should exist");
std::fs::create_dir_all(&other).expect("other workspace should exist");
let mut conversation = Conversation {
conversation_id: Uuid::new_v4(),
workspace: WorkspaceConfig {
working_dir: checkout.display().to_string(),
kind: "LocalWorkspace".to_owned(),
},
persistence_dir: root.path().join("state").display().to_string(),
max_iterations: 1,
stuck_detection: true,
execution_status: "idle".to_owned(),
confirmation_policy: ConfirmationPolicy {
kind: "NeverConfirm".to_owned(),
},
agent: AgentConfig {
kind: "Agent".to_owned(),
llm: LlmConfig {
model: "test".to_owned(),
api_key: None,
base_url: None,
usage_id: None,
extra_headers: None,
litellm_extra_body: None,
stream: None,
},
condenser: None,
tools: None,
mcp_config: None,
include_default_tools: None,
},
stats: None,
};
verify_conversation_workspace(&conversation, &checkout)
.await
.expect("matching canonical workspace should be accepted");
conversation.workspace.working_dir = other.display().to_string();
let error = verify_conversation_workspace(&conversation, &checkout)
.await
.expect_err("different canonical workspace should be rejected");
assert!(error.contains("does not match verified checkout"));
}
}