use std::{
collections::{HashSet, VecDeque},
path::Path,
};
use crate::opensymphony_control::{
AgentServerStatus, ConversationEvent, DaemonSnapshot, DaemonState, DaemonStatus,
IssueRuntimeState, IssueSnapshot, MemoryServerStatus, MetricsSnapshot, RecentEvent,
RecentEventKind, WorkerOutcome,
};
use crate::opensymphony_domain::{
HarnessInterruptStatus, HealthStatus, IssueIdentifier, IssueStateCategory,
OrchestratorSnapshot, SchedulerStatus, WorkerOutcomeKind,
};
use crate::opensymphony_openhands::LocalServerSupervisor;
use crate::opensymphony_workflow::ResolvedWorkflow;
use chrono::{DateTime, Utc};
use super::timestamp_to_datetime;
const RECENT_EVENT_LIMIT: usize = 24;
pub(super) fn map_snapshot(
snapshot: &OrchestratorSnapshot,
workspace_root: &Path,
terminal_states: &HashSet<String>,
agent_server: AgentServerStatus,
memory_server: MemoryServerStatus,
recent_events: &VecDeque<RecentEvent>,
) -> DaemonSnapshot {
let generated_at = timestamp_to_datetime(snapshot.generated_at);
let last_poll_at = snapshot
.daemon
.last_poll_at
.map(timestamp_to_datetime)
.unwrap_or(generated_at);
DaemonSnapshot {
generated_at,
daemon: DaemonStatus {
state: map_daemon_state(snapshot.daemon.health),
last_poll_at,
workspace_root: workspace_root.display().to_string(),
status_line: format!(
"poll={}ms, running={}, retry_queue={}",
snapshot.daemon.poll_interval_ms,
snapshot.daemon.running_issue_count,
snapshot.daemon.retry_queue_count
),
},
agent_server,
memory_server,
metrics: MetricsSnapshot {
running_issues: snapshot.daemon.running_issue_count as u32,
retry_queue_depth: snapshot.daemon.retry_queue_count as u32,
input_tokens: snapshot.daemon.usage.input_tokens,
output_tokens: snapshot.daemon.usage.output_tokens,
cache_read_tokens: snapshot.daemon.usage.cache_read_tokens,
total_tokens: snapshot.daemon.usage.total_tokens,
total_cost_micros: snapshot.daemon.usage.estimated_cost_usd_micros.unwrap_or(0),
},
issues: snapshot
.issues
.iter()
.map(|issue| {
map_issue(
issue,
terminal_states,
generated_at,
snapshot.hierarchy.get(issue.issue.id.as_str()),
&snapshot.operator_interactions,
)
})
.collect(),
recent_events: recent_events.iter().cloned().collect(),
}
}
pub(super) fn current_memory_server_status(
memory_server: Option<&super::super::memory::MemoryServerHandle>,
) -> MemoryServerStatus {
let Some(memory_server) = memory_server else {
return MemoryServerStatus::default();
};
let reachable = !memory_server.is_finished();
MemoryServerStatus {
enabled: true,
reachable,
endpoint: Some(memory_server.endpoint().to_string()),
status_line: if reachable {
"listening".to_string()
} else {
"degraded: memory server stopped; scoped worker reads are blocked".to_string()
},
}
}
fn map_issue(
issue: &crate::opensymphony_domain::IssueSnapshot,
terminal_states: &HashSet<String>,
generated_at: DateTime<Utc>,
hierarchy: Option<&crate::opensymphony_domain::HierarchyStateSnapshot>,
operator_interactions: &[crate::opensymphony_gateway_schema::approval::OperatorInteraction],
) -> IssueSnapshot {
let runtime_state = match issue.runtime.state {
SchedulerStatus::Running | SchedulerStatus::Claimed => IssueRuntimeState::Running,
SchedulerStatus::RetryQueued => IssueRuntimeState::RetryQueued,
SchedulerStatus::Released => {
if issue.issue.state.category == IssueStateCategory::Terminal {
IssueRuntimeState::Completed
} else {
match issue
.last_worker_outcome
.as_ref()
.map(|outcome| outcome.outcome)
{
Some(
WorkerOutcomeKind::Failed
| WorkerOutcomeKind::TimedOut
| WorkerOutcomeKind::Stalled
| WorkerOutcomeKind::Detached
| WorkerOutcomeKind::CancelFailed,
) => IssueRuntimeState::Failed,
_ if issue.runtime.release_reason
== Some(crate::opensymphony_domain::ReleaseReason::RetryExhausted) =>
{
IssueRuntimeState::Failed
}
None => IssueRuntimeState::Idle,
_ => IssueRuntimeState::Completed,
}
}
}
SchedulerStatus::Unclaimed => IssueRuntimeState::Idle,
};
let last_outcome = map_worker_outcome(issue, runtime_state);
let last_event_at = issue
.conversation
.as_ref()
.and_then(|conversation| conversation.last_event_at)
.map(timestamp_to_datetime)
.or_else(|| {
issue
.last_worker_outcome
.as_ref()
.map(|outcome| timestamp_to_datetime(outcome.finished_at))
})
.unwrap_or(generated_at);
let worker = issue.runtime.worker.as_ref();
let last_worker_outcome = issue.last_worker_outcome.as_ref();
let interrupt = issue.runtime.interrupt.as_ref();
let started_at = issue
.runtime
.started_at
.or_else(|| last_worker_outcome.map(|outcome| outcome.started_at))
.map(timestamp_to_datetime);
let finished_at = issue
.runtime
.released_at
.or_else(|| last_worker_outcome.map(|outcome| outcome.finished_at))
.map(timestamp_to_datetime);
let runtime_seconds = issue
.conversation
.as_ref()
.map(|conversation| conversation.runtime_seconds)
.unwrap_or(0)
.max(runtime_seconds_from_timestamps(
started_at,
finished_at,
generated_at,
runtime_state,
));
let repository_binding = issue.issue.repository_binding.clone();
let repository_binding_blocked = repository_binding
.as_ref()
.is_some_and(|binding| binding.resolved_binding().is_none());
let hierarchy_blocked = hierarchy.and_then(|state| state.blocked_reason.clone());
let operator = Some(crate::opensymphony_domain::ControlPlaneOperatorSnapshot {
routing_mode: None,
active_project_set: Vec::new(),
linear_project: issue
.issue
.project_slug
.clone()
.or_else(|| issue.issue.project_id.clone()),
binding_status: repository_binding.as_ref().map(|binding| match binding {
crate::opensymphony_domain::RepositoryBindingOutcome::Resolved(_) => "resolved",
crate::opensymphony_domain::RepositoryBindingOutcome::MissingBinding => "missing_binding",
crate::opensymphony_domain::RepositoryBindingOutcome::UnknownAlias(_) => "unknown_alias",
crate::opensymphony_domain::RepositoryBindingOutcome::MultipleBindings(_) => "multiple_bindings",
crate::opensymphony_domain::RepositoryBindingOutcome::RepositoryNotAllowedForProject(_, _) => "repository_not_allowed_for_project",
crate::opensymphony_domain::RepositoryBindingOutcome::ParentBindingNotAllowed => "parent_binding_not_allowed",
crate::opensymphony_domain::RepositoryBindingOutcome::ProjectOutsideActiveSet(_) => "project_outside_active_set",
}.to_owned()),
parent: hierarchy.map(
|state| crate::opensymphony_domain::ControlPlaneParentSnapshot {
parent_id: issue.issue.identifier.to_string(),
state: None,
hierarchy_generation: Some(state.generation),
blocked_reason: hierarchy_blocked.clone(),
descendant_repositories: Vec::new(),
checkout_handles: Vec::new(),
},
),
repository: repository_binding.as_ref().and_then(|binding| {
binding.resolved_binding().map(|binding| {
crate::opensymphony_domain::ControlPlaneRepositorySnapshot {
canonical_id: binding.repository.id.as_str().to_owned(),
display_alias: binding.alias.clone(),
safe_remote_fingerprint: Some(
binding
.repository
.safe_remote_fingerprint
.as_str()
.to_owned(),
),
config_generation: Some(binding.config_generation.clone()),
inventory_generation: Some(binding.inventory_generation.clone()),
checkout_generation: None,
target_branch: None,
target_commit: None,
instruction_source: None,
instruction_hash: None,
}
})
}),
leases: Vec::new(),
repairs: Vec::new(),
memory: None,
containment: None,
provider: None,
verification: None,
cleanup: None,
});
IssueSnapshot {
operator_interactions: operator_interactions
.iter()
.filter(|request| request.issue_id == issue.issue.id.as_str())
.cloned()
.collect(),
harness_capability: issue
.conversation
.as_ref()
.and_then(|c| c.harness_capability.as_deref().cloned()),
identifier: issue.issue.identifier.to_string(),
title: issue.issue.title.clone(),
tracker_state: issue.issue.state.name.clone(),
runtime_state,
last_outcome,
last_event_at,
conversation_id_suffix: issue
.conversation
.as_ref()
.map(|conversation| suffix(conversation.conversation_id.as_str()))
.unwrap_or_else(|| "-".to_string()),
codex_thread_id: issue.conversation.as_ref().and_then(|conversation| {
let is_codex = conversation.transport_target.as_deref()
== Some(crate::opensymphony_codex::CODEX_APP_SERVER_KIND)
|| conversation.runtime_contract_version.as_deref()
== Some(crate::opensymphony_codex::CODEX_APP_SERVER_CONTRACT);
let is_route_preview = conversation
.conversation_id
.as_str()
.starts_with(super::backends::ROUTE_PREVIEW_CONVERSATION_PREFIX);
(is_codex && !is_route_preview)
.then(|| conversation.conversation_id.as_str().to_string())
}),
workspace_path_suffix: issue
.workspace
.as_ref()
.map(|workspace| suffix_path(&workspace.path))
.unwrap_or_else(|| "-".to_string()),
branch_name: issue.issue.branch_name.clone(),
pr_url: issue.issue.pr_url.clone(),
project_id: issue.issue.project_id.clone(),
project_slug: issue.issue.project_slug.clone(),
project_name: issue.issue.project_name.clone(),
workspace_label: issue
.workspace
.as_ref()
.map(|workspace| suffix_path(&workspace.path))
.filter(|label| label != "-"),
retry_count: issue
.retry
.as_ref()
.map(|retry| retry.normal_retry_count)
.or(issue.retry_count_override)
.or_else(|| {
matches!(
issue.runtime.release_reason,
Some(
crate::opensymphony_domain::ReleaseReason::RetryExhausted
| crate::opensymphony_domain::ReleaseReason::Completed
)
)
.then(|| {
issue
.last_worker_outcome
.as_ref()
.and_then(|outcome| outcome.attempt.map(|attempt| attempt.get()))
.unwrap_or(0)
})
})
.unwrap_or(0),
release_reason: issue.runtime.release_reason,
claimed_at: issue.runtime.claimed_at.map(timestamp_to_datetime),
started_at,
finished_at,
turn_count: worker
.map(|worker| worker.turn_count)
.or_else(|| last_worker_outcome.map(|outcome| outcome.turn_count))
.unwrap_or(0),
max_turns: worker.map(|worker| worker.max_turns).unwrap_or(0),
runtime_seconds,
blocked: repository_binding_blocked
|| hierarchy_blocked.is_some()
|| issue.issue.blocked_by.iter().any(|blocker| {
blocker
.state
.as_deref()
.is_none_or(|state| !is_terminal_state(terminal_states, state))
})
|| (!issue.issue.sub_issues.is_empty()
&& issue
.issue
.sub_issues
.iter()
.any(|sub_issue| !is_terminal_state(terminal_states, &sub_issue.state))),
repository_binding,
hierarchy_generation: hierarchy.map(|state| state.generation),
hierarchy_blocked_reason: hierarchy_blocked,
blocked_by: issue
.issue
.blocked_by
.iter()
.filter_map(|blocker| blocker.identifier.as_ref())
.map(ToString::to_string)
.collect(),
server_base_url: issue
.conversation
.as_ref()
.and_then(|conversation| conversation.server_base_url.clone()),
transport_target: issue
.conversation
.as_ref()
.and_then(|conversation| conversation.transport_target.clone()),
http_auth_mode: issue
.conversation
.as_ref()
.and_then(|conversation| conversation.http_auth_mode.clone()),
websocket_auth_mode: issue
.conversation
.as_ref()
.and_then(|conversation| conversation.websocket_auth_mode.clone()),
websocket_query_param_name: issue
.conversation
.as_ref()
.and_then(|conversation| conversation.websocket_query_param_name.clone()),
recent_events: issue
.conversation
.as_ref()
.map(|conversation| {
conversation
.recent_activity
.iter()
.rev()
.map(|activity| ConversationEvent {
event_id: activity.event_id.clone(),
happened_at: timestamp_to_datetime(activity.happened_at),
kind: activity.kind.clone(),
summary: activity.summary.clone(),
payload: activity.payload.clone(),
sequence: activity.sequence,
})
.collect()
})
.unwrap_or_default(),
modified_files: Vec::new(),
input_tokens: issue
.conversation
.as_ref()
.map(|conversation| conversation.input_tokens)
.unwrap_or(0),
output_tokens: issue
.conversation
.as_ref()
.map(|conversation| conversation.output_tokens)
.unwrap_or(0),
cache_read_tokens: issue
.conversation
.as_ref()
.map(|conversation| conversation.cache_read_tokens)
.unwrap_or(0),
total_tokens: issue
.conversation
.as_ref()
.map(|conversation| conversation.effective_total_tokens())
.unwrap_or(0),
detached: false,
cancel_requested: matches!(
interrupt.map(|interrupt| interrupt.status),
Some(HarnessInterruptStatus::Requested)
),
cancel_acknowledged: matches!(
interrupt.map(|interrupt| interrupt.status),
Some(HarnessInterruptStatus::Acknowledged)
),
cancel_failed: matches!(
interrupt.map(|interrupt| interrupt.status),
Some(HarnessInterruptStatus::Failed)
),
cancel_timed_out: matches!(
interrupt.map(|interrupt| interrupt.status),
Some(HarnessInterruptStatus::TimedOut)
),
cancel_reason: interrupt.map(|interrupt| interrupt.command.reason.as_str().to_string()),
operator,
}
}
fn runtime_seconds_from_timestamps(
started_at: Option<DateTime<Utc>>,
finished_at: Option<DateTime<Utc>>,
generated_at: DateTime<Utc>,
runtime_state: IssueRuntimeState,
) -> u64 {
let Some(started_at) = started_at else {
return 0;
};
let end = match runtime_state {
IssueRuntimeState::Running | IssueRuntimeState::Releasing => generated_at,
IssueRuntimeState::Completed | IssueRuntimeState::Failed => match finished_at {
Some(finished_at) => finished_at,
None => return 0,
},
_ => return 0,
};
elapsed_seconds(started_at, end)
}
fn elapsed_seconds(started_at: DateTime<Utc>, ended_at: DateTime<Utc>) -> u64 {
ended_at
.signed_duration_since(started_at)
.num_seconds()
.max(0) as u64
}
fn map_worker_outcome(
issue: &crate::opensymphony_domain::IssueSnapshot,
runtime_state: IssueRuntimeState,
) -> WorkerOutcome {
match runtime_state {
IssueRuntimeState::Running => WorkerOutcome::Running,
IssueRuntimeState::Paused => WorkerOutcome::Unknown,
IssueRuntimeState::RetryQueued => match issue
.last_worker_outcome
.as_ref()
.map(|outcome| outcome.outcome)
{
Some(WorkerOutcomeKind::Succeeded) => WorkerOutcome::Continued,
Some(WorkerOutcomeKind::Cancelled) => WorkerOutcome::Canceled,
Some(
WorkerOutcomeKind::Failed
| WorkerOutcomeKind::TimedOut
| WorkerOutcomeKind::Stalled
| WorkerOutcomeKind::Detached
| WorkerOutcomeKind::CancelFailed,
) => WorkerOutcome::Failed,
None => WorkerOutcome::Continued,
},
IssueRuntimeState::Completed => match issue
.last_worker_outcome
.as_ref()
.map(|outcome| outcome.outcome)
{
Some(WorkerOutcomeKind::Cancelled) => WorkerOutcome::Canceled,
Some(
WorkerOutcomeKind::Failed
| WorkerOutcomeKind::TimedOut
| WorkerOutcomeKind::Stalled
| WorkerOutcomeKind::Detached
| WorkerOutcomeKind::CancelFailed,
) => WorkerOutcome::Failed,
_ => WorkerOutcome::Completed,
},
IssueRuntimeState::Failed => WorkerOutcome::Failed,
IssueRuntimeState::Idle => WorkerOutcome::Unknown,
IssueRuntimeState::Releasing => WorkerOutcome::Unknown,
}
}
pub(super) fn current_agent_server_status(
supervisor: &mut Option<LocalServerSupervisor>,
base_url: &str,
) -> AgentServerStatus {
if let Some(supervisor) = supervisor.as_mut()
&& let Ok(status) = supervisor.status()
{
return AgentServerStatus {
reachable: matches!(
status.state,
crate::opensymphony_openhands::ServerState::Ready
),
base_url: status.base_url,
conversation_count: 0,
status_line: format!("{:?}", status.state).to_ascii_lowercase(),
};
}
AgentServerStatus {
reachable: !base_url.is_empty(),
base_url: base_url.to_string(),
conversation_count: 0,
status_line: if base_url.is_empty() {
"not_selected"
} else {
"reachable"
}
.to_string(),
}
}
pub(super) fn push_recent_event(
recent_events: &mut VecDeque<RecentEvent>,
kind: RecentEventKind,
issue_identifier: Option<IssueIdentifier>,
summary: String,
happened_at: DateTime<Utc>,
) {
recent_events.push_front(RecentEvent {
happened_at,
issue_identifier: issue_identifier.map(|identifier| identifier.to_string()),
kind,
summary,
});
while recent_events.len() > RECENT_EVENT_LIMIT {
let _ = recent_events.pop_back();
}
}
pub(super) fn terminal_state_set(workflow: &ResolvedWorkflow) -> HashSet<String> {
workflow
.config
.tracker
.terminal_states
.iter()
.map(|state| state.trim().to_ascii_lowercase())
.collect()
}
fn is_terminal_state(terminal_states: &HashSet<String>, state: &str) -> bool {
terminal_states.contains(&state.trim().to_ascii_lowercase())
}
fn map_daemon_state(health: HealthStatus) -> DaemonState {
match health {
HealthStatus::Unknown | HealthStatus::Starting => DaemonState::Starting,
HealthStatus::Healthy => DaemonState::Ready,
HealthStatus::Degraded | HealthStatus::Failed => DaemonState::Degraded,
}
}
fn suffix(value: &str) -> String {
if value.len() <= 8 {
value.to_string()
} else {
value[value.len() - 8..].to_string()
}
}
fn suffix_path(path: &Path) -> String {
path.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| path.display().to_string())
}
#[cfg(test)]
mod tests {
use std::{collections::BTreeMap, path::PathBuf};
use crate::opensymphony_domain::{
BlockerRef, ComponentHealthSnapshot, ConversationId, ConversationMetadata, DaemonSnapshot,
HealthStatus, IssueId, IssueIdentifier, IssueRef, IssueSnapshot as DomainIssueSnapshot,
IssueState, IssueStateCategory, NormalizedIssue, OrchestratorSnapshot, RetryAttempt,
RuntimeStateSnapshot, RuntimeStreamState, RuntimeUsageTotals, SchedulerStatus, TimestampMs,
WorkerAttemptSnapshot, WorkerId, WorkerOutcomeKind, WorkerOutcomeRecord, WorkspaceKey,
WorkspaceRecord,
};
use serde_json::json;
use super::{map_snapshot, terminal_state_set, timestamp_to_datetime};
fn must<T, E: std::fmt::Display>(result: Result<T, E>) -> T {
match result {
Ok(value) => value,
Err(error) => panic!("{error}"),
}
}
fn ts(value: u64) -> TimestampMs {
TimestampMs::new(value)
}
fn resolved_workflow_for_tests() -> crate::opensymphony_workflow::ResolvedWorkflow {
let workflow = crate::opensymphony_workflow::WorkflowDefinition::parse(
r#"---
tracker:
kind: linear
project_slug: sample-project
active_states:
- In Progress
terminal_states:
- Done
---
{{ issue.identifier }}
"#,
)
.expect("workflow should parse");
workflow
.resolve(
std::path::Path::new("/tmp"),
&BTreeMap::from([("LINEAR_API_KEY".to_owned(), "linear-token".to_owned())]),
)
.expect("workflow should resolve")
}
#[test]
fn map_snapshot_preserves_recent_events_and_run_metrics() {
let recent_activity = (0..12)
.map(
|index| crate::opensymphony_domain::ConversationActivityEvent {
event_id: format!("evt-{index}"),
happened_at: ts(1_000 + index),
kind: "ActionEvent".to_owned(),
summary: format!("summary {index}"),
payload: (index == 0).then(|| json!({"command": "npm test"})),
sequence: index,
},
)
.collect();
let mut snapshot = OrchestratorSnapshot::new(
ts(2_000),
DaemonSnapshot::new(
HealthStatus::Healthy,
1_000,
4,
Some(ts(2_000)),
ComponentHealthSnapshot::default(),
RuntimeUsageTotals::default(),
),
vec![DomainIssueSnapshot {
issue: NormalizedIssue {
id: must(IssueId::new("lin_352")),
identifier: must(IssueIdentifier::new("COE-352")),
title: "Render media pipeline".to_owned(),
description: None,
priority: None,
state: IssueState {
id: None,
name: "In Progress".to_owned(),
category: IssueStateCategory::Active,
},
branch_name: None,
pr_url: None,
pr_urls: Vec::new(),
url: None,
labels: Vec::new(),
project_id: Some("proj-open".to_owned()),
project_slug: Some("opensymphony-bootstrap".to_owned()),
project_name: Some("OpenSymphony".to_owned()),
parent_id: None,
repository_binding: None,
blocked_by: Vec::<BlockerRef>::new(),
sub_issues: Vec::<IssueRef>::new(),
created_at: None,
updated_at: None,
},
runtime: RuntimeStateSnapshot {
state: SchedulerStatus::Running,
claimed_at: Some(ts(900)),
started_at: Some(ts(1_000)),
released_at: None,
release_reason: None,
worker: Some(WorkerAttemptSnapshot {
worker_id: must(WorkerId::new("worker-352")),
attempt: None,
normal_retry_count: 0,
turn_count: 3,
max_turns: 8,
}),
last_event_at: Some(ts(1_011)),
stalled_at: None,
interrupt: None,
},
workspace: Some(WorkspaceRecord {
path: PathBuf::from("/tmp/workspaces/COE-352"),
workspace_key: must(WorkspaceKey::new("COE-352")),
created_now: false,
created_at: None,
updated_at: None,
last_seen_tracker_refresh_at: None,
}),
conversation: Some(ConversationMetadata {
harness_capability: None,
conversation_id: must(ConversationId::new("conv_352")),
server_base_url: Some("http://127.0.0.1:3000".to_owned()),
transport_target: Some("loopback".to_owned()),
http_auth_mode: Some("none".to_owned()),
websocket_auth_mode: Some("none".to_owned()),
websocket_query_param_name: None,
fresh_conversation: false,
runtime_contract_version: Some("openhands-sdk-agent-server-v1".to_owned()),
stream_state: RuntimeStreamState::Ready,
last_event_id: Some("evt-11".to_owned()),
last_event_kind: Some("ActionEvent".to_owned()),
last_event_at: Some(ts(1_011)),
last_event_summary: Some("summary 11".to_owned()),
recent_activity,
input_tokens: 0,
output_tokens: 0,
cache_read_tokens: 0,
total_tokens: 0,
runtime_seconds: 0,
next_activity_sequence: 0,
}),
retry: None,
retry_count_override: None,
last_worker_outcome: None,
recent_worker_outcomes: Vec::new(),
}],
);
snapshot.hierarchy.insert(
"lin_352".to_owned(),
crate::opensymphony_domain::HierarchyStateSnapshot {
generation: 3,
blocked_reason: Some("HierarchyChanged".to_owned()),
},
);
let mapped = map_snapshot(
&snapshot,
PathBuf::from("/tmp/workspaces").as_path(),
&terminal_state_set(&resolved_workflow_for_tests()),
crate::opensymphony_control::AgentServerStatus {
reachable: true,
base_url: "http://127.0.0.1:3000".to_owned(),
conversation_count: 1,
status_line: "healthy".to_owned(),
},
crate::opensymphony_control::MemoryServerStatus::default(),
&std::collections::VecDeque::new(),
);
assert_eq!(mapped.issues[0].recent_events.len(), 12);
assert_eq!(mapped.issues[0].recent_events[0].summary, "summary 11");
assert_eq!(mapped.issues[0].recent_events[11].summary, "summary 0");
assert_eq!(
mapped.issues[0].recent_events[11].payload,
Some(json!({"command": "npm test"}))
);
assert_eq!(
mapped.issues[0].claimed_at,
Some(timestamp_to_datetime(ts(900)))
);
assert_eq!(
mapped.issues[0].started_at,
Some(timestamp_to_datetime(ts(1_000)))
);
assert_eq!(mapped.issues[0].turn_count, 3);
assert_eq!(mapped.issues[0].max_turns, 8);
assert_eq!(mapped.issues[0].runtime_seconds, 1);
assert_eq!(mapped.issues[0].project_id.as_deref(), Some("proj-open"));
assert_eq!(
mapped.issues[0].project_slug.as_deref(),
Some("opensymphony-bootstrap")
);
assert_eq!(
mapped.issues[0].project_name.as_deref(),
Some("OpenSymphony")
);
assert_eq!(mapped.issues[0].workspace_label.as_deref(), Some("COE-352"));
assert!(mapped.issues[0].blocked);
assert_eq!(mapped.issues[0].hierarchy_generation, Some(3));
assert_eq!(
mapped.issues[0].hierarchy_blocked_reason.as_deref(),
Some("HierarchyChanged")
);
}
#[test]
fn map_snapshot_exposes_blocked_repository_binding() {
let mut issue = released_issue_snapshot(
"In Progress",
IssueStateCategory::Active,
crate::opensymphony_domain::ReleaseReason::Completed,
);
issue.issue.repository_binding = Some(
crate::opensymphony_domain::RepositoryBindingOutcome::UnknownAlias(
"missing".to_string(),
),
);
let mapped = map_single_issue(issue);
assert!(mapped.blocked);
assert!(matches!(
mapped.repository_binding,
Some(crate::opensymphony_domain::RepositoryBindingOutcome::UnknownAlias(alias))
if alias == "missing"
));
}
fn released_issue_snapshot(
state_name: &str,
category: IssueStateCategory,
release_reason: crate::opensymphony_domain::ReleaseReason,
) -> DomainIssueSnapshot {
DomainIssueSnapshot {
issue: NormalizedIssue {
id: must(IssueId::new("lin_532")),
identifier: must(IssueIdentifier::new("COE-532")),
title: "Symbol identity container chain".to_owned(),
description: None,
priority: None,
state: IssueState {
id: None,
name: state_name.to_owned(),
category,
},
branch_name: None,
pr_url: None,
pr_urls: Vec::new(),
url: None,
labels: Vec::new(),
project_id: None,
project_slug: None,
project_name: None,
parent_id: None,
repository_binding: None,
blocked_by: Vec::<BlockerRef>::new(),
sub_issues: Vec::<IssueRef>::new(),
created_at: None,
updated_at: None,
},
runtime: RuntimeStateSnapshot {
state: SchedulerStatus::Released,
claimed_at: None,
started_at: None,
released_at: Some(ts(1_500)),
release_reason: Some(release_reason),
worker: None,
last_event_at: None,
stalled_at: None,
interrupt: None,
},
workspace: None,
conversation: None,
retry: None,
retry_count_override: None,
last_worker_outcome: None,
recent_worker_outcomes: Vec::new(),
}
}
fn map_single_issue(issue: DomainIssueSnapshot) -> crate::opensymphony_control::IssueSnapshot {
let snapshot = OrchestratorSnapshot::new(
ts(2_000),
DaemonSnapshot::new(
HealthStatus::Healthy,
1_000,
4,
Some(ts(2_000)),
ComponentHealthSnapshot::default(),
RuntimeUsageTotals::default(),
),
vec![issue],
);
let mut mapped = map_snapshot(
&snapshot,
PathBuf::from("/tmp/workspaces").as_path(),
&terminal_state_set(&resolved_workflow_for_tests()),
crate::opensymphony_control::AgentServerStatus {
reachable: true,
base_url: "http://127.0.0.1:3000".to_owned(),
conversation_count: 0,
status_line: "healthy".to_owned(),
},
crate::opensymphony_control::MemoryServerStatus::default(),
&std::collections::VecDeque::new(),
);
mapped.issues.remove(0)
}
fn map_single_issue_runtime_state(
issue: DomainIssueSnapshot,
) -> crate::opensymphony_control::IssueRuntimeState {
map_single_issue(issue).runtime_state
}
fn codex_conversation(thread_id: &str) -> ConversationMetadata {
ConversationMetadata {
harness_capability: None,
conversation_id: must(ConversationId::new(thread_id.to_owned())),
server_base_url: None,
transport_target: Some("codex_app_server".to_owned()),
http_auth_mode: None,
websocket_auth_mode: None,
websocket_query_param_name: None,
fresh_conversation: false,
runtime_contract_version: Some("codex-app-server-json-rpc-v2".to_owned()),
stream_state: RuntimeStreamState::Ready,
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 running_issue_with_conversation(conversation: ConversationMetadata) -> DomainIssueSnapshot {
let mut issue = released_issue_snapshot(
"In Progress",
IssueStateCategory::Active,
crate::opensymphony_domain::ReleaseReason::Completed,
);
issue.conversation = Some(conversation);
issue
}
#[test]
fn codex_conversation_records_full_thread_id() {
let issue = map_single_issue(running_issue_with_conversation(codex_conversation(
"019f3979-3aa3-71f3-86b1-18e92c71fbc9",
)));
assert_eq!(
issue.codex_thread_id.as_deref(),
Some("019f3979-3aa3-71f3-86b1-18e92c71fbc9")
);
}
#[test]
fn openhands_conversation_records_no_codex_thread_id() {
let mut conversation = codex_conversation("ignored");
conversation.transport_target = Some("loopback".to_owned());
conversation.runtime_contract_version = Some("openhands-sdk-agent-server-v1".to_owned());
let issue = map_single_issue(running_issue_with_conversation(conversation));
assert_eq!(issue.codex_thread_id, None);
}
#[test]
fn dry_run_route_preview_records_no_codex_thread_id() {
let mut conversation = codex_conversation("route-preview-worker-1");
conversation.transport_target =
Some(crate::opensymphony_codex::CODEX_APP_SERVER_KIND.to_owned());
conversation.runtime_contract_version = Some("opensymphony-routing-alpha-v1".to_owned());
let issue = map_single_issue(running_issue_with_conversation(conversation));
assert_eq!(issue.codex_thread_id, None);
}
#[test]
fn legacy_codex_manifest_without_contract_still_records_thread_id() {
let mut conversation = codex_conversation("019f3979-3aa3-71f3-86b1-18e92c71fbc9");
conversation.transport_target =
Some(crate::opensymphony_codex::CODEX_APP_SERVER_KIND.to_owned());
conversation.runtime_contract_version = None;
let issue = map_single_issue(running_issue_with_conversation(conversation));
assert_eq!(
issue.codex_thread_id.as_deref(),
Some("019f3979-3aa3-71f3-86b1-18e92c71fbc9")
);
}
#[test]
fn tracker_inactive_release_without_a_run_maps_to_idle_not_completed() {
let state = map_single_issue_runtime_state(released_issue_snapshot(
"Backlog",
IssueStateCategory::NonActive,
crate::opensymphony_domain::ReleaseReason::TrackerInactive,
));
assert_eq!(state, crate::opensymphony_control::IssueRuntimeState::Idle);
}
#[test]
fn tracker_terminal_release_without_a_run_still_maps_to_completed() {
let state = map_single_issue_runtime_state(released_issue_snapshot(
"Done",
IssueStateCategory::Terminal,
crate::opensymphony_domain::ReleaseReason::TrackerTerminal,
));
assert_eq!(
state,
crate::opensymphony_control::IssueRuntimeState::Completed
);
}
#[test]
fn retry_exhausted_release_preserves_explicit_reason() {
let mut domain_issue = released_issue_snapshot(
"In Progress",
IssueStateCategory::Active,
crate::opensymphony_domain::ReleaseReason::RetryExhausted,
);
domain_issue.last_worker_outcome = Some(WorkerOutcomeRecord {
worker_id: must(WorkerId::new("worker-532")),
attempt: Some(must(RetryAttempt::new(3))),
outcome: WorkerOutcomeKind::Failed,
started_at: ts(1_000),
finished_at: ts(1_400),
turn_count: 1,
summary: None,
error: Some("historical failure".to_owned()),
harness_stopped: false,
parent_verification: None,
});
let issue = map_single_issue(domain_issue);
assert_eq!(
issue.release_reason,
Some(crate::opensymphony_domain::ReleaseReason::RetryExhausted)
);
assert_eq!(issue.retry_count, 3);
assert_eq!(
issue.runtime_state,
crate::opensymphony_control::IssueRuntimeState::Failed
);
}
#[test]
fn completed_final_retry_preserves_attempt_count() {
let mut domain_issue = released_issue_snapshot(
"In Progress",
IssueStateCategory::Active,
crate::opensymphony_domain::ReleaseReason::Completed,
);
domain_issue.last_worker_outcome = Some(WorkerOutcomeRecord {
worker_id: must(WorkerId::new("worker-532")),
attempt: Some(must(RetryAttempt::new(3))),
outcome: WorkerOutcomeKind::Succeeded,
started_at: ts(1_000),
finished_at: ts(1_400),
turn_count: 1,
summary: None,
error: None,
harness_stopped: false,
parent_verification: None,
});
assert_eq!(map_single_issue(domain_issue).retry_count, 3);
}
#[test]
fn stale_tracker_inactive_reason_with_terminal_state_still_maps_to_completed() {
let state = map_single_issue_runtime_state(released_issue_snapshot(
"Done",
IssueStateCategory::Terminal,
crate::opensymphony_domain::ReleaseReason::TrackerInactive,
));
assert_eq!(
state,
crate::opensymphony_control::IssueRuntimeState::Completed
);
}
#[test]
fn detached_active_run_is_failed_while_its_cleanup_remains_fenced() {
let mut issue = released_issue_snapshot(
"In Progress",
IssueStateCategory::Active,
crate::opensymphony_domain::ReleaseReason::Completed,
);
issue.last_worker_outcome = Some(WorkerOutcomeRecord {
worker_id: must(WorkerId::new("worker-acp-uncertain")),
attempt: None,
outcome: WorkerOutcomeKind::Detached,
started_at: ts(1_000),
finished_at: ts(1_400),
turn_count: 1,
summary: Some("ACP prompt outcome is uncertain".into()),
error: None,
harness_stopped: false,
parent_verification: None,
});
let issue = map_single_issue(issue);
assert_eq!(
issue.runtime_state,
crate::opensymphony_control::IssueRuntimeState::Failed
);
assert_eq!(
issue.last_outcome,
crate::opensymphony_control::WorkerOutcome::Failed
);
}
#[test]
fn terminal_tracker_state_overrides_a_failed_worker_outcome() {
let mut issue = released_issue_snapshot(
"Done",
IssueStateCategory::Terminal,
crate::opensymphony_domain::ReleaseReason::RetryExhausted,
);
issue.last_worker_outcome = Some(WorkerOutcomeRecord {
worker_id: must(WorkerId::new("worker-532")),
attempt: None,
outcome: WorkerOutcomeKind::Failed,
started_at: ts(1_000),
finished_at: ts(1_400),
turn_count: 1,
summary: None,
error: Some("historical failure".to_owned()),
harness_stopped: false,
parent_verification: None,
});
assert_eq!(
map_single_issue_runtime_state(issue),
crate::opensymphony_control::IssueRuntimeState::Completed
);
}
}