Download codex-rs/protocol/src/protocol.rs from SaylorTwift/codex: direct link, hf CLI and curl.
- Browser
- Download file 233 kB
-
https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/protocol/src/protocol.rs
- Command line
-
hf download hf://SaylorTwift/codex/codex-rs/protocol/src/protocol.rs
-
curl -L -o protocol.rs https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/protocol/src/protocol.rs
233 kB
| //! Defines the protocol for a Codex session between a client and an agent. | |
| //! | |
| //! Uses a SQ (Submission Queue) / EQ (Event Queue) pattern to asynchronously communicate | |
| //! between user and agent. | |
| use std::collections::BTreeMap; | |
| use std::collections::HashMap; | |
| use std::fmt; | |
| use std::ops::Mul; | |
| use std::path::Path; | |
| use std::path::PathBuf; | |
| use std::str::FromStr; | |
| use std::time::Duration; | |
| use strum_macros::EnumIter; | |
| use crate::AgentPath; | |
| use crate::ResponseItemId; | |
| use crate::SanitizedGitUrl; | |
| use crate::SessionId; | |
| use crate::ThreadId; | |
| use crate::approvals::ElicitationRequestEvent; | |
| use crate::capabilities::SelectedCapabilityRoot; | |
| use crate::config_types::ApprovalsReviewer; | |
| use crate::config_types::CollaborationMode; | |
| use crate::config_types::ModeKind; | |
| use crate::config_types::MultiAgentMode; | |
| use crate::config_types::Personality; | |
| use crate::config_types::ReasoningSummary as ReasoningSummaryConfig; | |
| use crate::config_types::WindowsSandboxLevel; | |
| use crate::dynamic_tools::DynamicToolCallOutputContentItem; | |
| use crate::dynamic_tools::DynamicToolCallRequest; | |
| use crate::dynamic_tools::DynamicToolResponse; | |
| use crate::dynamic_tools::DynamicToolSpec; | |
| use crate::error::Result as CodexResult; | |
| use crate::items::AgentMessageDelivery; | |
| use crate::items::AsyncUserInputQuestion; | |
| use crate::items::TurnItem; | |
| use crate::mcp::CallToolResult; | |
| use crate::mcp::RequestId; | |
| use crate::memory_citation::MemoryCitation; | |
| use crate::models::ActivePermissionProfile; | |
| use crate::models::AgentMessageInputContent; | |
| use crate::models::BaseInstructions; | |
| use crate::models::ContentItem; | |
| use crate::models::ImageDetail; | |
| use crate::models::InternalChatMessageMetadataPassthrough; | |
| use crate::models::MessagePhase; | |
| use crate::models::PermissionProfile; | |
| use crate::models::ProfileWorkspaceRoot; | |
| use crate::models::ResponseInputItem; | |
| use crate::models::ResponseItem; | |
| use crate::models::SandboxEnforcement; | |
| use crate::models::WebSearchAction; | |
| use crate::num_format::format_with_separators; | |
| use crate::openai_models::ReasoningEffort as ReasoningEffortConfig; | |
| use crate::parse_command::ParsedCommand; | |
| use crate::plan_tool::UpdatePlanArgs; | |
| use crate::request_permissions::RequestPermissionsEvent; | |
| use crate::request_permissions::RequestPermissionsResponse; | |
| use crate::request_user_input::RequestUserInputResponse; | |
| use crate::turn_input::CyberAccessProgram; | |
| use crate::turn_input::SuspendTurnOutcome; | |
| use crate::turn_input::TurnInputMode; | |
| use crate::turn_input::TurnInputRequest; | |
| use crate::turn_input::TurnInputSubmission; | |
| use crate::turn_input::TurnStartOptions; | |
| use codex_extension_items::image_generation::ImageGenerationFailure; | |
| use codex_utils_absolute_path::AbsolutePathBuf; | |
| use codex_utils_path_uri::PathUri; | |
| use schemars::JsonSchema; | |
| use serde::Deserialize; | |
| use serde::Deserializer; | |
| use serde::Serialize; | |
| use serde::de::Error as _; | |
| use serde_json::Map; | |
| use serde_json::Value; | |
| use serde_with::serde_as; | |
| use strum_macros::Display; | |
| use tokio::sync::oneshot; | |
| use tracing::error; | |
| use ts_rs::TS; | |
| pub use crate::approvals::ApplyPatchApprovalRequestEvent; | |
| pub use crate::approvals::ElicitationAction; | |
| pub use crate::approvals::ExecApprovalRequestEvent; | |
| pub use crate::approvals::ExecPolicyAmendment; | |
| pub use crate::approvals::GuardianAssessmentAction; | |
| pub use crate::approvals::GuardianAssessmentDecisionSource; | |
| pub use crate::approvals::GuardianAssessmentEvent; | |
| pub use crate::approvals::GuardianAssessmentOutcome; | |
| pub use crate::approvals::GuardianAssessmentStatus; | |
| pub use crate::approvals::GuardianCommandSource; | |
| pub use crate::approvals::GuardianRiskLevel; | |
| pub use crate::approvals::GuardianUserAuthorization; | |
| pub use crate::approvals::NetworkApprovalContext; | |
| pub use crate::approvals::NetworkApprovalProtocol; | |
| pub use crate::approvals::NetworkPolicyAmendment; | |
| pub use crate::approvals::NetworkPolicyRuleAction; | |
| pub use crate::environment::EnvironmentConfig; | |
| pub use crate::environment::EnvironmentConfigState; | |
| pub use crate::environment::has_full_access; | |
| pub use crate::legacy_events::HasLegacyEvent; | |
| pub use crate::permissions::FileSystemAccessMode; | |
| pub use crate::permissions::FileSystemPath; | |
| pub use crate::permissions::FileSystemSandboxEntry; | |
| pub use crate::permissions::FileSystemSandboxKind; | |
| pub use crate::permissions::FileSystemSandboxPolicy; | |
| pub use crate::permissions::FileSystemSpecialPath; | |
| pub use crate::permissions::NetworkSandboxPolicy; | |
| pub use crate::permissions::RawFileSystemSandboxPolicy; | |
| use crate::permissions::default_read_only_subpaths_for_writable_root; | |
| pub use crate::request_permissions::RequestPermissionsArgs; | |
| pub use crate::request_user_input::RequestUserInputEvent; | |
| /// Open/close tags for special context blocks. Used across crates to avoid duplicated hardcoded | |
| /// strings. | |
| pub const USER_INSTRUCTIONS_OPEN_TAG: &str = "<user_instructions>"; | |
| pub const USER_INSTRUCTIONS_CLOSE_TAG: &str = "</user_instructions>"; | |
| pub const ENVIRONMENT_CONTEXT_OPEN_TAG: &str = "<environment_context>"; | |
| pub const ENVIRONMENT_CONTEXT_CLOSE_TAG: &str = "</environment_context>"; | |
| pub const ENVIRONMENTS_INSTRUCTIONS_OPEN_TAG: &str = "<environments_instructions>"; | |
| pub const ENVIRONMENTS_INSTRUCTIONS_CLOSE_TAG: &str = "</environments_instructions>"; | |
| pub const APPS_INSTRUCTIONS_OPEN_TAG: &str = "<apps_instructions>"; | |
| pub const APPS_INSTRUCTIONS_CLOSE_TAG: &str = "</apps_instructions>"; | |
| pub const SKILLS_INSTRUCTIONS_OPEN_TAG: &str = "<skills_instructions>"; | |
| pub const SKILLS_INSTRUCTIONS_CLOSE_TAG: &str = "</skills_instructions>"; | |
| pub const PLUGINS_INSTRUCTIONS_OPEN_TAG: &str = "<plugins_instructions>"; | |
| pub const PLUGINS_INSTRUCTIONS_CLOSE_TAG: &str = "</plugins_instructions>"; | |
| pub const TOOLS_OPEN_TAG: &str = "<tools>"; | |
| pub const TOOLS_CLOSE_TAG: &str = "</tools>"; | |
| pub const COLLABORATION_MODE_OPEN_TAG: &str = "<collaboration_mode>"; | |
| pub const COLLABORATION_MODE_CLOSE_TAG: &str = "</collaboration_mode>"; | |
| pub const MULTI_AGENT_MODE_OPEN_TAG: &str = "<multi_agent_mode>"; | |
| pub const MULTI_AGENT_MODE_CLOSE_TAG: &str = "</multi_agent_mode>"; | |
| pub const REALTIME_CONVERSATION_OPEN_TAG: &str = "<realtime_conversation>"; | |
| pub const REALTIME_CONVERSATION_CLOSE_TAG: &str = "</realtime_conversation>"; | |
| pub const CONTEXT_WINDOW_OPEN_TAG: &str = "<context_window>"; | |
| pub const CONTEXT_WINDOW_CLOSE_TAG: &str = "</context_window>"; | |
| pub const CONTEXT_WINDOW_GUIDANCE_OPEN_TAG: &str = "<context_window_guidance>"; | |
| pub const CONTEXT_WINDOW_GUIDANCE_CLOSE_TAG: &str = "</context_window_guidance>"; | |
| pub const USER_MESSAGE_BEGIN: &str = "## My request for Codex:"; | |
| /// Removes the model-context prefix from a user message before displaying it. | |
| pub fn strip_user_message_prefix(text: &str) -> &str { | |
| match text.find(USER_MESSAGE_BEGIN) { | |
| Some(idx) => text[idx + USER_MESSAGE_BEGIN.len()..].trim(), | |
| None => text.trim(), | |
| } | |
| } | |
| // TODO(anp): Replace `TurnEnvironmentSelection` with `PathUri` once path URIs carry environment | |
| // identifiers. | |
| pub struct TurnEnvironmentSelection { | |
| pub environment_id: String, | |
| pub cwd: PathUri, | |
| pub workspace_roots: Vec<PathUri>, | |
| pub config: EnvironmentConfigState, | |
| } | |
| pub struct TurnEnvironmentSelections { | |
| pub legacy_fallback_cwd: AbsolutePathBuf, | |
| pub environments: Vec<TurnEnvironmentSelection>, | |
| } | |
| impl TurnEnvironmentSelections { | |
| pub fn new( | |
| legacy_fallback_cwd: AbsolutePathBuf, | |
| environments: Vec<TurnEnvironmentSelection>, | |
| ) -> Self { | |
| Self { | |
| legacy_fallback_cwd, | |
| environments, | |
| } | |
| } | |
| } | |
| pub struct GitSha(pub String); | |
| impl GitSha { | |
| pub fn new(sha: &str) -> Self { | |
| Self(sha.to_string()) | |
| } | |
| } | |
| /// Submission Queue Entry - requests from user | |
| pub struct Submission { | |
| /// Unique id for this Submission to correlate with Events | |
| pub id: String, | |
| /// Payload | |
| pub op: Op, | |
| /// Optional W3C trace carrier propagated across async submission handoffs. | |
| pub trace: Option<W3cTraceContext>, | |
| /// Core-provided ID of the parent turn that directly initiated this submission. | |
| /// | |
| /// This is only used for inter-agent communication. | |
| pub parent_turn_id: Option<String>, | |
| /// Core-provided ID of the top-level turn that causally initiated this submission. | |
| pub root_turn_id: Option<String>, | |
| } | |
| pub struct W3cTraceContext { | |
| pub traceparent: Option<String>, | |
| pub tracestate: Option<String>, | |
| } | |
| pub struct ConversationStartParams { | |
| /// Whether Codex response handoffs are managed through explicit client append calls. | |
| pub client_managed_handoffs: bool, | |
| /// Whether a realtime V3 delegation produces an acknowledgement filler. | |
| /// `None` preserves the Realtime API's default behavior. | |
| pub delegation_ack_filler: Option<bool>, | |
| /// Whether to route any remaining transcript tail through Codex when the session ends. | |
| /// TODO: Remove this rollout knob once transcript-tail flushing is always enabled. | |
| pub flush_transcript_tail_on_session_end: bool, | |
| /// Sends automatic Codex responses as realtime conversation items instead of handoff appends. | |
| pub codex_responses_as_items: bool, | |
| /// Optional prefix added to automatic Codex response items when `codex_responses_as_items` is set. | |
| pub codex_response_item_prefix: Option<String>, | |
| /// Selects how automatic Codex handoffs are routed in Frameless Bidi sessions. | |
| /// Realtime V1 and V2 ignore this setting. | |
| pub codex_response_handoff_mode: CodexResponseHandoffMode, | |
| /// Optional client-selected BEM prefixes keyed by `analysis`, `commentary`, and `final`. | |
| pub codex_response_handoff_channel_prefixes: Option<BTreeMap<String, Vec<String>>>, | |
| /// Overrides the configured realtime model for this session only. | |
| pub model: Option<String>, | |
| /// Selects whether the realtime session should produce text or audio output. | |
| pub output_modality: RealtimeOutputModality, | |
| /// Whether to append Codex's startup context to the realtime backend prompt. | |
| pub include_startup_context: bool, | |
| /// Complete role-bearing text items to include in the initial realtime session history. | |
| pub initial_items: Vec<ConversationTextParams>, | |
| /// Developer instructions given to Codex when this realtime session starts. | |
| pub realtime_start_instructions: Option<String>, | |
| /// Developer instructions given to Codex when this realtime session ends. | |
| pub realtime_end_instructions: Option<String>, | |
| pub prompt: Option<Option<String>>, | |
| pub realtime_session_id: Option<String>, | |
| pub transport: Option<ConversationStartTransport>, | |
| /// Overrides the configured realtime protocol version for this session only. | |
| pub version: Option<RealtimeConversationVersion>, | |
| pub voice: Option<RealtimeVoice>, | |
| } | |
| pub enum ConversationStartTransport { | |
| Websocket, | |
| Webrtc { | |
| sdp: String, | |
| }, | |
| ExistingCall { | |
| call_id: String, | |
| /// Endpoint selected by the embedding runtime for this call's sideband. | |
| /// This is an in-process override, not a client-supplied API parameter. | |
| /// `None` uses the configured endpoint or the default public API. | |
| sideband_base_url: Option<String>, | |
| }, | |
| } | |
| pub enum RealtimeOutputModality { | |
| Text, | |
| Audio, | |
| } | |
| pub enum RealtimeVoice { | |
| Alloy, | |
| Arbor, | |
| Ash, | |
| Ballad, | |
| Breeze, | |
| Cedar, | |
| Coral, | |
| Cove, | |
| Echo, | |
| Ember, | |
| Juniper, | |
| Maple, | |
| Marin, | |
| Sage, | |
| Shimmer, | |
| Sol, | |
| Spruce, | |
| Vale, | |
| Verse, | |
| } | |
| impl RealtimeVoice { | |
| pub fn wire_name(self) -> &'static str { | |
| match self { | |
| Self::Alloy => "alloy", | |
| Self::Arbor => "arbor", | |
| Self::Ash => "ash", | |
| Self::Ballad => "ballad", | |
| Self::Breeze => "breeze", | |
| Self::Cedar => "cedar", | |
| Self::Coral => "coral", | |
| Self::Cove => "cove", | |
| Self::Echo => "echo", | |
| Self::Ember => "ember", | |
| Self::Juniper => "juniper", | |
| Self::Maple => "maple", | |
| Self::Marin => "marin", | |
| Self::Sage => "sage", | |
| Self::Shimmer => "shimmer", | |
| Self::Sol => "sol", | |
| Self::Spruce => "spruce", | |
| Self::Vale => "vale", | |
| Self::Verse => "verse", | |
| } | |
| } | |
| } | |
| pub struct RealtimeVoicesList { | |
| pub v1: Vec<RealtimeVoice>, | |
| pub v2: Vec<RealtimeVoice>, | |
| pub default_v1: RealtimeVoice, | |
| pub default_v2: RealtimeVoice, | |
| } | |
| impl RealtimeVoicesList { | |
| pub fn builtin() -> Self { | |
| Self { | |
| v1: vec![ | |
| RealtimeVoice::Juniper, | |
| RealtimeVoice::Maple, | |
| RealtimeVoice::Spruce, | |
| RealtimeVoice::Ember, | |
| RealtimeVoice::Vale, | |
| RealtimeVoice::Breeze, | |
| RealtimeVoice::Arbor, | |
| RealtimeVoice::Sol, | |
| RealtimeVoice::Cove, | |
| ], | |
| v2: vec![ | |
| RealtimeVoice::Alloy, | |
| RealtimeVoice::Ash, | |
| RealtimeVoice::Ballad, | |
| RealtimeVoice::Coral, | |
| RealtimeVoice::Echo, | |
| RealtimeVoice::Sage, | |
| RealtimeVoice::Shimmer, | |
| RealtimeVoice::Verse, | |
| RealtimeVoice::Marin, | |
| RealtimeVoice::Cedar, | |
| ], | |
| default_v1: RealtimeVoice::Cove, | |
| default_v2: RealtimeVoice::Marin, | |
| } | |
| } | |
| } | |
| pub struct RealtimeAudioFrame { | |
| pub data: String, | |
| pub sample_rate: u32, | |
| pub num_channels: u16, | |
| pub samples_per_channel: Option<u32>, | |
| pub item_id: Option<String>, | |
| } | |
| pub struct RealtimeTranscriptDelta { | |
| pub delta: String, | |
| } | |
| pub struct RealtimeTranscriptDone { | |
| pub text: String, | |
| } | |
| pub struct RealtimeTranscriptEntry { | |
| pub role: String, | |
| pub text: String, | |
| } | |
| pub struct RealtimeHandoffRequested { | |
| pub handoff_id: String, | |
| pub item_id: String, | |
| pub input_transcript: String, | |
| pub active_transcript: Vec<RealtimeTranscriptEntry>, | |
| } | |
| pub struct RealtimeNoopRequested { | |
| pub call_id: String, | |
| pub item_id: String, | |
| } | |
| pub struct RealtimeInputAudioSpeechStarted { | |
| pub item_id: Option<String>, | |
| } | |
| pub struct RealtimeResponseCancelled { | |
| pub response_id: Option<String>, | |
| } | |
| pub struct RealtimeResponseCreated { | |
| pub response_id: Option<String>, | |
| } | |
| pub struct RealtimeResponseDone { | |
| pub response_id: Option<String>, | |
| } | |
| pub enum RealtimeEvent { | |
| SessionUpdated { | |
| realtime_session_id: String, | |
| instructions: Option<String>, | |
| }, | |
| InputAudioSpeechStarted(RealtimeInputAudioSpeechStarted), | |
| InputTranscriptDelta(RealtimeTranscriptDelta), | |
| InputTranscriptDone(RealtimeTranscriptDone), | |
| OutputTranscriptDelta(RealtimeTranscriptDelta), | |
| OutputTranscriptDone(RealtimeTranscriptDone), | |
| AudioOut(RealtimeAudioFrame), | |
| ResponseCreated(RealtimeResponseCreated), | |
| ResponseCancelled(RealtimeResponseCancelled), | |
| ResponseDone(RealtimeResponseDone), | |
| ConversationItemAdded(Value), | |
| ConversationItemDone { | |
| item_id: String, | |
| }, | |
| /// Canonical display history produced by Core, separate from provider events. | |
| HistoryItemStarted(crate::realtime::RealtimeItem), | |
| HistoryTranscriptDelta { | |
| item_id: String, | |
| delta: String, | |
| }, | |
| HistoryItemCompleted(crate::realtime::RealtimeItem), | |
| HandoffRequested(RealtimeHandoffRequested), | |
| NoopRequested(RealtimeNoopRequested), | |
| Error(String), | |
| } | |
| pub struct ConversationAudioParams { | |
| pub frame: RealtimeAudioFrame, | |
| } | |
| pub struct ConversationTextParams { | |
| pub text: String, | |
| pub role: ConversationTextRole, | |
| } | |
| pub enum ConversationTextRole { | |
| User, | |
| Developer, | |
| Assistant, | |
| } | |
| pub struct ConversationSpeechParams { | |
| pub text: String, | |
| } | |
| /// Supported sparse changes to one live task's current settings, regardless of | |
| /// task kind. Child sessions and consumers of frozen initial settings are unchanged. | |
| pub struct TurnSettingsUpdate { | |
| /// Changes the reviewer for subsequent approval requests, not pending reviews. | |
| pub approvals_reviewer: Option<ApprovalsReviewer>, | |
| pub model: Option<String>, | |
| /// `None` preserves the selection; `Some(None)` clears it. | |
| pub effort: Option<Option<ReasoningEffortConfig>>, | |
| pub summary: Option<ReasoningSummaryConfig>, | |
| /// `None` preserves the requested tier; `Some(None)` clears it. | |
| pub service_tier: Option<Option<String>>, | |
| } | |
| /// The result of processing a turn-settings update, not merely queueing it. | |
| pub enum TurnSettingsUpdateOutcome { | |
| /// Published for subsequent captures; already captured steps are unchanged. | |
| /// The task need not sample or consume every selected preference. | |
| Applied, | |
| /// The named live task was absent or lost before publication. | |
| TargetUnavailable, | |
| Rejected { | |
| reason: String, | |
| }, | |
| } | |
| /// Thread-settings overrides that can be applied before user input or on their | |
| /// own. Standalone updates change the settings inherited by future turns. | |
| pub struct ThreadSettingsOverrides { | |
| /// Updated fallback `cwd` and environments supplied together as a complete pair. | |
| pub environments: Option<TurnEnvironmentSelections>, | |
| /// Updated top-level runtime workspace roots for default environments. | |
| /// Explicit environment selections own their roots separately. | |
| pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>, | |
| /// Updated profile-defined workspace roots for status summaries and | |
| /// per-turn config reconstruction. | |
| pub profile_workspace_roots: Option<Vec<ProfileWorkspaceRoot>>, | |
| /// Updated command approval policy. | |
| pub approval_policy: Option<AskForApproval>, | |
| /// Updated approval reviewer for future approval prompts. | |
| pub approvals_reviewer: Option<ApprovalsReviewer>, | |
| /// Updated sandbox policy for tool calls. | |
| pub sandbox_policy: Option<SandboxPolicy>, | |
| /// Updated permissions profile for tool calls. | |
| pub permission_profile: Option<PermissionProfile>, | |
| /// Named or built-in profile that produced `permission_profile`, if the | |
| /// update selected a profile rather than supplying raw permissions. | |
| pub active_permission_profile: Option<ActivePermissionProfile>, | |
| /// Updated Windows sandbox mode for tool execution. | |
| pub windows_sandbox_level: Option<WindowsSandboxLevel>, | |
| /// Updated model slug. When set, the model info is derived automatically. | |
| pub model: Option<String>, | |
| /// Updated reasoning effort (honored only for reasoning-capable models). | |
| /// | |
| /// Use `Some(Some(_))` to set a specific effort, `Some(None)` to clear the | |
| /// effort, or `None` to leave the existing value unchanged. | |
| pub effort: Option<Option<ReasoningEffortConfig>>, | |
| /// Updated reasoning summary preference (honored only for reasoning-capable models). | |
| pub summary: Option<ReasoningSummaryConfig>, | |
| /// Updated service tier preference for future turns. | |
| /// | |
| /// Use `Some(Some(_))` to set a specific tier, `Some(None)` to clear the | |
| /// preference, or `None` to leave the existing value unchanged. | |
| pub service_tier: Option<Option<String>>, | |
| /// EXPERIMENTAL - set a pre-set collaboration mode. | |
| /// Takes precedence over model, effort, and developer instructions if set. | |
| pub collaboration_mode: Option<CollaborationMode>, | |
| /// Updated personality preference. | |
| pub personality: Option<Personality>, | |
| /// Replace the thread's disabled plugin IDs. Omission preserves the current | |
| /// selection, and an empty list clears it. | |
| pub disabled_plugin_ids: Option<Vec<String>>, | |
| } | |
| /// Source classification for client-supplied context. | |
| pub enum AdditionalContextKind { | |
| Untrusted, | |
| Application, | |
| } | |
| /// Client-supplied context keyed by an opaque source identifier. | |
| pub struct AdditionalContextEntry { | |
| pub value: String, | |
| pub kind: AdditionalContextKind, | |
| } | |
| /// Submission operation | |
| pub enum Op { | |
| /// Abort current task without terminating background terminal processes. | |
| /// This server sends [`EventMsg::TurnAborted`] in response. | |
| Interrupt, | |
| /// Terminate all running background terminal processes for this thread. | |
| /// Use this when callers intentionally want to stop long-lived background shells. | |
| CleanBackgroundTerminals, | |
| /// Start a realtime conversation stream. | |
| RealtimeConversationStart(ConversationStartParams), | |
| /// Send audio input to the running realtime conversation stream. | |
| RealtimeConversationAudio(ConversationAudioParams), | |
| /// Send text input to the running realtime conversation stream. | |
| RealtimeConversationText(ConversationTextParams), | |
| /// Append speakable text to the running realtime conversation stream. | |
| RealtimeConversationSpeech(ConversationSpeechParams), | |
| /// Close the running realtime conversation stream. | |
| RealtimeConversationClose, | |
| /// Request the list of voices supported by realtime conversation streams. | |
| RealtimeConversationListVoices, | |
| /// Submit turn input using the requested routing behavior. | |
| TurnInput { | |
| request: Box<TurnInputRequest>, | |
| mode: TurnInputMode, | |
| reply: oneshot::Sender<CodexResult<TurnInputSubmission>>, | |
| }, | |
| /// Resume an interrupted regular turn. | |
| RecoverTurn { | |
| thread_settings: ThreadSettingsOverrides, | |
| start_options: TurnStartOptions, | |
| reply: oneshot::Sender<CodexResult<TurnInputSubmission>>, | |
| }, | |
| /// Stop the active root turn without recording a terminal turn event. | |
| SuspendTurnAndShutdown { | |
| reply: oneshot::Sender<CodexResult<SuspendTurnOutcome>>, | |
| }, | |
| /// Apply thread-settings overrides without starting a turn. | |
| /// | |
| /// This uses the same submission queue as turn starts so app-server can | |
| /// preserve caller order between both kinds of mutation. | |
| ThreadSettings { | |
| /// Sparse thread-settings overrides to apply. | |
| thread_settings: ThreadSettingsOverrides, | |
| }, | |
| /// Update only the named running turn, without changing future settings. | |
| /// The reply reports the actual publication or why it did not occur. | |
| TurnSettings { | |
| turn_id: String, | |
| update: TurnSettingsUpdate, | |
| reply: oneshot::Sender<TurnSettingsUpdateOutcome>, | |
| }, | |
| /// Inter-agent communication that should be recorded as agent-message history | |
| /// while still using the normal thread submission lifecycle. | |
| InterAgentCommunication { | |
| communication: InterAgentCommunication, | |
| start_options: TurnStartOptions, | |
| }, | |
| /// Approve a command execution | |
| ExecApproval { | |
| /// The id of the submission we are approving | |
| id: String, | |
| /// Turn id associated with the approval event, when available. | |
| turn_id: Option<String>, | |
| /// The user's decision in response to the request. | |
| decision: ReviewDecision, | |
| }, | |
| /// Approve a code patch | |
| PatchApproval { | |
| /// The id of the submission we are approving | |
| id: String, | |
| /// The user's decision in response to the request. | |
| decision: ReviewDecision, | |
| }, | |
| /// Resolve an MCP elicitation request. | |
| ResolveElicitation { | |
| /// Name of the MCP server that issued the request. | |
| server_name: String, | |
| /// Request identifier from the MCP server. | |
| request_id: RequestId, | |
| /// User's decision for the request. | |
| decision: ElicitationAction, | |
| /// Structured user input supplied for accepted elicitations. | |
| content: Option<Value>, | |
| /// Optional client metadata associated with the elicitation response. | |
| meta: Option<Value>, | |
| }, | |
| /// Resolve a request_user_input tool call. | |
| UserInputAnswer { | |
| /// Turn id for the in-flight request. | |
| id: String, | |
| /// User-provided answers. | |
| response: RequestUserInputResponse, | |
| }, | |
| /// Resolve a request_permissions tool call. | |
| RequestPermissionsResponse { | |
| /// Call id for the in-flight request. | |
| id: String, | |
| /// User-granted permissions. | |
| response: RequestPermissionsResponse, | |
| }, | |
| /// Resolve a dynamic tool call request. | |
| DynamicToolResponse { | |
| /// Call id for the in-flight request. | |
| id: String, | |
| /// Tool output payload. | |
| response: DynamicToolResponse, | |
| }, | |
| /// Request MCP servers to reinitialize and refresh cached tool lists. | |
| RefreshMcpServers, | |
| /// Reload user config layer overrides for the active session. | |
| /// | |
| /// This updates runtime config-derived behavior (for example app | |
| /// enable/disable state) without restarting the thread. | |
| ReloadUserConfig, | |
| /// Request the agent to summarize the current conversation context. | |
| /// The agent will use its existing context (either conversation history or previous response id) | |
| /// to generate a summary which will be returned as an AgentMessage event. | |
| Compact, | |
| /// Set whether the thread remains eligible for memory generation. | |
| /// | |
| /// This persists thread-level memory mode metadata without involving the | |
| /// model. | |
| SetThreadMemoryMode { mode: ThreadMemoryMode }, | |
| /// Request a code review from the agent. | |
| Review { review_request: ReviewRequest }, | |
| /// Record that the user approved one retry of a concrete Guardian-denied action. | |
| ApproveGuardianDeniedAction { event: GuardianAssessmentEvent }, | |
| /// Request to shut down codex instance. | |
| Shutdown, | |
| /// Execute a user-initiated one-off shell command (triggered by "!cmd"). | |
| /// | |
| /// The command string is executed using the user's default shell and may | |
| /// include shell syntax (pipes, redirects, etc.). Output is streamed via | |
| /// `ExecCommand*` events and the UI regains control upon `TurnComplete`. | |
| RunUserShellCommand { | |
| /// The raw command string after '!' | |
| command: String, | |
| /// Maximum execution time in milliseconds. Defaults to one hour. | |
| timeout_ms: Option<u64>, | |
| }, | |
| } | |
| pub enum ThreadMemoryMode { | |
| Enabled, | |
| Disabled, | |
| } | |
| pub enum ThreadHistoryMode { | |
| Legacy, | |
| Paginated, | |
| } | |
| impl ThreadHistoryMode { | |
| pub const fn as_str(self) -> &'static str { | |
| match self { | |
| Self::Legacy => "legacy", | |
| Self::Paginated => "paginated", | |
| } | |
| } | |
| } | |
| impl FromStr for ThreadHistoryMode { | |
| type Err = String; | |
| fn from_str(value: &str) -> Result<Self, Self::Err> { | |
| match value { | |
| "legacy" => Ok(Self::Legacy), | |
| "paginated" => Ok(Self::Paginated), | |
| _ => Err(format!("unknown thread history mode `{value}`")), | |
| } | |
| } | |
| } | |
| pub struct InterAgentCommunication { | |
| pub id: Option<ResponseItemId>, | |
| pub author: AgentPath, | |
| pub recipient: AgentPath, | |
| pub other_recipients: Vec<AgentPath>, | |
| pub content: String, | |
| pub encrypted_content: Option<String>, | |
| pub internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>, | |
| pub trigger_turn: bool, | |
| } | |
| impl InterAgentCommunication { | |
| pub fn new( | |
| author: AgentPath, | |
| recipient: AgentPath, | |
| other_recipients: Vec<AgentPath>, | |
| content: String, | |
| trigger_turn: bool, | |
| ) -> Self { | |
| Self { | |
| id: None, | |
| author, | |
| recipient, | |
| other_recipients, | |
| content, | |
| encrypted_content: None, | |
| internal_chat_message_metadata_passthrough: None, | |
| trigger_turn, | |
| } | |
| } | |
| pub fn new_encrypted( | |
| author: AgentPath, | |
| recipient: AgentPath, | |
| other_recipients: Vec<AgentPath>, | |
| encrypted_content: String, | |
| trigger_turn: bool, | |
| ) -> Self { | |
| Self { | |
| id: None, | |
| author, | |
| recipient, | |
| other_recipients, | |
| content: String::new(), | |
| encrypted_content: Some(encrypted_content), | |
| internal_chat_message_metadata_passthrough: None, | |
| trigger_turn, | |
| } | |
| } | |
| pub fn set_turn_id_if_missing(&mut self, turn_id: &str) { | |
| InternalChatMessageMetadataPassthrough::set_turn_id_if_missing( | |
| &mut self.internal_chat_message_metadata_passthrough, | |
| turn_id, | |
| ); | |
| } | |
| pub fn to_response_input_item(&self) -> ResponseInputItem { | |
| let mut communication = self.clone(); | |
| communication.id = None; | |
| communication.internal_chat_message_metadata_passthrough = None; | |
| ResponseInputItem::Message { | |
| role: "assistant".to_string(), | |
| content: vec![ContentItem::OutputText { | |
| text: serde_json::to_string(&communication).unwrap_or_default(), | |
| }], | |
| phase: Some(MessagePhase::Commentary), | |
| } | |
| } | |
| pub fn to_model_input_item(&self) -> ResponseItem { | |
| let content = match &self.encrypted_content { | |
| Some(encrypted_content) => { | |
| let message_type = if self.trigger_turn { | |
| "NEW_TASK" | |
| } else { | |
| "MESSAGE" | |
| }; | |
| vec![ | |
| AgentMessageInputContent::InputText { | |
| text: format!( | |
| "Message Type: {message_type}\nTask name: {}\nSender: {}\nPayload:\n", | |
| self.recipient, self.author | |
| ), | |
| }, | |
| AgentMessageInputContent::EncryptedContent { | |
| encrypted_content: encrypted_content.clone(), | |
| }, | |
| ] | |
| } | |
| None => vec![AgentMessageInputContent::InputText { | |
| text: self.content.clone(), | |
| }], | |
| }; | |
| ResponseItem::AgentMessage { | |
| id: self.id.clone(), | |
| author: self.author.to_string(), | |
| recipient: self.recipient.to_string(), | |
| content, | |
| internal_chat_message_metadata_passthrough: self | |
| .internal_chat_message_metadata_passthrough | |
| .clone(), | |
| } | |
| } | |
| pub fn is_message_content(content: &[ContentItem]) -> bool { | |
| Self::from_message_content(content).is_some() | |
| } | |
| pub fn from_message_content(content: &[ContentItem]) -> Option<Self> { | |
| match content { | |
| [ContentItem::InputText { text }] | [ContentItem::OutputText { text }] => { | |
| serde_json::from_str(text).ok() | |
| } | |
| _ => None, | |
| } | |
| } | |
| } | |
| impl Op { | |
| pub fn kind(&self) -> &'static str { | |
| match self { | |
| Self::Interrupt => "interrupt", | |
| Self::CleanBackgroundTerminals => "clean_background_terminals", | |
| Self::RealtimeConversationStart(_) => "realtime_conversation_start", | |
| Self::RealtimeConversationAudio(_) => "realtime_conversation_audio", | |
| Self::RealtimeConversationText(_) => "realtime_conversation_text", | |
| Self::RealtimeConversationSpeech(_) => "realtime_conversation_speech", | |
| Self::RealtimeConversationClose => "realtime_conversation_close", | |
| Self::RealtimeConversationListVoices => "realtime_conversation_list_voices", | |
| Self::TurnInput { .. } => "turn_input", | |
| Self::RecoverTurn { .. } => "recover_turn", | |
| Self::SuspendTurnAndShutdown { .. } => "suspend_turn_and_shutdown", | |
| Self::ThreadSettings { .. } => "thread_settings", | |
| Self::TurnSettings { .. } => "turn_settings", | |
| Self::InterAgentCommunication { .. } => "inter_agent_communication", | |
| Self::ExecApproval { .. } => "exec_approval", | |
| Self::PatchApproval { .. } => "patch_approval", | |
| Self::ResolveElicitation { .. } => "resolve_elicitation", | |
| Self::UserInputAnswer { .. } => "user_input_answer", | |
| Self::RequestPermissionsResponse { .. } => "request_permissions_response", | |
| Self::DynamicToolResponse { .. } => "dynamic_tool_response", | |
| Self::RefreshMcpServers => "refresh_mcp_servers", | |
| Self::ReloadUserConfig => "reload_user_config", | |
| Self::Compact => "compact", | |
| Self::SetThreadMemoryMode { .. } => "set_thread_memory_mode", | |
| Self::Review { .. } => "review", | |
| Self::ApproveGuardianDeniedAction { .. } => "approve_guardian_denied_action", | |
| Self::Shutdown => "shutdown", | |
| Self::RunUserShellCommand { .. } => "run_user_shell_command", | |
| } | |
| } | |
| } | |
| /// Determines the conditions under which the user is consulted to approve | |
| /// running the command proposed by Codex. | |
| pub enum AskForApproval { | |
| /// Internal policy for projects marked untrusted. Commands require | |
| /// approval unless an explicit exec policy rule allows them. | |
| UnlessTrusted, | |
| /// The model decides when to ask the user for approval. | |
| OnRequest, | |
| /// Fine-grained controls for individual approval flows. | |
| /// | |
| /// When a field is `true`, commands in that category are allowed. When it | |
| /// is `false`, those requests are automatically rejected instead of shown | |
| /// to the user. | |
| Granular(GranularApprovalConfig), | |
| /// Never ask the user to approve commands. Failures are immediately returned | |
| /// to the model, and never escalated to the user for approval. | |
| Never, | |
| } | |
| pub struct GranularApprovalConfig { | |
| /// Whether to allow shell command approval requests, including inline | |
| /// `with_additional_permissions` and `require_escalated` requests. | |
| pub sandbox_approval: bool, | |
| /// Whether to allow prompts triggered by execpolicy `prompt` rules. | |
| pub rules: bool, | |
| /// Whether to allow approval prompts triggered by skill script execution. | |
| pub skill_approval: bool, | |
| /// Whether to allow prompts triggered by the `request_permissions` tool. | |
| pub request_permissions: bool, | |
| /// Whether to allow MCP elicitation prompts. | |
| pub mcp_elicitations: bool, | |
| } | |
| impl GranularApprovalConfig { | |
| pub const fn allows_sandbox_approval(self) -> bool { | |
| self.sandbox_approval | |
| } | |
| pub const fn allows_rules_approval(self) -> bool { | |
| self.rules | |
| } | |
| pub const fn allows_skill_approval(self) -> bool { | |
| self.skill_approval | |
| } | |
| pub const fn allows_request_permissions(self) -> bool { | |
| self.request_permissions | |
| } | |
| pub const fn allows_mcp_elicitations(self) -> bool { | |
| self.mcp_elicitations | |
| } | |
| } | |
| /// Represents whether outbound network access is available to the agent. | |
| pub enum NetworkAccess { | |
| Restricted, | |
| Enabled, | |
| } | |
| impl NetworkAccess { | |
| pub fn is_enabled(self) -> bool { | |
| matches!(self, NetworkAccess::Enabled) | |
| } | |
| } | |
| /// Determines execution restrictions for model shell commands. | |
| pub enum SandboxPolicy { | |
| /// No restrictions whatsoever. Use with caution. | |
| DangerFullAccess, | |
| /// Read-only access configuration. | |
| ReadOnly { | |
| /// When set to `true`, outbound network access is allowed. `false` by | |
| /// default. | |
| network_access: bool, | |
| }, | |
| /// Indicates the process is already in an external sandbox. Allows full | |
| /// disk access while honoring the provided network setting. | |
| ExternalSandbox { | |
| /// Whether the external sandbox permits outbound network traffic. | |
| network_access: NetworkAccess, | |
| }, | |
| /// Same as `ReadOnly` but additionally grants write access to the current | |
| /// working directory ("workspace"). | |
| WorkspaceWrite { | |
| /// Additional folders (beyond cwd and possibly TMPDIR) that should be | |
| /// writable from within the sandbox. | |
| writable_roots: Vec<AbsolutePathBuf>, | |
| /// When set to `true`, outbound network access is allowed. `false` by | |
| /// default. | |
| network_access: bool, | |
| /// When set to `true`, will NOT include the per-user `TMPDIR` | |
| /// environment variable among the default writable roots. Defaults to | |
| /// `false`. | |
| exclude_tmpdir_env_var: bool, | |
| /// When set to `true`, will NOT include the `/tmp` among the default | |
| /// writable roots on UNIX. Defaults to `false`. | |
| exclude_slash_tmp: bool, | |
| }, | |
| } | |
| /// A writable root path accompanied by a list of subpaths that should remain | |
| /// read‑only even when the root is writable. This is primarily used to ensure | |
| /// that folders containing files that could be modified to escalate the | |
| /// privileges of the agent (e.g. `.codex`, `.git`, notably `.git/hooks`) under | |
| /// a writable root are not modified by the agent. | |
| pub struct WritableRoot { | |
| pub root: AbsolutePathBuf, | |
| /// By construction, these subpaths are all under `root`. | |
| pub read_only_subpaths: Vec<AbsolutePathBuf>, | |
| /// Workspace metadata path names that must not be created or replaced under | |
| /// `root` unless the policy grants an explicit write rule for that metadata | |
| /// path. | |
| pub protected_metadata_names: Vec<String>, | |
| } | |
| impl WritableRoot { | |
| pub fn is_path_writable(&self, path: &Path) -> bool { | |
| // Check if the path is under the root. | |
| if !path.starts_with(&self.root) { | |
| return false; | |
| } | |
| // Check if the path is under any of the read-only subpaths. | |
| for subpath in &self.read_only_subpaths { | |
| if path.starts_with(subpath) { | |
| return false; | |
| } | |
| } | |
| if self.path_contains_protected_metadata_name(path) { | |
| return false; | |
| } | |
| true | |
| } | |
| fn path_contains_protected_metadata_name(&self, path: &Path) -> bool { | |
| let Ok(relative_path) = path.strip_prefix(&self.root) else { | |
| return false; | |
| }; | |
| let Some(first_component) = relative_path.components().next() else { | |
| return false; | |
| }; | |
| self.protected_metadata_names | |
| .iter() | |
| .any(|name| first_component.as_os_str() == std::ffi::OsStr::new(name)) | |
| } | |
| } | |
| impl FromStr for SandboxPolicy { | |
| type Err = serde_json::Error; | |
| fn from_str(s: &str) -> Result<Self, Self::Err> { | |
| serde_json::from_str(s) | |
| } | |
| } | |
| impl FromStr for FileSystemSandboxPolicy { | |
| type Err = serde_json::Error; | |
| fn from_str(s: &str) -> Result<Self, Self::Err> { | |
| serde_json::from_str::<RawFileSystemSandboxPolicy>(s)? | |
| .try_into() | |
| .map_err(serde_json::Error::custom) | |
| } | |
| } | |
| impl FromStr for NetworkSandboxPolicy { | |
| type Err = serde_json::Error; | |
| fn from_str(s: &str) -> Result<Self, Self::Err> { | |
| serde_json::from_str(s) | |
| } | |
| } | |
| impl SandboxPolicy { | |
| /// Returns a policy with read-only disk access and no network. | |
| pub fn new_read_only_policy() -> Self { | |
| SandboxPolicy::ReadOnly { | |
| network_access: false, | |
| } | |
| } | |
| /// Returns a policy that can read the entire disk, but can only write to | |
| /// the current working directory and the per-user tmp dir on macOS. It does | |
| /// not allow network access. | |
| pub fn new_workspace_write_policy() -> Self { | |
| SandboxPolicy::WorkspaceWrite { | |
| writable_roots: vec![], | |
| network_access: false, | |
| exclude_tmpdir_env_var: false, | |
| exclude_slash_tmp: false, | |
| } | |
| } | |
| pub fn has_full_disk_read_access(&self) -> bool { | |
| true | |
| } | |
| pub fn has_full_disk_write_access(&self) -> bool { | |
| match self { | |
| SandboxPolicy::DangerFullAccess => true, | |
| SandboxPolicy::ExternalSandbox { .. } => true, | |
| SandboxPolicy::ReadOnly { .. } => false, | |
| SandboxPolicy::WorkspaceWrite { .. } => false, | |
| } | |
| } | |
| pub fn has_full_network_access(&self) -> bool { | |
| match self { | |
| SandboxPolicy::DangerFullAccess => true, | |
| SandboxPolicy::ExternalSandbox { network_access } => network_access.is_enabled(), | |
| SandboxPolicy::ReadOnly { network_access, .. } => *network_access, | |
| SandboxPolicy::WorkspaceWrite { network_access, .. } => *network_access, | |
| } | |
| } | |
| /// Returns the list of writable roots (tailored to the current working | |
| /// directory) together with subpaths that should remain read‑only under | |
| /// each writable root. | |
| pub fn get_writable_roots_with_cwd(&self, cwd: &Path) -> Vec<WritableRoot> { | |
| match self { | |
| SandboxPolicy::DangerFullAccess => Vec::new(), | |
| SandboxPolicy::ExternalSandbox { .. } => Vec::new(), | |
| SandboxPolicy::ReadOnly { .. } => Vec::new(), | |
| SandboxPolicy::WorkspaceWrite { | |
| writable_roots, | |
| exclude_tmpdir_env_var, | |
| exclude_slash_tmp, | |
| network_access: _, | |
| } => { | |
| // Start from explicitly configured writable roots. | |
| let mut roots: Vec<AbsolutePathBuf> = writable_roots.clone(); | |
| // Always include defaults: cwd, /tmp (if present on Unix), and | |
| // on macOS, the per-user TMPDIR unless explicitly excluded. | |
| // TODO(mbolin): cwd param should be AbsolutePathBuf. | |
| let cwd_absolute = AbsolutePathBuf::from_absolute_path(cwd); | |
| match cwd_absolute { | |
| Ok(cwd) => { | |
| roots.push(cwd); | |
| } | |
| Err(e) => { | |
| error!( | |
| "Ignoring invalid cwd {:?} for sandbox writable root: {}", | |
| cwd, e | |
| ); | |
| } | |
| } | |
| // Include /tmp on Unix unless explicitly excluded. | |
| if cfg!(unix) && !exclude_slash_tmp { | |
| match AbsolutePathBuf::from_absolute_path("/tmp") { | |
| Ok(slash_tmp) => { | |
| if slash_tmp.as_path().is_dir() { | |
| roots.push(slash_tmp); | |
| } | |
| } | |
| Err(e) => { | |
| error!("Ignoring invalid /tmp for sandbox writable root: {e}"); | |
| } | |
| } | |
| } | |
| // Include $TMPDIR unless explicitly excluded. On macOS, TMPDIR | |
| // is per-user, so writes to TMPDIR should not be readable by | |
| // other users on the system. | |
| // | |
| // By comparison, TMPDIR is not guaranteed to be defined on | |
| // Linux or Windows, but supporting it here gives users a way to | |
| // provide the model with their own temporary directory without | |
| // having to hardcode it in the config. | |
| if !exclude_tmpdir_env_var | |
| && let Some(tmpdir) = std::env::var_os("TMPDIR") | |
| && !tmpdir.is_empty() | |
| { | |
| match AbsolutePathBuf::from_absolute_path(PathBuf::from(&tmpdir)) { | |
| Ok(tmpdir_path) => { | |
| roots.push(tmpdir_path); | |
| } | |
| Err(e) => { | |
| error!( | |
| "Ignoring invalid TMPDIR value {tmpdir:?} for sandbox writable root: {e}", | |
| ); | |
| } | |
| } | |
| } | |
| // For each root, compute subpaths that should remain read-only. | |
| let cwd_root = AbsolutePathBuf::from_absolute_path(cwd).ok(); | |
| roots | |
| .into_iter() | |
| .map(|writable_root| { | |
| let protect_missing_dot_codex = cwd_root | |
| .as_ref() | |
| .is_some_and(|cwd_root| cwd_root == &writable_root); | |
| WritableRoot { | |
| read_only_subpaths: default_read_only_subpaths_for_writable_root( | |
| &writable_root, | |
| protect_missing_dot_codex, | |
| ), | |
| protected_metadata_names: Vec::new(), | |
| root: writable_root, | |
| } | |
| }) | |
| .collect() | |
| } | |
| } | |
| } | |
| } | |
| /// Event Queue Entry - events from agent | |
| pub struct Event { | |
| /// Submission `id` that this event is correlated with. | |
| pub id: String, | |
| /// Payload | |
| pub msg: EventMsg, | |
| } | |
| pub struct EnvironmentConnectionEvent { | |
| pub environment_id: String, | |
| } | |
| /// Response event from the agent | |
| /// NOTE: Make sure none of these values have optional types, as it will mess up the extension code-gen. | |
| pub enum EventMsg { | |
| /// Error while executing a submission | |
| Error(ErrorEvent), | |
| /// Warning issued while processing a submission. Unlike `Error`, this | |
| /// indicates the turn continued but the user should still be notified. | |
| Warning(WarningEvent), | |
| /// Provider-owned authentication recovery has started for the current turn. | |
| AuthRecoveryStarted(AuthRecoveryEvent), | |
| /// Provider-owned authentication recovery has completed for the current turn. | |
| AuthRecoveryCompleted(AuthRecoveryEvent), | |
| /// Warning issued by the guardian automatic approval reviewer. | |
| GuardianWarning(WarningEvent), | |
| /// Realtime conversation lifecycle start event. | |
| RealtimeConversationStarted(RealtimeConversationStartedEvent), | |
| /// Realtime conversation streaming payload event. | |
| RealtimeConversationRealtime(RealtimeConversationRealtimeEvent), | |
| /// Realtime conversation lifecycle close event. | |
| RealtimeConversationClosed(RealtimeConversationClosedEvent), | |
| /// Realtime session description protocol payload. | |
| RealtimeConversationSdp(RealtimeConversationSdpEvent), | |
| /// Model routing changed from the requested model to a different model. | |
| ModelReroute(ModelRerouteEvent), | |
| /// Backend recommends additional account verification for this turn. | |
| ModelVerification(ModelVerificationEvent), | |
| /// Backend moderation metadata intended for first-party turn presentation. | |
| TurnModerationMetadata(TurnModerationMetadataEvent), | |
| /// Backend indicates that response output is waiting on a safety review. | |
| SafetyBuffering(SafetyBufferingEvent), | |
| /// Conversation history was compacted (either automatically or manually). | |
| ContextCompacted(ContextCompactedEvent), | |
| /// Legacy persisted marker for dropping the last N user turns. | |
| /// Retained for replay of existing rollouts; live rollback operations are unsupported. | |
| ThreadRolledBack(ThreadRolledBackEvent), | |
| /// Agent has started a turn. | |
| /// v1 wire format uses `task_started`; accept `turn_started` for v2 interop. | |
| TurnStarted(TurnStartedEvent), | |
| /// Persistent thread-settings overrides from the correlated submission have | |
| /// been applied to the session configuration. | |
| ThreadSettingsApplied(ThreadSettingsAppliedEvent), | |
| /// Agent has completed all actions. | |
| /// v1 wire format uses `task_complete`; accept `turn_complete` for v2 interop. | |
| TurnComplete(TurnCompleteEvent), | |
| /// Usage update for the current session, including totals and last turn. | |
| /// Optional means unknown — UIs should not display when `None`. | |
| TokenCount(TokenCountEvent), | |
| /// Agent text output message | |
| AgentMessage(AgentMessageEvent), | |
| /// User/system input message (what was sent to the model) | |
| UserMessage(UserMessageEvent), | |
| /// Reasoning event from agent. | |
| AgentReasoning(AgentReasoningEvent), | |
| /// Raw chain-of-thought from agent. | |
| AgentReasoningRawContent(AgentReasoningRawContentEvent), | |
| /// Signaled when the model begins a new reasoning summary section (e.g., a new titled block). | |
| AgentReasoningSectionBreak(AgentReasoningSectionBreakEvent), | |
| /// Ack the client's configure message. | |
| SessionConfigured(SessionConfiguredEvent), | |
| /// A selected environment completed its connection handshake. | |
| EnvironmentConnected(EnvironmentConnectionEvent), | |
| /// A selected environment lost its established connection. | |
| EnvironmentDisconnected(EnvironmentConnectionEvent), | |
| /// Updated long-running goal metadata for the thread. | |
| ThreadGoalUpdated(ThreadGoalUpdatedEvent), | |
| /// A durable thread-scoped user-message queue changed. | |
| ThreadQueueChanged(ThreadQueueChangedEvent), | |
| /// Incremental MCP startup progress updates. | |
| McpStartupUpdate(McpStartupUpdateEvent), | |
| /// Aggregate MCP startup completion summary. | |
| McpStartupComplete(McpStartupCompleteEvent), | |
| McpToolCallBegin(McpToolCallBeginEvent), | |
| McpToolCallEnd(McpToolCallEndEvent), | |
| WebSearchBegin(WebSearchBeginEvent), | |
| WebSearchEnd(WebSearchEndEvent), | |
| ImageGenerationBegin(ImageGenerationBeginEvent), | |
| ImageGenerationEnd(ImageGenerationEndEvent), | |
| /// Notification that the server is about to execute a command. | |
| ExecCommandBegin(ExecCommandBeginEvent), | |
| /// Incremental chunk of output from a running command. | |
| ExecCommandOutputDelta(ExecCommandOutputDeltaEvent), | |
| /// Terminal interaction for an in-progress command (stdin sent and stdout observed). | |
| TerminalInteraction(TerminalInteractionEvent), | |
| ExecCommandEnd(ExecCommandEndEvent), | |
| /// Notification that the agent attached a local image via the view_image tool. | |
| ViewImageToolCall(ViewImageToolCallEvent), | |
| ExecApprovalRequest(ExecApprovalRequestEvent), | |
| RequestPermissions(RequestPermissionsEvent), | |
| RequestUserInput(RequestUserInputEvent), | |
| DynamicToolCallRequest(DynamicToolCallRequest), | |
| DynamicToolCallResponse(DynamicToolCallResponseEvent), | |
| ElicitationRequest(ElicitationRequestEvent), | |
| ApplyPatchApprovalRequest(ApplyPatchApprovalRequestEvent), | |
| /// Structured lifecycle event for a guardian-reviewed approval request. | |
| GuardianAssessment(GuardianAssessmentEvent), | |
| /// Notification advising the user that something they are using has been | |
| /// deprecated and should be phased out. | |
| DeprecationNotice(DeprecationNoticeEvent), | |
| /// Notification that a model stream experienced an error or disconnect | |
| /// and the system is handling it (e.g., retrying with backoff). | |
| StreamError(StreamErrorEvent), | |
| /// Notification that the agent is about to apply a code patch. Mirrors | |
| /// `ExecCommandBegin` so front‑ends can show progress indicators. | |
| PatchApplyBegin(PatchApplyBeginEvent), | |
| /// Latest model-generated structured changes for an `apply_patch` call. | |
| PatchApplyUpdated(PatchApplyUpdatedEvent), | |
| /// Notification that a patch application has finished. | |
| PatchApplyEnd(PatchApplyEndEvent), | |
| TurnDiff(TurnDiffEvent), | |
| /// List of voices supported by realtime conversation streams. | |
| RealtimeConversationListVoicesResponse(RealtimeConversationListVoicesResponseEvent), | |
| PlanUpdate(UpdatePlanArgs), | |
| TurnAborted(TurnAbortedEvent), | |
| /// Notification that the agent is shutting down. | |
| ShutdownComplete, | |
| /// Entered review mode. | |
| EnteredReviewMode(EnteredReviewModeEvent), | |
| /// Exited review mode with an optional final result to apply. | |
| ExitedReviewMode(ExitedReviewModeEvent), | |
| RawResponseItem(RawResponseItemEvent), | |
| RawResponseCompleted(RawResponseCompletedEvent), | |
| ItemStarted(ItemStartedEvent), | |
| ItemCompleted(ItemCompletedEvent), | |
| HookStarted(HookStartedEvent), | |
| HookCompleted(HookCompletedEvent), | |
| AgentMessageContentDelta(AgentMessageContentDeltaEvent), | |
| PlanDelta(PlanDeltaEvent), | |
| ReasoningContentDelta(ReasoningContentDeltaEvent), | |
| ReasoningRawContentDelta(ReasoningRawContentDeltaEvent), | |
| /// Collab interaction: agent spawn begin. | |
| CollabAgentSpawnBegin(CollabAgentSpawnBeginEvent), | |
| /// Collab interaction: agent spawn end. | |
| CollabAgentSpawnEnd(CollabAgentSpawnEndEvent), | |
| /// Collab interaction: agent interaction begin. | |
| CollabAgentInteractionBegin(CollabAgentInteractionBeginEvent), | |
| /// Collab interaction: agent interaction end. | |
| CollabAgentInteractionEnd(CollabAgentInteractionEndEvent), | |
| /// Collab interaction: waiting begin. | |
| CollabWaitingBegin(CollabWaitingBeginEvent), | |
| /// Collab interaction: waiting end. | |
| CollabWaitingEnd(CollabWaitingEndEvent), | |
| /// Collab interaction: close begin. | |
| CollabCloseBegin(CollabCloseBeginEvent), | |
| /// Collab interaction: close end. | |
| CollabCloseEnd(CollabCloseEndEvent), | |
| /// Collab interaction: resume begin. | |
| CollabResumeBegin(CollabResumeBeginEvent), | |
| /// Collab interaction: resume end. | |
| CollabResumeEnd(CollabResumeEndEvent), | |
| /// Path-based v2 sub-agent activity. | |
| SubAgentActivity(SubAgentActivityEvent), | |
| } | |
| pub enum HookEventName { | |
| PreToolUse, | |
| PermissionRequest, | |
| PostToolUse, | |
| PreCompact, | |
| PostCompact, | |
| SessionStart, | |
| SessionEnd, | |
| UserPromptSubmit, | |
| SubagentStart, | |
| SubagentStop, | |
| Stop, | |
| Interrupt, | |
| } | |
| pub enum HookHandlerType { | |
| Command, | |
| McpTool, | |
| Prompt, | |
| Agent, | |
| } | |
| pub enum HookExecutionMode { | |
| Sync, | |
| Async, | |
| } | |
| pub enum HookScope { | |
| Thread, | |
| Turn, | |
| } | |
| pub enum HookSource { | |
| System, | |
| User, | |
| Project, | |
| Mdm, | |
| SessionFlags, | |
| Plugin, | |
| CloudRequirements, | |
| CloudManagedConfig, | |
| LegacyManagedConfigFile, | |
| LegacyManagedConfigMdm, | |
| Unknown, | |
| } | |
| pub enum HookTrustStatus { | |
| Managed, | |
| Untrusted, | |
| Trusted, | |
| Modified, | |
| } | |
| pub enum HookRunStatus { | |
| Running, | |
| Completed, | |
| Failed, | |
| Blocked, | |
| Stopped, | |
| } | |
| pub enum HookOutputEntryKind { | |
| Warning, | |
| Stop, | |
| Feedback, | |
| Context, | |
| Error, | |
| } | |
| pub struct HookOutputEntry { | |
| pub kind: HookOutputEntryKind, | |
| pub text: String, | |
| } | |
| pub struct HookRunSummary { | |
| /// Internal classification used to suppress lifecycle notifications without losing telemetry. | |
| pub builtin: bool, | |
| pub id: String, | |
| pub event_name: HookEventName, | |
| pub handler_type: HookHandlerType, | |
| pub execution_mode: HookExecutionMode, | |
| pub scope: HookScope, | |
| pub source_path: AbsolutePathBuf, | |
| pub source: HookSource, | |
| pub display_order: i64, | |
| pub status: HookRunStatus, | |
| pub status_message: Option<String>, | |
| pub started_at: i64, | |
| pub completed_at: Option<i64>, | |
| pub duration_ms: Option<i64>, | |
| pub entries: Vec<HookOutputEntry>, | |
| } | |
| pub struct HookStartedEvent { | |
| pub turn_id: Option<String>, | |
| pub run: HookRunSummary, | |
| } | |
| pub struct HookCompletedEvent { | |
| pub turn_id: Option<String>, | |
| pub run: HookRunSummary, | |
| } | |
| pub enum RealtimeConversationVersion { | |
| V1, | |
| V2, | |
| V3, | |
| } | |
| pub enum CodexResponseHandoffMode { | |
| Thinking, | |
| Commentary, | |
| BemTags, | |
| } | |
| pub struct RealtimeConversationStartedEvent { | |
| pub realtime_session_id: Option<String>, | |
| pub version: RealtimeConversationVersion, | |
| } | |
| pub struct RealtimeConversationRealtimeEvent { | |
| pub payload: RealtimeEvent, | |
| } | |
| pub struct RealtimeConversationClosedEvent { | |
| pub reason: Option<String>, | |
| } | |
| pub struct RealtimeConversationSdpEvent { | |
| pub sdp: String, | |
| } | |
| impl From<CollabAgentSpawnBeginEvent> for EventMsg { | |
| fn from(event: CollabAgentSpawnBeginEvent) -> Self { | |
| EventMsg::CollabAgentSpawnBegin(event) | |
| } | |
| } | |
| impl From<CollabAgentSpawnEndEvent> for EventMsg { | |
| fn from(event: CollabAgentSpawnEndEvent) -> Self { | |
| EventMsg::CollabAgentSpawnEnd(event) | |
| } | |
| } | |
| impl From<CollabAgentInteractionBeginEvent> for EventMsg { | |
| fn from(event: CollabAgentInteractionBeginEvent) -> Self { | |
| EventMsg::CollabAgentInteractionBegin(event) | |
| } | |
| } | |
| impl From<CollabAgentInteractionEndEvent> for EventMsg { | |
| fn from(event: CollabAgentInteractionEndEvent) -> Self { | |
| EventMsg::CollabAgentInteractionEnd(event) | |
| } | |
| } | |
| impl From<CollabWaitingBeginEvent> for EventMsg { | |
| fn from(event: CollabWaitingBeginEvent) -> Self { | |
| EventMsg::CollabWaitingBegin(event) | |
| } | |
| } | |
| impl From<CollabWaitingEndEvent> for EventMsg { | |
| fn from(event: CollabWaitingEndEvent) -> Self { | |
| EventMsg::CollabWaitingEnd(event) | |
| } | |
| } | |
| impl From<CollabCloseBeginEvent> for EventMsg { | |
| fn from(event: CollabCloseBeginEvent) -> Self { | |
| EventMsg::CollabCloseBegin(event) | |
| } | |
| } | |
| impl From<CollabCloseEndEvent> for EventMsg { | |
| fn from(event: CollabCloseEndEvent) -> Self { | |
| EventMsg::CollabCloseEnd(event) | |
| } | |
| } | |
| impl From<CollabResumeBeginEvent> for EventMsg { | |
| fn from(event: CollabResumeBeginEvent) -> Self { | |
| EventMsg::CollabResumeBegin(event) | |
| } | |
| } | |
| impl From<CollabResumeEndEvent> for EventMsg { | |
| fn from(event: CollabResumeEndEvent) -> Self { | |
| EventMsg::CollabResumeEnd(event) | |
| } | |
| } | |
| impl From<SubAgentActivityEvent> for EventMsg { | |
| fn from(event: SubAgentActivityEvent) -> Self { | |
| EventMsg::SubAgentActivity(event) | |
| } | |
| } | |
| /// Agent lifecycle status, derived from emitted events. | |
| pub enum AgentStatus { | |
| /// Agent is waiting for initialization. | |
| PendingInit, | |
| /// Agent is currently running. | |
| Running, | |
| /// Agent's current turn was interrupted and it may receive more input. | |
| Interrupted, | |
| /// Agent is done. Contains the final assistant message. | |
| Completed(Option<String>), | |
| /// Agent encountered an error. | |
| Errored(String), | |
| /// Agent has been shutdown. | |
| Shutdown, | |
| /// Agent is not found. | |
| NotFound, | |
| } | |
| /// Turn kinds that reject same-turn steering. | |
| pub enum NonSteerableTurnKind { | |
| Review, | |
| Compact, | |
| } | |
| /// Codex errors that we expose to clients. | |
| pub enum CodexErrorInfo { | |
| ContextWindowExceeded, | |
| SessionBudgetExceeded, | |
| UsageLimitExceeded, | |
| RateLimitExceeded, | |
| ServerOverloaded, | |
| CyberPolicy, | |
| MisalignmentPolicyViolation, | |
| HttpConnectionFailed { | |
| http_status_code: Option<u16>, | |
| }, | |
| /// Failed to connect to the response SSE stream. | |
| ResponseStreamConnectionFailed { | |
| http_status_code: Option<u16>, | |
| }, | |
| InternalServerError, | |
| Unauthorized, | |
| BadRequest, | |
| SandboxError, | |
| /// The response SSE stream disconnected in the middle of a turnbefore completion. | |
| ResponseStreamDisconnected { | |
| http_status_code: Option<u16>, | |
| }, | |
| /// Reached the retry limit for responses. | |
| ResponseTooManyFailedAttempts { | |
| http_status_code: Option<u16>, | |
| }, | |
| /// Returned when `turn/start` or `turn/steer` is submitted while the current active turn | |
| /// cannot accept same-turn steering, for example `/review` or manual `/compact`. | |
| ActiveTurnNotSteerable { | |
| turn_kind: NonSteerableTurnKind, | |
| }, | |
| // Retained to deserialize errors recorded in legacy rollouts. | |
| ThreadRollbackFailed, | |
| Other, | |
| } | |
| impl CodexErrorInfo { | |
| /// Whether this error should mark the current turn as failed when replaying history. | |
| pub fn affects_turn_status(&self) -> bool { | |
| match self { | |
| Self::ThreadRollbackFailed | Self::ActiveTurnNotSteerable { .. } => false, | |
| Self::ContextWindowExceeded | |
| | Self::SessionBudgetExceeded | |
| | Self::UsageLimitExceeded | |
| | Self::RateLimitExceeded | |
| | Self::ServerOverloaded | |
| | Self::CyberPolicy | |
| | Self::MisalignmentPolicyViolation | |
| | Self::HttpConnectionFailed { .. } | |
| | Self::ResponseStreamConnectionFailed { .. } | |
| | Self::InternalServerError | |
| | Self::Unauthorized | |
| | Self::BadRequest | |
| | Self::SandboxError | |
| | Self::ResponseStreamDisconnected { .. } | |
| | Self::ResponseTooManyFailedAttempts { .. } | |
| | Self::Other => true, | |
| } | |
| } | |
| } | |
| pub struct RawResponseItemEvent { | |
| pub item: ResponseItem, | |
| } | |
| /// Exact usage and metadata reported by one upstream Responses API completion. | |
| /// | |
| /// Unlike TokenCountEvent, this is not accumulated, estimated, or replayed. | |
| pub struct RawResponseCompletedEvent { | |
| pub response_id: String, | |
| pub token_usage: Option<TokenUsage>, | |
| pub usage_metadata: Option<crate::ResponseUsageMetadata>, | |
| } | |
| pub struct ItemStartedEvent { | |
| pub thread_id: ThreadId, | |
| pub turn_id: String, | |
| pub item: TurnItem, | |
| pub started_at_ms: i64, | |
| } | |
| pub struct ItemCompletedEvent { | |
| pub thread_id: ThreadId, | |
| pub turn_id: String, | |
| pub item: TurnItem, | |
| pub started_at_ms: Option<i64>, | |
| // Old rollout files may contain ItemCompleted events for PlanItem without | |
| // this field. Default to 0 so those persisted rollouts still deserialize | |
| // after tightening the core event contract. | |
| pub completed_at_ms: i64, | |
| } | |
| const fn default_item_completed_at_ms() -> i64 { | |
| 0 | |
| } | |
| pub struct AgentMessageContentDeltaEvent { | |
| pub thread_id: String, | |
| pub turn_id: String, | |
| pub item_id: String, | |
| pub delta: String, | |
| } | |
| pub struct PlanDeltaEvent { | |
| pub thread_id: String, | |
| pub turn_id: String, | |
| pub item_id: String, | |
| pub delta: String, | |
| } | |
| pub struct ReasoningContentDeltaEvent { | |
| pub thread_id: String, | |
| pub turn_id: String, | |
| pub item_id: String, | |
| pub delta: String, | |
| // load with default value so it's backward compatible with the old format. | |
| pub summary_index: i64, | |
| } | |
| pub struct ReasoningRawContentDeltaEvent { | |
| pub thread_id: String, | |
| pub turn_id: String, | |
| pub item_id: String, | |
| pub delta: String, | |
| // load with default value so it's backward compatible with the old format. | |
| pub content_index: i64, | |
| } | |
| pub struct EnteredReviewModeEvent { | |
| pub target: ReviewTarget, | |
| pub user_facing_hint: Option<String>, | |
| pub turn_id: Option<String>, | |
| pub item_id: Option<String>, | |
| } | |
| pub struct ExitedReviewModeEvent { | |
| pub turn_id: Option<String>, | |
| pub item_id: Option<String>, | |
| pub review_output: Option<ReviewOutputEvent>, | |
| } | |
| // Individual event payload types matching each `EventMsg` variant. | |
| /// Public, customer-facing details supplied by the Responses API for a misalignment block. | |
| pub struct MisalignmentErrorDetails { | |
| /// Open-ended classification; new values must not prevent the error from being surfaced. | |
| pub error_type: Option<String>, | |
| /// A localized explanation is required before a client may offer continuation. | |
| pub detailed_explanation: Option<String>, | |
| /// Model-visible instruction to submit if the user elects to continue. | |
| pub steer: Option<MisalignmentSteer>, | |
| } | |
| impl fmt::Debug for MisalignmentErrorDetails { | |
| fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| formatter | |
| .debug_struct("MisalignmentErrorDetails") | |
| .field("error_type", &self.error_type) | |
| .field( | |
| "has_detailed_explanation", | |
| &self.detailed_explanation.is_some(), | |
| ) | |
| .field("has_steer", &self.steer.is_some()) | |
| .finish() | |
| } | |
| } | |
| /// Public steering instruction returned alongside a resumable misalignment block. | |
| pub struct MisalignmentSteer { | |
| pub message: String, | |
| } | |
| impl fmt::Debug for MisalignmentSteer { | |
| fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| formatter | |
| .debug_struct("MisalignmentSteer") | |
| .field("message", &"[REDACTED]") | |
| .finish() | |
| } | |
| } | |
| pub struct ErrorEvent { | |
| pub message: String, | |
| pub codex_error_info: Option<CodexErrorInfo>, | |
| /// Sensitive explanation and steering are delivered live but never enter rollout storage. | |
| pub misalignment: Option<MisalignmentErrorDetails>, | |
| } | |
| impl ErrorEvent { | |
| /// Whether this error should mark the current turn as failed when replaying history. | |
| pub fn affects_turn_status(&self) -> bool { | |
| self.codex_error_info | |
| .as_ref() | |
| .is_none_or(CodexErrorInfo::affects_turn_status) | |
| } | |
| } | |
| pub struct WarningEvent { | |
| pub message: String, | |
| } | |
| /// User-facing progress for provider-owned authentication recovery. | |
| pub struct AuthRecoveryEvent { | |
| /// Display name of the model provider whose authentication is recovering. | |
| pub provider: String, | |
| /// User-facing description of the authentication recovery stage. | |
| pub message: String, | |
| } | |
| pub enum ModelRerouteReason { | |
| HighRiskCyberActivity, | |
| } | |
| pub struct ModelRerouteEvent { | |
| pub from_model: String, | |
| pub to_model: String, | |
| pub reason: ModelRerouteReason, | |
| } | |
| pub enum ModelVerification { | |
| TrustedAccessForCyber, | |
| } | |
| pub struct ModelVerificationEvent { | |
| pub verifications: Vec<ModelVerification>, | |
| } | |
| pub struct TurnModerationMetadataEvent { | |
| pub metadata: Value, | |
| } | |
| pub struct SafetyBufferingEvent { | |
| pub model: String, | |
| pub use_cases: Vec<String>, | |
| pub reasons: Vec<String>, | |
| pub show_buffering_ui: bool, | |
| pub faster_model: Option<String>, | |
| } | |
| pub struct ContextCompactedEvent; | |
| pub struct TurnCompleteEvent { | |
| pub turn_id: String, | |
| pub last_agent_message: Option<String>, | |
| /// Terminal error details when the turn completed unsuccessfully. | |
| pub error: Option<ErrorEvent>, | |
| /// Unix timestamp (in seconds) when the turn started. | |
| pub started_at: Option<i64>, | |
| /// Unix timestamp (in seconds) when the turn completed. | |
| pub completed_at: Option<i64>, | |
| /// Duration between turn start and completion in milliseconds, if known. | |
| pub duration_ms: Option<i64>, | |
| /// Duration between turn start and the first model token in milliseconds, if known. | |
| pub time_to_first_token_ms: Option<i64>, | |
| } | |
| pub struct TurnStartedEvent { | |
| pub turn_id: String, | |
| /// ID of the originating turn in the root thread; equals `turn_id` for root turns. | |
| pub root_turn_id: Option<String>, | |
| // Persist for rollout consumers that correlate turns with telemetry traces. | |
| pub trace_id: Option<String>, | |
| /// Unix timestamp (in seconds) when the turn started. | |
| pub started_at: Option<i64>, | |
| // TODO(aibrahim): make this not optional | |
| pub model_context_window: Option<i64>, | |
| pub collaboration_mode_kind: ModeKind, | |
| } | |
| pub struct ThreadSettingsAppliedEvent { | |
| /// Logical task that owns this snapshot, independent of the physical rollout file. | |
| /// Absent in older histories; copied snapshots retain their original owner's ID. | |
| pub thread_id: Option<ThreadId>, | |
| pub thread_settings: ThreadSettingsSnapshot, | |
| } | |
| pub struct ThreadSettingsSnapshot { | |
| pub model: String, | |
| pub model_provider_id: String, | |
| pub service_tier: Option<String>, | |
| pub approval_policy: AskForApproval, | |
| pub approvals_reviewer: ApprovalsReviewer, | |
| pub permission_profile: PermissionProfile, | |
| pub active_permission_profile: Option<ActivePermissionProfile>, | |
| pub cwd: AbsolutePathBuf, | |
| /// Top-level runtime workspace roots for default environments, excluding roots | |
| /// supplied by explicit environment selections or permission profiles. | |
| /// An absent value means unknown; an empty list means no roots. | |
| pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>, | |
| pub reasoning_effort: Option<ReasoningEffortConfig>, | |
| pub reasoning_summary: Option<ReasoningSummaryConfig>, | |
| pub personality: Option<Personality>, | |
| pub collaboration_mode: CollaborationMode, | |
| /// Thread-owned plugin selection, retained even when a plugin is unavailable. | |
| pub disabled_plugin_ids: Vec<String>, | |
| } | |
| pub struct TokenUsage { | |
| pub input_tokens: i64, | |
| pub cached_input_tokens: i64, | |
| pub cache_write_input_tokens: i64, | |
| pub output_tokens: i64, | |
| pub reasoning_output_tokens: i64, | |
| pub total_tokens: i64, | |
| /// Provider-reported units consumed from the shared rollout budget. | |
| pub codex_rollout_budget_units: Option<serde_json::Number>, | |
| } | |
| /// Best-effort Responses API usage observed for one completed response. | |
| pub struct TokenUsageRecord { | |
| pub thread_id: ThreadId, | |
| pub turn_id: String, | |
| pub session_id: SessionId, | |
| pub root_turn_id: String, | |
| pub response_id: String, | |
| pub usage: TokenUsage, | |
| pub turn_token_usage: TokenUsage, | |
| pub thread_token_usage: TokenUsage, | |
| } | |
| pub struct TokenUsageInfo { | |
| pub total_token_usage: TokenUsage, | |
| pub last_token_usage: TokenUsage, | |
| // TODO(aibrahim): make this not optional | |
| pub model_context_window: Option<i64>, | |
| } | |
| impl TokenUsageInfo { | |
| pub fn new_or_append( | |
| info: &Option<TokenUsageInfo>, | |
| last: &Option<TokenUsage>, | |
| model_context_window: Option<i64>, | |
| ) -> Option<Self> { | |
| if info.is_none() && last.is_none() { | |
| return None; | |
| } | |
| let mut info = match info { | |
| Some(info) => info.clone(), | |
| None => Self { | |
| total_token_usage: TokenUsage::default(), | |
| last_token_usage: TokenUsage::default(), | |
| model_context_window, | |
| }, | |
| }; | |
| if let Some(last) = last { | |
| info.append_last_usage(last); | |
| } | |
| if let Some(model_context_window) = model_context_window { | |
| info.model_context_window = Some(model_context_window); | |
| } | |
| Some(info) | |
| } | |
| pub fn append_last_usage(&mut self, last: &TokenUsage) { | |
| self.total_token_usage.add_assign(last); | |
| self.last_token_usage = last.clone(); | |
| } | |
| pub fn fill_to_context_window(&mut self, context_window: i64) { | |
| let previous_total = self.total_token_usage.total_tokens; | |
| let delta = (context_window - previous_total).max(0); | |
| self.model_context_window = Some(context_window); | |
| self.total_token_usage = TokenUsage { | |
| total_tokens: context_window, | |
| ..TokenUsage::default() | |
| }; | |
| self.last_token_usage = TokenUsage { | |
| total_tokens: delta, | |
| ..TokenUsage::default() | |
| }; | |
| } | |
| pub fn full_context_window(context_window: i64) -> Self { | |
| let mut info = Self { | |
| total_token_usage: TokenUsage::default(), | |
| last_token_usage: TokenUsage::default(), | |
| model_context_window: Some(context_window), | |
| }; | |
| info.fill_to_context_window(context_window); | |
| info | |
| } | |
| } | |
| pub struct TokenCountEvent { | |
| pub info: Option<TokenUsageInfo>, | |
| pub rate_limits: Option<RateLimitSnapshot>, | |
| } | |
| pub struct RateLimitSnapshot { | |
| pub limit_id: Option<String>, | |
| pub limit_name: Option<String>, | |
| /// Normal model metadata for a quota alias; never a replacement for the request model. | |
| pub normal_model_slug: Option<String>, | |
| pub primary: Option<RateLimitWindow>, | |
| pub secondary: Option<RateLimitWindow>, | |
| pub credits: Option<CreditsSnapshot>, | |
| pub individual_limit: Option<SpendControlLimitSnapshot>, | |
| /// Backend-reported spend-control state. `None` is unavailable, not a sparse-update recovery. | |
| pub spend_control_reached: Option<bool>, | |
| pub plan_type: Option<crate::account::PlanType>, | |
| pub rate_limit_reached_type: Option<RateLimitReachedType>, | |
| } | |
| pub enum RateLimitReachedType { | |
| RateLimitReached, | |
| WorkspaceOwnerCreditsDepleted, | |
| WorkspaceMemberCreditsDepleted, | |
| WorkspaceOwnerUsageLimitReached, | |
| WorkspaceMemberUsageLimitReached, | |
| } | |
| impl FromStr for RateLimitReachedType { | |
| type Err = String; | |
| fn from_str(value: &str) -> Result<Self, Self::Err> { | |
| match value { | |
| "rate_limit_reached" => Ok(Self::RateLimitReached), | |
| "workspace_owner_credits_depleted" => Ok(Self::WorkspaceOwnerCreditsDepleted), | |
| "workspace_member_credits_depleted" => Ok(Self::WorkspaceMemberCreditsDepleted), | |
| "workspace_owner_usage_limit_reached" => Ok(Self::WorkspaceOwnerUsageLimitReached), | |
| "workspace_member_usage_limit_reached" => Ok(Self::WorkspaceMemberUsageLimitReached), | |
| other => Err(format!("unknown rate limit reached type: {other}")), | |
| } | |
| } | |
| } | |
| pub struct RateLimitWindow { | |
| /// Percentage (0-100) of the window that has been consumed. | |
| pub used_percent: f64, | |
| /// Rolling window duration, in minutes. | |
| pub window_minutes: Option<i64>, | |
| /// Unix timestamp (seconds since epoch) when the window resets. | |
| pub resets_at: Option<i64>, | |
| } | |
| pub struct CreditsSnapshot { | |
| pub has_credits: bool, | |
| pub unlimited: bool, | |
| pub balance: Option<String>, | |
| } | |
| pub struct SpendControlLimitSnapshot { | |
| pub limit: String, | |
| pub used: String, | |
| pub remaining_percent: i32, | |
| pub resets_at: i64, | |
| } | |
| // Includes prompts, tools and space to call compact. | |
| const BASELINE_TOKENS: i64 = 12000; | |
| impl TokenUsage { | |
| pub fn is_zero(&self) -> bool { | |
| self.total_tokens == 0 | |
| } | |
| pub fn cached_input(&self) -> i64 { | |
| self.cached_input_tokens.max(0) | |
| } | |
| pub fn non_cached_input(&self) -> i64 { | |
| (self.input_tokens - self.cached_input()).max(0) | |
| } | |
| /// Primary count for display as a single absolute value: non-cached input + output. | |
| pub fn blended_total(&self) -> i64 { | |
| (self.non_cached_input() + self.output_tokens.max(0)).max(0) | |
| } | |
| pub fn tokens_in_context_window(&self) -> i64 { | |
| self.total_tokens | |
| } | |
| /// Estimate the remaining user-controllable percentage of the model's context window. | |
| /// | |
| /// `context_window` is the total size of the model's context window. | |
| /// `BASELINE_TOKENS` should capture tokens that are always present in | |
| /// the context (e.g., system prompt and fixed tool instructions) so that | |
| /// the percentage reflects the portion the user can influence. | |
| /// | |
| /// This normalizes both the numerator and denominator by subtracting the | |
| /// baseline, so immediately after the first prompt the UI shows 100% left | |
| /// and trends toward 0% as the user fills the effective window. | |
| pub fn percent_of_context_window_remaining(&self, context_window: i64) -> i64 { | |
| if context_window <= BASELINE_TOKENS { | |
| return 0; | |
| } | |
| let effective_window = context_window - BASELINE_TOKENS; | |
| let used = (self.tokens_in_context_window() - BASELINE_TOKENS).max(0); | |
| let remaining = (effective_window - used).max(0); | |
| ((remaining as f64 / effective_window as f64) * 100.0) | |
| .clamp(0.0, 100.0) | |
| .round() as i64 | |
| } | |
| /// In-place element-wise sum of token counts. | |
| pub fn add_assign(&mut self, other: &TokenUsage) { | |
| self.input_tokens += other.input_tokens; | |
| self.cached_input_tokens += other.cached_input_tokens; | |
| self.cache_write_input_tokens += other.cache_write_input_tokens; | |
| self.output_tokens += other.output_tokens; | |
| self.reasoning_output_tokens += other.reasoning_output_tokens; | |
| self.total_tokens += other.total_tokens; | |
| } | |
| } | |
| pub struct FinalOutput { | |
| pub token_usage: TokenUsage, | |
| } | |
| impl From<TokenUsage> for FinalOutput { | |
| fn from(token_usage: TokenUsage) -> Self { | |
| Self { token_usage } | |
| } | |
| } | |
| impl fmt::Display for FinalOutput { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| let token_usage = &self.token_usage; | |
| write!( | |
| f, | |
| "Token usage: total={} input={}{} output={}{}", | |
| format_with_separators(token_usage.blended_total()), | |
| format_with_separators(token_usage.non_cached_input()), | |
| if token_usage.cached_input() > 0 { | |
| format!( | |
| " (+ {} cached)", | |
| format_with_separators(token_usage.cached_input()) | |
| ) | |
| } else { | |
| String::new() | |
| }, | |
| format_with_separators(token_usage.output_tokens), | |
| if token_usage.reasoning_output_tokens > 0 { | |
| format!( | |
| " (reasoning {})", | |
| format_with_separators(token_usage.reasoning_output_tokens) | |
| ) | |
| } else { | |
| String::new() | |
| } | |
| ) | |
| } | |
| } | |
| pub struct AgentMessageEvent { | |
| pub message: String, | |
| pub phase: Option<MessagePhase>, | |
| pub memory_citation: Option<MemoryCitation>, | |
| pub delivery: Option<AgentMessageDelivery>, | |
| pub questions: Option<Vec<AsyncUserInputQuestion>>, | |
| } | |
| pub enum UserMessageImageKind { | |
| Inline, | |
| File, | |
| } | |
| pub struct UserMessageEvent { | |
| pub client_id: Option<String>, | |
| pub message: String, | |
| /// Image URLs sourced from `UserInput::Image`. These are safe | |
| /// to replay in legacy UI history events and correspond to images sent to | |
| /// the model. | |
| pub images: Option<Vec<String>>, | |
| /// Detail hints for `images`, indexed in parallel. Missing entries imply | |
| /// default image detail behavior. | |
| pub image_details: Vec<Option<ImageDetail>>, | |
| /// File IDs sourced from `UserInput::Image`. These are passed through as | |
| /// opaque references and are not created by image preparation. | |
| pub file_ids: Option<Vec<String>>, | |
| /// Detail hints for `file_ids`, indexed in parallel. Missing entries imply | |
| /// default image detail behavior. | |
| pub file_id_details: Vec<Option<ImageDetail>>, | |
| /// Inline and file-backed image kinds in their original input order. | |
| /// New producers populate this alongside `images` and `file_ids`; when it | |
| /// is absent, consumers retain the legacy inline-then-file ordering. | |
| pub image_order: Vec<UserMessageImageKind>, | |
| /// Local file paths sourced from `UserInput::LocalImage`. These are kept so | |
| /// the UI can reattach images when editing history. Local image prompts may | |
| /// include a display form of the path, but these should not be treated as | |
| /// API-ready URLs. | |
| pub local_images: Vec<std::path::PathBuf>, | |
| /// Detail hints for `local_images`, indexed in parallel. Missing entries | |
| /// imply default image detail behavior. | |
| pub local_image_details: Vec<Option<ImageDetail>>, | |
| /// Audio URLs sourced from `UserInput::Audio`. These are safe to replay in | |
| /// legacy UI history events and correspond to audio sent to the model. | |
| pub audio: Option<Vec<String>>, | |
| /// Local file paths sourced from `UserInput::LocalAudio`. These are kept so | |
| /// clients can reattach audio when editing history and should not be | |
| /// treated as API-ready URLs. | |
| pub local_audio: Vec<std::path::PathBuf>, | |
| /// UI-defined spans within `message` used to render or persist special elements. | |
| pub text_elements: Vec<crate::user_input::TextElement>, | |
| } | |
| impl UserMessageEvent { | |
| /// Returns whether `image_order` accounts for every split image reference exactly once. | |
| pub fn has_complete_image_order(&self) -> bool { | |
| if self.image_order.is_empty() { | |
| return false; | |
| } | |
| let mut inline_count = 0; | |
| let mut file_count = 0; | |
| for image_kind in &self.image_order { | |
| match image_kind { | |
| UserMessageImageKind::Inline => inline_count += 1, | |
| UserMessageImageKind::File => file_count += 1, | |
| } | |
| } | |
| inline_count == self.images.as_ref().map_or(0, Vec::len) | |
| && file_count == self.file_ids.as_ref().map_or(0, Vec::len) | |
| } | |
| } | |
| /// Returns the user-facing preview text for a user message. | |
| pub fn user_message_preview(user: &UserMessageEvent) -> Option<String> { | |
| let message = strip_user_message_prefix(user.message.as_str()); | |
| if !message.is_empty() { | |
| return Some(message.to_string()); | |
| } | |
| if user | |
| .images | |
| .as_ref() | |
| .is_some_and(|images| !images.is_empty()) | |
| || user | |
| .file_ids | |
| .as_ref() | |
| .is_some_and(|file_ids| !file_ids.is_empty()) | |
| || !user.local_images.is_empty() | |
| { | |
| return Some("[Image]".to_string()); | |
| } | |
| if user.audio.as_ref().is_some_and(|audio| !audio.is_empty()) || !user.local_audio.is_empty() { | |
| return Some("[Audio]".to_string()); | |
| } | |
| None | |
| } | |
| pub struct AgentReasoningEvent { | |
| pub text: String, | |
| } | |
| pub struct AgentReasoningRawContentEvent { | |
| pub text: String, | |
| } | |
| pub struct AgentReasoningSectionBreakEvent { | |
| // load with default value so it's backward compatible with the old format. | |
| pub item_id: String, | |
| pub summary_index: i64, | |
| } | |
| pub struct McpInvocation { | |
| /// Name of the MCP server as defined in the config. | |
| pub server: String, | |
| /// Name of the tool as given by the MCP server. | |
| pub tool: String, | |
| /// Arguments to the tool call. | |
| pub arguments: Option<serde_json::Value>, | |
| } | |
| pub struct McpToolCallBeginEvent { | |
| /// Identifier so this can be paired with the McpToolCallEnd event. | |
| pub call_id: String, | |
| /// Originating turn; absent in older rollout records. | |
| pub turn_id: String, | |
| pub invocation: McpInvocation, | |
| pub connector_id: Option<String>, | |
| pub mcp_app_resource_uri: Option<String>, | |
| pub mcp_app_ui: Option<crate::items::McpAppUi>, | |
| pub link_id: Option<String>, | |
| pub app_name: Option<String>, | |
| pub action_name: Option<String>, | |
| pub plugin_id: Option<String>, | |
| /// Whether the selected tool is annotated as read-only, not its execution outcome. | |
| pub read_only_hint: Option<bool>, | |
| } | |
| pub struct McpToolCallEndEvent { | |
| /// Identifier for the corresponding McpToolCallBegin that finished. | |
| pub call_id: String, | |
| /// Originating turn; absent in older rollout records. | |
| pub turn_id: String, | |
| pub invocation: McpInvocation, | |
| pub connector_id: Option<String>, | |
| pub mcp_app_resource_uri: Option<String>, | |
| pub mcp_app_ui: Option<crate::items::McpAppUi>, | |
| pub link_id: Option<String>, | |
| pub app_name: Option<String>, | |
| pub action_name: Option<String>, | |
| pub plugin_id: Option<String>, | |
| pub read_only_hint: Option<bool>, | |
| pub duration: Duration, | |
| /// Result of the tool call. Note this could be an error. | |
| pub result: Result<CallToolResult, String>, | |
| } | |
| pub struct DynamicToolCallResponseEvent { | |
| /// Identifier for the corresponding DynamicToolCallRequest. | |
| pub call_id: String, | |
| /// Turn ID that this dynamic tool call belongs to. | |
| pub turn_id: String, | |
| pub completed_at_ms: i64, | |
| /// Dynamic tool namespace, when one was provided. | |
| pub namespace: Option<String>, | |
| /// Dynamic tool name. | |
| pub tool: String, | |
| /// Dynamic tool call arguments. | |
| pub arguments: serde_json::Value, | |
| /// Dynamic tool response content items. | |
| pub content_items: Vec<DynamicToolCallOutputContentItem>, | |
| /// Whether the tool call succeeded. | |
| pub success: bool, | |
| /// Optional error text when the tool call failed before producing a response. | |
| pub error: Option<String>, | |
| /// The duration of the dynamic tool call. | |
| pub duration: Duration, | |
| } | |
| impl McpToolCallEndEvent { | |
| pub fn is_success(&self) -> bool { | |
| match &self.result { | |
| Ok(result) => !result.is_error.unwrap_or(false), | |
| Err(_) => false, | |
| } | |
| } | |
| } | |
| pub struct WebSearchBeginEvent { | |
| pub call_id: String, | |
| } | |
| pub struct WebSearchEndEvent { | |
| pub call_id: String, | |
| pub query: String, | |
| pub action: WebSearchAction, | |
| /// Structured search results returned out-of-band by standalone web search. | |
| pub results: Option<Vec<Value>>, | |
| } | |
| pub struct ImageGenerationBeginEvent { | |
| pub call_id: String, | |
| } | |
| pub struct ImageGenerationEndEvent { | |
| pub call_id: String, | |
| pub status: String, | |
| pub revised_prompt: Option<String>, | |
| pub result: String, | |
| pub transparent_background: Option<bool>, | |
| pub failure: Option<ImageGenerationFailure>, | |
| pub saved_path: Option<AbsolutePathBuf>, | |
| } | |
| // Conversation kept for backward compatibility. | |
| /// Response payload for `Op::GetHistory` containing the current session's | |
| /// in-memory transcript. | |
| pub struct ConversationPathResponseEvent { | |
| pub conversation_id: ThreadId, | |
| pub path: PathBuf, | |
| } | |
| pub enum SessionSource { | |
| Cli, | |
| VSCode, | |
| Exec, | |
| Mcp, | |
| Custom(String), | |
| Internal(InternalSessionSource), | |
| SubAgent(SubAgentSource), | |
| Unknown, | |
| } | |
| pub enum ThreadSource { | |
| User, | |
| Subagent, | |
| GuardianReview, | |
| Feature(String), | |
| MemoryConsolidation, | |
| } | |
| impl ThreadSource { | |
| pub fn as_str(&self) -> &str { | |
| match self { | |
| ThreadSource::User => "user", | |
| ThreadSource::Subagent => "subagent", | |
| ThreadSource::GuardianReview => "guardian_review", | |
| ThreadSource::Feature(feature) => feature, | |
| ThreadSource::MemoryConsolidation => "memory_consolidation", | |
| } | |
| } | |
| } | |
| impl fmt::Display for ThreadSource { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| f.write_str(self.as_str()) | |
| } | |
| } | |
| impl TryFrom<String> for ThreadSource { | |
| type Error = String; | |
| fn try_from(value: String) -> Result<Self, Self::Error> { | |
| value.parse() | |
| } | |
| } | |
| impl From<ThreadSource> for String { | |
| fn from(value: ThreadSource) -> Self { | |
| value.to_string() | |
| } | |
| } | |
| impl FromStr for ThreadSource { | |
| type Err = String; | |
| fn from_str(value: &str) -> Result<Self, Self::Err> { | |
| match value { | |
| "user" => Ok(ThreadSource::User), | |
| "subagent" => Ok(ThreadSource::Subagent), | |
| "guardian_review" => Ok(ThreadSource::GuardianReview), | |
| "memory_consolidation" => Ok(ThreadSource::MemoryConsolidation), | |
| other => Ok(ThreadSource::Feature(other.to_string())), | |
| } | |
| } | |
| } | |
| pub enum InternalSessionSource { | |
| MemoryConsolidation, | |
| Guardian, | |
| } | |
| pub enum SubAgentSource { | |
| Review, | |
| Compact, | |
| ThreadSpawn { | |
| parent_thread_id: ThreadId, | |
| depth: i32, | |
| agent_path: Option<AgentPath>, | |
| agent_nickname: Option<String>, | |
| agent_role: Option<String>, | |
| }, | |
| MemoryConsolidation, | |
| Other(String), | |
| } | |
| impl fmt::Display for SessionSource { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| match self { | |
| SessionSource::Cli => f.write_str("cli"), | |
| SessionSource::VSCode => f.write_str("vscode"), | |
| SessionSource::Exec => f.write_str("exec"), | |
| SessionSource::Mcp => f.write_str("mcp"), | |
| SessionSource::Custom(source) => f.write_str(source), | |
| SessionSource::Internal(source) => write!(f, "internal_{source}"), | |
| SessionSource::SubAgent(sub_source) => write!(f, "subagent_{sub_source}"), | |
| SessionSource::Unknown => f.write_str("unknown"), | |
| } | |
| } | |
| } | |
| impl SessionSource { | |
| pub fn from_startup_arg(value: &str) -> Result<Self, &'static str> { | |
| let trimmed = value.trim(); | |
| if trimmed.is_empty() { | |
| return Err("session source must not be empty"); | |
| } | |
| let normalized = trimmed.to_ascii_lowercase(); | |
| Ok(match normalized.as_str() { | |
| "cli" => SessionSource::Cli, | |
| "vscode" => SessionSource::VSCode, | |
| "exec" => SessionSource::Exec, | |
| "mcp" | "appserver" | "app-server" | "app_server" => SessionSource::Mcp, | |
| "unknown" => SessionSource::Unknown, | |
| _ => SessionSource::Custom(normalized), | |
| }) | |
| } | |
| pub fn is_internal(&self) -> bool { | |
| matches!(self, SessionSource::Internal(_)) | |
| } | |
| pub fn is_non_root_agent(&self) -> bool { | |
| matches!( | |
| self, | |
| SessionSource::Internal(_) | SessionSource::SubAgent(_) | |
| ) | |
| } | |
| pub fn get_nickname(&self) -> Option<String> { | |
| match self { | |
| SessionSource::SubAgent(SubAgentSource::ThreadSpawn { agent_nickname, .. }) => { | |
| agent_nickname.clone() | |
| } | |
| _ => None, | |
| } | |
| } | |
| pub fn get_agent_role(&self) -> Option<String> { | |
| match self { | |
| SessionSource::SubAgent(SubAgentSource::ThreadSpawn { agent_role, .. }) => { | |
| agent_role.clone() | |
| } | |
| _ => None, | |
| } | |
| } | |
| pub fn get_agent_path(&self) -> Option<AgentPath> { | |
| match self { | |
| SessionSource::SubAgent(SubAgentSource::ThreadSpawn { agent_path, .. }) => { | |
| agent_path.clone() | |
| } | |
| _ => None, | |
| } | |
| } | |
| pub fn restriction_product(&self) -> Option<Product> { | |
| match self { | |
| SessionSource::Custom(source) => Product::from_session_source_name(source), | |
| SessionSource::Cli | |
| | SessionSource::VSCode | |
| | SessionSource::Exec | |
| | SessionSource::Mcp | |
| | SessionSource::Unknown => Some(Product::Codex), | |
| SessionSource::Internal(_) | SessionSource::SubAgent(_) => None, | |
| } | |
| } | |
| pub fn matches_product_restriction(&self, products: &[Product]) -> bool { | |
| products.is_empty() | |
| || self | |
| .restriction_product() | |
| .is_some_and(|product| product.matches_product_restriction(products)) | |
| } | |
| pub fn parent_thread_id(&self) -> Option<ThreadId> { | |
| match self { | |
| SessionSource::SubAgent(subagent_source) => subagent_source.parent_thread_id(), | |
| SessionSource::Cli | |
| | SessionSource::VSCode | |
| | SessionSource::Exec | |
| | SessionSource::Mcp | |
| | SessionSource::Custom(_) | |
| | SessionSource::Internal(_) | |
| | SessionSource::Unknown => None, | |
| } | |
| } | |
| } | |
| impl fmt::Display for SubAgentSource { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| match self { | |
| SubAgentSource::Review => f.write_str("review"), | |
| SubAgentSource::Compact => f.write_str("compact"), | |
| SubAgentSource::MemoryConsolidation => f.write_str("memory_consolidation"), | |
| SubAgentSource::ThreadSpawn { | |
| parent_thread_id, | |
| depth, | |
| .. | |
| } => { | |
| write!(f, "thread_spawn_{parent_thread_id}_d{depth}") | |
| } | |
| SubAgentSource::Other(other) => f.write_str(other), | |
| } | |
| } | |
| } | |
| impl SubAgentSource { | |
| pub fn kind(&self) -> &str { | |
| match self { | |
| SubAgentSource::Review => "review", | |
| SubAgentSource::Compact => "compact", | |
| SubAgentSource::ThreadSpawn { .. } => "thread_spawn", | |
| SubAgentSource::MemoryConsolidation => "memory_consolidation", | |
| SubAgentSource::Other(other) => other, | |
| } | |
| } | |
| pub fn parent_thread_id(&self) -> Option<ThreadId> { | |
| match self { | |
| SubAgentSource::ThreadSpawn { | |
| parent_thread_id, .. | |
| } => Some(*parent_thread_id), | |
| SubAgentSource::Review | |
| | SubAgentSource::Compact | |
| | SubAgentSource::MemoryConsolidation | |
| | SubAgentSource::Other(_) => None, | |
| } | |
| } | |
| } | |
| impl fmt::Display for InternalSessionSource { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| match self { | |
| InternalSessionSource::MemoryConsolidation => f.write_str("memory_consolidation"), | |
| InternalSessionSource::Guardian => f.write_str("guardian"), | |
| } | |
| } | |
| } | |
| pub enum MultiAgentVersion { | |
| Disabled, | |
| V1, | |
| V2, | |
| } | |
| pub struct SessionContextWindow { | |
| /// UUIDv7 identity of this context window. | |
| pub window_id: String, | |
| } | |
| impl SessionContextWindow { | |
| pub fn new(window_id: String) -> Self { | |
| Self { window_id } | |
| } | |
| } | |
| /// Exclusive position in another rollout's paginated history. | |
| pub struct HistoryPosition { | |
| /// Rollout ID for the immutable prefix file. | |
| /// | |
| /// `HistoryPosition` predates `thread/revert`, so this field is named `thread_id`. Treat its | |
| /// value as a `rollout_id`: ordinary rollouts use the thread ID as their rollout ID, while a | |
| /// reverted thread's filename carries a distinct rollout ID. It is not necessarily | |
| /// [`SessionMeta::id`], which remains the stable thread ID across revert. | |
| pub thread_id: ThreadId, | |
| /// First rollout ordinal not included from the prefix file. | |
| pub end_ordinal_exclusive: u64, | |
| /// Byte offset immediately after the last included JSONL record from the prefix file. | |
| pub end_byte_offset: u64, | |
| } | |
| /// SessionMeta contains session-level data that doesn't correspond to a specific turn. | |
| /// | |
| /// NOTE: There used to be an `instructions` field here, which stored user_instructions, but we | |
| /// now save that on TurnContext. base_instructions stores the base instructions for the session, | |
| /// and should be used when there is no config override. | |
| pub struct SessionMeta { | |
| /// session_id is equal to the root thread's ID. | |
| pub session_id: SessionId, | |
| pub id: ThreadId, | |
| pub forked_from_id: Option<ThreadId>, | |
| /// Exclusive ordinal inherited from the logical fork parent, independent of `history_base`. | |
| /// Revert may replace the physical history base while retaining this fork boundary. | |
| pub forked_from_ordinal_exclusive: Option<u64>, | |
| pub parent_thread_id: Option<ThreadId>, | |
| pub timestamp: String, | |
| pub cwd: PathBuf, | |
| /// Top-level runtime workspace roots at creation for default environments, | |
| /// excluding roots supplied by explicit environment selections or permission profiles. | |
| /// An absent value means unknown; an empty list means no roots. | |
| /// Keep native paths parseable across hosts; validate them when restoring settings. | |
| pub runtime_workspace_roots: Option<Vec<PathBuf>>, | |
| pub originator: String, | |
| pub cli_version: String, | |
| pub source: SessionSource, | |
| /// Optional analytics source classification for this thread. | |
| pub thread_source: Option<ThreadSource>, | |
| /// Optional random unique nickname assigned to an AgentControl-spawned sub-agent. | |
| pub agent_nickname: Option<String>, | |
| /// Optional role (agent_role) assigned to an AgentControl-spawned sub-agent. | |
| pub agent_role: Option<String>, | |
| /// Optional canonical agent path assigned to an AgentControl-spawned sub-agent. | |
| pub agent_path: Option<String>, | |
| pub model_provider: Option<String>, | |
| /// base_instructions for the session. This *should* always be present when creating a new session, | |
| /// but may be missing for older sessions. If not present, fall back to rendering the base_instructions | |
| /// from ModelsManager. | |
| pub base_instructions: Option<BaseInstructions>, | |
| pub dynamic_tools: Option<Vec<DynamicToolSpec>>, | |
| /// Capability roots selected for this thread by the hosting platform. | |
| pub selected_capability_roots: Vec<SelectedCapabilityRoot>, | |
| pub memory_mode: Option<String>, | |
| pub history_mode: ThreadHistoryMode, | |
| /// Exclusive prefix of another paginated rollout inherited by this thread. | |
| pub history_base: Option<HistoryPosition>, | |
| /// First rollout ordinal that belongs to this subagent's own projected history. | |
| /// | |
| /// Earlier rollout records are inherited model context and stay out of child | |
| /// turn/item projection. | |
| pub subagent_history_start_ordinal: Option<u64>, | |
| pub multi_agent_version: Option<MultiAgentVersion>, | |
| /// Initial context-window identity for consumers that tail rollout JSONL before compaction. | |
| pub context_window: Option<SessionContextWindow>, | |
| } | |
| impl Default for SessionMeta { | |
| fn default() -> Self { | |
| let id = ThreadId::default(); | |
| SessionMeta { | |
| session_id: id.into(), | |
| id, | |
| forked_from_id: None, | |
| forked_from_ordinal_exclusive: None, | |
| parent_thread_id: None, | |
| timestamp: String::new(), | |
| cwd: PathBuf::new(), | |
| runtime_workspace_roots: None, | |
| originator: String::new(), | |
| cli_version: String::new(), | |
| source: SessionSource::default(), | |
| thread_source: None, | |
| agent_nickname: None, | |
| agent_role: None, | |
| agent_path: None, | |
| model_provider: None, | |
| base_instructions: None, | |
| dynamic_tools: None, | |
| selected_capability_roots: Vec::new(), | |
| memory_mode: None, | |
| history_mode: ThreadHistoryMode::default(), | |
| history_base: None, | |
| subagent_history_start_ordinal: None, | |
| multi_agent_version: None, | |
| context_window: None, | |
| } | |
| } | |
| } | |
| pub struct SessionMetaLine { | |
| pub meta: SessionMeta, | |
| pub git: Option<GitInfo>, | |
| } | |
| impl<'de> Deserialize<'de> for SessionMetaLine { | |
| fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> | |
| where | |
| D: Deserializer<'de>, | |
| { | |
| struct SessionMetaLineFields { | |
| meta: SessionMeta, | |
| git: Option<GitInfo>, | |
| } | |
| let mut value = Value::deserialize(deserializer)?; | |
| let fields = value | |
| .as_object_mut() | |
| .ok_or_else(|| D::Error::custom("session metadata must be an object"))?; | |
| if !fields.contains_key("session_id") { | |
| let thread_id = fields | |
| .get("id") | |
| .cloned() | |
| .ok_or_else(|| D::Error::missing_field("id"))?; | |
| fields.insert("session_id".to_string(), thread_id); | |
| } | |
| let SessionMetaLineFields { meta, git } = | |
| serde_json::from_value(value).map_err(D::Error::custom)?; | |
| Ok(Self { meta, git }) | |
| } | |
| } | |
| /// Persisted comparison state used to resume model-visible world-state diffing. | |
| pub struct WorldStateItem { | |
| /// Full snapshots establish a new baseline; patches update the current baseline. | |
| pub full: bool, | |
| pub state: Map<String, Value>, | |
| } | |
| impl WorldStateItem { | |
| pub fn full(state: Map<String, Value>) -> Self { | |
| Self { full: true, state } | |
| } | |
| pub fn patch(state: Map<String, Value>) -> Self { | |
| Self { full: false, state } | |
| } | |
| } | |
| pub struct TurnContextNetworkItem { | |
| pub allowed_domains: Vec<String>, | |
| pub denied_domains: Vec<String>, | |
| } | |
| /// Persist once per real user turn after computing that turn's model-visible | |
| /// context updates, and again after mid-turn compaction when replacement | |
| /// history re-establishes full context, so resume/fork replay can recover the | |
| /// latest durable baseline. | |
| pub struct TurnContextItem { | |
| pub turn_id: Option<String>, | |
| /// Root turn that owns this subagent turn's attribution. | |
| /// Only set for subagent turns; persisted so resume keeps the scope frozen at turn start. | |
| pub root_turn_id: Option<String>, | |
| /// Plugin selection captured for this turn. Absent in older histories. | |
| pub disabled_plugin_ids: Option<Vec<String>>, | |
| pub cwd: AbsolutePathBuf, | |
| /// Effective workspace roots used to materialize symbolic | |
| /// `:workspace_roots` filesystem permissions in `permission_profile`. | |
| pub workspace_roots: Option<Vec<AbsolutePathBuf>>, | |
| pub current_date: Option<String>, | |
| pub timezone: Option<String>, | |
| pub approval_policy: AskForApproval, | |
| pub approvals_reviewer: Option<ApprovalsReviewer>, | |
| pub sandbox_policy: SandboxPolicy, | |
| pub permission_profile: Option<PermissionProfile>, | |
| /// Built-in or named profile that produced `permission_profile`, when known. | |
| pub active_permission_profile: Option<ActivePermissionProfile>, | |
| pub network: Option<TurnContextNetworkItem>, | |
| pub file_system_sandbox_policy: Option<RawFileSystemSandboxPolicy>, | |
| pub model: String, | |
| pub comp_hash: Option<String>, | |
| pub personality: Option<Personality>, | |
| pub collaboration_mode: Option<CollaborationMode>, | |
| pub multi_agent_version: Option<MultiAgentVersion>, | |
| /// Legacy effective model-visible mode retained to deserialize older rollouts. | |
| pub multi_agent_mode: Option<MultiAgentMode>, | |
| pub realtime_active: Option<bool>, | |
| pub cyber_access_program: Option<CyberAccessProgram>, | |
| pub effort: Option<ReasoningEffortConfig>, | |
| // Compatibility-only field written with a default value so older Codex | |
| // versions can deserialize turn-context rollout items. It is no longer | |
| // read by context reconstruction and should be removed in a future schema | |
| // cleanup. | |
| pub summary: ReasoningSummaryConfig, | |
| } | |
| impl TurnContextItem { | |
| pub fn permission_profile(&self) -> PermissionProfile { | |
| self.permission_profile.clone().unwrap_or_else(|| { | |
| let file_system_sandbox_policy = self | |
| .file_system_sandbox_policy | |
| .clone() | |
| .map(TryInto::try_into) | |
| .transpose() | |
| .unwrap_or_else(|_| Some(FileSystemSandboxPolicy::restricted(Vec::new()))) | |
| .unwrap_or_else(|| { | |
| FileSystemSandboxPolicy::from_legacy_sandbox_policy_for_cwd( | |
| &self.sandbox_policy, | |
| self.cwd.as_path(), | |
| ) | |
| }); | |
| PermissionProfile::from_runtime_permissions_with_enforcement( | |
| SandboxEnforcement::from_legacy_sandbox_policy(&self.sandbox_policy), | |
| &file_system_sandbox_policy, | |
| NetworkSandboxPolicy::from(&self.sandbox_policy), | |
| ) | |
| }) | |
| } | |
| } | |
| pub enum TruncationPolicy { | |
| Bytes(usize), | |
| Tokens(usize), | |
| } | |
| impl From<crate::openai_models::TruncationPolicyConfig> for TruncationPolicy { | |
| fn from(config: crate::openai_models::TruncationPolicyConfig) -> Self { | |
| match config.mode { | |
| crate::openai_models::TruncationMode::Bytes => Self::Bytes(config.limit as usize), | |
| crate::openai_models::TruncationMode::Tokens => Self::Tokens(config.limit as usize), | |
| } | |
| } | |
| } | |
| impl TruncationPolicy { | |
| pub fn token_budget(&self) -> usize { | |
| match self { | |
| TruncationPolicy::Bytes(bytes) => { | |
| usize::try_from(codex_utils_string::approx_tokens_from_byte_count(*bytes)) | |
| .unwrap_or(usize::MAX) | |
| } | |
| TruncationPolicy::Tokens(tokens) => *tokens, | |
| } | |
| } | |
| pub fn byte_budget(&self) -> usize { | |
| match self { | |
| TruncationPolicy::Bytes(bytes) => *bytes, | |
| TruncationPolicy::Tokens(tokens) => { | |
| codex_utils_string::approx_bytes_for_tokens(*tokens) | |
| } | |
| } | |
| } | |
| } | |
| impl Mul<f64> for TruncationPolicy { | |
| type Output = Self; | |
| fn mul(self, multiplier: f64) -> Self::Output { | |
| match self { | |
| TruncationPolicy::Bytes(bytes) => { | |
| TruncationPolicy::Bytes((bytes as f64 * multiplier).ceil() as usize) | |
| } | |
| TruncationPolicy::Tokens(tokens) => { | |
| TruncationPolicy::Tokens((tokens as f64 * multiplier).ceil() as usize) | |
| } | |
| } | |
| } | |
| } | |
| pub struct GitInfo { | |
| /// Current commit hash (SHA) | |
| pub commit_hash: Option<GitSha>, | |
| /// Current branch name | |
| pub branch: Option<String>, | |
| /// Repository URL (if available from remote) | |
| pub repository_url: Option<SanitizedGitUrl>, | |
| } | |
| pub enum ReviewDelivery { | |
| Inline, | |
| Detached, | |
| } | |
| pub enum ReviewTarget { | |
| /// Review the working tree: staged, unstaged, and untracked files. | |
| UncommittedChanges, | |
| /// Review changes between the current branch and the given base branch. | |
| BaseBranch { branch: String }, | |
| /// Review the changes introduced by a specific commit. | |
| Commit { | |
| sha: String, | |
| /// Optional human-readable label (e.g., commit subject) for UIs. | |
| title: Option<String>, | |
| }, | |
| /// Arbitrary instructions provided by the user. | |
| Custom { instructions: String }, | |
| } | |
| /// Review request sent to the review session. | |
| pub struct ReviewRequest { | |
| pub target: ReviewTarget, | |
| pub user_facing_hint: Option<String>, | |
| } | |
| /// Structured review result produced by a child review session. | |
| pub struct ReviewOutputEvent { | |
| pub findings: Vec<ReviewFinding>, | |
| pub overall_correctness: String, | |
| pub overall_explanation: String, | |
| pub overall_confidence_score: f32, | |
| } | |
| impl Default for ReviewOutputEvent { | |
| fn default() -> Self { | |
| Self { | |
| findings: Vec::new(), | |
| overall_correctness: String::default(), | |
| overall_explanation: String::default(), | |
| overall_confidence_score: 0.0, | |
| } | |
| } | |
| } | |
| /// A single review finding describing an observed issue or recommendation. | |
| pub struct ReviewFinding { | |
| pub title: String, | |
| pub body: String, | |
| pub confidence_score: f32, | |
| pub priority: i32, | |
| pub code_location: ReviewCodeLocation, | |
| } | |
| /// Location of the code related to a review finding. | |
| pub struct ReviewCodeLocation { | |
| pub absolute_file_path: PathBuf, | |
| pub line_range: ReviewLineRange, | |
| } | |
| /// Inclusive line range in a file associated with the finding. | |
| pub struct ReviewLineRange { | |
| pub start: u32, | |
| pub end: u32, | |
| } | |
| pub enum ExecCommandSource { | |
| Agent, | |
| UserShell, | |
| UnifiedExecStartup, | |
| UnifiedExecInteraction, | |
| } | |
| pub enum ExecCommandStatus { | |
| Completed, | |
| Failed, | |
| Declined, | |
| } | |
| pub struct ExecCommandBeginEvent { | |
| /// Identifier so this can be paired with the ExecCommandEnd event. | |
| pub call_id: String, | |
| /// Trusted first-party plugin attributed to this command, when known. | |
| pub plugin_id: Option<String>, | |
| /// Safe plugin-relative path attributed to this command, when known. | |
| pub script_path: Option<String>, | |
| /// Identifier for the underlying PTY process (when available). | |
| pub process_id: Option<String>, | |
| /// Turn ID that this command belongs to. | |
| pub turn_id: String, | |
| pub started_at_ms: i64, | |
| /// The command to be executed. | |
| pub command: Vec<String>, | |
| /// The command's working directory if not the default cwd for the agent. | |
| pub cwd: PathUri, | |
| pub parsed_cmd: Vec<ParsedCommand>, | |
| /// Where the command originated. Defaults to Agent for backward compatibility. | |
| pub source: ExecCommandSource, | |
| /// Raw input sent to a unified exec session (if this is an interaction event). | |
| pub interaction_input: Option<String>, | |
| } | |
| pub struct ExecCommandEndEvent { | |
| /// Identifier for the ExecCommandBegin that finished. | |
| pub call_id: String, | |
| /// Trusted first-party plugin attributed to this command, when known. | |
| pub plugin_id: Option<String>, | |
| /// Safe plugin-relative path attributed to this command, when known. | |
| pub script_path: Option<String>, | |
| /// Identifier for the underlying PTY process (when available). | |
| pub process_id: Option<String>, | |
| /// Turn ID that this command belongs to. | |
| pub turn_id: String, | |
| pub completed_at_ms: i64, | |
| /// The command that was executed. | |
| pub command: Vec<String>, | |
| /// The command's working directory if not the default cwd for the agent. | |
| pub cwd: PathUri, | |
| pub parsed_cmd: Vec<ParsedCommand>, | |
| /// Where the command originated. Defaults to Agent for backward compatibility. | |
| pub source: ExecCommandSource, | |
| /// Raw input sent to a unified exec session (if this is an interaction event). | |
| pub interaction_input: Option<String>, | |
| /// Captured stdout | |
| pub stdout: String, | |
| /// Captured stderr | |
| pub stderr: String, | |
| /// Captured aggregated output | |
| pub aggregated_output: String, | |
| /// The command's exit code. | |
| pub exit_code: i32, | |
| /// The duration of the command execution. | |
| pub duration: Duration, | |
| /// Formatted output from the command, as seen by the model. | |
| pub formatted_output: String, | |
| /// Completion status for this command execution. | |
| pub status: ExecCommandStatus, | |
| } | |
| pub struct ViewImageToolCallEvent { | |
| /// Identifier for the originating tool call. | |
| pub call_id: String, | |
| /// Filesystem path resolved for the selected environment. | |
| /// | |
| /// This core event is not exposed directly in the app-server API. App-server | |
| /// converts the path to `LegacyAppPathString` when building its public item. | |
| pub path: PathUri, | |
| } | |
| pub enum ExecOutputStream { | |
| Stdout, | |
| Stderr, | |
| } | |
| pub struct ExecCommandOutputDeltaEvent { | |
| /// Identifier for the ExecCommandBegin that produced this chunk. | |
| pub call_id: String, | |
| /// Which stream produced this chunk. | |
| pub stream: ExecOutputStream, | |
| /// Raw bytes from the stream (may not be valid UTF-8). | |
| pub chunk: Vec<u8>, | |
| } | |
| pub struct TerminalInteractionEvent { | |
| /// Identifier for the ExecCommandBegin that produced this chunk. | |
| pub call_id: String, | |
| /// Process id associated with the running command. | |
| pub process_id: String, | |
| /// Stdin sent to the running session. | |
| pub stdin: String, | |
| } | |
| pub struct DeprecationNoticeEvent { | |
| /// Concise summary of what is deprecated. | |
| pub summary: String, | |
| /// Optional extra guidance, such as migration steps or rationale. | |
| pub details: Option<String>, | |
| } | |
| pub struct ThreadRolledBackEvent { | |
| /// Number of user turns that were removed from context. | |
| pub num_turns: u32, | |
| } | |
| pub struct StreamErrorEvent { | |
| pub message: String, | |
| pub codex_error_info: Option<CodexErrorInfo>, | |
| /// Optional details about the underlying stream failure (often the same | |
| /// human-readable message that is surfaced as the terminal error if retries | |
| /// are exhausted). | |
| pub additional_details: Option<String>, | |
| } | |
| pub struct StreamInfoEvent { | |
| pub message: String, | |
| } | |
| pub struct PatchApplyBeginEvent { | |
| /// Identifier so this can be paired with the PatchApplyEnd event. | |
| pub call_id: String, | |
| /// Turn ID that this patch belongs to. | |
| /// Uses `#[serde(default)]` for backwards compatibility. | |
| pub turn_id: String, | |
| /// If true, there was no ApplyPatchApprovalRequest for this patch. | |
| pub auto_approved: bool, | |
| /// The changes to be applied. | |
| pub changes: HashMap<PathBuf, FileChange>, | |
| } | |
| pub struct PatchApplyUpdatedEvent { | |
| /// Identifier for the originating `apply_patch` tool call. | |
| pub call_id: String, | |
| /// Structured file changes parsed from the model-generated patch input so far. | |
| pub changes: HashMap<PathBuf, FileChange>, | |
| } | |
| pub struct PatchApplyEndEvent { | |
| /// Identifier for the PatchApplyBegin that finished. | |
| pub call_id: String, | |
| /// Turn ID that this patch belongs to. | |
| /// Uses `#[serde(default)]` for backwards compatibility. | |
| pub turn_id: String, | |
| /// Captured stdout (summary printed by apply_patch). | |
| pub stdout: String, | |
| /// Captured stderr (parser errors, IO failures, etc.). | |
| pub stderr: String, | |
| /// Whether the patch was applied successfully. | |
| pub success: bool, | |
| /// The changes that were applied (mirrors PatchApplyBeginEvent::changes). | |
| pub changes: HashMap<PathBuf, FileChange>, | |
| /// Completion status for this patch application. | |
| pub status: PatchApplyStatus, | |
| } | |
| pub enum PatchApplyStatus { | |
| Completed, | |
| Failed, | |
| Declined, | |
| } | |
| pub struct TurnDiffEvent { | |
| pub unified_diff: String, | |
| } | |
| pub struct McpStartupUpdateEvent { | |
| /// Server name being started. | |
| pub server: String, | |
| /// Current startup status. | |
| pub status: McpStartupStatus, | |
| } | |
| pub enum McpStartupStatus { | |
| Starting, | |
| Ready, | |
| Failed { | |
| error: String, | |
| reason: Option<McpStartupFailureReason>, | |
| }, | |
| Cancelled, | |
| } | |
| pub enum McpStartupFailureReason { | |
| ReauthenticationRequired, | |
| } | |
| pub struct McpStartupCompleteEvent { | |
| pub ready: Vec<String>, | |
| pub failed: Vec<McpStartupFailure>, | |
| pub cancelled: Vec<String>, | |
| } | |
| pub struct McpStartupFailure { | |
| pub server: String, | |
| pub error: String, | |
| } | |
| pub enum McpAuthStatus { | |
| Unknown, | |
| Unsupported, | |
| NotLoggedIn, | |
| BearerToken, | |
| OAuth, | |
| } | |
| impl fmt::Display for McpAuthStatus { | |
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |
| let text = match self { | |
| McpAuthStatus::Unknown => "Unknown", | |
| McpAuthStatus::Unsupported => "Unsupported", | |
| McpAuthStatus::NotLoggedIn => "Not logged in", | |
| McpAuthStatus::BearerToken => "Bearer token", | |
| McpAuthStatus::OAuth => "OAuth", | |
| }; | |
| f.write_str(text) | |
| } | |
| } | |
| pub struct RealtimeConversationListVoicesResponseEvent { | |
| pub voices: RealtimeVoicesList, | |
| } | |
| pub enum Product { | |
| Chatgpt, | |
| Codex, | |
| Atlas, | |
| } | |
| impl Product { | |
| pub fn to_app_platform(self) -> &'static str { | |
| match self { | |
| Self::Chatgpt => "chat", | |
| Self::Codex => "codex", | |
| Self::Atlas => "atlas", | |
| } | |
| } | |
| pub fn from_session_source_name(value: &str) -> Option<Self> { | |
| let normalized = value.trim().to_ascii_lowercase(); | |
| match normalized.as_str() { | |
| "chatgpt" => Some(Self::Chatgpt), | |
| "codex" => Some(Self::Codex), | |
| "atlas" => Some(Self::Atlas), | |
| _ => None, | |
| } | |
| } | |
| pub fn matches_product_restriction(&self, products: &[Product]) -> bool { | |
| products.is_empty() || products.contains(self) | |
| } | |
| } | |
| pub enum SkillScope { | |
| User, | |
| Repo, | |
| System, | |
| Admin, | |
| } | |
| pub struct SkillMetadata { | |
| pub name: String, | |
| pub description: String, | |
| /// Legacy short_description from SKILL.md. Prefer SKILL.json interface.short_description. | |
| pub short_description: Option<String>, | |
| pub interface: Option<SkillInterface>, | |
| pub dependencies: Option<SkillDependencies>, | |
| pub path: AbsolutePathBuf, | |
| pub scope: SkillScope, | |
| pub enabled: bool, | |
| } | |
| pub struct SkillInterface { | |
| pub display_name: Option<String>, | |
| pub short_description: Option<String>, | |
| pub icon_small: Option<AbsolutePathBuf>, | |
| pub icon_large: Option<AbsolutePathBuf>, | |
| pub brand_color: Option<String>, | |
| pub default_prompt: Option<String>, | |
| } | |
| pub struct SkillDependencies { | |
| pub tools: Vec<SkillToolDependency>, | |
| } | |
| pub struct SkillToolDependency { | |
| pub r#type: String, | |
| pub value: String, | |
| pub description: Option<String>, | |
| pub transport: Option<String>, | |
| pub command: Option<String>, | |
| pub url: Option<String>, | |
| } | |
| pub struct SessionNetworkProxyRuntime { | |
| pub http_addr: String, | |
| pub socks_addr: String, | |
| } | |
| pub struct SessionConfiguredEvent { | |
| pub session_id: SessionId, | |
| pub thread_id: ThreadId, | |
| pub forked_from_id: Option<ThreadId>, | |
| pub parent_thread_id: Option<ThreadId>, | |
| /// Optional analytics source classification for this thread. | |
| pub thread_source: Option<ThreadSource>, | |
| /// Optional user-facing thread name (may be unset). | |
| pub thread_name: Option<String>, | |
| /// Tell the client what model is being queried. | |
| pub model: String, | |
| pub model_provider_id: String, | |
| pub service_tier: Option<String>, | |
| /// When to escalate for approval for execution | |
| pub approval_policy: AskForApproval, | |
| /// Configures who approval requests are routed to for review once they have | |
| /// been escalated. This does not disable separate safety checks such as | |
| /// ARC. | |
| pub approvals_reviewer: ApprovalsReviewer, | |
| /// Canonical effective permissions for commands executed in the session. | |
| pub permission_profile: PermissionProfile, | |
| /// Named or implicit built-in profile that produced `permission_profile`, | |
| /// when known. | |
| pub active_permission_profile: Option<ActivePermissionProfile>, | |
| /// Working directory that should be treated as the *root* of the | |
| /// session. | |
| pub cwd: AbsolutePathBuf, | |
| /// The effort the model is putting into reasoning about the user's request. | |
| pub reasoning_effort: Option<ReasoningEffortConfig>, | |
| /// Optional initial messages (as events) for resumed sessions. | |
| /// When present, UIs can use these to seed the history. | |
| pub initial_messages: Option<Vec<EventMsg>>, | |
| /// Runtime proxy bind addresses, when the managed proxy was started for this session. | |
| pub network_proxy: Option<SessionNetworkProxyRuntime>, | |
| /// Path in which the rollout is stored. Can be `None` for ephemeral threads | |
| pub rollout_path: Option<PathBuf>, | |
| } | |
| impl<'de> Deserialize<'de> for SessionConfiguredEvent { | |
| fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> | |
| where | |
| D: serde::Deserializer<'de>, | |
| { | |
| struct Wire { | |
| session_id: SessionId, | |
| thread_id: Option<ThreadId>, | |
| forked_from_id: Option<ThreadId>, | |
| parent_thread_id: Option<ThreadId>, | |
| thread_source: Option<ThreadSource>, | |
| thread_name: Option<String>, | |
| model: String, | |
| model_provider_id: String, | |
| service_tier: Option<String>, | |
| approval_policy: AskForApproval, | |
| approvals_reviewer: ApprovalsReviewer, | |
| // `SessionConfiguredEvent` is persisted into rollout history. Older | |
| // rollouts only have `sandbox_policy`, so accept it on deserialize | |
| // and immediately project it into the canonical `permission_profile`. | |
| sandbox_policy: Option<SandboxPolicy>, | |
| permission_profile: Option<PermissionProfile>, | |
| active_permission_profile: Option<ActivePermissionProfile>, | |
| cwd: AbsolutePathBuf, | |
| reasoning_effort: Option<ReasoningEffortConfig>, | |
| initial_messages: Option<Vec<EventMsg>>, | |
| network_proxy: Option<SessionNetworkProxyRuntime>, | |
| rollout_path: Option<PathBuf>, | |
| } | |
| let wire = Wire::deserialize(deserializer)?; | |
| let permission_profile = match (wire.permission_profile, wire.sandbox_policy) { | |
| (Some(permission_profile), _) => permission_profile, | |
| (None, Some(sandbox_policy)) => PermissionProfile::from_legacy_sandbox_policy_for_cwd( | |
| &sandbox_policy, | |
| wire.cwd.as_path(), | |
| ), | |
| (None, None) => { | |
| return Err(serde::de::Error::missing_field("permission_profile")); | |
| } | |
| }; | |
| Ok(Self { | |
| session_id: wire.session_id, | |
| thread_id: wire.thread_id.unwrap_or_else(|| wire.session_id.into()), | |
| forked_from_id: wire.forked_from_id, | |
| parent_thread_id: wire.parent_thread_id, | |
| thread_source: wire.thread_source, | |
| thread_name: wire.thread_name, | |
| model: wire.model, | |
| model_provider_id: wire.model_provider_id, | |
| service_tier: wire.service_tier, | |
| approval_policy: wire.approval_policy, | |
| approvals_reviewer: wire.approvals_reviewer, | |
| permission_profile, | |
| active_permission_profile: wire.active_permission_profile, | |
| cwd: wire.cwd, | |
| reasoning_effort: wire.reasoning_effort, | |
| initial_messages: wire.initial_messages, | |
| network_proxy: wire.network_proxy, | |
| rollout_path: wire.rollout_path, | |
| }) | |
| } | |
| } | |
| pub enum ThreadGoalStatus { | |
| Active, | |
| Paused, | |
| Blocked, | |
| UsageLimited, | |
| BudgetLimited, | |
| Complete, | |
| } | |
| pub const MAX_THREAD_GOAL_OBJECTIVE_CHARS: usize = 4_000; | |
| pub fn validate_thread_goal_objective(value: &str) -> Result<(), String> { | |
| if value.is_empty() { | |
| return Err("goal objective must not be empty".to_string()); | |
| } | |
| if value.chars().count() > MAX_THREAD_GOAL_OBJECTIVE_CHARS { | |
| return Err(format!( | |
| "goal objective must be at most {MAX_THREAD_GOAL_OBJECTIVE_CHARS} characters" | |
| )); | |
| } | |
| Ok(()) | |
| } | |
| pub struct ThreadGoal { | |
| pub thread_id: ThreadId, | |
| pub objective: String, | |
| pub status: ThreadGoalStatus, | |
| pub token_budget: Option<i64>, | |
| pub tokens_used: i64, | |
| pub time_used_seconds: i64, | |
| pub created_at: i64, | |
| pub updated_at: i64, | |
| } | |
| pub struct ThreadGoalUpdatedEvent { | |
| pub thread_id: ThreadId, | |
| pub turn_id: Option<String>, | |
| pub goal: ThreadGoal, | |
| } | |
| pub struct ThreadQueueChangedEvent { | |
| pub thread_id: ThreadId, | |
| } | |
| /// User's decision in response to an ExecApprovalRequest. | |
| pub enum ReviewDecision { | |
| /// User has approved this command and the agent should execute it. | |
| Approved, | |
| /// User has approved this command and wants to apply the proposed execpolicy | |
| /// amendment so future matching commands are permitted. | |
| ApprovedExecpolicyAmendment { | |
| proposed_execpolicy_amendment: ExecPolicyAmendment, | |
| }, | |
| /// User has approved this request and wants future prompts in the same | |
| /// session-scoped approval cache to be automatically approved for the | |
| /// remainder of the session. | |
| ApprovedForSession, | |
| /// User has approved this MCP tool call and wants to amend its policy so | |
| /// matching future calls are automatically approved across sessions. | |
| ApprovedMcpPolicyAmendment, | |
| /// User chose to persist a network policy rule (allow/deny) for future | |
| /// requests to the same host. | |
| NetworkPolicyAmendment { | |
| network_policy_amendment: NetworkPolicyAmendment, | |
| }, | |
| /// User has denied this command and the agent should not execute it, but | |
| /// it should continue the session and try something else. | |
| Denied { rejection: String }, | |
| /// Automatic approval review timed out before reaching a decision. | |
| TimedOut, | |
| /// User has denied this command and the agent should not do anything until | |
| /// the user's next command. | |
| Abort, | |
| } | |
| impl Default for ReviewDecision { | |
| fn default() -> Self { | |
| Self::Denied { | |
| rejection: "denied".to_string(), | |
| } | |
| } | |
| } | |
| impl ReviewDecision { | |
| pub fn denied(rejection: impl Into<String>) -> Self { | |
| Self::Denied { | |
| rejection: rejection.into(), | |
| } | |
| } | |
| /// Returns an opaque version of the decision without PII. We can't use an ignored flag | |
| /// on `serde` because the serialization is required by some surfaces. | |
| pub fn to_opaque_string(&self) -> &'static str { | |
| match self { | |
| ReviewDecision::Approved => "approved", | |
| ReviewDecision::ApprovedExecpolicyAmendment { .. } => "approved_with_amendment", | |
| ReviewDecision::ApprovedForSession => "approved_for_session", | |
| ReviewDecision::ApprovedMcpPolicyAmendment => "approved_mcp_policy_amendment", | |
| ReviewDecision::NetworkPolicyAmendment { | |
| network_policy_amendment, | |
| } => match network_policy_amendment.action { | |
| NetworkPolicyRuleAction::Allow => "approved_with_network_policy_allow", | |
| NetworkPolicyRuleAction::Deny => "denied_with_network_policy_deny", | |
| }, | |
| ReviewDecision::Denied { .. } => "denied", | |
| ReviewDecision::TimedOut => "timed_out", | |
| ReviewDecision::Abort => "abort", | |
| } | |
| } | |
| } | |
| pub enum FileChange { | |
| Add { | |
| content: String, | |
| }, | |
| Delete { | |
| content: String, | |
| }, | |
| Update { | |
| unified_diff: String, | |
| move_path: Option<PathBuf>, | |
| }, | |
| } | |
| pub struct Chunk { | |
| /// 1-based line index of the first line in the original file | |
| pub orig_index: u32, | |
| pub deleted_lines: Vec<String>, | |
| pub inserted_lines: Vec<String>, | |
| } | |
| pub struct TurnAbortedEvent { | |
| pub turn_id: Option<String>, | |
| pub reason: TurnAbortReason, | |
| /// Unix timestamp (in seconds) when the turn started. | |
| pub started_at: Option<i64>, | |
| /// Unix timestamp (in seconds) when the turn was aborted. | |
| pub completed_at: Option<i64>, | |
| /// Duration between turn start and abort in milliseconds, if known. | |
| pub duration_ms: Option<i64>, | |
| } | |
| pub enum TurnAbortReason { | |
| Interrupted, | |
| Replaced, | |
| ReviewEnded, | |
| BudgetLimited, | |
| } | |
| pub struct CollabAgentSpawnBeginEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub started_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Initial prompt sent to the agent. Can be empty to prevent CoT leaking at the | |
| /// beginning. | |
| pub prompt: String, | |
| pub model: String, | |
| pub reasoning_effort: ReasoningEffortConfig, | |
| } | |
| pub struct CollabAgentRef { | |
| /// Thread ID of the receiver/new agent. | |
| pub thread_id: ThreadId, | |
| /// Optional nickname assigned to an AgentControl-spawned sub-agent. | |
| pub agent_nickname: Option<String>, | |
| /// Optional role (agent_role) assigned to an AgentControl-spawned sub-agent. | |
| pub agent_role: Option<String>, | |
| } | |
| pub struct CollabAgentStatusEntry { | |
| /// Thread ID of the receiver/new agent. | |
| pub thread_id: ThreadId, | |
| /// Optional nickname assigned to an AgentControl-spawned sub-agent. | |
| pub agent_nickname: Option<String>, | |
| /// Optional role (agent_role) assigned to an AgentControl-spawned sub-agent. | |
| pub agent_role: Option<String>, | |
| /// Last known status of the agent. | |
| pub status: AgentStatus, | |
| } | |
| pub struct CollabAgentSpawnEndEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub completed_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the newly spawned agent, if it was created. | |
| pub new_thread_id: Option<ThreadId>, | |
| /// Optional nickname assigned to the new agent. | |
| pub new_agent_nickname: Option<String>, | |
| /// Optional role assigned to the new agent. | |
| pub new_agent_role: Option<String>, | |
| /// Initial prompt sent to the agent. Can be empty to prevent CoT leaking at the | |
| /// beginning. | |
| pub prompt: String, | |
| /// Effective model used by the spawned agent after inheritance and role overrides. | |
| pub model: String, | |
| /// Effective reasoning effort used by the spawned agent after inheritance and role overrides. | |
| pub reasoning_effort: ReasoningEffortConfig, | |
| /// Last known status of the new agent reported to the sender agent. | |
| pub status: AgentStatus, | |
| } | |
| pub struct CollabAgentInteractionBeginEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub started_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| /// Prompt sent from the sender to the receiver. Can be empty to prevent CoT | |
| /// leaking at the beginning. | |
| pub prompt: String, | |
| } | |
| pub struct CollabAgentInteractionEndEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub completed_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| /// Optional nickname assigned to the receiver agent. | |
| pub receiver_agent_nickname: Option<String>, | |
| /// Optional role assigned to the receiver agent. | |
| pub receiver_agent_role: Option<String>, | |
| /// Prompt sent from the sender to the receiver. Can be empty to prevent CoT | |
| /// leaking at the beginning. | |
| pub prompt: String, | |
| /// Last known status of the receiver agent reported to the sender agent. | |
| pub status: AgentStatus, | |
| } | |
| pub enum SubAgentActivityKind { | |
| Started, | |
| Interacted, | |
| Interrupted, | |
| Completed, | |
| } | |
| pub struct SubAgentActivityEvent { | |
| pub event_id: String, | |
| pub occurred_at_ms: i64, | |
| /// Thread ID of the affected sub-agent. | |
| pub agent_thread_id: ThreadId, | |
| /// Canonical v2 path of the affected sub-agent. | |
| pub agent_path: AgentPath, | |
| pub kind: SubAgentActivityKind, | |
| } | |
| pub struct CollabWaitingBeginEvent { | |
| pub started_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receivers. | |
| pub receiver_thread_ids: Vec<ThreadId>, | |
| /// Optional nicknames/roles for receivers. | |
| pub receiver_agents: Vec<CollabAgentRef>, | |
| /// ID of the waiting call. | |
| pub call_id: String, | |
| } | |
| pub struct CollabWaitingEndEvent { | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// ID of the waiting call. | |
| pub call_id: String, | |
| pub completed_at_ms: i64, | |
| /// Optional receiver metadata paired with final statuses. | |
| pub agent_statuses: Vec<CollabAgentStatusEntry>, | |
| /// Last known status of the receiver agents reported to the sender agent. | |
| pub statuses: HashMap<ThreadId, AgentStatus>, | |
| } | |
| pub struct CollabCloseBeginEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub started_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| } | |
| pub struct CollabCloseEndEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub completed_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| /// Optional nickname assigned to the receiver agent. | |
| pub receiver_agent_nickname: Option<String>, | |
| /// Optional role assigned to the receiver agent. | |
| pub receiver_agent_role: Option<String>, | |
| /// Last known status of the receiver agent reported to the sender agent before | |
| /// the close. | |
| pub status: AgentStatus, | |
| } | |
| pub struct CollabResumeBeginEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub started_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| /// Optional nickname assigned to the receiver agent. | |
| pub receiver_agent_nickname: Option<String>, | |
| /// Optional role assigned to the receiver agent. | |
| pub receiver_agent_role: Option<String>, | |
| } | |
| pub struct CollabResumeEndEvent { | |
| /// Identifier for the collab tool call. | |
| pub call_id: String, | |
| pub completed_at_ms: i64, | |
| /// Thread ID of the sender. | |
| pub sender_thread_id: ThreadId, | |
| /// Thread ID of the receiver. | |
| pub receiver_thread_id: ThreadId, | |
| /// Optional nickname assigned to the receiver agent. | |
| pub receiver_agent_nickname: Option<String>, | |
| /// Optional role assigned to the receiver agent. | |
| pub receiver_agent_role: Option<String>, | |
| /// Last known status of the receiver agent reported to the sender agent after | |
| /// resume. | |
| pub status: AgentStatus, | |
| } | |
| mod tests { | |
| use super::*; | |
| use crate::items::CommandExecutionItem; | |
| use crate::items::CommandExecutionStatus; | |
| use crate::items::DynamicToolCallItem; | |
| use crate::items::DynamicToolCallStatus; | |
| use crate::items::EnteredReviewModeItem; | |
| use crate::items::ExitedReviewModeItem; | |
| use crate::items::FileChangeItem; | |
| use crate::items::ImageGenerationItem; | |
| use crate::items::McpToolCallItem; | |
| use crate::items::McpToolCallStatus; | |
| use crate::items::UserMessageItem; | |
| use crate::items::WebSearchItem; | |
| use crate::mcp::CallToolResult; | |
| use crate::permissions::FileSystemAccessMode; | |
| use crate::permissions::FileSystemPath; | |
| use crate::permissions::FileSystemSandboxEntry; | |
| use crate::permissions::FileSystemSandboxPolicy; | |
| use crate::permissions::FileSystemSpecialPath; | |
| use crate::permissions::NetworkSandboxPolicy; | |
| use anyhow::Result; | |
| use codex_utils_absolute_path::AbsolutePathBuf; | |
| use codex_utils_absolute_path::test_support::PathBufExt; | |
| use codex_utils_absolute_path::test_support::test_path_buf; | |
| use pretty_assertions::assert_eq; | |
| use serde_json::json; | |
| use std::path::PathBuf; | |
| use tempfile::NamedTempFile; | |
| use tempfile::TempDir; | |
| fn old_turn_started_records_have_no_root_attribution() { | |
| let event: TurnStartedEvent = serde_json::from_value(serde_json::json!({ | |
| "turn_id": "old-turn", | |
| "model_context_window": null | |
| })) | |
| .unwrap(); | |
| assert_eq!(event.root_turn_id, None); | |
| } | |
| fn review_decision_denied_round_trip() -> Result<()> { | |
| let decision = ReviewDecision::Denied { | |
| rejection: "denied reason".to_string(), | |
| }; | |
| let value = json!({"denied": {"rejection": "denied reason"}}); | |
| assert_eq!(serde_json::to_value(&decision)?, value); | |
| assert_eq!(serde_json::from_value::<ReviewDecision>(value)?, decision); | |
| Ok(()) | |
| } | |
| fn hook_builtin_classification_stays_internal() -> Result<()> { | |
| let wire = json!({ | |
| "id": "cleanup-hook", | |
| "event_name": "stop", | |
| "handler_type": "mcp_tool", | |
| "execution_mode": "sync", | |
| "scope": "turn", | |
| "source_path": test_path_buf("/tmp/hooks.json").abs(), | |
| "source": "plugin", | |
| "display_order": 0, | |
| "status": "completed", | |
| "status_message": null, | |
| "started_at": 10, | |
| "completed_at": 11, | |
| "duration_ms": 1000, | |
| "entries": [], | |
| }); | |
| let mut run: HookRunSummary = serde_json::from_value(wire.clone())?; | |
| assert!(!run.builtin); | |
| run.builtin = true; | |
| assert_eq!(serde_json::to_value(run)?, wire); | |
| let mut untrusted_wire = wire; | |
| untrusted_wire["builtin"] = json!(true); | |
| assert!(!serde_json::from_value::<HookRunSummary>(untrusted_wire)?.builtin); | |
| let schema = serde_json::to_value(schemars::schema_for!(HookRunSummary))?; | |
| assert!( | |
| !schema["properties"] | |
| .as_object() | |
| .expect("hook properties") | |
| .contains_key("builtin") | |
| ); | |
| assert!(!HookRunSummary::decl().contains("builtin:")); | |
| Ok(()) | |
| } | |
| fn feature_thread_source_serializes_as_its_app_owned_label() -> Result<()> { | |
| let source = ThreadSource::Feature("automation".to_string()); | |
| assert_eq!(serde_json::to_value(&source)?, json!("automation")); | |
| assert_eq!( | |
| serde_json::from_value::<ThreadSource>(json!("automation"))?, | |
| source | |
| ); | |
| Ok(()) | |
| } | |
| fn session_meta_normalizes_legacy_dynamic_tools() -> Result<()> { | |
| let mut value = serde_json::to_value(SessionMeta::default())?; | |
| value["dynamic_tools"] = json!([ | |
| { | |
| "namespace": "legacy_app", | |
| "name": "lookup_ticket", | |
| "description": "Look up a ticket", | |
| "inputSchema": {"type": "object", "properties": {}}, | |
| "exposeToContext": false | |
| }, | |
| { | |
| "namespace": "legacy_app", | |
| "name": "update_ticket", | |
| "description": "Update a ticket", | |
| "inputSchema": {"type": "object", "properties": {}}, | |
| "deferLoading": false, | |
| "exposeToContext": false | |
| } | |
| ]); | |
| let meta: SessionMeta = serde_json::from_value(value)?; | |
| assert_eq!( | |
| meta.dynamic_tools, | |
| Some(vec![DynamicToolSpec::Namespace( | |
| crate::dynamic_tools::DynamicToolNamespaceSpec { | |
| name: "legacy_app".to_string(), | |
| description: String::new(), | |
| tools: vec![ | |
| crate::dynamic_tools::DynamicToolNamespaceTool::Function( | |
| crate::dynamic_tools::DynamicToolFunctionSpec { | |
| name: "lookup_ticket".to_string(), | |
| description: "Look up a ticket".to_string(), | |
| input_schema: json!({"type": "object", "properties": {}}), | |
| defer_loading: true, | |
| }, | |
| ), | |
| crate::dynamic_tools::DynamicToolNamespaceTool::Function( | |
| crate::dynamic_tools::DynamicToolFunctionSpec { | |
| name: "update_ticket".to_string(), | |
| description: "Update a ticket".to_string(), | |
| input_schema: json!({"type": "object", "properties": {}}), | |
| defer_loading: false, | |
| }, | |
| ), | |
| ], | |
| }, | |
| )]) | |
| ); | |
| Ok(()) | |
| } | |
| fn sorted_writable_roots(roots: Vec<WritableRoot>) -> Vec<(PathBuf, Vec<PathBuf>)> { | |
| let mut sorted_roots: Vec<(PathBuf, Vec<PathBuf>)> = roots | |
| .into_iter() | |
| .map(|root| { | |
| let mut read_only_subpaths: Vec<PathBuf> = root | |
| .read_only_subpaths | |
| .into_iter() | |
| .map(|path| path.to_path_buf()) | |
| .collect(); | |
| read_only_subpaths.sort(); | |
| (root.root.to_path_buf(), read_only_subpaths) | |
| }) | |
| .collect(); | |
| sorted_roots.sort_by(|left, right| left.0.cmp(&right.0)); | |
| sorted_roots | |
| } | |
| fn sandbox_policy_allows_read(policy: &SandboxPolicy, _path: &Path, _cwd: &Path) -> bool { | |
| policy.has_full_disk_read_access() | |
| } | |
| fn sandbox_policy_allows_write(policy: &SandboxPolicy, path: &Path, cwd: &Path) -> bool { | |
| if policy.has_full_disk_write_access() { | |
| return true; | |
| } | |
| policy | |
| .get_writable_roots_with_cwd(cwd) | |
| .iter() | |
| .any(|root| root.is_path_writable(path)) | |
| } | |
| fn session_source_from_startup_arg_maps_known_values() { | |
| assert_eq!( | |
| SessionSource::from_startup_arg("vscode").unwrap(), | |
| SessionSource::VSCode | |
| ); | |
| assert_eq!( | |
| SessionSource::from_startup_arg("app-server").unwrap(), | |
| SessionSource::Mcp | |
| ); | |
| } | |
| fn inter_agent_communication_response_input_item_preserves_commentary_phase() { | |
| let mut communication = InterAgentCommunication { | |
| id: Some(ResponseItemId::with_suffix("amsg", "1")), | |
| author: AgentPath::root(), | |
| recipient: AgentPath::root().join("reviewer").expect("recipient path"), | |
| other_recipients: vec![AgentPath::root().join("worker").expect("recipient path")], | |
| content: "review the diff".to_string(), | |
| encrypted_content: None, | |
| internal_chat_message_metadata_passthrough: None, | |
| trigger_turn: true, | |
| }; | |
| communication.set_turn_id_if_missing("turn-1"); | |
| let mut serialized_communication = communication.clone(); | |
| serialized_communication.id = None; | |
| serialized_communication.internal_chat_message_metadata_passthrough = None; | |
| assert_eq!( | |
| communication.to_response_input_item(), | |
| ResponseInputItem::Message { | |
| role: "assistant".to_string(), | |
| content: vec![ContentItem::OutputText { | |
| text: serde_json::to_string(&serialized_communication) | |
| .expect("serialize communication"), | |
| }], | |
| phase: Some(MessagePhase::Commentary), | |
| } | |
| ); | |
| } | |
| fn queued_encrypted_inter_agent_communication_renders_message_envelope() { | |
| let communication = InterAgentCommunication::new_encrypted( | |
| AgentPath::root().join("worker").expect("author path"), | |
| AgentPath::root(), | |
| Vec::new(), | |
| "encrypted payload".to_string(), | |
| /*trigger_turn*/ false, | |
| ); | |
| assert_eq!( | |
| communication.to_model_input_item(), | |
| ResponseItem::AgentMessage { | |
| id: None, | |
| author: "/root/worker".to_string(), | |
| recipient: "/root".to_string(), | |
| content: vec![ | |
| AgentMessageInputContent::InputText { | |
| text: "Message Type: MESSAGE\nTask name: /root\nSender: /root/worker\nPayload:\n" | |
| .to_string(), | |
| }, | |
| AgentMessageInputContent::EncryptedContent { | |
| encrypted_content: "encrypted payload".to_string(), | |
| }, | |
| ], | |
| internal_chat_message_metadata_passthrough: None, | |
| } | |
| ); | |
| } | |
| fn session_source_from_startup_arg_normalizes_custom_values() { | |
| assert_eq!( | |
| SessionSource::from_startup_arg("atlas").unwrap(), | |
| SessionSource::Custom("atlas".to_string()) | |
| ); | |
| assert_eq!( | |
| SessionSource::from_startup_arg(" Atlas ").unwrap(), | |
| SessionSource::Custom("atlas".to_string()) | |
| ); | |
| } | |
| fn session_source_restriction_product_defaults_non_subagent_sources_to_codex() { | |
| assert_eq!( | |
| SessionSource::Cli.restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| assert_eq!( | |
| SessionSource::VSCode.restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| assert_eq!( | |
| SessionSource::Exec.restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| assert_eq!( | |
| SessionSource::Mcp.restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| assert_eq!( | |
| SessionSource::Unknown.restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| } | |
| fn session_source_restriction_product_does_not_guess_subagent_products() { | |
| assert_eq!( | |
| SessionSource::SubAgent(SubAgentSource::Review).restriction_product(), | |
| None | |
| ); | |
| assert_eq!( | |
| SessionSource::Internal(InternalSessionSource::MemoryConsolidation) | |
| .restriction_product(), | |
| None | |
| ); | |
| } | |
| fn session_source_restriction_product_maps_custom_sources_to_products() { | |
| assert_eq!( | |
| SessionSource::Custom("chatgpt".to_string()).restriction_product(), | |
| Some(Product::Chatgpt) | |
| ); | |
| assert_eq!( | |
| SessionSource::Custom("ATLAS".to_string()).restriction_product(), | |
| Some(Product::Atlas) | |
| ); | |
| assert_eq!( | |
| SessionSource::Custom("codex".to_string()).restriction_product(), | |
| Some(Product::Codex) | |
| ); | |
| assert_eq!( | |
| SessionSource::Custom("atlas-dev".to_string()).restriction_product(), | |
| None | |
| ); | |
| } | |
| fn session_source_matches_product_restriction() { | |
| assert!( | |
| SessionSource::Custom("chatgpt".to_string()) | |
| .matches_product_restriction(&[Product::Chatgpt]) | |
| ); | |
| assert!( | |
| !SessionSource::Custom("chatgpt".to_string()) | |
| .matches_product_restriction(&[Product::Codex]) | |
| ); | |
| assert!(SessionSource::VSCode.matches_product_restriction(&[Product::Codex])); | |
| assert!( | |
| !SessionSource::Custom("atlas-dev".to_string()) | |
| .matches_product_restriction(&[Product::Atlas]) | |
| ); | |
| assert!(SessionSource::Custom("atlas-dev".to_string()).matches_product_restriction(&[])); | |
| } | |
| fn sandbox_policy_probe_paths(policy: &SandboxPolicy, cwd: &Path) -> Vec<PathBuf> { | |
| let mut paths = vec![cwd.to_path_buf()]; | |
| for root in policy.get_writable_roots_with_cwd(cwd) { | |
| paths.push(root.root.to_path_buf()); | |
| paths.extend( | |
| root.read_only_subpaths | |
| .into_iter() | |
| .map(|path| path.to_path_buf()), | |
| ); | |
| } | |
| paths.sort(); | |
| paths.dedup(); | |
| paths | |
| } | |
| fn assert_same_sandbox_policy_semantics( | |
| expected: &SandboxPolicy, | |
| actual: &SandboxPolicy, | |
| cwd: &Path, | |
| ) { | |
| assert_eq!( | |
| actual.has_full_disk_read_access(), | |
| expected.has_full_disk_read_access() | |
| ); | |
| assert_eq!( | |
| actual.has_full_disk_write_access(), | |
| expected.has_full_disk_write_access() | |
| ); | |
| assert_eq!( | |
| actual.has_full_network_access(), | |
| expected.has_full_network_access() | |
| ); | |
| let mut probe_paths = sandbox_policy_probe_paths(expected, cwd); | |
| probe_paths.extend(sandbox_policy_probe_paths(actual, cwd)); | |
| probe_paths.sort(); | |
| probe_paths.dedup(); | |
| for path in probe_paths { | |
| assert_eq!( | |
| sandbox_policy_allows_read(actual, &path, cwd), | |
| sandbox_policy_allows_read(expected, &path, cwd), | |
| "read access mismatch for {}", | |
| path.display() | |
| ); | |
| assert_eq!( | |
| sandbox_policy_allows_write(actual, &path, cwd), | |
| sandbox_policy_allows_write(expected, &path, cwd), | |
| "write access mismatch for {}", | |
| path.display() | |
| ); | |
| } | |
| } | |
| fn external_sandbox_reports_full_access_flags() { | |
| let restricted = SandboxPolicy::ExternalSandbox { | |
| network_access: NetworkAccess::Restricted, | |
| }; | |
| assert!(restricted.has_full_disk_write_access()); | |
| assert!(!restricted.has_full_network_access()); | |
| let enabled = SandboxPolicy::ExternalSandbox { | |
| network_access: NetworkAccess::Enabled, | |
| }; | |
| assert!(enabled.has_full_disk_write_access()); | |
| assert!(enabled.has_full_network_access()); | |
| } | |
| fn read_only_reports_network_access_flags() { | |
| let restricted = SandboxPolicy::new_read_only_policy(); | |
| assert!(!restricted.has_full_network_access()); | |
| let enabled = SandboxPolicy::ReadOnly { | |
| network_access: true, | |
| }; | |
| assert!(enabled.has_full_network_access()); | |
| } | |
| fn granular_approval_config_mcp_elicitation_flag_is_field_driven() { | |
| assert!( | |
| GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: false, | |
| mcp_elicitations: true, | |
| } | |
| .allows_mcp_elicitations() | |
| ); | |
| assert!( | |
| !GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: false, | |
| mcp_elicitations: false, | |
| } | |
| .allows_mcp_elicitations() | |
| ); | |
| } | |
| fn granular_approval_config_skill_approval_flag_is_field_driven() { | |
| assert!( | |
| GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: true, | |
| request_permissions: false, | |
| mcp_elicitations: false, | |
| } | |
| .allows_skill_approval() | |
| ); | |
| assert!( | |
| !GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: false, | |
| mcp_elicitations: false, | |
| } | |
| .allows_skill_approval() | |
| ); | |
| } | |
| fn granular_approval_config_request_permissions_flag_is_field_driven() { | |
| assert!( | |
| GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: true, | |
| mcp_elicitations: false, | |
| } | |
| .allows_request_permissions() | |
| ); | |
| assert!( | |
| !GranularApprovalConfig { | |
| sandbox_approval: false, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: false, | |
| mcp_elicitations: false, | |
| } | |
| .allows_request_permissions() | |
| ); | |
| } | |
| fn granular_approval_config_defaults_missing_optional_flags_to_false() { | |
| let decoded = serde_json::from_value::<GranularApprovalConfig>(serde_json::json!({ | |
| "sandbox_approval": true, | |
| "rules": false, | |
| "mcp_elicitations": true, | |
| })) | |
| .expect("granular approval config should deserialize"); | |
| assert_eq!( | |
| decoded, | |
| GranularApprovalConfig { | |
| sandbox_approval: true, | |
| rules: false, | |
| skill_approval: false, | |
| request_permissions: false, | |
| mcp_elicitations: true, | |
| } | |
| ); | |
| } | |
| fn restricted_file_system_policy_reports_full_access_from_root_entries() { | |
| let read_only = FileSystemSandboxPolicy::restricted(vec![FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::Root, | |
| }, | |
| access: FileSystemAccessMode::Read, | |
| missing_path_behavior: None, | |
| }]); | |
| assert!(read_only.has_full_disk_read_access()); | |
| assert!(!read_only.has_full_disk_write_access()); | |
| assert!(!read_only.include_platform_defaults()); | |
| let writable = FileSystemSandboxPolicy::restricted(vec![FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::Root, | |
| }, | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }]); | |
| assert!(writable.has_full_disk_read_access()); | |
| assert!(writable.has_full_disk_write_access()); | |
| } | |
| fn restricted_file_system_policy_treats_root_with_carveouts_as_scoped_access() { | |
| let cwd = TempDir::new().expect("tempdir"); | |
| let canonical_cwd = codex_utils_absolute_path::canonicalize_preserving_symlinks(cwd.path()) | |
| .expect("canonicalize cwd"); | |
| let root = AbsolutePathBuf::from_absolute_path(&canonical_cwd) | |
| .expect("absolute canonical tempdir") | |
| .as_path() | |
| .ancestors() | |
| .last() | |
| .and_then(|path| AbsolutePathBuf::from_absolute_path(path).ok()) | |
| .expect("filesystem root"); | |
| let blocked = AbsolutePathBuf::resolve_path_against_base("blocked", cwd.path()); | |
| let expected_blocked = AbsolutePathBuf::from_absolute_path( | |
| codex_utils_absolute_path::canonicalize_preserving_symlinks(cwd.path()) | |
| .expect("canonicalize cwd") | |
| .join("blocked"), | |
| ) | |
| .expect("canonical blocked"); | |
| let policy = FileSystemSandboxPolicy::restricted(vec![ | |
| FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::Root, | |
| }, | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }, | |
| FileSystemSandboxEntry { | |
| path: blocked.into(), | |
| access: FileSystemAccessMode::Deny, | |
| missing_path_behavior: None, | |
| }, | |
| ]); | |
| assert!(!policy.has_full_disk_read_access()); | |
| assert!(!policy.has_full_disk_write_access()); | |
| assert_eq!( | |
| policy.get_readable_roots_with_cwd(cwd.path()), | |
| vec![root.clone()] | |
| ); | |
| assert_eq!( | |
| policy.get_unreadable_roots_with_cwd(cwd.path()), | |
| vec![expected_blocked.clone()] | |
| ); | |
| let writable_roots = policy.get_writable_roots_with_cwd(cwd.path()); | |
| assert_eq!(writable_roots.len(), 1); | |
| assert_eq!(writable_roots[0].root, root); | |
| assert!( | |
| writable_roots[0] | |
| .read_only_subpaths | |
| .iter() | |
| .any(|path| path.as_path() == expected_blocked.as_path()) | |
| ); | |
| } | |
| fn restricted_file_system_policy_derives_effective_paths() { | |
| let cwd = TempDir::new().expect("tempdir"); | |
| std::fs::create_dir_all(cwd.path().join(".agents")).expect("create .agents"); | |
| std::fs::create_dir_all(cwd.path().join(".codex")).expect("create .codex"); | |
| let canonical_cwd = codex_utils_absolute_path::canonicalize_preserving_symlinks(cwd.path()) | |
| .expect("canonicalize cwd"); | |
| let cwd_absolute = | |
| AbsolutePathBuf::from_absolute_path(&canonical_cwd).expect("absolute tempdir"); | |
| let secret = AbsolutePathBuf::resolve_path_against_base("secret", cwd.path()); | |
| let expected_secret = AbsolutePathBuf::from_absolute_path(canonical_cwd.join("secret")) | |
| .expect("canonical secret"); | |
| let expected_agents = AbsolutePathBuf::from_absolute_path(canonical_cwd.join(".agents")) | |
| .expect("canonical .agents"); | |
| let expected_codex = AbsolutePathBuf::from_absolute_path(canonical_cwd.join(".codex")) | |
| .expect("canonical .codex"); | |
| let policy = FileSystemSandboxPolicy::restricted(vec![ | |
| FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::Minimal, | |
| }, | |
| access: FileSystemAccessMode::Read, | |
| missing_path_behavior: None, | |
| }, | |
| FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::project_roots(/*subpath*/ None), | |
| }, | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }, | |
| FileSystemSandboxEntry { | |
| path: secret.into(), | |
| access: FileSystemAccessMode::Deny, | |
| missing_path_behavior: None, | |
| }, | |
| ]); | |
| assert!(!policy.has_full_disk_read_access()); | |
| assert!(!policy.has_full_disk_write_access()); | |
| assert!(policy.include_platform_defaults()); | |
| assert_eq!( | |
| policy.get_readable_roots_with_cwd(cwd.path()), | |
| vec![cwd_absolute.clone()] | |
| ); | |
| assert_eq!( | |
| policy.get_unreadable_roots_with_cwd(cwd.path()), | |
| vec![expected_secret.clone()] | |
| ); | |
| let writable_roots = policy.get_writable_roots_with_cwd(cwd.path()); | |
| assert_eq!(writable_roots.len(), 1); | |
| assert_eq!(writable_roots[0].root, cwd_absolute); | |
| assert!( | |
| writable_roots[0] | |
| .read_only_subpaths | |
| .iter() | |
| .any(|path| path.as_path() == expected_secret.as_path()) | |
| ); | |
| assert!( | |
| writable_roots[0] | |
| .read_only_subpaths | |
| .iter() | |
| .any(|path| path.as_path() == expected_agents.as_path()) | |
| ); | |
| assert!( | |
| writable_roots[0] | |
| .read_only_subpaths | |
| .iter() | |
| .any(|path| path.as_path() == expected_codex.as_path()) | |
| ); | |
| } | |
| fn restricted_file_system_policy_treats_read_entries_as_read_only_subpaths() { | |
| let cwd = TempDir::new().expect("tempdir"); | |
| let canonical_cwd = codex_utils_absolute_path::canonicalize_preserving_symlinks(cwd.path()) | |
| .expect("canonicalize cwd"); | |
| let docs = AbsolutePathBuf::resolve_path_against_base("docs", cwd.path()); | |
| let docs_public = AbsolutePathBuf::resolve_path_against_base("docs/public", cwd.path()); | |
| let expected_docs = AbsolutePathBuf::from_absolute_path(canonical_cwd.join("docs")) | |
| .expect("canonical docs"); | |
| let expected_docs_public = | |
| AbsolutePathBuf::from_absolute_path(canonical_cwd.join("docs/public")) | |
| .expect("canonical docs/public"); | |
| let expected_dot_codex = AbsolutePathBuf::from_absolute_path(canonical_cwd.join(".codex")) | |
| .expect("canonical .codex"); | |
| let policy = FileSystemSandboxPolicy::restricted(vec![ | |
| FileSystemSandboxEntry { | |
| path: FileSystemPath::Special { | |
| value: FileSystemSpecialPath::project_roots(/*subpath*/ None), | |
| }, | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }, | |
| FileSystemSandboxEntry { | |
| path: docs.into(), | |
| access: FileSystemAccessMode::Read, | |
| missing_path_behavior: None, | |
| }, | |
| FileSystemSandboxEntry { | |
| path: docs_public.into(), | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }, | |
| ]); | |
| assert!(!policy.has_full_disk_write_access()); | |
| assert_eq!( | |
| sorted_writable_roots(policy.get_writable_roots_with_cwd(cwd.path())), | |
| vec![ | |
| ( | |
| canonical_cwd, | |
| vec![ | |
| expected_dot_codex.to_path_buf(), | |
| expected_docs.to_path_buf() | |
| ], | |
| ), | |
| (expected_docs_public.to_path_buf(), Vec::new()), | |
| ] | |
| ); | |
| } | |
| fn file_system_policy_rejects_legacy_bridge_for_non_workspace_writes() { | |
| let cwd = if cfg!(windows) { | |
| Path::new(r"C:\workspace") | |
| } else { | |
| Path::new("/tmp/workspace") | |
| }; | |
| let external_write_path = if cfg!(windows) { | |
| AbsolutePathBuf::from_absolute_path(r"C:\temp").expect("absolute windows temp path") | |
| } else { | |
| AbsolutePathBuf::from_absolute_path("/tmp").expect("absolute tmp path") | |
| }; | |
| let policy = FileSystemSandboxPolicy::restricted(vec![FileSystemSandboxEntry { | |
| path: FileSystemPath::Path { | |
| path: external_write_path.into(), | |
| }, | |
| access: FileSystemAccessMode::Write, | |
| missing_path_behavior: None, | |
| }]); | |
| let err = policy | |
| .to_legacy_sandbox_policy(NetworkSandboxPolicy::Restricted, cwd) | |
| .expect_err("non-workspace writes should be rejected"); | |
| assert!( | |
| err.to_string() | |
| .contains("filesystem writes outside the workspace root"), | |
| "{err}" | |
| ); | |
| } | |
| fn legacy_sandbox_policy_semantics_survive_split_bridge() { | |
| let cwd = TempDir::new().expect("tempdir"); | |
| let writable_root = AbsolutePathBuf::resolve_path_against_base("writable", cwd.path()); | |
| let policies = [ | |
| SandboxPolicy::DangerFullAccess, | |
| SandboxPolicy::ExternalSandbox { | |
| network_access: NetworkAccess::Restricted, | |
| }, | |
| SandboxPolicy::ExternalSandbox { | |
| network_access: NetworkAccess::Enabled, | |
| }, | |
| SandboxPolicy::ReadOnly { | |
| network_access: false, | |
| }, | |
| SandboxPolicy::WorkspaceWrite { | |
| writable_roots: vec![], | |
| network_access: false, | |
| exclude_tmpdir_env_var: true, | |
| exclude_slash_tmp: true, | |
| }, | |
| SandboxPolicy::WorkspaceWrite { | |
| writable_roots: vec![writable_root], | |
| network_access: true, | |
| exclude_tmpdir_env_var: false, | |
| exclude_slash_tmp: true, | |
| }, | |
| ]; | |
| for expected in policies { | |
| let actual = | |
| FileSystemSandboxPolicy::from_legacy_sandbox_policy_for_cwd(&expected, cwd.path()) | |
| .to_legacy_sandbox_policy(NetworkSandboxPolicy::from(&expected), cwd.path()) | |
| .expect("legacy bridge should preserve legacy policy semantics"); | |
| assert_same_sandbox_policy_semantics(&expected, &actual, cwd.path()); | |
| } | |
| } | |
| fn item_started_event_from_web_search_emits_begin_event() { | |
| let event = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::WebSearch(WebSearchItem { | |
| id: "search-1".into(), | |
| query: "find docs".into(), | |
| action: WebSearchAction::Search { | |
| query: Some("find docs".into()), | |
| queries: None, | |
| }, | |
| results: None, | |
| }), | |
| started_at_ms: 0, | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::WebSearchBegin(event) => assert_eq!(event.call_id, "search-1"), | |
| _ => panic!("expected WebSearchBegin event"), | |
| } | |
| } | |
| fn item_started_event_from_non_web_search_emits_no_legacy_events() { | |
| let event = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::UserMessage(UserMessageItem::new(&[])), | |
| started_at_ms: 0, | |
| }; | |
| assert!( | |
| event | |
| .as_legacy_events(/*show_raw_agent_reasoning*/ false) | |
| .is_empty() | |
| ); | |
| } | |
| fn item_started_event_from_image_generation_emits_begin_event() { | |
| let event = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::ImageGeneration(ImageGenerationItem { | |
| id: "ig-1".into(), | |
| status: "in_progress".into(), | |
| revised_prompt: None, | |
| result: String::new(), | |
| saved_path: None, | |
| }), | |
| started_at_ms: 0, | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::ImageGenerationBegin(event) => assert_eq!(event.call_id, "ig-1"), | |
| _ => panic!("expected ImageGenerationBegin event"), | |
| } | |
| } | |
| fn item_started_event_from_file_change_emits_patch_begin_event() { | |
| let event = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: 0, | |
| item: TurnItem::FileChange(FileChangeItem { | |
| id: "patch-1".into(), | |
| changes: [( | |
| PathBuf::from("new.txt"), | |
| FileChange::Add { | |
| content: "hello".into(), | |
| }, | |
| )] | |
| .into_iter() | |
| .collect(), | |
| status: None, | |
| auto_approved: Some(true), | |
| stdout: None, | |
| stderr: None, | |
| }), | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::PatchApplyBegin(event) => { | |
| assert_eq!(event.call_id, "patch-1"); | |
| assert_eq!(event.turn_id, "turn-1"); | |
| assert!(event.auto_approved); | |
| assert!(event.changes.contains_key(&PathBuf::from("new.txt"))); | |
| } | |
| _ => panic!("expected PatchApplyBegin event"), | |
| } | |
| } | |
| fn item_started_event_from_mcp_tool_call_emits_begin_event() { | |
| let event = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: 0, | |
| item: TurnItem::McpToolCall(McpToolCallItem { | |
| id: "mcp-1".into(), | |
| server: "server".into(), | |
| tool: "tool".into(), | |
| arguments: json!({"arg": "value"}), | |
| connector_id: Some("connector".into()), | |
| mcp_app_resource_uri: Some("app://connector".into()), | |
| mcp_app_ui: None, | |
| link_id: Some("link_123".into()), | |
| app_name: Some("Calendar".into()), | |
| action_name: Some("create_event".into()), | |
| plugin_id: Some("sample@test".into()), | |
| read_only_hint: Some(false), | |
| status: McpToolCallStatus::InProgress, | |
| result: None, | |
| error: None, | |
| duration: None, | |
| }), | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::McpToolCallBegin(event) => { | |
| assert_eq!(event.turn_id, "turn-1"); | |
| assert_eq!(event.call_id, "mcp-1"); | |
| assert_eq!(event.invocation.server, "server"); | |
| assert_eq!(event.invocation.tool, "tool"); | |
| assert_eq!(event.connector_id.as_deref(), Some("connector")); | |
| assert_eq!( | |
| event.mcp_app_resource_uri.as_deref(), | |
| Some("app://connector") | |
| ); | |
| assert_eq!(event.link_id.as_deref(), Some("link_123")); | |
| assert_eq!(event.app_name.as_deref(), Some("Calendar")); | |
| assert_eq!(event.action_name.as_deref(), Some("create_event")); | |
| assert_eq!(event.plugin_id.as_deref(), Some("sample@test")); | |
| assert_eq!(event.read_only_hint, Some(false)); | |
| } | |
| _ => panic!("expected McpToolCallBegin event"), | |
| } | |
| } | |
| fn item_completed_event_from_image_generation_emits_end_event() { | |
| let event = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::ImageGeneration(ImageGenerationItem { | |
| id: "ig-1".into(), | |
| status: "completed".into(), | |
| revised_prompt: Some("A tiny blue square".into()), | |
| result: "Zm9v".into(), | |
| saved_path: Some(test_path_buf("/tmp/ig-1.png").abs()), | |
| }), | |
| started_at_ms: Some(0), | |
| completed_at_ms: 0, | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::ImageGenerationEnd(event) => { | |
| assert_eq!(event.call_id, "ig-1"); | |
| assert_eq!(event.status, "completed"); | |
| assert_eq!(event.revised_prompt.as_deref(), Some("A tiny blue square")); | |
| assert_eq!(event.result, "Zm9v"); | |
| assert_eq!( | |
| event.saved_path.as_ref().map(AbsolutePathBuf::as_path), | |
| Some(test_path_buf("/tmp/ig-1.png").as_path()) | |
| ); | |
| } | |
| _ => panic!("expected ImageGenerationEnd event"), | |
| } | |
| } | |
| fn item_completed_event_from_file_change_emits_patch_end_event() { | |
| let event = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(0), | |
| completed_at_ms: 0, | |
| item: TurnItem::FileChange(FileChangeItem { | |
| id: "patch-1".into(), | |
| changes: [( | |
| PathBuf::from("new.txt"), | |
| FileChange::Add { | |
| content: "hello".into(), | |
| }, | |
| )] | |
| .into_iter() | |
| .collect(), | |
| status: Some(PatchApplyStatus::Completed), | |
| auto_approved: None, | |
| stdout: Some("Done!".into()), | |
| stderr: Some(String::new()), | |
| }), | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::PatchApplyEnd(event) => { | |
| assert_eq!(event.call_id, "patch-1"); | |
| assert_eq!(event.turn_id, "turn-1"); | |
| assert_eq!(event.stdout, "Done!"); | |
| assert!(event.success); | |
| assert_eq!(event.status, PatchApplyStatus::Completed); | |
| assert!(event.changes.contains_key(&PathBuf::from("new.txt"))); | |
| } | |
| _ => panic!("expected PatchApplyEnd event"), | |
| } | |
| } | |
| fn item_completed_event_from_mcp_tool_call_emits_end_event() { | |
| let event = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(0), | |
| completed_at_ms: 0, | |
| item: TurnItem::McpToolCall(McpToolCallItem { | |
| id: "mcp-1".into(), | |
| server: "server".into(), | |
| tool: "tool".into(), | |
| arguments: json!({"arg": "value"}), | |
| connector_id: Some("connector".into()), | |
| mcp_app_resource_uri: Some("app://connector".into()), | |
| mcp_app_ui: None, | |
| link_id: Some("link_123".into()), | |
| app_name: Some("Calendar".into()), | |
| action_name: Some("create_event".into()), | |
| plugin_id: Some("sample@test".into()), | |
| read_only_hint: None, | |
| status: McpToolCallStatus::Completed, | |
| result: Some(CallToolResult { | |
| content: vec![json!({"type": "text", "text": "ok"})], | |
| structured_content: None, | |
| is_error: Some(false), | |
| meta: None, | |
| }), | |
| error: None, | |
| duration: Some(Duration::from_millis(42)), | |
| }), | |
| }; | |
| let legacy_events = event.as_legacy_events(/*show_raw_agent_reasoning*/ false); | |
| assert_eq!(legacy_events.len(), 1); | |
| match &legacy_events[0] { | |
| EventMsg::McpToolCallEnd(event) => { | |
| assert_eq!(event.turn_id, "turn-1"); | |
| assert_eq!(event.call_id, "mcp-1"); | |
| assert_eq!(event.invocation.server, "server"); | |
| assert_eq!(event.invocation.tool, "tool"); | |
| assert_eq!(event.connector_id.as_deref(), Some("connector")); | |
| assert_eq!( | |
| event.mcp_app_resource_uri.as_deref(), | |
| Some("app://connector") | |
| ); | |
| assert_eq!(event.link_id.as_deref(), Some("link_123")); | |
| assert_eq!(event.app_name.as_deref(), Some("Calendar")); | |
| assert_eq!(event.action_name.as_deref(), Some("create_event")); | |
| assert_eq!(event.plugin_id.as_deref(), Some("sample@test")); | |
| assert_eq!(event.duration, Duration::from_millis(42)); | |
| assert!(event.is_success()); | |
| } | |
| _ => panic!("expected McpToolCallEnd event"), | |
| } | |
| } | |
| fn command_execution_item_lifecycle_emits_legacy_exec_events() { | |
| let cwd = PathUri::from_abs_path(&test_path_buf("/tmp").abs()); | |
| let started = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: 10, | |
| item: TurnItem::CommandExecution(CommandExecutionItem { | |
| model_context: None, | |
| id: "exec-1".into(), | |
| plugin_id: Some("sample@openai-curated".into()), | |
| script_path: Some("scripts/run.py".into()), | |
| process_id: Some("pid-1".into()), | |
| command: vec!["echo".into(), "done".into()], | |
| cwd: cwd.clone(), | |
| parsed_cmd: vec![ParsedCommand::Unknown { | |
| cmd: "echo done".into(), | |
| }], | |
| source: ExecCommandSource::Agent, | |
| interaction_input: None, | |
| status: CommandExecutionStatus::InProgress, | |
| stdout: None, | |
| stderr: None, | |
| aggregated_output: None, | |
| exit_code: None, | |
| duration: None, | |
| formatted_output: None, | |
| }), | |
| }; | |
| let completed = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(10), | |
| completed_at_ms: 20, | |
| item: TurnItem::CommandExecution(CommandExecutionItem { | |
| model_context: None, | |
| id: "exec-1".into(), | |
| plugin_id: Some("sample@openai-curated".into()), | |
| script_path: Some("scripts/run.py".into()), | |
| process_id: Some("pid-1".into()), | |
| command: vec!["echo".into(), "done".into()], | |
| cwd, | |
| parsed_cmd: vec![ParsedCommand::Unknown { | |
| cmd: "echo done".into(), | |
| }], | |
| source: ExecCommandSource::Agent, | |
| interaction_input: None, | |
| status: CommandExecutionStatus::Completed, | |
| stdout: Some("done\n".into()), | |
| stderr: Some(String::new()), | |
| aggregated_output: Some("done\n".into()), | |
| exit_code: Some(0), | |
| duration: Some(Duration::from_millis(5)), | |
| formatted_output: Some("done\n".into()), | |
| }), | |
| }; | |
| assert!(matches!( | |
| started.as_legacy_events(/*show_raw_agent_reasoning*/ false).as_slice(), | |
| [EventMsg::ExecCommandBegin(ExecCommandBeginEvent { | |
| call_id, | |
| plugin_id, | |
| script_path, | |
| turn_id, | |
| started_at_ms: 10, | |
| .. | |
| })] if call_id == "exec-1" | |
| && plugin_id.as_deref() == Some("sample@openai-curated") | |
| && script_path.as_deref() == Some("scripts/run.py") | |
| && turn_id == "turn-1" | |
| )); | |
| assert!(matches!( | |
| completed | |
| .as_legacy_events(/*show_raw_agent_reasoning*/ false) | |
| .as_slice(), | |
| [EventMsg::ExecCommandEnd(ExecCommandEndEvent { | |
| call_id, | |
| plugin_id, | |
| script_path, | |
| turn_id, | |
| completed_at_ms: 20, | |
| aggregated_output, | |
| .. | |
| })] if call_id == "exec-1" | |
| && plugin_id.as_deref() == Some("sample@openai-curated") | |
| && script_path.as_deref() == Some("scripts/run.py") | |
| && turn_id == "turn-1" | |
| && aggregated_output == "done\n" | |
| )); | |
| } | |
| fn dynamic_tool_call_item_lifecycle_emits_legacy_dynamic_tool_events() { | |
| let started = ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: 10, | |
| item: TurnItem::DynamicToolCall(DynamicToolCallItem { | |
| id: "dynamic-1".into(), | |
| namespace: Some("apps".into()), | |
| tool: "lookup".into(), | |
| arguments: json!({"id": "123"}), | |
| status: DynamicToolCallStatus::InProgress, | |
| content_items: None, | |
| success: None, | |
| error: None, | |
| duration: None, | |
| }), | |
| }; | |
| let completed = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(10), | |
| completed_at_ms: 20, | |
| item: TurnItem::DynamicToolCall(DynamicToolCallItem { | |
| id: "dynamic-1".into(), | |
| namespace: Some("apps".into()), | |
| tool: "lookup".into(), | |
| arguments: json!({"id": "123"}), | |
| status: DynamicToolCallStatus::Completed, | |
| content_items: Some(vec![DynamicToolCallOutputContentItem::InputText { | |
| text: "ok".into(), | |
| }]), | |
| success: Some(true), | |
| error: None, | |
| duration: Some(Duration::from_millis(5)), | |
| }), | |
| }; | |
| assert!(matches!( | |
| started.as_legacy_events(/*show_raw_agent_reasoning*/ false).as_slice(), | |
| [EventMsg::DynamicToolCallRequest(DynamicToolCallRequest { | |
| call_id, | |
| turn_id, | |
| started_at_ms: 10, | |
| .. | |
| })] if call_id == "dynamic-1" && turn_id == "turn-1" | |
| )); | |
| assert!(matches!( | |
| completed | |
| .as_legacy_events(/*show_raw_agent_reasoning*/ false) | |
| .as_slice(), | |
| [EventMsg::DynamicToolCallResponse(DynamicToolCallResponseEvent { | |
| call_id, | |
| turn_id, | |
| completed_at_ms: 20, | |
| success: true, | |
| .. | |
| })] if call_id == "dynamic-1" && turn_id == "turn-1" | |
| )); | |
| } | |
| fn review_mode_item_completion_emits_legacy_events_with_ids() { | |
| let entered = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(0), | |
| completed_at_ms: 0, | |
| item: TurnItem::EnteredReviewMode(EnteredReviewModeItem { | |
| id: "entered-review".into(), | |
| target: ReviewTarget::Custom { | |
| instructions: "review this".into(), | |
| }, | |
| user_facing_hint: "Review requested.".into(), | |
| }), | |
| }; | |
| let exited = ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| started_at_ms: Some(0), | |
| completed_at_ms: 0, | |
| item: TurnItem::ExitedReviewMode(ExitedReviewModeItem { | |
| id: "exited-review".into(), | |
| review_output: Some(ReviewOutputEvent { | |
| overall_explanation: "Looks good.".into(), | |
| ..Default::default() | |
| }), | |
| }), | |
| }; | |
| assert!(matches!( | |
| entered | |
| .as_legacy_events(/*show_raw_agent_reasoning*/ false) | |
| .as_slice(), | |
| [EventMsg::EnteredReviewMode(EnteredReviewModeEvent { | |
| target: ReviewTarget::Custom { instructions }, | |
| user_facing_hint: Some(user_facing_hint), | |
| turn_id: Some(turn_id), | |
| item_id: Some(item_id), | |
| })] | |
| if instructions == "review this" | |
| && user_facing_hint == "Review requested." | |
| && turn_id == "turn-1" | |
| && item_id == "entered-review" | |
| )); | |
| assert!(matches!( | |
| exited | |
| .as_legacy_events(/*show_raw_agent_reasoning*/ false) | |
| .as_slice(), | |
| [EventMsg::ExitedReviewMode(ExitedReviewModeEvent { | |
| turn_id: Some(turn_id), | |
| item_id: Some(item_id), | |
| review_output: Some(review_output), | |
| })] | |
| if turn_id == "turn-1" | |
| && item_id == "exited-review" | |
| && review_output.overall_explanation == "Looks good." | |
| )); | |
| } | |
| fn item_started_event_requires_started_at_ms() { | |
| let mut value = serde_json::to_value(ItemStartedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::UserMessage(UserMessageItem::new(&[])), | |
| started_at_ms: 123, | |
| }) | |
| .unwrap(); | |
| value.as_object_mut().unwrap().remove("started_at_ms"); | |
| assert!(serde_json::from_value::<ItemStartedEvent>(value).is_err()); | |
| } | |
| fn item_completed_event_defaults_missing_completed_at_ms() { | |
| let mut value = serde_json::to_value(ItemCompletedEvent { | |
| thread_id: ThreadId::new(), | |
| turn_id: "turn-1".into(), | |
| item: TurnItem::UserMessage(UserMessageItem::new(&[])), | |
| started_at_ms: None, | |
| completed_at_ms: 123, | |
| }) | |
| .unwrap(); | |
| value.as_object_mut().unwrap().remove("completed_at_ms"); | |
| let event = serde_json::from_value::<ItemCompletedEvent>(value).unwrap(); | |
| assert_eq!(event.started_at_ms, None); | |
| assert_eq!(event.completed_at_ms, 0); | |
| } | |
| fn review_mode_events_deserialize_legacy_payloads() { | |
| let entered = serde_json::from_value::<EnteredReviewModeEvent>(json!({ | |
| "target": { | |
| "type": "custom", | |
| "instructions": "review this" | |
| }, | |
| "user_facing_hint": "hint" | |
| })) | |
| .unwrap(); | |
| assert_eq!(entered.turn_id, None); | |
| assert_eq!(entered.item_id, None); | |
| let exited = serde_json::from_value::<ExitedReviewModeEvent>(json!({ | |
| "review_output": null | |
| })) | |
| .unwrap(); | |
| assert_eq!(exited.turn_id, None); | |
| assert_eq!(exited.item_id, None); | |
| } | |
| fn rollback_failed_error_does_not_affect_turn_status() { | |
| let event = ErrorEvent { | |
| misalignment: None, | |
| message: "rollback failed".into(), | |
| codex_error_info: Some(CodexErrorInfo::ThreadRollbackFailed), | |
| }; | |
| assert!(!event.affects_turn_status()); | |
| } | |
| fn active_turn_not_steerable_error_does_not_affect_turn_status() { | |
| let event = ErrorEvent { | |
| misalignment: None, | |
| message: "cannot steer a review turn".into(), | |
| codex_error_info: Some(CodexErrorInfo::ActiveTurnNotSteerable { | |
| turn_kind: NonSteerableTurnKind::Review, | |
| }), | |
| }; | |
| assert!(!event.affects_turn_status()); | |
| } | |
| fn generic_error_affects_turn_status() { | |
| let event = ErrorEvent { | |
| misalignment: None, | |
| message: "generic".into(), | |
| codex_error_info: Some(CodexErrorInfo::Other), | |
| }; | |
| assert!(event.affects_turn_status()); | |
| } | |
| fn misalignment_explanation_and_steer_are_never_serialized_into_error_events() { | |
| let event = ErrorEvent { | |
| message: "This request violated the misalignment policy.".to_string(), | |
| codex_error_info: Some(CodexErrorInfo::MisalignmentPolicyViolation), | |
| misalignment: Some(MisalignmentErrorDetails { | |
| error_type: Some("unauthorized_data_transfer".to_string()), | |
| detailed_explanation: Some("Sensitive customer explanation".to_string()), | |
| steer: Some(MisalignmentSteer { | |
| message: "Sensitive customer steering".to_string(), | |
| }), | |
| }), | |
| }; | |
| let serialized = serde_json::to_value(&event).expect("serialize error event"); | |
| assert_eq!( | |
| serialized, | |
| json!({ | |
| "message": "This request violated the misalignment policy.", | |
| "codex_error_info": "misalignment_policy_violation" | |
| }) | |
| ); | |
| let restored: ErrorEvent = | |
| serde_json::from_value(serialized).expect("deserialize persisted error event"); | |
| assert_eq!(restored.misalignment, None); | |
| let debug = format!("{event:?}"); | |
| assert!(!debug.contains("Sensitive customer explanation")); | |
| assert!(!debug.contains("Sensitive customer steering")); | |
| } | |
| fn realtime_conversation_started_event_uses_realtime_session_id() { | |
| let event = RealtimeConversationStartedEvent { | |
| realtime_session_id: Some("conv_1".to_string()), | |
| version: RealtimeConversationVersion::V2, | |
| }; | |
| assert_eq!( | |
| serde_json::to_value(&event).unwrap(), | |
| json!({ | |
| "realtime_session_id": "conv_1", | |
| "version": "v2" | |
| }) | |
| ); | |
| } | |
| fn realtime_voice_list_is_stable() { | |
| assert_eq!( | |
| RealtimeVoicesList::builtin(), | |
| RealtimeVoicesList { | |
| v1: vec![ | |
| RealtimeVoice::Juniper, | |
| RealtimeVoice::Maple, | |
| RealtimeVoice::Spruce, | |
| RealtimeVoice::Ember, | |
| RealtimeVoice::Vale, | |
| RealtimeVoice::Breeze, | |
| RealtimeVoice::Arbor, | |
| RealtimeVoice::Sol, | |
| RealtimeVoice::Cove, | |
| ], | |
| v2: vec![ | |
| RealtimeVoice::Alloy, | |
| RealtimeVoice::Ash, | |
| RealtimeVoice::Ballad, | |
| RealtimeVoice::Coral, | |
| RealtimeVoice::Echo, | |
| RealtimeVoice::Sage, | |
| RealtimeVoice::Shimmer, | |
| RealtimeVoice::Verse, | |
| RealtimeVoice::Marin, | |
| RealtimeVoice::Cedar, | |
| ], | |
| default_v1: RealtimeVoice::Cove, | |
| default_v2: RealtimeVoice::Marin, | |
| } | |
| ); | |
| } | |
| fn user_input_text_serializes_empty_text_elements() -> Result<()> { | |
| let input = crate::user_input::UserInput::Text { | |
| text: "hello".to_string(), | |
| text_elements: Vec::new(), | |
| }; | |
| let json_input = serde_json::to_value(input)?; | |
| assert_eq!( | |
| json_input, | |
| json!({ | |
| "type": "text", | |
| "text": "hello", | |
| "text_elements": [], | |
| }) | |
| ); | |
| Ok(()) | |
| } | |
| fn user_message_event_serializes_empty_metadata_vectors() -> Result<()> { | |
| let event = UserMessageEvent { | |
| client_id: None, | |
| message: "hello".to_string(), | |
| images: None, | |
| local_images: Vec::new(), | |
| text_elements: Vec::new(), | |
| ..Default::default() | |
| }; | |
| let json_event = serde_json::to_value(event)?; | |
| assert_eq!( | |
| json_event, | |
| json!({ | |
| "message": "hello", | |
| "local_images": [], | |
| "local_audio": [], | |
| "text_elements": [], | |
| }) | |
| ); | |
| Ok(()) | |
| } | |
| fn user_message_event_deserializes_without_image_detail_fields() -> Result<()> { | |
| let event: UserMessageEvent = serde_json::from_value(json!({ | |
| "message": "hello", | |
| "images": ["https://example.com/image.png"], | |
| "local_images": ["/tmp/local.png"], | |
| "text_elements": [], | |
| }))?; | |
| assert_eq!(event.message, "hello"); | |
| assert_eq!( | |
| event.images, | |
| Some(vec!["https://example.com/image.png".to_string()]) | |
| ); | |
| assert_eq!(event.image_details, Vec::<Option<ImageDetail>>::new()); | |
| assert_eq!(event.file_ids, None); | |
| assert_eq!(event.file_id_details, Vec::<Option<ImageDetail>>::new()); | |
| assert_eq!(event.image_order, Vec::<UserMessageImageKind>::new()); | |
| assert_eq!(event.local_images, vec![PathBuf::from("/tmp/local.png")]); | |
| assert_eq!(event.local_image_details, Vec::<Option<ImageDetail>>::new()); | |
| assert_eq!(event.audio, None); | |
| assert_eq!(event.local_audio, Vec::<PathBuf>::new()); | |
| assert_eq!(event.text_elements, Vec::new()); | |
| Ok(()) | |
| } | |
| fn user_message_item_legacy_event_preserves_attachments() { | |
| let local_path = PathBuf::from("/tmp/local.png"); | |
| let local_audio_path = PathBuf::from("/tmp/local.wav"); | |
| let mut item = UserMessageItem::new(&[ | |
| crate::user_input::UserInput::Image { | |
| image: crate::models::ImageReference::Inline { | |
| image_url: "https://example.com/first.png".to_string(), | |
| }, | |
| detail: Some(ImageDetail::Original), | |
| }, | |
| crate::user_input::UserInput::Image { | |
| image: crate::models::ImageReference::File { | |
| file_id: "file_123".to_string(), | |
| }, | |
| detail: Some(ImageDetail::Low), | |
| }, | |
| crate::user_input::UserInput::Image { | |
| image: crate::models::ImageReference::Inline { | |
| image_url: "https://example.com/second.png".to_string(), | |
| }, | |
| detail: None, | |
| }, | |
| crate::user_input::UserInput::LocalImage { | |
| path: local_path.clone(), | |
| detail: Some(ImageDetail::Original), | |
| }, | |
| crate::user_input::UserInput::Audio { | |
| audio_url: "https://example.com/remote.mp3".to_string(), | |
| }, | |
| crate::user_input::UserInput::LocalAudio { | |
| path: local_audio_path.clone(), | |
| }, | |
| ]); | |
| item.client_id = Some("client-message-1".to_string()); | |
| let EventMsg::UserMessage(event) = item.as_legacy_event() else { | |
| panic!("expected user message event"); | |
| }; | |
| let event_json = serde_json::to_value(&event).expect("serialize user message event"); | |
| assert_eq!( | |
| event.images, | |
| Some(vec![ | |
| "https://example.com/first.png".to_string(), | |
| "https://example.com/second.png".to_string(), | |
| ]) | |
| ); | |
| assert_eq!(event.client_id, Some("client-message-1".to_string())); | |
| assert_eq!(event.image_details, vec![Some(ImageDetail::Original)]); | |
| assert_eq!(event.file_ids, Some(vec!["file_123".to_string()])); | |
| assert_eq!(event.file_id_details, vec![Some(ImageDetail::Low)]); | |
| assert_eq!( | |
| event.image_order, | |
| vec![ | |
| UserMessageImageKind::Inline, | |
| UserMessageImageKind::File, | |
| UserMessageImageKind::Inline, | |
| ] | |
| ); | |
| assert_eq!(event_json["file_ids"], json!(["file_123"])); | |
| assert_eq!(event_json["file_id_details"], json!(["low"])); | |
| assert_eq!( | |
| event_json["image_order"], | |
| json!(["inline", "file", "inline"]) | |
| ); | |
| assert_eq!(event.local_images, vec![local_path]); | |
| assert_eq!(event.local_image_details, vec![Some(ImageDetail::Original)]); | |
| assert_eq!( | |
| event.audio, | |
| Some(vec!["https://example.com/remote.mp3".to_string()]) | |
| ); | |
| assert_eq!(event.local_audio, vec![local_audio_path]); | |
| } | |
| fn audio_only_user_message_has_placeholder_preview() { | |
| let event = UserMessageEvent { | |
| audio: Some(vec!["https://example.com/remote.mp3".to_string()]), | |
| ..Default::default() | |
| }; | |
| assert_eq!(user_message_preview(&event), Some("[Audio]".to_string())); | |
| } | |
| fn file_only_user_message_has_placeholder_preview() { | |
| let event = UserMessageEvent { | |
| file_ids: Some(vec!["file_123".to_string()]), | |
| ..Default::default() | |
| }; | |
| assert_eq!(user_message_preview(&event), Some("[Image]".to_string())); | |
| } | |
| fn turn_aborted_event_deserializes_without_turn_id() -> Result<()> { | |
| let event: EventMsg = serde_json::from_value(json!({ | |
| "type": "turn_aborted", | |
| "reason": "interrupted", | |
| }))?; | |
| match event { | |
| EventMsg::TurnAborted(TurnAbortedEvent { | |
| turn_id, reason, .. | |
| }) => { | |
| assert_eq!(turn_id, None); | |
| assert_eq!(reason, TurnAbortReason::Interrupted); | |
| } | |
| _ => panic!("expected turn_aborted event"), | |
| } | |
| Ok(()) | |
| } | |
| fn session_meta_defaults_legacy_history_mode() -> Result<()> { | |
| let session_meta: SessionMeta = serde_json::from_value(json!({ | |
| "session_id": "00000000-0000-0000-0000-000000000001", | |
| "id": "00000000-0000-0000-0000-000000000001", | |
| "timestamp": "2026-01-01T00:00:00Z", | |
| "cwd": "/tmp", | |
| "originator": "codex", | |
| "cli_version": "0.0.0", | |
| "model_provider": null, | |
| "base_instructions": null | |
| }))?; | |
| assert_eq!(session_meta.history_mode, ThreadHistoryMode::Legacy); | |
| assert_eq!(session_meta.history_base, None); | |
| assert_eq!(session_meta.forked_from_ordinal_exclusive, None); | |
| let serialized = serde_json::to_value(&session_meta)?; | |
| assert!(serialized.get("forked_from_ordinal_exclusive").is_none()); | |
| assert_eq!(serialized["history_mode"], json!("legacy")); | |
| let mut unknown = serialized; | |
| unknown["history_mode"] = json!("future"); | |
| assert!(serde_json::from_value::<SessionMeta>(unknown).is_err()); | |
| Ok(()) | |
| } | |
| fn turn_context_item_deserializes_without_network() -> Result<()> { | |
| let item: TurnContextItem = serde_json::from_value(json!({ | |
| "cwd": test_path_buf("/tmp"), | |
| "approval_policy": "never", | |
| "sandbox_policy": { "type": "danger-full-access" }, | |
| "model": "gpt-5", | |
| "summary": "auto", | |
| }))?; | |
| assert_eq!(item.network, None); | |
| assert_eq!(item.file_system_sandbox_policy, None); | |
| assert_eq!(item.comp_hash, None); | |
| Ok(()) | |
| } | |
| fn turn_context_item_deserializes_legacy_on_failure_as_on_request() -> Result<()> { | |
| let item: TurnContextItem = serde_json::from_value(json!({ | |
| "cwd": test_path_buf("/tmp"), | |
| "approval_policy": "on-failure", | |
| "sandbox_policy": { "type": "danger-full-access" }, | |
| "model": "gpt-5", | |
| "summary": "auto", | |
| }))?; | |
| assert_eq!(item.approval_policy, AskForApproval::OnRequest); | |
| Ok(()) | |
| } | |
| fn turn_context_item_serializes_network_when_present() -> Result<()> { | |
| let item = TurnContextItem { | |
| turn_id: None, | |
| root_turn_id: None, | |
| disabled_plugin_ids: None, | |
| cwd: test_path_buf("/tmp").abs(), | |
| workspace_roots: None, | |
| current_date: None, | |
| timezone: None, | |
| approval_policy: AskForApproval::Never, | |
| approvals_reviewer: None, | |
| sandbox_policy: SandboxPolicy::DangerFullAccess, | |
| permission_profile: None, | |
| active_permission_profile: None, | |
| network: Some(TurnContextNetworkItem { | |
| allowed_domains: vec!["api.example.com".to_string()], | |
| denied_domains: vec!["blocked.example.com".to_string()], | |
| }), | |
| file_system_sandbox_policy: Some( | |
| FileSystemSandboxPolicy::restricted(vec![FileSystemSandboxEntry { | |
| path: FileSystemPath::GlobPattern { | |
| pattern: "/tmp/private/**/*.txt".to_string(), | |
| }, | |
| access: FileSystemAccessMode::Deny, | |
| missing_path_behavior: None, | |
| }]) | |
| .try_into() | |
| .expect("serializable split policy"), | |
| ), | |
| model: "gpt-5".to_string(), | |
| comp_hash: None, | |
| personality: None, | |
| collaboration_mode: None, | |
| multi_agent_version: None, | |
| multi_agent_mode: None, | |
| realtime_active: None, | |
| cyber_access_program: None, | |
| effort: None, | |
| summary: ReasoningSummaryConfig::Auto, | |
| }; | |
| let value = serde_json::to_value(item)?; | |
| assert_eq!( | |
| value["network"], | |
| json!({ | |
| "allowed_domains": ["api.example.com"], | |
| "denied_domains": ["blocked.example.com"], | |
| }) | |
| ); | |
| assert_eq!( | |
| value["file_system_sandbox_policy"], | |
| json!({ | |
| "kind": "restricted", | |
| "entries": [{ | |
| "path": { | |
| "type": "glob_pattern", | |
| "pattern": "/tmp/private/**/*.txt" | |
| }, | |
| "access": "deny" | |
| }] | |
| }) | |
| ); | |
| assert_eq!(value["summary"], json!("auto")); | |
| Ok(()) | |
| } | |
| /// Serialize Event to verify that its JSON representation has the expected | |
| /// amount of nesting. | |
| fn serialize_event() -> Result<()> { | |
| let session_id = SessionId::from_string("67e55044-10b1-426f-9247-bb680e5fe0c7")?; | |
| let thread_id = ThreadId::from_string("67e55044-10b1-426f-9247-bb680e5fe0c8")?; | |
| let rollout_file = NamedTempFile::new()?; | |
| let permission_profile = PermissionProfile::read_only(); | |
| let event = Event { | |
| id: "1234".to_string(), | |
| msg: EventMsg::SessionConfigured(SessionConfiguredEvent { | |
| session_id, | |
| thread_id, | |
| forked_from_id: None, | |
| parent_thread_id: None, | |
| thread_source: None, | |
| thread_name: None, | |
| model: "codex-mini-latest".to_string(), | |
| model_provider_id: "openai".to_string(), | |
| service_tier: None, | |
| approval_policy: AskForApproval::Never, | |
| approvals_reviewer: ApprovalsReviewer::User, | |
| permission_profile: permission_profile.clone(), | |
| active_permission_profile: None, | |
| cwd: test_path_buf("/home/user/project").abs(), | |
| reasoning_effort: Some(ReasoningEffortConfig::default()), | |
| initial_messages: None, | |
| network_proxy: None, | |
| rollout_path: Some(rollout_file.path().to_path_buf()), | |
| }), | |
| }; | |
| let expected = json!({ | |
| "id": "1234", | |
| "msg": { | |
| "type": "session_configured", | |
| "session_id": "67e55044-10b1-426f-9247-bb680e5fe0c7", | |
| "thread_id": "67e55044-10b1-426f-9247-bb680e5fe0c8", | |
| "model": "codex-mini-latest", | |
| "model_provider_id": "openai", | |
| "approval_policy": "never", | |
| "approvals_reviewer": "user", | |
| "permission_profile": permission_profile, | |
| "cwd": test_path_buf("/home/user/project"), | |
| "reasoning_effort": "medium", | |
| "rollout_path": format!("{}", rollout_file.path().display()), | |
| } | |
| }); | |
| assert_eq!(expected, serde_json::to_value(&event)?); | |
| Ok(()) | |
| } | |
| fn deserialize_legacy_session_configured_event_uses_sandbox_policy() -> Result<()> { | |
| let cwd = test_path_buf("/home/user/project"); | |
| let value = json!({ | |
| "session_id": "67e55044-10b1-426f-9247-bb680e5fe0c8", | |
| "model": "codex-mini-latest", | |
| "model_provider_id": "openai", | |
| "approval_policy": "never", | |
| "approvals_reviewer": "user", | |
| "sandbox_policy": { | |
| "type": "read-only" | |
| }, | |
| "cwd": cwd, | |
| }); | |
| let event: SessionConfiguredEvent = serde_json::from_value(value)?; | |
| assert_eq!(event.permission_profile, PermissionProfile::read_only()); | |
| Ok(()) | |
| } | |
| fn vec_u8_as_base64_serialization_and_deserialization() -> Result<()> { | |
| let event = ExecCommandOutputDeltaEvent { | |
| call_id: "call21".to_string(), | |
| stream: ExecOutputStream::Stdout, | |
| chunk: vec![1, 2, 3, 4, 5], | |
| }; | |
| let serialized = serde_json::to_string(&event)?; | |
| assert_eq!( | |
| r#"{"call_id":"call21","stream":"stdout","chunk":"AQIDBAU="}"#, | |
| serialized, | |
| ); | |
| let deserialized: ExecCommandOutputDeltaEvent = serde_json::from_str(&serialized)?; | |
| assert_eq!(deserialized, event); | |
| Ok(()) | |
| } | |
| fn serialize_mcp_startup_update_event() -> Result<()> { | |
| let event = Event { | |
| id: "init".to_string(), | |
| msg: EventMsg::McpStartupUpdate(McpStartupUpdateEvent { | |
| server: "srv".to_string(), | |
| status: McpStartupStatus::Failed { | |
| error: "boom".to_string(), | |
| reason: Some(McpStartupFailureReason::ReauthenticationRequired), | |
| }, | |
| }), | |
| }; | |
| let value = serde_json::to_value(&event)?; | |
| assert_eq!(value["msg"]["type"], "mcp_startup_update"); | |
| assert_eq!(value["msg"]["server"], "srv"); | |
| assert_eq!(value["msg"]["status"]["state"], "failed"); | |
| assert_eq!(value["msg"]["status"]["error"], "boom"); | |
| assert_eq!( | |
| value["msg"]["status"]["reason"], | |
| "reauthentication_required" | |
| ); | |
| Ok(()) | |
| } | |
| fn serialize_mcp_startup_complete_event() -> Result<()> { | |
| let event = Event { | |
| id: "init".to_string(), | |
| msg: EventMsg::McpStartupComplete(McpStartupCompleteEvent { | |
| ready: vec!["a".to_string()], | |
| failed: vec![McpStartupFailure { | |
| server: "b".to_string(), | |
| error: "bad".to_string(), | |
| }], | |
| cancelled: vec!["c".to_string()], | |
| }), | |
| }; | |
| let value = serde_json::to_value(&event)?; | |
| assert_eq!(value["msg"]["type"], "mcp_startup_complete"); | |
| assert_eq!(value["msg"]["ready"][0], "a"); | |
| assert_eq!(value["msg"]["failed"][0]["server"], "b"); | |
| assert_eq!(value["msg"]["failed"][0]["error"], "bad"); | |
| assert_eq!(value["msg"]["cancelled"][0], "c"); | |
| Ok(()) | |
| } | |
| fn token_usage_info_new_or_append_updates_context_window_when_provided() { | |
| let initial = Some(TokenUsageInfo { | |
| total_token_usage: TokenUsage::default(), | |
| last_token_usage: TokenUsage::default(), | |
| model_context_window: Some(258_400), | |
| }); | |
| let last = Some(TokenUsage { | |
| input_tokens: 10, | |
| cached_input_tokens: 0, | |
| cache_write_input_tokens: 0, | |
| output_tokens: 0, | |
| reasoning_output_tokens: 0, | |
| total_tokens: 10, | |
| codex_rollout_budget_units: None, | |
| }); | |
| let info = TokenUsageInfo::new_or_append(&initial, &last, Some(128_000)) | |
| .expect("new_or_append should return info"); | |
| assert_eq!(info.model_context_window, Some(128_000)); | |
| } | |
| fn token_usage_info_new_or_append_preserves_context_window_when_not_provided() { | |
| let initial = Some(TokenUsageInfo { | |
| total_token_usage: TokenUsage::default(), | |
| last_token_usage: TokenUsage::default(), | |
| model_context_window: Some(258_400), | |
| }); | |
| let last = Some(TokenUsage { | |
| input_tokens: 10, | |
| cached_input_tokens: 0, | |
| cache_write_input_tokens: 0, | |
| output_tokens: 0, | |
| reasoning_output_tokens: 0, | |
| total_tokens: 10, | |
| codex_rollout_budget_units: None, | |
| }); | |
| let info = | |
| TokenUsageInfo::new_or_append(&initial, &last, /*model_context_window*/ None) | |
| .expect("new_or_append should return info"); | |
| assert_eq!(info.model_context_window, Some(258_400)); | |
| } | |
| } | |