mirror of
https://github.com/harivansh-afk/sandbox-agent.git
synced 2026-04-15 07:04:48 +00:00
chore: fix bad merge
This commit is contained in:
parent
e72eb9f611
commit
b9efe971ff
32 changed files with 3654 additions and 15543 deletions
File diff suppressed because it is too large
Load diff
|
|
@ -1,20 +0,0 @@
|
|||
use std::env;
|
||||
|
||||
use reqwest::blocking::ClientBuilder;
|
||||
|
||||
const NO_SYSTEM_PROXY_ENV: &str = "SANDBOX_AGENT_NO_SYSTEM_PROXY";
|
||||
|
||||
fn disable_system_proxy() -> bool {
|
||||
env::var(NO_SYSTEM_PROXY_ENV)
|
||||
.map(|value| matches!(value.as_str(), "1" | "true" | "TRUE" | "yes" | "YES"))
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
pub(crate) fn blocking_client_builder() -> ClientBuilder {
|
||||
let builder = reqwest::blocking::Client::builder();
|
||||
if disable_system_proxy() {
|
||||
builder.no_proxy()
|
||||
} else {
|
||||
builder
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,3 @@
|
|||
pub mod agents;
|
||||
pub mod credentials;
|
||||
mod http_client;
|
||||
pub mod testing;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ use std::env;
|
|||
use std::path::PathBuf;
|
||||
use std::time::Duration;
|
||||
|
||||
use reqwest::blocking::Client;
|
||||
use reqwest::header::{HeaderMap, HeaderValue, AUTHORIZATION, CONTENT_TYPE};
|
||||
use reqwest::StatusCode;
|
||||
use thiserror::Error;
|
||||
|
|
@ -35,7 +36,6 @@ pub enum TestAgentConfigError {
|
|||
const AGENTS_ENV: &str = "SANDBOX_TEST_AGENTS";
|
||||
const ANTHROPIC_ENV: &str = "SANDBOX_TEST_ANTHROPIC_API_KEY";
|
||||
const OPENAI_ENV: &str = "SANDBOX_TEST_OPENAI_API_KEY";
|
||||
const PI_ENV: &str = "SANDBOX_TEST_PI";
|
||||
const ANTHROPIC_MODELS_URL: &str = "https://api.anthropic.com/v1/models";
|
||||
const OPENAI_MODELS_URL: &str = "https://api.openai.com/v1/models";
|
||||
const ANTHROPIC_VERSION: &str = "2023-06-01";
|
||||
|
|
@ -64,6 +64,7 @@ pub fn test_agents_from_env() -> Result<Vec<TestAgentConfig>, TestAgentConfigErr
|
|||
AgentId::Opencode,
|
||||
AgentId::Amp,
|
||||
AgentId::Pi,
|
||||
AgentId::Cursor,
|
||||
]);
|
||||
continue;
|
||||
}
|
||||
|
|
@ -74,12 +75,6 @@ pub fn test_agents_from_env() -> Result<Vec<TestAgentConfig>, TestAgentConfigErr
|
|||
agents
|
||||
};
|
||||
|
||||
let include_pi = pi_tests_enabled() && find_in_path(AgentId::Pi.binary_name());
|
||||
if !include_pi && agents.iter().any(|agent| *agent == AgentId::Pi) {
|
||||
eprintln!("Skipping Pi tests: set {PI_ENV}=1 and ensure pi is on PATH.");
|
||||
}
|
||||
agents.retain(|agent| *agent != AgentId::Pi || include_pi);
|
||||
|
||||
agents.sort_by(|a, b| a.as_str().cmp(b.as_str()));
|
||||
agents.dedup();
|
||||
|
||||
|
|
@ -144,22 +139,7 @@ pub fn test_agents_from_env() -> Result<Vec<TestAgentConfig>, TestAgentConfigErr
|
|||
}
|
||||
credentials_with(anthropic_cred.clone(), openai_cred.clone())
|
||||
}
|
||||
AgentId::Pi => {
|
||||
if anthropic_cred.is_none() && openai_cred.is_none() {
|
||||
return Err(TestAgentConfigError::MissingCredentials {
|
||||
agent,
|
||||
missing: format!("{ANTHROPIC_ENV} or {OPENAI_ENV}"),
|
||||
});
|
||||
}
|
||||
if let Some(cred) = anthropic_cred.as_ref() {
|
||||
ensure_anthropic_ok(&mut health_cache, cred)?;
|
||||
}
|
||||
if let Some(cred) = openai_cred.as_ref() {
|
||||
ensure_openai_ok(&mut health_cache, cred)?;
|
||||
}
|
||||
credentials_with(anthropic_cred.clone(), openai_cred.clone())
|
||||
}
|
||||
AgentId::Cursor => credentials_with(None, None),
|
||||
AgentId::Pi | AgentId::Cursor => credentials_with(None, None),
|
||||
AgentId::Mock => credentials_with(None, None),
|
||||
};
|
||||
configs.push(TestAgentConfig { agent, credentials });
|
||||
|
|
@ -195,7 +175,7 @@ fn ensure_openai_ok(
|
|||
fn health_check_anthropic(credentials: &ProviderCredentials) -> Result<(), TestAgentConfigError> {
|
||||
let credentials = credentials.clone();
|
||||
run_blocking_check("anthropic", move || {
|
||||
let client = crate::http_client::blocking_client_builder()
|
||||
let client = Client::builder()
|
||||
.timeout(Duration::from_secs(10))
|
||||
.build()
|
||||
.map_err(|err| TestAgentConfigError::HealthCheckFailed {
|
||||
|
|
@ -249,7 +229,7 @@ fn health_check_anthropic(credentials: &ProviderCredentials) -> Result<(), TestA
|
|||
fn health_check_openai(credentials: &ProviderCredentials) -> Result<(), TestAgentConfigError> {
|
||||
let credentials = credentials.clone();
|
||||
run_blocking_check("openai", move || {
|
||||
let client = crate::http_client::blocking_client_builder()
|
||||
let client = Client::builder()
|
||||
.timeout(Duration::from_secs(10))
|
||||
.build()
|
||||
.map_err(|err| TestAgentConfigError::HealthCheckFailed {
|
||||
|
|
@ -321,15 +301,14 @@ where
|
|||
}
|
||||
|
||||
fn detect_system_agents() -> Vec<AgentId> {
|
||||
let mut candidates = vec![
|
||||
let candidates = [
|
||||
AgentId::Claude,
|
||||
AgentId::Codex,
|
||||
AgentId::Opencode,
|
||||
AgentId::Amp,
|
||||
AgentId::Pi,
|
||||
AgentId::Cursor,
|
||||
];
|
||||
if pi_tests_enabled() && find_in_path(AgentId::Pi.binary_name()) {
|
||||
candidates.push(AgentId::Pi);
|
||||
}
|
||||
let install_dir = default_install_dir();
|
||||
candidates
|
||||
.into_iter()
|
||||
|
|
@ -371,15 +350,6 @@ fn read_env_key(name: &str) -> Option<String> {
|
|||
})
|
||||
}
|
||||
|
||||
fn pi_tests_enabled() -> bool {
|
||||
env::var(PI_ENV)
|
||||
.map(|value| {
|
||||
let value = value.trim().to_ascii_lowercase();
|
||||
value == "1" || value == "true" || value == "yes"
|
||||
})
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn credentials_with(
|
||||
anthropic_cred: Option<ProviderCredentials>,
|
||||
openai_cred: Option<ProviderCredentials>,
|
||||
|
|
|
|||
|
|
@ -774,7 +774,6 @@ enum CredentialAgent {
|
|||
Codex,
|
||||
Opencode,
|
||||
Amp,
|
||||
Pi,
|
||||
}
|
||||
|
||||
fn credentials_to_output(credentials: ExtractedCredentials, reveal: bool) -> CredentialsOutput {
|
||||
|
|
@ -877,31 +876,6 @@ fn select_token_for_agent(
|
|||
)))
|
||||
}
|
||||
}
|
||||
CredentialAgent::Pi => {
|
||||
if let Some(provider) = provider {
|
||||
return select_token_for_provider(credentials, provider);
|
||||
}
|
||||
if let Some(openai) = credentials.openai.as_ref() {
|
||||
return Ok(openai.api_key.clone());
|
||||
}
|
||||
if let Some(anthropic) = credentials.anthropic.as_ref() {
|
||||
return Ok(anthropic.api_key.clone());
|
||||
}
|
||||
if credentials.other.len() == 1 {
|
||||
if let Some((_, cred)) = credentials.other.iter().next() {
|
||||
return Ok(cred.api_key.clone());
|
||||
}
|
||||
}
|
||||
let available = available_providers(credentials);
|
||||
if available.is_empty() {
|
||||
Err(CliError::Server("no credentials found for pi".to_string()))
|
||||
} else {
|
||||
Err(CliError::Server(format!(
|
||||
"multiple providers available for pi: {} (use --provider)",
|
||||
available.join(", ")
|
||||
)))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,30 +0,0 @@
|
|||
use std::env;
|
||||
|
||||
use reqwest::blocking::ClientBuilder as BlockingClientBuilder;
|
||||
use reqwest::ClientBuilder;
|
||||
|
||||
const NO_SYSTEM_PROXY_ENV: &str = "SANDBOX_AGENT_NO_SYSTEM_PROXY";
|
||||
|
||||
fn disable_system_proxy() -> bool {
|
||||
env::var(NO_SYSTEM_PROXY_ENV)
|
||||
.map(|value| matches!(value.as_str(), "1" | "true" | "TRUE" | "yes" | "YES"))
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
pub fn client_builder() -> ClientBuilder {
|
||||
let builder = reqwest::Client::builder();
|
||||
if disable_system_proxy() {
|
||||
builder.no_proxy()
|
||||
} else {
|
||||
builder
|
||||
}
|
||||
}
|
||||
|
||||
pub fn blocking_client_builder() -> BlockingClientBuilder {
|
||||
let builder = reqwest::blocking::Client::builder();
|
||||
if disable_system_proxy() {
|
||||
builder.no_proxy()
|
||||
} else {
|
||||
builder
|
||||
}
|
||||
}
|
||||
|
|
@ -1,13 +1,8 @@
|
|||
//! Sandbox agent core utilities.
|
||||
|
||||
mod acp_runtime;
|
||||
mod agent_server_logs;
|
||||
mod opencode_session_manager;
|
||||
mod universal_events;
|
||||
mod acp_proxy_runtime;
|
||||
pub mod cli;
|
||||
pub mod daemon;
|
||||
pub mod http_client;
|
||||
pub mod opencode_compat;
|
||||
pub mod router;
|
||||
pub mod server_logs;
|
||||
pub mod telemetry;
|
||||
|
|
|
|||
|
|
@ -10,7 +10,6 @@ use std::str::FromStr;
|
|||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::sync::Mutex as StdMutex;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use axum::body::Body;
|
||||
use axum::extract::{Path, Query, State};
|
||||
|
|
@ -32,16 +31,16 @@ use crate::router::{
|
|||
is_question_tool_action, AgentModelInfo, AppState, CreateSessionRequest, PermissionReply,
|
||||
SessionInfo,
|
||||
};
|
||||
use sandbox_agent_agent_management::agents::AgentId;
|
||||
use sandbox_agent_agent_management::credentials::{
|
||||
extract_all_credentials, CredentialExtractionOptions, ExtractedCredentials,
|
||||
};
|
||||
use sandbox_agent_error::SandboxError;
|
||||
use sandbox_agent_universal_agent_schema::{
|
||||
use crate::universal_events::{
|
||||
ContentPart, FileAction, ItemDeltaData, ItemEventData, ItemKind, ItemRole, ItemStatus,
|
||||
PermissionEventData, PermissionStatus, QuestionEventData, QuestionStatus, UniversalEvent,
|
||||
UniversalEventData, UniversalEventType, UniversalItem,
|
||||
};
|
||||
use sandbox_agent_agent_credentials::{
|
||||
extract_all_credentials, CredentialExtractionOptions, ExtractedCredentials,
|
||||
};
|
||||
use sandbox_agent_agent_management::agents::AgentId;
|
||||
use sandbox_agent_error::SandboxError;
|
||||
|
||||
static SESSION_COUNTER: AtomicU64 = AtomicU64::new(1);
|
||||
static MESSAGE_COUNTER: AtomicU64 = AtomicU64::new(1);
|
||||
|
|
@ -53,7 +52,6 @@ const OPENCODE_EVENT_LOG_SIZE: usize = 4096;
|
|||
const OPENCODE_DEFAULT_MODEL_ID: &str = "mock";
|
||||
const OPENCODE_DEFAULT_PROVIDER_ID: &str = "mock";
|
||||
const OPENCODE_DEFAULT_AGENT_MODE: &str = "build";
|
||||
const OPENCODE_MODEL_CACHE_TTL: Duration = Duration::from_secs(30);
|
||||
const OPENCODE_MODEL_CHANGE_AFTER_SESSION_CREATE_ERROR: &str = "OpenCode compatibility currently does not support changing the model after creating a session. Export with /export and load in to a new session.";
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
|
|
@ -153,6 +151,9 @@ impl OpenCodeSessionRecord {
|
|||
if let Some(url) = &self.share_url {
|
||||
map.insert("share".to_string(), json!({"url": url}));
|
||||
}
|
||||
if let Some(permission_mode) = &self.permission_mode {
|
||||
map.insert("permissionMode".to_string(), json!(permission_mode));
|
||||
}
|
||||
Value::Object(map)
|
||||
}
|
||||
}
|
||||
|
|
@ -164,7 +165,7 @@ fn session_info_to_opencode_value(info: &SessionInfo, default_project_id: &str)
|
|||
.clone()
|
||||
.unwrap_or_else(|| format!("Session {}", info.session_id));
|
||||
let directory = info.directory.clone().unwrap_or_default();
|
||||
json!({
|
||||
let mut value = json!({
|
||||
"id": info.session_id,
|
||||
"slug": format!("session-{}", info.session_id),
|
||||
"projectID": default_project_id,
|
||||
|
|
@ -175,7 +176,15 @@ fn session_info_to_opencode_value(info: &SessionInfo, default_project_id: &str)
|
|||
"created": info.created_at,
|
||||
"updated": info.updated_at,
|
||||
}
|
||||
})
|
||||
});
|
||||
if let Some(obj) = value.as_object_mut() {
|
||||
obj.insert("agent".to_string(), json!(info.agent));
|
||||
obj.insert("permissionMode".to_string(), json!(info.permission_mode));
|
||||
if let Some(model) = &info.model {
|
||||
obj.insert("model".to_string(), json!(model));
|
||||
}
|
||||
}
|
||||
value
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
|
|
@ -303,8 +312,6 @@ struct OpenCodeModelCache {
|
|||
group_names: HashMap<String, String>,
|
||||
default_group: String,
|
||||
default_model: String,
|
||||
cached_at: Instant,
|
||||
had_discovery_errors: bool,
|
||||
/// Group IDs that have valid credentials available
|
||||
connected: Vec<String>,
|
||||
}
|
||||
|
|
@ -818,8 +825,6 @@ fn available_agent_ids() -> Vec<AgentId> {
|
|||
AgentId::Codex,
|
||||
AgentId::Opencode,
|
||||
AgentId::Amp,
|
||||
AgentId::Pi,
|
||||
AgentId::Cursor,
|
||||
AgentId::Mock,
|
||||
]
|
||||
}
|
||||
|
|
@ -838,30 +843,18 @@ async fn opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCache {
|
|||
// spawning duplicate provider/model fetches.
|
||||
let mut slot = state.opencode.model_cache.lock().await;
|
||||
if let Some(cache) = slot.as_ref() {
|
||||
if cache.cached_at.elapsed() < OPENCODE_MODEL_CACHE_TTL {
|
||||
info!(
|
||||
entries = cache.entries.len(),
|
||||
groups = cache.group_names.len(),
|
||||
connected = cache.connected.len(),
|
||||
"opencode model cache hit"
|
||||
);
|
||||
return cache.clone();
|
||||
}
|
||||
info!(
|
||||
entries = cache.entries.len(),
|
||||
groups = cache.group_names.len(),
|
||||
connected = cache.connected.len(),
|
||||
"opencode model cache hit"
|
||||
);
|
||||
return cache.clone();
|
||||
}
|
||||
let previous_cache = slot.clone();
|
||||
|
||||
let started = std::time::Instant::now();
|
||||
info!("opencode model cache miss; building cache");
|
||||
let mut cache = build_opencode_model_cache(state).await;
|
||||
if let Some(previous_cache) = previous_cache {
|
||||
if cache.had_discovery_errors
|
||||
&& cache.entries.is_empty()
|
||||
&& !previous_cache.entries.is_empty()
|
||||
{
|
||||
cache = previous_cache;
|
||||
cache.cached_at = Instant::now();
|
||||
}
|
||||
}
|
||||
let cache = build_opencode_model_cache(state).await;
|
||||
info!(
|
||||
elapsed_ms = started.elapsed().as_millis() as u64,
|
||||
entries = cache.entries.len(),
|
||||
|
|
@ -902,7 +895,6 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
let mut group_agents: HashMap<String, AgentId> = HashMap::new();
|
||||
let mut group_names: HashMap<String, String> = HashMap::new();
|
||||
let mut default_model: Option<String> = None;
|
||||
let mut had_discovery_errors = false;
|
||||
|
||||
let agents = available_agent_ids();
|
||||
let manager = state.inner.session_manager();
|
||||
|
|
@ -920,10 +912,6 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
let response = match response {
|
||||
Ok(response) => response,
|
||||
Err(err) => {
|
||||
had_discovery_errors = true;
|
||||
let (group_id, group_name) = fallback_group_for_agent(agent);
|
||||
group_agents.entry(group_id.clone()).or_insert(agent);
|
||||
group_names.entry(group_id).or_insert(group_name);
|
||||
warn!(
|
||||
agent = agent.as_str(),
|
||||
elapsed_ms = elapsed.as_millis() as u64,
|
||||
|
|
@ -941,12 +929,6 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
"opencode model cache fetched agent models"
|
||||
);
|
||||
|
||||
if response.models.is_empty() {
|
||||
let (group_id, group_name) = fallback_group_for_agent(agent);
|
||||
group_agents.entry(group_id.clone()).or_insert(agent);
|
||||
group_names.entry(group_id).or_insert(group_name);
|
||||
}
|
||||
|
||||
let first_model_id = response.models.first().map(|model| model.id.clone());
|
||||
for model in response.models {
|
||||
let model_id = model.id.clone();
|
||||
|
|
@ -1031,25 +1013,10 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
}
|
||||
}
|
||||
|
||||
// Build connected list based on credential availability
|
||||
// Build connected list conservatively for deterministic compat behavior.
|
||||
let mut connected = Vec::new();
|
||||
for group_id in group_names.keys() {
|
||||
let is_connected = match group_agents.get(group_id) {
|
||||
Some(AgentId::Claude) | Some(AgentId::Amp) => has_anthropic,
|
||||
Some(AgentId::Codex) => has_openai,
|
||||
Some(AgentId::Opencode) => {
|
||||
// Check the specific provider for opencode groups (e.g., "opencode:anthropic")
|
||||
match opencode_group_provider(group_id) {
|
||||
Some("anthropic") => has_anthropic,
|
||||
Some("openai") => has_openai,
|
||||
_ => has_anthropic || has_openai,
|
||||
}
|
||||
}
|
||||
Some(AgentId::Pi) => true,
|
||||
Some(AgentId::Cursor) => true,
|
||||
Some(AgentId::Mock) => true,
|
||||
None => false,
|
||||
};
|
||||
let is_connected = matches!(group_agents.get(group_id), Some(AgentId::Mock));
|
||||
if is_connected {
|
||||
connected.push(group_id.clone());
|
||||
}
|
||||
|
|
@ -1063,8 +1030,6 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
group_names,
|
||||
default_group,
|
||||
default_model,
|
||||
cached_at: Instant::now(),
|
||||
had_discovery_errors,
|
||||
connected,
|
||||
};
|
||||
info!(
|
||||
|
|
@ -1079,19 +1044,6 @@ async fn build_opencode_model_cache(state: &OpenCodeAppState) -> OpenCodeModelCa
|
|||
cache
|
||||
}
|
||||
|
||||
fn fallback_group_for_agent(agent: AgentId) -> (String, String) {
|
||||
if agent == AgentId::Opencode {
|
||||
return (
|
||||
"opencode".to_string(),
|
||||
agent_display_name(agent).to_string(),
|
||||
);
|
||||
}
|
||||
(
|
||||
agent.as_str().to_string(),
|
||||
agent_display_name(agent).to_string(),
|
||||
)
|
||||
}
|
||||
|
||||
fn resolve_agent_from_model(
|
||||
cache: &OpenCodeModelCache,
|
||||
provider_id: &str,
|
||||
|
|
@ -1205,8 +1157,6 @@ fn agent_display_name(agent: AgentId) -> &'static str {
|
|||
AgentId::Codex => "Codex",
|
||||
AgentId::Opencode => "OpenCode",
|
||||
AgentId::Amp => "Amp",
|
||||
AgentId::Pi => "Pi",
|
||||
AgentId::Cursor => "Cursor",
|
||||
AgentId::Mock => "Mock",
|
||||
}
|
||||
}
|
||||
|
|
@ -3295,9 +3245,6 @@ async fn oc_config_providers(State(state): State<Arc<OpenCodeAppState>>) -> impl
|
|||
.or_default()
|
||||
.push(entry);
|
||||
}
|
||||
for group_id in cache.group_names.keys() {
|
||||
grouped.entry(group_id.clone()).or_default();
|
||||
}
|
||||
let mut providers = Vec::new();
|
||||
let mut defaults = serde_json::Map::new();
|
||||
for (group_id, entries) in grouped {
|
||||
|
|
@ -4886,9 +4833,6 @@ async fn oc_provider_list(State(state): State<Arc<OpenCodeAppState>>) -> impl In
|
|||
.or_default()
|
||||
.push(entry);
|
||||
}
|
||||
for group_id in cache.group_names.keys() {
|
||||
grouped.entry(group_id.clone()).or_default();
|
||||
}
|
||||
let mut providers = Vec::new();
|
||||
let mut defaults = serde_json::Map::new();
|
||||
for (group_id, entries) in grouped {
|
||||
|
|
@ -5834,7 +5778,7 @@ pub struct OpenCodeApiDoc;
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use sandbox_agent_universal_agent_schema::ReasoningVisibility;
|
||||
use crate::universal_events::ReasoningVisibility;
|
||||
|
||||
fn assistant_item(content: Vec<ContentPart>) -> UniversalItem {
|
||||
UniversalItem {
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
|
|
@ -11,7 +11,6 @@ use serde::Serialize;
|
|||
use time::OffsetDateTime;
|
||||
use tokio::time::Instant;
|
||||
|
||||
use crate::http_client;
|
||||
static TELEMETRY_ENABLED: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
const TELEMETRY_URL: &str = "https://tc.rivet.dev";
|
||||
|
|
@ -83,7 +82,7 @@ pub fn log_enabled_message() {
|
|||
|
||||
pub fn spawn_telemetry_task() {
|
||||
tokio::spawn(async move {
|
||||
let client = match http_client::client_builder()
|
||||
let client = match Client::builder()
|
||||
.timeout(Duration::from_millis(TELEMETRY_TIMEOUT_MS))
|
||||
.build()
|
||||
{
|
||||
|
|
|
|||
|
|
@ -3,15 +3,6 @@ use serde::{Deserialize, Serialize};
|
|||
use serde_json::Value;
|
||||
use utoipa::ToSchema;
|
||||
|
||||
pub use sandbox_agent_extracted_agent_schemas::{amp, claude, codex, opencode, pi};
|
||||
|
||||
pub mod agents;
|
||||
|
||||
pub use agents::{
|
||||
amp as convert_amp, claude as convert_claude, codex as convert_codex,
|
||||
opencode as convert_opencode, pi as convert_pi,
|
||||
};
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
pub struct UniversalEvent {
|
||||
pub event_id: String,
|
||||
|
|
@ -87,13 +78,10 @@ pub struct SessionStartedData {
|
|||
pub struct SessionEndedData {
|
||||
pub reason: SessionEndReason,
|
||||
pub terminated_by: TerminatedBy,
|
||||
/// Error message when reason is Error
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub message: Option<String>,
|
||||
/// Process exit code when reason is Error
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub exit_code: Option<i32>,
|
||||
/// Agent stderr output when reason is Error
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub stderr: Option<StderrOutput>,
|
||||
}
|
||||
|
|
@ -116,15 +104,11 @@ pub enum TurnPhase {
|
|||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
pub struct StderrOutput {
|
||||
/// First N lines of stderr (if truncated) or full stderr (if not truncated)
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub head: Option<String>,
|
||||
/// Last N lines of stderr (only present if truncated)
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub tail: Option<String>,
|
||||
/// Whether the output was truncated
|
||||
pub truncated: bool,
|
||||
/// Total number of lines in stderr
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub total_lines: Option<usize>,
|
||||
}
|
||||
|
|
@ -226,7 +210,7 @@ pub enum ItemKind {
|
|||
Unknown,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ItemRole {
|
||||
User,
|
||||
|
|
@ -235,7 +219,7 @@ pub enum ItemRole {
|
|||
Tool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, ToSchema)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ItemStatus {
|
||||
InProgress,
|
||||
|
|
@ -294,93 +278,3 @@ pub enum ReasoningVisibility {
|
|||
Public,
|
||||
Private,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct EventConversion {
|
||||
pub event_type: UniversalEventType,
|
||||
pub data: UniversalEventData,
|
||||
pub native_session_id: Option<String>,
|
||||
pub source: EventSource,
|
||||
pub synthetic: bool,
|
||||
pub raw: Option<Value>,
|
||||
}
|
||||
|
||||
impl EventConversion {
|
||||
pub fn new(event_type: UniversalEventType, data: UniversalEventData) -> Self {
|
||||
Self {
|
||||
event_type,
|
||||
data,
|
||||
native_session_id: None,
|
||||
source: EventSource::Agent,
|
||||
synthetic: false,
|
||||
raw: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn with_native_session(mut self, session_id: Option<String>) -> Self {
|
||||
self.native_session_id = session_id;
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_raw(mut self, raw: Option<Value>) -> Self {
|
||||
self.raw = raw;
|
||||
self
|
||||
}
|
||||
|
||||
pub fn synthetic(mut self) -> Self {
|
||||
self.synthetic = true;
|
||||
self.source = EventSource::Daemon;
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_source(mut self, source: EventSource) -> Self {
|
||||
self.source = source;
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
pub fn turn_started_event(turn_id: Option<String>, metadata: Option<Value>) -> EventConversion {
|
||||
EventConversion::new(
|
||||
UniversalEventType::TurnStarted,
|
||||
UniversalEventData::Turn(TurnEventData {
|
||||
phase: TurnPhase::Started,
|
||||
turn_id,
|
||||
metadata,
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
pub fn turn_ended_event(turn_id: Option<String>, metadata: Option<Value>) -> EventConversion {
|
||||
EventConversion::new(
|
||||
UniversalEventType::TurnEnded,
|
||||
UniversalEventData::Turn(TurnEventData {
|
||||
phase: TurnPhase::Ended,
|
||||
turn_id,
|
||||
metadata,
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
pub fn item_from_text(role: ItemRole, text: String) -> UniversalItem {
|
||||
UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: None,
|
||||
parent_id: None,
|
||||
kind: ItemKind::Message,
|
||||
role: Some(role),
|
||||
content: vec![ContentPart::Text { text }],
|
||||
status: ItemStatus::Completed,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn item_from_parts(role: ItemRole, kind: ItemKind, parts: Vec<ContentPart>) -> UniversalItem {
|
||||
UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: None,
|
||||
parent_id: None,
|
||||
kind,
|
||||
role: Some(role),
|
||||
content: parts,
|
||||
status: ItemStatus::Completed,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,769 +0,0 @@
|
|||
use std::collections::{HashMap, HashSet};
|
||||
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::pi as schema;
|
||||
use crate::{
|
||||
ContentPart, EventConversion, ItemDeltaData, ItemEventData, ItemKind, ItemRole, ItemStatus,
|
||||
ReasoningVisibility, UniversalEventData, UniversalEventType, UniversalItem,
|
||||
};
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct PiEventConverter {
|
||||
tool_result_buffers: HashMap<String, String>,
|
||||
tool_result_started: HashSet<String>,
|
||||
message_completed: HashSet<String>,
|
||||
message_errors: HashSet<String>,
|
||||
message_reasoning: HashMap<String, String>,
|
||||
message_text: HashMap<String, String>,
|
||||
last_message_id: Option<String>,
|
||||
message_started: HashSet<String>,
|
||||
message_counter: u64,
|
||||
}
|
||||
|
||||
impl PiEventConverter {
|
||||
pub fn event_value_to_universal(
|
||||
&mut self,
|
||||
raw: &Value,
|
||||
) -> Result<Vec<EventConversion>, String> {
|
||||
let event_type = raw
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.ok_or_else(|| "missing event type".to_string())?;
|
||||
let native_session_id = extract_session_id(raw);
|
||||
|
||||
let conversions = match event_type {
|
||||
"message_start" => self.message_start(raw),
|
||||
"message_update" => self.message_update(raw),
|
||||
"message_end" => self.message_end(raw),
|
||||
"tool_execution_start" => self.tool_execution_start(raw),
|
||||
"tool_execution_update" => self.tool_execution_update(raw),
|
||||
"tool_execution_end" => self.tool_execution_end(raw),
|
||||
"agent_start"
|
||||
| "agent_end"
|
||||
| "turn_start"
|
||||
| "turn_end"
|
||||
| "auto_compaction"
|
||||
| "auto_compaction_start"
|
||||
| "auto_compaction_end"
|
||||
| "auto_retry"
|
||||
| "auto_retry_start"
|
||||
| "auto_retry_end"
|
||||
| "hook_error" => Ok(vec![status_event(event_type, raw)]),
|
||||
"extension_ui_request" | "extension_ui_response" | "extension_error" => {
|
||||
Ok(vec![status_event(event_type, raw)])
|
||||
}
|
||||
other => Err(format!("unsupported Pi event type: {other}")),
|
||||
}?;
|
||||
|
||||
Ok(conversions
|
||||
.into_iter()
|
||||
.map(|conversion| attach_metadata(conversion, &native_session_id, raw))
|
||||
.collect())
|
||||
}
|
||||
|
||||
fn next_synthetic_message_id(&mut self) -> String {
|
||||
self.message_counter += 1;
|
||||
format!("pi_msg_{}", self.message_counter)
|
||||
}
|
||||
|
||||
fn ensure_message_id(&mut self, message_id: Option<String>) -> String {
|
||||
if let Some(id) = message_id {
|
||||
self.last_message_id = Some(id.clone());
|
||||
return id;
|
||||
}
|
||||
if let Some(id) = self.last_message_id.clone() {
|
||||
return id;
|
||||
}
|
||||
let id = self.next_synthetic_message_id();
|
||||
self.last_message_id = Some(id.clone());
|
||||
id
|
||||
}
|
||||
|
||||
fn ensure_message_started(&mut self, message_id: &str) -> Option<EventConversion> {
|
||||
if !self.message_started.insert(message_id.to_string()) {
|
||||
return None;
|
||||
}
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(message_id.to_string()),
|
||||
parent_id: None,
|
||||
kind: ItemKind::Message,
|
||||
role: Some(ItemRole::Assistant),
|
||||
content: Vec::new(),
|
||||
status: ItemStatus::InProgress,
|
||||
};
|
||||
Some(
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemStarted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)
|
||||
.synthetic(),
|
||||
)
|
||||
}
|
||||
|
||||
fn clear_last_message_id(&mut self, message_id: Option<&str>) {
|
||||
if message_id.is_none() || self.last_message_id.as_deref() == message_id {
|
||||
self.last_message_id = None;
|
||||
}
|
||||
}
|
||||
|
||||
pub fn event_to_universal(
|
||||
&mut self,
|
||||
event: &schema::RpcEvent,
|
||||
) -> Result<Vec<EventConversion>, String> {
|
||||
let raw = serde_json::to_value(event).map_err(|err| err.to_string())?;
|
||||
self.event_value_to_universal(&raw)
|
||||
}
|
||||
|
||||
fn message_start(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let message = raw.get("message");
|
||||
if is_user_role(message) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let message_id = self.ensure_message_id(extract_message_id(raw));
|
||||
self.message_completed.remove(&message_id);
|
||||
self.message_started.insert(message_id.clone());
|
||||
let content = message.and_then(parse_message_content).unwrap_or_default();
|
||||
let entry = self.message_text.entry(message_id.clone()).or_default();
|
||||
for part in &content {
|
||||
if let ContentPart::Text { text } = part {
|
||||
entry.push_str(text);
|
||||
}
|
||||
}
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(message_id),
|
||||
parent_id: None,
|
||||
kind: ItemKind::Message,
|
||||
role: Some(ItemRole::Assistant),
|
||||
content,
|
||||
status: ItemStatus::InProgress,
|
||||
};
|
||||
Ok(vec![EventConversion::new(
|
||||
UniversalEventType::ItemStarted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)])
|
||||
}
|
||||
|
||||
fn message_update(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let assistant_event = raw
|
||||
.get("assistantMessageEvent")
|
||||
.or_else(|| raw.get("assistant_message_event"))
|
||||
.ok_or_else(|| "missing assistantMessageEvent".to_string())?;
|
||||
let event_type = assistant_event
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or("");
|
||||
let message_id = extract_message_id(raw)
|
||||
.or_else(|| extract_message_id(assistant_event))
|
||||
.or_else(|| self.last_message_id.clone());
|
||||
|
||||
match event_type {
|
||||
"start" => {
|
||||
if let Some(id) = message_id {
|
||||
self.last_message_id = Some(id);
|
||||
}
|
||||
Ok(Vec::new())
|
||||
}
|
||||
"text_start" | "text_delta" | "text_end" => {
|
||||
let Some(delta) = extract_delta_text(assistant_event) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let message_id = self.ensure_message_id(message_id);
|
||||
let entry = self.message_text.entry(message_id.clone()).or_default();
|
||||
entry.push_str(&delta);
|
||||
let mut conversions = Vec::new();
|
||||
if let Some(start) = self.ensure_message_started(&message_id) {
|
||||
conversions.push(start);
|
||||
}
|
||||
conversions.push(item_delta(Some(message_id), delta));
|
||||
Ok(conversions)
|
||||
}
|
||||
"thinking_start" | "thinking_delta" | "thinking_end" => {
|
||||
let Some(delta) = extract_delta_text(assistant_event) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let message_id = self.ensure_message_id(message_id);
|
||||
let entry = self
|
||||
.message_reasoning
|
||||
.entry(message_id.clone())
|
||||
.or_default();
|
||||
entry.push_str(&delta);
|
||||
let mut conversions = Vec::new();
|
||||
if let Some(start) = self.ensure_message_started(&message_id) {
|
||||
conversions.push(start);
|
||||
}
|
||||
conversions.push(item_delta(Some(message_id), delta));
|
||||
Ok(conversions)
|
||||
}
|
||||
"toolcall_start"
|
||||
| "toolcall_delta"
|
||||
| "toolcall_end"
|
||||
| "toolcall_args_start"
|
||||
| "toolcall_args_delta"
|
||||
| "toolcall_args_end" => Ok(Vec::new()),
|
||||
"done" => {
|
||||
let message_id = self.ensure_message_id(message_id);
|
||||
if self.message_errors.remove(&message_id) {
|
||||
self.message_text.remove(&message_id);
|
||||
self.message_reasoning.remove(&message_id);
|
||||
self.message_started.remove(&message_id);
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
if self.message_completed.contains(&message_id) {
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let message = raw
|
||||
.get("message")
|
||||
.or_else(|| assistant_event.get("message"));
|
||||
let conversion = self.complete_message(Some(message_id.clone()), message);
|
||||
self.message_completed.insert(message_id.clone());
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
Ok(vec![conversion])
|
||||
}
|
||||
"error" => {
|
||||
let message_id = self.ensure_message_id(message_id);
|
||||
if self.message_completed.contains(&message_id) {
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let error_text = assistant_event
|
||||
.get("error")
|
||||
.or_else(|| raw.get("error"))
|
||||
.map(value_to_string)
|
||||
.unwrap_or_else(|| "Pi message error".to_string());
|
||||
self.message_reasoning.remove(&message_id);
|
||||
self.message_text.remove(&message_id);
|
||||
self.message_errors.insert(message_id.clone());
|
||||
self.message_started.remove(&message_id);
|
||||
self.message_completed.insert(message_id.clone());
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(message_id),
|
||||
parent_id: None,
|
||||
kind: ItemKind::Message,
|
||||
role: Some(ItemRole::Assistant),
|
||||
content: vec![ContentPart::Text { text: error_text }],
|
||||
status: ItemStatus::Failed,
|
||||
};
|
||||
Ok(vec![EventConversion::new(
|
||||
UniversalEventType::ItemCompleted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)])
|
||||
}
|
||||
other => Err(format!("unsupported assistantMessageEvent: {other}")),
|
||||
}
|
||||
}
|
||||
|
||||
fn message_end(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let message = raw.get("message");
|
||||
if is_user_role(message) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let message_id = self
|
||||
.ensure_message_id(extract_message_id(raw).or_else(|| self.last_message_id.clone()));
|
||||
if self.message_errors.remove(&message_id) {
|
||||
self.message_text.remove(&message_id);
|
||||
self.message_reasoning.remove(&message_id);
|
||||
self.message_started.remove(&message_id);
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
if self.message_completed.contains(&message_id) {
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let conversion = self.complete_message(Some(message_id.clone()), message);
|
||||
self.message_completed.insert(message_id.clone());
|
||||
self.clear_last_message_id(Some(&message_id));
|
||||
Ok(vec![conversion])
|
||||
}
|
||||
|
||||
fn complete_message(
|
||||
&mut self,
|
||||
message_id: Option<String>,
|
||||
message: Option<&Value>,
|
||||
) -> EventConversion {
|
||||
let mut content = message.and_then(parse_message_content).unwrap_or_default();
|
||||
let failed = message_is_failed(message);
|
||||
let message_error_text = extract_message_error_text(message);
|
||||
|
||||
if let Some(id) = message_id.clone() {
|
||||
if content.is_empty() {
|
||||
if let Some(text) = self.message_text.remove(&id) {
|
||||
if !text.is_empty() {
|
||||
content.push(ContentPart::Text { text });
|
||||
}
|
||||
}
|
||||
} else {
|
||||
self.message_text.remove(&id);
|
||||
}
|
||||
|
||||
if let Some(reasoning) = self.message_reasoning.remove(&id) {
|
||||
if !reasoning.trim().is_empty()
|
||||
&& !content
|
||||
.iter()
|
||||
.any(|part| matches!(part, ContentPart::Reasoning { .. }))
|
||||
{
|
||||
content.push(ContentPart::Reasoning {
|
||||
text: reasoning,
|
||||
visibility: ReasoningVisibility::Private,
|
||||
});
|
||||
}
|
||||
}
|
||||
self.message_started.remove(&id);
|
||||
}
|
||||
|
||||
if failed && content.is_empty() {
|
||||
if let Some(text) = message_error_text {
|
||||
content.push(ContentPart::Text { text });
|
||||
}
|
||||
}
|
||||
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: message_id,
|
||||
parent_id: None,
|
||||
kind: ItemKind::Message,
|
||||
role: Some(ItemRole::Assistant),
|
||||
content,
|
||||
status: if failed {
|
||||
ItemStatus::Failed
|
||||
} else {
|
||||
ItemStatus::Completed
|
||||
},
|
||||
};
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemCompleted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)
|
||||
}
|
||||
|
||||
fn tool_execution_start(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let tool_call_id =
|
||||
extract_tool_call_id(raw).ok_or_else(|| "missing toolCallId".to_string())?;
|
||||
let tool_name = extract_tool_name(raw).unwrap_or_else(|| "tool".to_string());
|
||||
let arguments = raw
|
||||
.get("args")
|
||||
.or_else(|| raw.get("arguments"))
|
||||
.map(value_to_string)
|
||||
.unwrap_or_else(|| "{}".to_string());
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(tool_call_id.clone()),
|
||||
parent_id: None,
|
||||
kind: ItemKind::ToolCall,
|
||||
role: Some(ItemRole::Assistant),
|
||||
content: vec![ContentPart::ToolCall {
|
||||
name: tool_name,
|
||||
arguments,
|
||||
call_id: tool_call_id,
|
||||
}],
|
||||
status: ItemStatus::InProgress,
|
||||
};
|
||||
Ok(vec![EventConversion::new(
|
||||
UniversalEventType::ItemStarted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)])
|
||||
}
|
||||
|
||||
fn tool_execution_update(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let tool_call_id = match extract_tool_call_id(raw) {
|
||||
Some(id) => id,
|
||||
None => return Ok(Vec::new()),
|
||||
};
|
||||
let partial = match raw
|
||||
.get("partialResult")
|
||||
.or_else(|| raw.get("partial_result"))
|
||||
{
|
||||
Some(value) => value_to_string(value),
|
||||
None => return Ok(Vec::new()),
|
||||
};
|
||||
let prior = self
|
||||
.tool_result_buffers
|
||||
.get(&tool_call_id)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let delta = delta_from_partial(&prior, &partial);
|
||||
self.tool_result_buffers
|
||||
.insert(tool_call_id.clone(), partial);
|
||||
|
||||
let mut conversions = Vec::new();
|
||||
if self.tool_result_started.insert(tool_call_id.clone()) {
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(tool_call_id.clone()),
|
||||
parent_id: None,
|
||||
kind: ItemKind::ToolResult,
|
||||
role: Some(ItemRole::Tool),
|
||||
content: vec![ContentPart::ToolResult {
|
||||
call_id: tool_call_id.clone(),
|
||||
output: String::new(),
|
||||
}],
|
||||
status: ItemStatus::InProgress,
|
||||
};
|
||||
conversions.push(
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemStarted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)
|
||||
.synthetic(),
|
||||
);
|
||||
}
|
||||
|
||||
if !delta.is_empty() {
|
||||
conversions.push(
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemDelta,
|
||||
UniversalEventData::ItemDelta(ItemDeltaData {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(tool_call_id.clone()),
|
||||
delta,
|
||||
}),
|
||||
)
|
||||
.synthetic(),
|
||||
);
|
||||
}
|
||||
|
||||
Ok(conversions)
|
||||
}
|
||||
|
||||
fn tool_execution_end(&mut self, raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
let tool_call_id =
|
||||
extract_tool_call_id(raw).ok_or_else(|| "missing toolCallId".to_string())?;
|
||||
self.tool_result_buffers.remove(&tool_call_id);
|
||||
self.tool_result_started.remove(&tool_call_id);
|
||||
|
||||
let output = raw
|
||||
.get("result")
|
||||
.and_then(extract_result_content)
|
||||
.unwrap_or_default();
|
||||
let is_error = raw.get("isError").and_then(Value::as_bool).unwrap_or(false);
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: Some(tool_call_id.clone()),
|
||||
parent_id: None,
|
||||
kind: ItemKind::ToolResult,
|
||||
role: Some(ItemRole::Tool),
|
||||
content: vec![ContentPart::ToolResult {
|
||||
call_id: tool_call_id,
|
||||
output,
|
||||
}],
|
||||
status: if is_error {
|
||||
ItemStatus::Failed
|
||||
} else {
|
||||
ItemStatus::Completed
|
||||
},
|
||||
};
|
||||
Ok(vec![EventConversion::new(
|
||||
UniversalEventType::ItemCompleted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)])
|
||||
}
|
||||
}
|
||||
|
||||
pub fn event_to_universal(event: &schema::RpcEvent) -> Result<Vec<EventConversion>, String> {
|
||||
PiEventConverter::default().event_to_universal(event)
|
||||
}
|
||||
|
||||
pub fn event_value_to_universal(raw: &Value) -> Result<Vec<EventConversion>, String> {
|
||||
PiEventConverter::default().event_value_to_universal(raw)
|
||||
}
|
||||
|
||||
fn attach_metadata(
|
||||
conversion: EventConversion,
|
||||
native_session_id: &Option<String>,
|
||||
raw: &Value,
|
||||
) -> EventConversion {
|
||||
conversion
|
||||
.with_native_session(native_session_id.clone())
|
||||
.with_raw(Some(raw.clone()))
|
||||
}
|
||||
|
||||
fn status_event(label: &str, raw: &Value) -> EventConversion {
|
||||
let detail = raw
|
||||
.get("error")
|
||||
.or_else(|| raw.get("message"))
|
||||
.map(value_to_string);
|
||||
let item = UniversalItem {
|
||||
item_id: String::new(),
|
||||
native_item_id: None,
|
||||
parent_id: None,
|
||||
kind: ItemKind::Status,
|
||||
role: Some(ItemRole::System),
|
||||
content: vec![ContentPart::Status {
|
||||
label: pi_status_label(label),
|
||||
detail,
|
||||
}],
|
||||
status: ItemStatus::Completed,
|
||||
};
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemCompleted,
|
||||
UniversalEventData::Item(ItemEventData { item }),
|
||||
)
|
||||
}
|
||||
|
||||
fn pi_status_label(label: &str) -> String {
|
||||
match label {
|
||||
"turn_end" => "turn.completed".to_string(),
|
||||
"agent_end" => "session.idle".to_string(),
|
||||
_ => format!("pi.{label}"),
|
||||
}
|
||||
}
|
||||
|
||||
fn item_delta(message_id: Option<String>, delta: String) -> EventConversion {
|
||||
EventConversion::new(
|
||||
UniversalEventType::ItemDelta,
|
||||
UniversalEventData::ItemDelta(ItemDeltaData {
|
||||
item_id: String::new(),
|
||||
native_item_id: message_id,
|
||||
delta,
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
fn is_user_role(message: Option<&Value>) -> bool {
|
||||
message
|
||||
.and_then(|msg| msg.get("role"))
|
||||
.and_then(Value::as_str)
|
||||
.is_some_and(|role| role == "user")
|
||||
}
|
||||
|
||||
fn extract_session_id(value: &Value) -> Option<String> {
|
||||
extract_string(value, &["sessionId"])
|
||||
.or_else(|| extract_string(value, &["session_id"]))
|
||||
.or_else(|| extract_string(value, &["session", "id"]))
|
||||
.or_else(|| extract_string(value, &["message", "sessionId"]))
|
||||
}
|
||||
|
||||
fn extract_message_id(value: &Value) -> Option<String> {
|
||||
extract_string(value, &["messageId"])
|
||||
.or_else(|| extract_string(value, &["message_id"]))
|
||||
.or_else(|| extract_string(value, &["message", "id"]))
|
||||
.or_else(|| extract_string(value, &["message", "messageId"]))
|
||||
.or_else(|| extract_string(value, &["assistantMessageEvent", "messageId"]))
|
||||
}
|
||||
|
||||
fn extract_tool_call_id(value: &Value) -> Option<String> {
|
||||
extract_string(value, &["toolCallId"]).or_else(|| extract_string(value, &["tool_call_id"]))
|
||||
}
|
||||
|
||||
fn extract_tool_name(value: &Value) -> Option<String> {
|
||||
extract_string(value, &["toolName"]).or_else(|| extract_string(value, &["tool_name"]))
|
||||
}
|
||||
|
||||
fn extract_string(value: &Value, path: &[&str]) -> Option<String> {
|
||||
let mut current = value;
|
||||
for key in path {
|
||||
current = current.get(*key)?;
|
||||
}
|
||||
current.as_str().map(|value| value.to_string())
|
||||
}
|
||||
|
||||
fn extract_delta_text(event: &Value) -> Option<String> {
|
||||
if let Some(value) = event.get("delta") {
|
||||
return Some(value_to_string(value));
|
||||
}
|
||||
if let Some(value) = event.get("text") {
|
||||
return Some(value_to_string(value));
|
||||
}
|
||||
if let Some(value) = event.get("partial") {
|
||||
return extract_text_from_value(value);
|
||||
}
|
||||
if let Some(value) = event.get("content") {
|
||||
return extract_text_from_value(value);
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn extract_text_from_value(value: &Value) -> Option<String> {
|
||||
if let Some(text) = value.as_str() {
|
||||
return Some(text.to_string());
|
||||
}
|
||||
if let Some(text) = value.get("text").and_then(Value::as_str) {
|
||||
return Some(text.to_string());
|
||||
}
|
||||
if let Some(text) = value.get("content").and_then(Value::as_str) {
|
||||
return Some(text.to_string());
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn extract_result_content(value: &Value) -> Option<String> {
|
||||
let content = value.get("content").and_then(Value::as_str);
|
||||
let text = value.get("text").and_then(Value::as_str);
|
||||
content
|
||||
.or(text)
|
||||
.map(|value| value.to_string())
|
||||
.or_else(|| Some(value_to_string(value)))
|
||||
}
|
||||
|
||||
fn parse_message_content(message: &Value) -> Option<Vec<ContentPart>> {
|
||||
if let Some(text) = message.as_str() {
|
||||
return Some(vec![ContentPart::Text {
|
||||
text: text.to_string(),
|
||||
}]);
|
||||
}
|
||||
let content_value = message
|
||||
.get("content")
|
||||
.or_else(|| message.get("text"))
|
||||
.or_else(|| message.get("value"))?;
|
||||
let mut parts = Vec::new();
|
||||
match content_value {
|
||||
Value::String(text) => parts.push(ContentPart::Text { text: text.clone() }),
|
||||
Value::Array(items) => {
|
||||
for item in items {
|
||||
if let Some(part) = content_part_from_value(item) {
|
||||
parts.push(part);
|
||||
}
|
||||
}
|
||||
}
|
||||
Value::Object(_) => {
|
||||
if let Some(part) = content_part_from_value(content_value) {
|
||||
parts.push(part);
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
Some(parts)
|
||||
}
|
||||
|
||||
fn message_is_failed(message: Option<&Value>) -> bool {
|
||||
message
|
||||
.and_then(|value| {
|
||||
value
|
||||
.get("stopReason")
|
||||
.or_else(|| value.get("stop_reason"))
|
||||
.and_then(Value::as_str)
|
||||
})
|
||||
.is_some_and(|reason| reason == "error" || reason == "aborted")
|
||||
}
|
||||
|
||||
fn extract_message_error_text(message: Option<&Value>) -> Option<String> {
|
||||
let value = message?;
|
||||
|
||||
if let Some(text) = value
|
||||
.get("errorMessage")
|
||||
.or_else(|| value.get("error_message"))
|
||||
.and_then(Value::as_str)
|
||||
{
|
||||
let trimmed = text.trim();
|
||||
if !trimmed.is_empty() {
|
||||
return Some(trimmed.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
let error = value.get("error")?;
|
||||
if let Some(text) = error.as_str() {
|
||||
let trimmed = text.trim();
|
||||
if !trimmed.is_empty() {
|
||||
return Some(trimmed.to_string());
|
||||
}
|
||||
}
|
||||
if let Some(text) = error
|
||||
.get("errorMessage")
|
||||
.or_else(|| error.get("error_message"))
|
||||
.or_else(|| error.get("message"))
|
||||
.and_then(Value::as_str)
|
||||
{
|
||||
let trimmed = text.trim();
|
||||
if !trimmed.is_empty() {
|
||||
return Some(trimmed.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
fn content_part_from_value(value: &Value) -> Option<ContentPart> {
|
||||
if let Some(text) = value.as_str() {
|
||||
return Some(ContentPart::Text {
|
||||
text: text.to_string(),
|
||||
});
|
||||
}
|
||||
let part_type = value.get("type").and_then(Value::as_str);
|
||||
match part_type {
|
||||
Some("text") | Some("markdown") => {
|
||||
extract_text_from_value(value).map(|text| ContentPart::Text { text })
|
||||
}
|
||||
Some("thinking") | Some("reasoning") => {
|
||||
extract_text_from_value(value).map(|text| ContentPart::Reasoning {
|
||||
text,
|
||||
visibility: ReasoningVisibility::Private,
|
||||
})
|
||||
}
|
||||
Some("image") => value
|
||||
.get("path")
|
||||
.or_else(|| value.get("url"))
|
||||
.and_then(|path| {
|
||||
path.as_str().map(|path| ContentPart::Image {
|
||||
path: path.to_string(),
|
||||
mime: value
|
||||
.get("mime")
|
||||
.or_else(|| value.get("mimeType"))
|
||||
.and_then(Value::as_str)
|
||||
.map(|mime| mime.to_string()),
|
||||
})
|
||||
}),
|
||||
Some("tool_call") | Some("toolcall") => {
|
||||
let name = value
|
||||
.get("name")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or("tool")
|
||||
.to_string();
|
||||
let arguments = value
|
||||
.get("arguments")
|
||||
.or_else(|| value.get("args"))
|
||||
.map(value_to_string)
|
||||
.unwrap_or_else(|| "{}".to_string());
|
||||
let call_id = value
|
||||
.get("call_id")
|
||||
.or_else(|| value.get("callId"))
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or_default()
|
||||
.to_string();
|
||||
Some(ContentPart::ToolCall {
|
||||
name,
|
||||
arguments,
|
||||
call_id,
|
||||
})
|
||||
}
|
||||
Some("tool_result") => {
|
||||
let call_id = value
|
||||
.get("call_id")
|
||||
.or_else(|| value.get("callId"))
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or_default()
|
||||
.to_string();
|
||||
let output = value
|
||||
.get("output")
|
||||
.or_else(|| value.get("content"))
|
||||
.map(value_to_string)
|
||||
.unwrap_or_default();
|
||||
Some(ContentPart::ToolResult { call_id, output })
|
||||
}
|
||||
_ => Some(ContentPart::Json {
|
||||
json: value.clone(),
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
fn value_to_string(value: &Value) -> String {
|
||||
if let Some(text) = value.as_str() {
|
||||
text.to_string()
|
||||
} else {
|
||||
value.to_string()
|
||||
}
|
||||
}
|
||||
|
||||
fn delta_from_partial(previous: &str, next: &str) -> String {
|
||||
if next.starts_with(previous) {
|
||||
next[previous.len()..].to_string()
|
||||
} else {
|
||||
next.to_string()
|
||||
}
|
||||
}
|
||||
|
|
@ -1,414 +0,0 @@
|
|||
use sandbox_agent_universal_agent_schema::convert_pi::PiEventConverter;
|
||||
use sandbox_agent_universal_agent_schema::pi as pi_schema;
|
||||
use sandbox_agent_universal_agent_schema::{
|
||||
ContentPart, ItemKind, ItemRole, ItemStatus, UniversalEventData, UniversalEventType,
|
||||
};
|
||||
use serde_json::json;
|
||||
|
||||
fn parse_event(value: serde_json::Value) -> pi_schema::RpcEvent {
|
||||
serde_json::from_value(value).expect("pi event")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_message_flow_converts() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let start_event = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello" }]
|
||||
}
|
||||
}));
|
||||
let start_events = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("start conversions");
|
||||
assert_eq!(start_events[0].event_type, UniversalEventType::ItemStarted);
|
||||
if let UniversalEventData::Item(item) = &start_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::Message);
|
||||
assert_eq!(item.item.role, Some(ItemRole::Assistant));
|
||||
assert_eq!(item.item.status, ItemStatus::InProgress);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
|
||||
let update_event = parse_event(json!({
|
||||
"type": "message_update",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"assistantMessageEvent": { "type": "text_delta", "delta": " world" }
|
||||
}));
|
||||
let update_events = converter
|
||||
.event_to_universal(&update_event)
|
||||
.expect("update conversions");
|
||||
assert_eq!(update_events[0].event_type, UniversalEventType::ItemDelta);
|
||||
|
||||
let end_event = parse_event(json!({
|
||||
"type": "message_end",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello world" }]
|
||||
}
|
||||
}));
|
||||
let end_events = converter
|
||||
.event_to_universal(&end_event)
|
||||
.expect("end conversions");
|
||||
assert_eq!(end_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
if let UniversalEventData::Item(item) = &end_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::Message);
|
||||
assert_eq!(item.item.role, Some(ItemRole::Assistant));
|
||||
assert_eq!(item.item.status, ItemStatus::Completed);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_user_message_echo_is_skipped() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
// Pi may echo the user message as a message_start with role "user".
|
||||
// The daemon already records synthetic user events, so the converter
|
||||
// must skip these to avoid a duplicate assistant-looking bubble.
|
||||
let start_event = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "user-msg-1",
|
||||
"message": {
|
||||
"role": "user",
|
||||
"content": [{ "type": "text", "text": "hello!" }]
|
||||
}
|
||||
}));
|
||||
let events = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("user message_start should not error");
|
||||
assert!(
|
||||
events.is_empty(),
|
||||
"user message_start should produce no events, got {}",
|
||||
events.len()
|
||||
);
|
||||
|
||||
let end_event = parse_event(json!({
|
||||
"type": "message_end",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "user-msg-1",
|
||||
"message": {
|
||||
"role": "user",
|
||||
"content": [{ "type": "text", "text": "hello!" }]
|
||||
}
|
||||
}));
|
||||
let events = converter
|
||||
.event_to_universal(&end_event)
|
||||
.expect("user message_end should not error");
|
||||
assert!(
|
||||
events.is_empty(),
|
||||
"user message_end should produce no events, got {}",
|
||||
events.len()
|
||||
);
|
||||
|
||||
// A subsequent assistant message should still work normally.
|
||||
let assistant_start = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello! How can I help?" }]
|
||||
}
|
||||
}));
|
||||
let events = converter
|
||||
.event_to_universal(&assistant_start)
|
||||
.expect("assistant message_start");
|
||||
assert_eq!(events.len(), 1);
|
||||
assert_eq!(events[0].event_type, UniversalEventType::ItemStarted);
|
||||
if let UniversalEventData::Item(item) = &events[0].data {
|
||||
assert_eq!(item.item.role, Some(ItemRole::Assistant));
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_tool_execution_converts_with_partial_deltas() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let start_event = parse_event(json!({
|
||||
"type": "tool_execution_start",
|
||||
"sessionId": "session-1",
|
||||
"toolCallId": "call-1",
|
||||
"toolName": "bash",
|
||||
"args": { "command": "ls" }
|
||||
}));
|
||||
let start_events = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("tool start");
|
||||
assert_eq!(start_events[0].event_type, UniversalEventType::ItemStarted);
|
||||
if let UniversalEventData::Item(item) = &start_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::ToolCall);
|
||||
assert_eq!(item.item.role, Some(ItemRole::Assistant));
|
||||
match &item.item.content[0] {
|
||||
ContentPart::ToolCall { name, .. } => assert_eq!(name, "bash"),
|
||||
_ => panic!("expected tool call content"),
|
||||
}
|
||||
}
|
||||
|
||||
let update_event = parse_event(json!({
|
||||
"type": "tool_execution_update",
|
||||
"sessionId": "session-1",
|
||||
"toolCallId": "call-1",
|
||||
"partialResult": "foo"
|
||||
}));
|
||||
let update_events = converter
|
||||
.event_to_universal(&update_event)
|
||||
.expect("tool update");
|
||||
assert!(update_events
|
||||
.iter()
|
||||
.any(|event| event.event_type == UniversalEventType::ItemDelta));
|
||||
|
||||
let update_event2 = parse_event(json!({
|
||||
"type": "tool_execution_update",
|
||||
"sessionId": "session-1",
|
||||
"toolCallId": "call-1",
|
||||
"partialResult": "foobar"
|
||||
}));
|
||||
let update_events2 = converter
|
||||
.event_to_universal(&update_event2)
|
||||
.expect("tool update 2");
|
||||
let delta = update_events2
|
||||
.iter()
|
||||
.find_map(|event| match &event.data {
|
||||
UniversalEventData::ItemDelta(data) => Some(data.delta.clone()),
|
||||
_ => None,
|
||||
})
|
||||
.unwrap_or_default();
|
||||
assert_eq!(delta, "bar");
|
||||
|
||||
let end_event = parse_event(json!({
|
||||
"type": "tool_execution_end",
|
||||
"sessionId": "session-1",
|
||||
"toolCallId": "call-1",
|
||||
"result": { "type": "text", "content": "done" },
|
||||
"isError": false
|
||||
}));
|
||||
let end_events = converter.event_to_universal(&end_event).expect("tool end");
|
||||
assert_eq!(end_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
if let UniversalEventData::Item(item) = &end_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::ToolResult);
|
||||
assert_eq!(item.item.role, Some(ItemRole::Tool));
|
||||
match &item.item.content[0] {
|
||||
ContentPart::ToolResult { output, .. } => assert_eq!(output, "done"),
|
||||
_ => panic!("expected tool result content"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_unknown_event_returns_error() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
let event = parse_event(json!({
|
||||
"type": "unknown_event",
|
||||
"sessionId": "session-1"
|
||||
}));
|
||||
assert!(converter.event_to_universal(&event).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_turn_and_agent_end_emit_terminal_status_labels() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let turn_end = parse_event(json!({
|
||||
"type": "turn_end",
|
||||
"sessionId": "session-1"
|
||||
}));
|
||||
let turn_events = converter
|
||||
.event_to_universal(&turn_end)
|
||||
.expect("turn_end conversions");
|
||||
assert_eq!(turn_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
if let UniversalEventData::Item(item) = &turn_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::Status);
|
||||
assert!(
|
||||
matches!(
|
||||
item.item.content.first(),
|
||||
Some(ContentPart::Status { label, .. }) if label == "turn.completed"
|
||||
),
|
||||
"turn_end should map to turn.completed status"
|
||||
);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
|
||||
let agent_end = parse_event(json!({
|
||||
"type": "agent_end",
|
||||
"sessionId": "session-1"
|
||||
}));
|
||||
let agent_events = converter
|
||||
.event_to_universal(&agent_end)
|
||||
.expect("agent_end conversions");
|
||||
assert_eq!(
|
||||
agent_events[0].event_type,
|
||||
UniversalEventType::ItemCompleted
|
||||
);
|
||||
if let UniversalEventData::Item(item) = &agent_events[0].data {
|
||||
assert_eq!(item.item.kind, ItemKind::Status);
|
||||
assert!(
|
||||
matches!(
|
||||
item.item.content.first(),
|
||||
Some(ContentPart::Status { label, .. }) if label == "session.idle"
|
||||
),
|
||||
"agent_end should map to session.idle status"
|
||||
);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_message_done_completes_without_message_end() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let start_event = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello" }]
|
||||
}
|
||||
}));
|
||||
let _start_events = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("start conversions");
|
||||
|
||||
let update_event = parse_event(json!({
|
||||
"type": "message_update",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"assistantMessageEvent": { "type": "text_delta", "delta": " world" }
|
||||
}));
|
||||
let _update_events = converter
|
||||
.event_to_universal(&update_event)
|
||||
.expect("update conversions");
|
||||
|
||||
let done_event = parse_event(json!({
|
||||
"type": "message_update",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"assistantMessageEvent": { "type": "done" }
|
||||
}));
|
||||
let done_events = converter
|
||||
.event_to_universal(&done_event)
|
||||
.expect("done conversions");
|
||||
assert_eq!(done_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
if let UniversalEventData::Item(item) = &done_events[0].data {
|
||||
assert_eq!(item.item.status, ItemStatus::Completed);
|
||||
assert!(
|
||||
matches!(item.item.content.get(0), Some(ContentPart::Text { text }) if text == "Hello world")
|
||||
);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_message_done_then_message_end_does_not_double_complete() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let start_event = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello" }]
|
||||
}
|
||||
}));
|
||||
let _ = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("start conversions");
|
||||
|
||||
let update_event = parse_event(json!({
|
||||
"type": "message_update",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"assistantMessageEvent": { "type": "text_delta", "delta": " world" }
|
||||
}));
|
||||
let _ = converter
|
||||
.event_to_universal(&update_event)
|
||||
.expect("update conversions");
|
||||
|
||||
let done_event = parse_event(json!({
|
||||
"type": "message_update",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"assistantMessageEvent": { "type": "done" }
|
||||
}));
|
||||
let done_events = converter
|
||||
.event_to_universal(&done_event)
|
||||
.expect("done conversions");
|
||||
assert_eq!(done_events.len(), 1);
|
||||
assert_eq!(done_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
|
||||
let end_event = parse_event(json!({
|
||||
"type": "message_end",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-1",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{ "type": "text", "text": "Hello world" }]
|
||||
}
|
||||
}));
|
||||
let end_events = converter
|
||||
.event_to_universal(&end_event)
|
||||
.expect("end conversions");
|
||||
assert!(
|
||||
end_events.is_empty(),
|
||||
"message_end after done should not emit a second completion"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pi_message_end_error_surfaces_failed_status_and_error_text() {
|
||||
let mut converter = PiEventConverter::default();
|
||||
|
||||
let start_event = parse_event(json!({
|
||||
"type": "message_start",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-err",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": []
|
||||
}
|
||||
}));
|
||||
let _ = converter
|
||||
.event_to_universal(&start_event)
|
||||
.expect("start conversions");
|
||||
|
||||
let end_raw = json!({
|
||||
"type": "message_end",
|
||||
"sessionId": "session-1",
|
||||
"messageId": "msg-err",
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [],
|
||||
"stopReason": "error",
|
||||
"errorMessage": "Connection error."
|
||||
}
|
||||
});
|
||||
let end_events = converter
|
||||
.event_value_to_universal(&end_raw)
|
||||
.expect("end conversions");
|
||||
|
||||
assert_eq!(end_events[0].event_type, UniversalEventType::ItemCompleted);
|
||||
if let UniversalEventData::Item(item) = &end_events[0].data {
|
||||
assert_eq!(item.item.status, ItemStatus::Failed);
|
||||
assert!(
|
||||
matches!(item.item.content.first(), Some(ContentPart::Text { text }) if text == "Connection error.")
|
||||
);
|
||||
} else {
|
||||
panic!("expected item event");
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue