Remove core protocol dependency [2/2] (#20325)

## Why

With the local model layer and app-server routing in place from PR1,
this PR moves the active TUI runtime onto app-server notifications. The
affected pieces share the same event flow, so the command surface,
session state, bottom-pane prompts, chat rendering, history/status
views, and tests move together to keep the stacked branch buildable.

This PR also removes the obsolete compatibility surface that is no
longer used after the migration. The proposed protocol-boundary verifier
layer was dropped from the stack; enforcing that final boundary will be
simpler once `codex-tui` no longer needs any `codex_protocol`
references.

This PR is part 2 of a 2-PR stack:

1. Add TUI-owned replacement models and extract app-server event
routing.
2. Move the active TUI flow to app-server notifications and delete
obsolete adapter code.

## What changed

- Rewired app command and session handling to use app-server request and
notification shapes.
- Moved approval overlays, request-user-input flows, MCP elicitation,
realtime events, and review commands onto the app-server-facing model
surface.
- Updated chat rendering, history cells, status views, multi-agent UI,
replay state, and TUI tests to use app-server notifications plus the
local models introduced in PR1.
- Deleted `codex-rs/tui/src/app/app_server_adapter.rs` and the
superseded `chatwidget/tests/background_events.rs` fixture path.

## Verification

- `cargo check -p codex-tui --tests`
- Top of stack: `cargo test -p codex-tui`
This commit is contained in:
Eric Traut
2026-04-30 11:34:34 -07:00
committed by GitHub
Unverified
parent 5cc5f12efc
commit f2bc2f26a9
76 changed files with 5078 additions and 9951 deletions
File diff suppressed because it is too large Load Diff
+2 -24
View File
@@ -7,11 +7,9 @@ use super::app_server_event_targets::server_request_thread_id;
use crate::app_command::AppCommand;
use crate::app_event::AppEvent;
use crate::app_server_session::AppServerSession;
use crate::app_server_session::app_server_rate_limit_snapshot_to_core;
use crate::app_server_session::status_account_display_from_auth_mode;
use codex_app_server_client::AppServerEvent;
use codex_app_server_protocol::AuthMode;
use codex_app_server_protocol::JSONRPCErrorError;
use codex_app_server_protocol::ServerNotification;
use codex_app_server_protocol::ServerRequest;
@@ -77,9 +75,8 @@ impl App {
self.refresh_mcp_startup_expected_servers_from_config();
}
ServerNotification::AccountRateLimitsUpdated(notification) => {
self.chat_widget.on_rate_limit_snapshot(Some(
app_server_rate_limit_snapshot_to_core(notification.rate_limits.clone()),
));
self.chat_widget
.on_rate_limit_snapshot(Some(notification.rate_limits.clone()));
return;
}
ServerNotification::AccountUpdated(notification) => {
@@ -186,23 +183,4 @@ impl App {
tracing::warn!("failed to enqueue app-server request: {err}");
}
}
async fn reject_app_server_request(
&self,
app_server_client: &AppServerSession,
request_id: codex_app_server_protocol::RequestId,
reason: String,
) -> std::result::Result<(), String> {
app_server_client
.reject_server_request(
request_id,
JSONRPCErrorError {
code: -32000,
message: reason,
data: None,
},
)
.await
.map_err(|err| format!("failed to reject app-server request: {err}"))
}
}
+62 -94
View File
@@ -1,20 +1,38 @@
use std::collections::HashMap;
use std::collections::VecDeque;
use super::App;
use crate::app_command::AppCommand;
use crate::app_command::AppCommandView;
use crate::app_server_approval_conversions::granted_permission_profile_from_request;
use crate::app_server_session::AppServerSession;
use codex_app_server_protocol::CommandExecutionRequestApprovalResponse;
use codex_app_server_protocol::FileChangeApprovalDecision;
use codex_app_server_protocol::FileChangeRequestApprovalResponse;
use codex_app_server_protocol::McpServerElicitationAction;
use codex_app_server_protocol::JSONRPCErrorError;
use codex_app_server_protocol::McpServerElicitationRequestResponse;
use codex_app_server_protocol::PermissionsRequestApprovalResponse;
use codex_app_server_protocol::RequestId as AppServerRequestId;
use codex_app_server_protocol::ServerRequest;
use codex_app_server_protocol::ToolRequestUserInputResponse;
use codex_protocol::mcp::RequestId as McpRequestId;
use codex_protocol::protocol::ReviewDecision;
impl App {
pub(super) async fn reject_app_server_request(
&self,
app_server_client: &AppServerSession,
request_id: AppServerRequestId,
reason: String,
) -> std::result::Result<(), String> {
app_server_client
.reject_server_request(
request_id,
JSONRPCErrorError {
code: -32000,
message: reason,
data: None,
},
)
.await
.map_err(|err| format!("failed to reject app-server request: {err}"))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct AppServerRequestResolution {
@@ -44,7 +62,7 @@ pub(crate) enum ResolvedAppServerRequest {
},
McpElicitation {
server_name: String,
request_id: McpRequestId,
request_id: AppServerRequestId,
},
}
@@ -54,7 +72,7 @@ pub(super) struct PendingAppServerRequests {
file_change_approvals: HashMap<String, AppServerRequestId>,
permissions_approvals: HashMap<String, AppServerRequestId>,
user_inputs: HashMap<String, VecDeque<PendingUserInputRequest>>,
mcp_requests: HashMap<McpLegacyRequestKey, AppServerRequestId>,
mcp_requests: HashMap<McpRequestKey, AppServerRequestId>,
}
impl PendingAppServerRequests {
@@ -101,9 +119,9 @@ impl PendingAppServerRequests {
}
ServerRequest::McpServerElicitationRequest { request_id, params } => {
self.mcp_requests.insert(
McpLegacyRequestKey {
McpRequestKey {
server_name: params.server_name.clone(),
request_id: app_server_request_id_to_mcp_request_id(request_id),
request_id: request_id.clone(),
},
request_id.clone(),
);
@@ -141,30 +159,32 @@ impl PendingAppServerRequests {
T: Into<AppCommand>,
{
let op: AppCommand = op.into();
let resolution = match op.view() {
AppCommandView::ExecApproval { id, decision, .. } => self
let resolution = match &op {
AppCommand::ExecApproval { id, decision, .. } => self
.exec_approvals
.remove(id)
.map(|request_id| {
Ok::<AppServerRequestResolution, String>(AppServerRequestResolution {
request_id,
result: serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: decision.clone().into(),
decision: decision.clone(),
})
.map_err(|err| {
format!("failed to serialize command execution approval response: {err}")
format!(
"failed to serialize command execution approval response: {err}"
)
})?,
})
})
.transpose()?,
AppCommandView::PatchApproval { id, decision } => self
AppCommand::PatchApproval { id, decision } => self
.file_change_approvals
.remove(id)
.map(|request_id| {
Ok::<AppServerRequestResolution, String>(AppServerRequestResolution {
request_id,
result: serde_json::to_value(FileChangeRequestApprovalResponse {
decision: file_change_decision(decision)?,
decision: decision.clone(),
})
.map_err(|err| {
format!("failed to serialize file change approval response: {err}")
@@ -172,7 +192,7 @@ impl PendingAppServerRequests {
})
})
.transpose()?,
AppCommandView::RequestPermissionsResponse { id, response } => self
AppCommand::RequestPermissionsResponse { id, response } => self
.permissions_approvals
.remove(id)
.map(|request_id| {
@@ -191,30 +211,18 @@ impl PendingAppServerRequests {
})
})
.transpose()?,
AppCommandView::UserInputAnswer { id, response } => self
AppCommand::UserInputAnswer { id, response } => self
.pop_user_input_request_for_turn(id)
.map(|pending| {
Ok::<AppServerRequestResolution, String>(AppServerRequestResolution {
request_id: pending.request_id,
result: serde_json::to_value(
serde_json::from_value::<ToolRequestUserInputResponse>(
serde_json::to_value(response).map_err(|err| {
format!("failed to encode request_user_input response: {err}")
})?,
)
.map_err(|err| {
format!(
"failed to decode request_user_input response for app-server: {err}"
)
})?,
)
.map_err(|err| {
result: serde_json::to_value(response).map_err(|err| {
format!("failed to serialize request_user_input response: {err}")
})?,
})
})
.transpose()?,
AppCommandView::ResolveElicitation {
AppCommand::ResolveElicitation {
server_name,
request_id,
decision,
@@ -222,7 +230,7 @@ impl PendingAppServerRequests {
meta,
} => self
.mcp_requests
.remove(&McpLegacyRequestKey {
.remove(&McpRequestKey {
server_name: server_name.to_string(),
request_id: request_id.clone(),
})
@@ -230,17 +238,7 @@ impl PendingAppServerRequests {
Ok::<AppServerRequestResolution, String>(AppServerRequestResolution {
request_id,
result: serde_json::to_value(McpServerElicitationRequestResponse {
action: match decision {
codex_protocol::approvals::ElicitationAction::Accept => {
McpServerElicitationAction::Accept
}
codex_protocol::approvals::ElicitationAction::Decline => {
McpServerElicitationAction::Decline
}
codex_protocol::approvals::ElicitationAction::Cancel => {
McpServerElicitationAction::Cancel
}
},
action: *decision,
content: content.clone(),
meta: meta.clone(),
})
@@ -383,44 +381,25 @@ struct PendingUserInputRequest {
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct McpLegacyRequestKey {
struct McpRequestKey {
server_name: String,
request_id: McpRequestId,
}
fn app_server_request_id_to_mcp_request_id(request_id: &AppServerRequestId) -> McpRequestId {
match request_id {
AppServerRequestId::String(value) => McpRequestId::String(value.clone()),
AppServerRequestId::Integer(value) => McpRequestId::Integer(*value),
}
}
fn file_change_decision(decision: &ReviewDecision) -> Result<FileChangeApprovalDecision, String> {
match decision {
ReviewDecision::Approved => Ok(FileChangeApprovalDecision::Accept),
ReviewDecision::ApprovedForSession => Ok(FileChangeApprovalDecision::AcceptForSession),
ReviewDecision::Denied => Ok(FileChangeApprovalDecision::Decline),
ReviewDecision::TimedOut => Ok(FileChangeApprovalDecision::Decline),
ReviewDecision::Abort => Ok(FileChangeApprovalDecision::Cancel),
ReviewDecision::ApprovedExecpolicyAmendment { .. } => {
Err("execpolicy amendment is not a valid file change approval decision".to_string())
}
ReviewDecision::NetworkPolicyAmendment { .. } => {
Err("network policy amendment is not a valid file change approval decision".to_string())
}
}
request_id: AppServerRequestId,
}
#[cfg(test)]
mod tests {
use super::PendingAppServerRequests;
use super::ResolvedAppServerRequest;
use crate::app_command::AppCommand as Op;
use codex_app_server_protocol::AdditionalFileSystemPermissions;
use codex_app_server_protocol::AdditionalNetworkPermissions;
use codex_app_server_protocol::CommandExecutionApprovalDecision;
use codex_app_server_protocol::CommandExecutionRequestApprovalParams;
use codex_app_server_protocol::FileChangeApprovalDecision;
use codex_app_server_protocol::FileChangeRequestApprovalParams;
use codex_app_server_protocol::McpElicitationObjectType;
use codex_app_server_protocol::McpElicitationSchema;
use codex_app_server_protocol::McpServerElicitationAction;
use codex_app_server_protocol::McpServerElicitationRequest;
use codex_app_server_protocol::McpServerElicitationRequestParams;
use codex_app_server_protocol::PermissionGrantScope;
@@ -431,13 +410,8 @@ mod tests {
use codex_app_server_protocol::ToolRequestUserInputAnswer;
use codex_app_server_protocol::ToolRequestUserInputParams;
use codex_app_server_protocol::ToolRequestUserInputResponse;
use codex_protocol::approvals::ElicitationAction;
use codex_protocol::approvals::ExecPolicyAmendment;
use codex_protocol::mcp::RequestId as McpRequestId;
use codex_protocol::models::FileSystemPermissions;
use codex_protocol::models::NetworkPermissions;
use codex_protocol::protocol::Op;
use codex_protocol::protocol::ReviewDecision;
use codex_protocol::request_permissions::RequestPermissionProfile;
use codex_utils_absolute_path::AbsolutePathBuf;
use pretty_assertions::assert_eq;
@@ -474,7 +448,7 @@ mod tests {
.take_resolution(&Op::ExecApproval {
id: "approval-1".to_string(),
turn_id: None,
decision: ReviewDecision::Approved,
decision: CommandExecutionApprovalDecision::Accept,
})
.expect("resolution should serialize")
.expect("request should be pending");
@@ -586,10 +560,10 @@ mod tests {
let user_input = pending
.take_resolution(&Op::UserInputAnswer {
id: "turn-2".to_string(),
response: codex_protocol::request_user_input::RequestUserInputResponse {
response: ToolRequestUserInputResponse {
answers: std::iter::once((
"question".to_string(),
codex_protocol::request_user_input::RequestUserInputAnswer {
ToolRequestUserInputAnswer {
answers: vec!["yes".to_string()],
},
))
@@ -643,8 +617,8 @@ mod tests {
let resolution = pending
.take_resolution(&Op::ResolveElicitation {
server_name: "example".to_string(),
request_id: McpRequestId::Integer(12),
decision: ElicitationAction::Accept,
request_id: AppServerRequestId::Integer(12),
decision: McpServerElicitationAction::Accept,
content: Some(json!({ "answer": "yes" })),
meta: Some(json!({ "source": "tui" })),
})
@@ -703,7 +677,7 @@ mod tests {
}
#[test]
fn rejects_invalid_patch_decisions_for_file_change_requests() {
fn resolves_patch_approval_through_app_server_request_id() {
let mut pending = PendingAppServerRequests::default();
assert_eq!(
pending.note_server_request(&ServerRequest::FileChangeRequestApproval {
@@ -719,22 +693,16 @@ mod tests {
None
);
let error = pending
let resolution = pending
.take_resolution(&Op::PatchApproval {
id: "patch-1".to_string(),
decision: ReviewDecision::ApprovedExecpolicyAmendment {
proposed_execpolicy_amendment: ExecPolicyAmendment::new(vec![
"echo".to_string(),
"hi".to_string(),
]),
},
decision: FileChangeApprovalDecision::Cancel,
})
.expect_err("invalid patch decision should fail");
.expect("resolution should serialize")
.expect("request should be pending");
assert_eq!(
error,
"execpolicy amendment is not a valid file change approval decision"
);
assert_eq!(resolution.request_id, AppServerRequestId::Integer(13));
assert_eq!(resolution.result, json!({ "decision": "cancel" }));
}
#[test]
@@ -803,7 +771,7 @@ mod tests {
pending.resolve_notification(&AppServerRequestId::Integer(12)),
Some(ResolvedAppServerRequest::McpElicitation {
server_name: "example".to_string(),
request_id: McpRequestId::Integer(12),
request_id: AppServerRequestId::Integer(12),
})
);
}
@@ -844,7 +812,7 @@ mod tests {
});
}
let response = codex_protocol::request_user_input::RequestUserInputResponse {
let response = ToolRequestUserInputResponse {
answers: HashMap::new(),
};
let first_response = pending
+2 -1
View File
@@ -9,6 +9,7 @@ use codex_app_server_protocol::MarketplaceAddParams;
use codex_app_server_protocol::MarketplaceAddResponse;
use codex_app_server_protocol::MarketplaceRemoveParams;
use codex_app_server_protocol::MarketplaceRemoveResponse;
use codex_app_server_protocol::RequestId;
use codex_utils_absolute_path::AbsolutePathBuf;
impl App {
@@ -483,7 +484,7 @@ pub(super) async fn fetch_account_rate_limits(
.await
.wrap_err("account/rateLimits/read failed in TUI")?;
Ok(app_server_rate_limit_snapshots_to_core(response))
Ok(app_server_rate_limit_snapshots(response))
}
pub(super) async fn send_add_credits_nudge_email(
+13 -16
View File
@@ -48,7 +48,7 @@ impl App {
match self.rebuild_config_for_cwd(resume_cwd.clone()).await {
Ok(config) => Ok(config),
Err(err) => {
if crate::cwds_differ(current_cwd, &resume_cwd) {
if crate::session_resume::cwds_differ(current_cwd, &resume_cwd) {
Err(err)
} else {
let resume_cwd_display = resume_cwd.display().to_string();
@@ -65,7 +65,7 @@ impl App {
pub(super) fn apply_runtime_policy_overrides(&mut self, config: &mut Config) {
if let Some(policy) = self.runtime_approval_policy_override.as_ref()
&& let Err(err) = config.permissions.approval_policy.set(*policy)
&& let Err(err) = config.permissions.approval_policy.set(policy.to_core())
{
tracing::warn!(%err, "failed to carry forward approval policy override");
self.chat_widget.add_error_message(format!(
@@ -94,7 +94,7 @@ impl App {
user_message_prefix: &str,
log_message: &str,
) -> bool {
if let Err(err) = config.permissions.approval_policy.set(policy) {
if let Err(err) = config.permissions.approval_policy.set(policy.to_core()) {
tracing::warn!(error = %err, "{log_message}");
self.chat_widget
.add_error_message(format!("{user_message_prefix}: {err}"));
@@ -294,8 +294,9 @@ impl App {
self.set_approvals_reviewer_in_app_and_widget(self.config.approvals_reviewer);
}
if approval_policy_override.is_some() {
self.chat_widget
.set_approval_policy(self.config.permissions.approval_policy.value());
self.chat_widget.set_approval_policy(AskForApproval::from(
self.config.permissions.approval_policy.value(),
));
}
if permission_profile_override.is_some()
&& let Err(err) = self
@@ -493,7 +494,7 @@ impl App {
(!model.starts_with("codex-auto-")).then(|| Self::reasoning_label(reasoning_effort))
}
pub(crate) fn token_usage(&self) -> codex_protocol::protocol::TokenUsage {
pub(crate) fn token_usage(&self) -> crate::token_usage::TokenUsage {
self.chat_widget.token_usage()
}
@@ -548,9 +549,6 @@ mod tests {
use crate::app::test_support::make_test_app;
use crate::test_support::PathBufExt;
use codex_protocol::models::PermissionProfile;
use codex_protocol::protocol::Event;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::SessionConfiguredEvent;
use pretty_assertions::assert_eq;
use tempfile::tempdir;
@@ -634,11 +632,11 @@ mod tests {
let next_cwd_tmp = tempdir()?;
let next_cwd = next_cwd_tmp.path().to_path_buf();
app.chat_widget.handle_codex_event(Event {
id: String::new(),
msg: EventMsg::SessionConfigured(SessionConfiguredEvent {
session_id: ThreadId::new(),
app.chat_widget
.handle_thread_session(crate::session_state::ThreadSessionState {
thread_id: ThreadId::new(),
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -648,14 +646,13 @@ mod tests {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: next_cwd.clone().abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
}),
});
});
assert_eq!(app.chat_widget.config_ref().cwd.to_path_buf(), next_cwd);
assert_eq!(app.config.cwd, original_cwd);
+8 -7
View File
@@ -1055,7 +1055,7 @@ impl App {
self.app_event_tx.send(AppEvent::CodexOp(
AppCommand::override_turn_context(
/*cwd*/ None,
Some(preset.approval),
Some(AskForApproval::from(preset.approval)),
Some(self.config.approvals_reviewer),
Some(preset.permission_profile.clone()),
#[cfg(target_os = "windows")]
@@ -1068,8 +1068,9 @@ impl App {
/*personality*/ None,
),
));
self.app_event_tx
.send(AppEvent::UpdateAskForApprovalPolicy(preset.approval));
self.app_event_tx.send(AppEvent::UpdateAskForApprovalPolicy(
AskForApproval::from(preset.approval),
));
self.app_event_tx.send(AppEvent::UpdatePermissionProfile(
preset.permission_profile.clone(),
));
@@ -1317,10 +1318,10 @@ impl App {
return Ok(AppRunControl::Continue);
}
self.config = config;
self.runtime_approval_policy_override =
Some(self.config.permissions.approval_policy.value());
self.chat_widget
.set_approval_policy(self.config.permissions.approval_policy.value());
let approval_policy =
AskForApproval::from(self.config.permissions.approval_policy.value());
self.runtime_approval_policy_override = Some(approval_policy);
self.chat_widget.set_approval_policy(approval_policy);
self.sync_active_thread_permission_settings_to_cached_session()
.await;
}
+36 -29
View File
@@ -14,10 +14,8 @@
//! `SessionSource::SubAgent(ThreadSpawn { parent_thread_id, .. })` edges until no new children are
//! found. The primary thread itself is never included in the output.
use codex_app_server_protocol::SessionSource;
use codex_app_server_protocol::Thread;
use codex_protocol::ThreadId;
use codex_protocol::protocol::SubAgentSource;
use std::collections::HashMap;
use std::collections::HashSet;
@@ -63,15 +61,12 @@ pub(crate) fn find_loaded_subagent_threads_for_primary(
continue;
}
let SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id: source_parent_thread_id,
..
}) = &thread.source
let Some(source_parent_thread_id) = thread_spawn_parent_thread_id(&thread.source)
else {
continue;
};
if *source_parent_thread_id != parent_thread_id {
if source_parent_thread_id != parent_thread_id {
continue;
}
@@ -96,6 +91,18 @@ pub(crate) fn find_loaded_subagent_threads_for_primary(
loaded_threads
}
fn thread_spawn_parent_thread_id(
source: &codex_app_server_protocol::SessionSource,
) -> Option<ThreadId> {
let value = serde_json::to_value(source).ok()?;
let parent_thread_id = value
.get("subAgent")?
.get("thread_spawn")?
.get("parent_thread_id")?
.as_str()?;
ThreadId::from_string(parent_thread_id).ok()
}
#[cfg(test)]
mod tests {
use super::LoadedSubagentThread;
@@ -104,7 +111,6 @@ mod tests {
use codex_app_server_protocol::Thread;
use codex_app_server_protocol::ThreadStatus;
use codex_protocol::ThreadId;
use codex_protocol::protocol::SubAgentSource;
use codex_utils_absolute_path::test_support::PathBufExt;
use codex_utils_absolute_path::test_support::test_path_buf;
use pretty_assertions::assert_eq;
@@ -131,6 +137,25 @@ mod tests {
}
}
fn thread_spawn_source(
parent_thread_id: ThreadId,
depth: i32,
agent_nickname: &str,
agent_role: &str,
) -> SessionSource {
serde_json::from_value(serde_json::json!({
"subAgent": {
"thread_spawn": {
"parent_thread_id": parent_thread_id.to_string(),
"depth": depth,
"agent_nickname": agent_nickname,
"agent_role": agent_role,
}
}
}))
.expect("valid subagent source")
}
#[test]
fn finds_loaded_subagent_tree_for_primary_thread() {
let primary_thread_id =
@@ -146,39 +171,21 @@ mod tests {
let mut child = test_thread(
child_thread_id,
SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id: primary_thread_id,
depth: 1,
agent_path: None,
agent_nickname: Some("Scout".to_string()),
agent_role: Some("explorer".to_string()),
}),
thread_spawn_source(primary_thread_id, /*depth*/ 1, "Scout", "explorer"),
);
child.agent_nickname = Some("Scout".to_string());
child.agent_role = Some("explorer".to_string());
let mut grandchild = test_thread(
grandchild_thread_id,
SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id: child_thread_id,
depth: 2,
agent_path: None,
agent_nickname: Some("Atlas".to_string()),
agent_role: Some("worker".to_string()),
}),
thread_spawn_source(child_thread_id, /*depth*/ 2, "Atlas", "worker"),
);
grandchild.agent_nickname = Some("Atlas".to_string());
grandchild.agent_role = Some("worker".to_string());
let unrelated_child = test_thread(
unrelated_child_id,
SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id: unrelated_parent_id,
depth: 1,
agent_path: None,
agent_nickname: Some("Other".to_string()),
agent_role: Some("researcher".to_string()),
}),
thread_spawn_source(unrelated_parent_id, /*depth*/ 1, "Other", "researcher"),
);
let loaded = find_loaded_subagent_threads_for_primary(
@@ -1,5 +1,4 @@
use crate::app_command::AppCommand;
use crate::app_command::AppCommandView;
use codex_app_server_protocol::RequestId as AppServerRequestId;
use codex_app_server_protocol::ServerNotification;
use codex_app_server_protocol::ServerRequest;
@@ -10,11 +9,11 @@ use std::collections::HashSet;
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct ElicitationRequestKey {
server_name: String,
request_id: codex_protocol::mcp::RequestId,
request_id: AppServerRequestId,
}
impl ElicitationRequestKey {
fn new(server_name: String, request_id: codex_protocol::mcp::RequestId) -> Self {
fn new(server_name: String, request_id: AppServerRequestId) -> Self {
Self {
server_name,
request_id,
@@ -33,9 +32,9 @@ impl ElicitationRequestKey {
// - buffer eviction (`note_evicted_event`)
//
// We keep both fast lookup sets (for snapshot filtering by call_id/request key) and
// turn-indexed queues/vectors so `TurnComplete`/`TurnAborted` can clear stale prompts tied
// to a turn. `request_user_input` removal is FIFO because the overlay answers queued prompts
// in FIFO order for a shared `turn_id`.
// turn-indexed queues/vectors so turn completion or interruption can clear
// stale prompts tied to a turn. `request_user_input` removal is FIFO because
// the overlay answers queued prompts in FIFO order for a shared `turn_id`.
pub(super) struct PendingInteractiveReplayState {
exec_approval_call_ids: HashSet<String>,
exec_approval_call_ids_by_turn_id: HashMap<String, Vec<String>>,
@@ -77,13 +76,13 @@ impl PendingInteractiveReplayState {
{
let op: AppCommand = op.into();
matches!(
op.view(),
AppCommandView::ExecApproval { .. }
| AppCommandView::PatchApproval { .. }
| AppCommandView::ResolveElicitation { .. }
| AppCommandView::RequestPermissionsResponse { .. }
| AppCommandView::UserInputAnswer { .. }
| AppCommandView::Shutdown
&op,
AppCommand::ExecApproval { .. }
| AppCommand::PatchApproval { .. }
| AppCommand::ResolveElicitation { .. }
| AppCommand::RequestPermissionsResponse { .. }
| AppCommand::UserInputAnswer { .. }
| AppCommand::Shutdown
)
}
@@ -92,8 +91,8 @@ impl PendingInteractiveReplayState {
T: Into<AppCommand>,
{
let op: AppCommand = op.into();
match op.view() {
AppCommandView::ExecApproval { id, turn_id, .. } => {
match &op {
AppCommand::ExecApproval { id, turn_id, .. } => {
self.exec_approval_call_ids.remove(id);
if let Some(turn_id) = turn_id {
Self::remove_call_id_from_turn_map_entry(
@@ -105,7 +104,7 @@ impl PendingInteractiveReplayState {
self.pending_requests_by_request_id
.retain(|_, pending| !matches!(pending, PendingInteractiveRequest::ExecApproval { approval_id, .. } if approval_id == id));
}
AppCommandView::PatchApproval { id, .. } => {
AppCommand::PatchApproval { id, .. } => {
self.patch_approval_call_ids.remove(id);
Self::remove_call_id_from_turn_map(
&mut self.patch_approval_call_ids_by_turn_id,
@@ -114,7 +113,7 @@ impl PendingInteractiveReplayState {
self.pending_requests_by_request_id
.retain(|_, pending| !matches!(pending, PendingInteractiveRequest::PatchApproval { item_id, .. } if item_id == id));
}
AppCommandView::ResolveElicitation {
AppCommand::ResolveElicitation {
server_name,
request_id,
..
@@ -130,7 +129,7 @@ impl PendingInteractiveReplayState {
},
);
}
AppCommandView::RequestPermissionsResponse { id, .. } => {
AppCommand::RequestPermissionsResponse { id, .. } => {
self.request_permissions_call_ids.remove(id);
Self::remove_call_id_from_turn_map(
&mut self.request_permissions_call_ids_by_turn_id,
@@ -145,7 +144,7 @@ impl PendingInteractiveReplayState {
// `Op::UserInputAnswer` identifies the turn, not the prompt call_id. The UI
// answers queued prompts for the same turn in FIFO order, so remove the oldest
// queued call_id for that turn.
AppCommandView::UserInputAnswer { id, .. } => {
AppCommand::UserInputAnswer { id, .. } => {
let mut remove_turn_entry = false;
if let Some(call_ids) = self.request_user_input_call_ids_by_turn_id.get_mut(id) {
if !call_ids.is_empty() {
@@ -165,7 +164,7 @@ impl PendingInteractiveReplayState {
self.request_user_input_call_ids_by_turn_id.remove(id);
}
}
AppCommandView::Shutdown => self.clear(),
AppCommand::Shutdown => self.clear(),
_ => {}
}
}
@@ -208,10 +207,8 @@ impl PendingInteractiveReplayState {
);
}
ServerRequest::McpServerElicitationRequest { request_id, params } => {
let key = ElicitationRequestKey::new(
params.server_name.clone(),
app_server_request_id_to_mcp_request_id(request_id),
);
let key =
ElicitationRequestKey::new(params.server_name.clone(), request_id.clone());
self.elicitation_requests.insert(key.clone());
self.pending_requests_by_request_id.insert(
request_id.clone(),
@@ -311,7 +308,7 @@ impl PendingInteractiveReplayState {
self.elicitation_requests
.remove(&ElicitationRequestKey::new(
params.server_name.clone(),
app_server_request_id_to_mcp_request_id(request_id),
request_id.clone(),
));
}
ServerRequest::ToolRequestUserInput { params, .. } => {
@@ -366,7 +363,7 @@ impl PendingInteractiveReplayState {
.elicitation_requests
.contains(&ElicitationRequestKey::new(
params.server_name.clone(),
app_server_request_id_to_mcp_request_id(request_id),
request_id.clone(),
)),
ServerRequest::ToolRequestUserInput { params, .. } => {
self.request_user_input_call_ids.contains(&params.item_id)
@@ -549,10 +546,7 @@ impl PendingInteractiveReplayState {
(
PendingInteractiveRequest::Elicitation(key),
ServerRequest::McpServerElicitationRequest { request_id, params },
) => {
key.server_name == params.server_name
&& key.request_id == app_server_request_id_to_mcp_request_id(request_id)
}
) => key.server_name == params.server_name && key.request_id == *request_id,
(
PendingInteractiveRequest::RequestPermissions { turn_id, item_id },
ServerRequest::PermissionsRequestApproval { params, .. },
@@ -566,23 +560,17 @@ impl PendingInteractiveReplayState {
}
}
fn app_server_request_id_to_mcp_request_id(
request_id: &AppServerRequestId,
) -> codex_protocol::mcp::RequestId {
match request_id {
AppServerRequestId::String(value) => codex_protocol::mcp::RequestId::String(value.clone()),
AppServerRequestId::Integer(value) => codex_protocol::mcp::RequestId::Integer(*value),
}
}
#[cfg(test)]
mod tests {
use super::super::ThreadBufferedEvent;
use super::super::ThreadEventStore;
use crate::app_command::AppCommand as Op;
use codex_app_server_protocol::CommandExecutionApprovalDecision;
use codex_app_server_protocol::CommandExecutionRequestApprovalParams;
use codex_app_server_protocol::FileChangeRequestApprovalParams;
use codex_app_server_protocol::McpElicitationObjectType;
use codex_app_server_protocol::McpElicitationSchema;
use codex_app_server_protocol::McpServerElicitationAction;
use codex_app_server_protocol::McpServerElicitationRequest;
use codex_app_server_protocol::McpServerElicitationRequestParams;
use codex_app_server_protocol::RequestId as AppServerRequestId;
@@ -591,11 +579,10 @@ mod tests {
use codex_app_server_protocol::ServerRequestResolvedNotification;
use codex_app_server_protocol::ThreadClosedNotification;
use codex_app_server_protocol::ToolRequestUserInputParams;
use codex_app_server_protocol::ToolRequestUserInputResponse;
use codex_app_server_protocol::Turn;
use codex_app_server_protocol::TurnCompletedNotification;
use codex_app_server_protocol::TurnStatus;
use codex_protocol::protocol::Op;
use codex_protocol::protocol::ReviewDecision;
use codex_utils_absolute_path::test_support::PathBufExt;
use codex_utils_absolute_path::test_support::test_path_buf;
use pretty_assertions::assert_eq;
@@ -724,7 +711,7 @@ mod tests {
store.note_outbound_op(&Op::UserInputAnswer {
id: "turn-1".to_string(),
response: codex_protocol::request_user_input::RequestUserInputResponse {
response: ToolRequestUserInputResponse {
answers: HashMap::new(),
},
});
@@ -767,7 +754,7 @@ mod tests {
store.note_outbound_op(&Op::ExecApproval {
id: "approval-1".to_string(),
turn_id: Some("turn-1".to_string()),
decision: ReviewDecision::Approved,
decision: CommandExecutionApprovalDecision::Accept,
});
let snapshot = store.snapshot();
@@ -809,7 +796,7 @@ mod tests {
store.note_outbound_op(&Op::UserInputAnswer {
id: "turn-1".to_string(),
response: codex_protocol::request_user_input::RequestUserInputResponse {
response: ToolRequestUserInputResponse {
answers: HashMap::new(),
},
});
@@ -833,7 +820,7 @@ mod tests {
store.note_outbound_op(&Op::UserInputAnswer {
id: "turn-1".to_string(),
response: codex_protocol::request_user_input::RequestUserInputResponse {
response: ToolRequestUserInputResponse {
answers: HashMap::new(),
},
});
@@ -854,7 +841,7 @@ mod tests {
store.note_outbound_op(&Op::PatchApproval {
id: "call-1".to_string(),
decision: ReviewDecision::Approved,
decision: codex_app_server_protocol::FileChangeApprovalDecision::Accept,
});
let snapshot = store.snapshot();
@@ -888,13 +875,13 @@ mod tests {
#[test]
fn thread_event_snapshot_drops_resolved_elicitation_after_outbound_resolution() {
let mut store = ThreadEventStore::new(/*capacity*/ 8);
let request_id = codex_protocol::mcp::RequestId::String("request-1".to_string());
let request_id = AppServerRequestId::String("request-1".to_string());
store.push_request(elicitation_request("server-1", "request-1", "turn-1"));
store.note_outbound_op(&Op::ResolveElicitation {
server_name: "server-1".to_string(),
request_id,
decision: codex_protocol::approvals::ElicitationAction::Accept,
decision: McpServerElicitationAction::Accept,
content: None,
meta: None,
});
@@ -920,7 +907,7 @@ mod tests {
store.note_outbound_op(&Op::ExecApproval {
id: "call-1".to_string(),
turn_id: Some("turn-1".to_string()),
decision: ReviewDecision::Approved,
decision: CommandExecutionApprovalDecision::Accept,
});
assert_eq!(store.has_pending_thread_approvals(), false);
+4 -4
View File
@@ -636,7 +636,7 @@ impl App {
let resume_cwd = if self.remote_app_server_url.is_some() {
current_cwd.clone()
} else {
match crate::resolve_cwd_for_resume_or_fork(
match crate::session_resume::resolve_cwd_for_resume_or_fork(
tui,
&self.config,
&current_cwd,
@@ -647,9 +647,9 @@ impl App {
)
.await?
{
crate::ResolveCwdOutcome::Continue(Some(cwd)) => cwd,
crate::ResolveCwdOutcome::Continue(None) => current_cwd.clone(),
crate::ResolveCwdOutcome::Exit => {
crate::session_resume::ResolveCwdOutcome::Continue(Some(cwd)) => cwd,
crate::session_resume::ResolveCwdOutcome::Continue(None) => current_cwd.clone(),
crate::session_resume::ResolveCwdOutcome::Exit => {
return Ok(AppRunControl::Exit(ExitReason::UserRequested));
}
}
+2 -1
View File
@@ -74,7 +74,8 @@ fn test_session_telemetry(config: &Config, model: &str) -> SessionTelemetry {
"test_originator".to_string(),
/*log_user_prompts*/ false,
"test".to_string(),
SessionSource::Cli,
serde_json::from_value(serde_json::json!("cli"))
.expect("cli session source should deserialize"),
)
}
+264 -258
View File
@@ -6,7 +6,6 @@ use super::*;
use crate::app_backtrack::BacktrackSelection;
use crate::app_backtrack::BacktrackState;
use crate::app_backtrack::user_count;
use crate::app_command::AppCommand;
use crate::chatwidget::ChatWidgetInit;
use crate::chatwidget::create_initial_user_message;
@@ -22,6 +21,8 @@ use crate::history_cell::new_session_info;
use crate::multi_agents::AgentPickerThreadEntry;
use assert_matches::assert_matches;
use crate::app_command::AppCommand as Op;
use crate::diff_model::FileChange;
use crate::legacy_core::config::ConfigBuilder;
use crate::legacy_core::config::ConfigOverrides;
use crate::legacy_core::config::TerminalResizeReflowMaxRows;
@@ -29,6 +30,7 @@ use codex_app_server_protocol::AdditionalFileSystemPermissions;
use codex_app_server_protocol::AdditionalNetworkPermissions;
use codex_app_server_protocol::AdditionalPermissionProfile;
use codex_app_server_protocol::AgentMessageDeltaNotification;
use codex_app_server_protocol::AskForApproval;
use codex_app_server_protocol::CommandExecutionRequestApprovalParams;
use codex_app_server_protocol::ConfigWarningNotification;
use codex_app_server_protocol::FileChangeRequestApprovalParams;
@@ -47,6 +49,7 @@ use codex_app_server_protocol::PermissionsRequestApprovalParams;
use codex_app_server_protocol::RequestId as AppServerRequestId;
use codex_app_server_protocol::ServerNotification;
use codex_app_server_protocol::ServerRequest;
use codex_app_server_protocol::SessionSource;
use codex_app_server_protocol::Thread;
use codex_app_server_protocol::ThreadClosedNotification;
use codex_app_server_protocol::ThreadItem;
@@ -60,6 +63,7 @@ use codex_app_server_protocol::TurnCompletedNotification;
use codex_app_server_protocol::TurnError as AppServerTurnError;
use codex_app_server_protocol::TurnStartedNotification;
use codex_app_server_protocol::TurnStatus;
use codex_app_server_protocol::UserInput;
use codex_app_server_protocol::UserInput as AppServerUserInput;
use codex_app_server_protocol::WarningNotification;
use codex_otel::SessionTelemetry;
@@ -68,24 +72,11 @@ use codex_protocol::config_types::CollaborationMode;
use codex_protocol::config_types::CollaborationModeMask;
use codex_protocol::config_types::ModeKind;
use codex_protocol::config_types::Settings;
use codex_protocol::models::AdditionalPermissionProfile as CoreAdditionalPermissionProfile;
use codex_protocol::models::FileSystemPermissions;
use codex_protocol::models::NetworkPermissions;
use codex_protocol::models::PermissionProfile;
use codex_protocol::protocol::AskForApproval;
use codex_protocol::protocol::Event;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::FileChange;
use codex_protocol::protocol::NetworkApprovalContext;
use codex_protocol::protocol::NetworkApprovalProtocol;
use codex_protocol::protocol::RolloutItem;
use codex_protocol::protocol::RolloutLine;
use codex_protocol::protocol::SessionConfiguredEvent;
use codex_protocol::protocol::SessionSource;
use codex_protocol::protocol::TurnContextItem;
use codex_protocol::request_permissions::RequestPermissionProfile;
use codex_protocol::user_input::TextElement;
use codex_protocol::user_input::UserInput;
use codex_utils_absolute_path::AbsolutePathBuf;
use crossterm::event::KeyModifiers;
use insta::assert_snapshot;
@@ -442,12 +433,12 @@ async fn enqueue_primary_thread_session_replays_turns_before_initial_prompt_subm
}
AppEvent::SubmitThreadOp {
thread_id: op_thread_id,
op: AppCommand::UserTurn { items, .. },
op: Op::UserTurn { items, .. },
} => {
assert_eq!(op_thread_id, thread_id);
submitted_items = Some(items);
}
AppEvent::CodexOp(AppCommand::UserTurn { items, .. }) => {
AppEvent::CodexOp(Op::UserTurn { items, .. }) => {
submitted_items = Some(items);
}
_ => {}
@@ -514,8 +505,7 @@ async fn history_lookup_response_is_routed_to_requesting_thread() -> Result<()>
&Op::GetHistoryEntryRequest {
offset: 0,
log_id: 1,
}
.into(),
},
)
.await?;
@@ -1338,32 +1328,43 @@ async fn open_agent_picker_marks_terminal_read_errors_closed() -> Result<()> {
Ok(())
}
#[tokio::test]
async fn open_agent_picker_marks_loaded_threads_open() -> Result<()> {
let mut app = Box::pin(make_test_app()).await;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("embedded app server");
let started = app_server
.start_thread(app.chat_widget.config_ref())
.await?;
let thread_id = started.session.thread_id;
app.thread_event_channels
.insert(thread_id, ThreadEventChannel::new(/*capacity*/ 1));
#[test]
fn open_agent_picker_marks_loaded_threads_open() -> Result<()> {
const WORKER_THREADS: usize = 1;
const TEST_STACK_SIZE_BYTES: usize = 8 * 1024 * 1024;
Box::pin(app.open_agent_picker(&mut app_server)).await;
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(WORKER_THREADS)
.thread_stack_size(TEST_STACK_SIZE_BYTES)
.enable_all()
.build()?;
assert_eq!(
app.agent_navigation.get(&thread_id),
Some(&AgentPickerThreadEntry {
agent_nickname: None,
agent_role: None,
is_closed: false,
})
);
Ok(())
runtime.block_on(async {
let mut app = Box::pin(make_test_app()).await;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("embedded app server");
let started = app_server
.start_thread(app.chat_widget.config_ref())
.await?;
let thread_id = started.session.thread_id;
app.thread_event_channels
.insert(thread_id, ThreadEventChannel::new(/*capacity*/ 1));
Box::pin(app.open_agent_picker(&mut app_server)).await;
assert_eq!(
app.agent_navigation.get(&thread_id),
Some(&AgentPickerThreadEntry {
agent_nickname: None,
agent_role: None,
is_closed: false,
})
);
Ok(())
})
}
#[test]
@@ -1575,68 +1576,84 @@ async fn update_memory_settings_persists_and_updates_widget_config() -> Result<(
Ok(())
}
#[tokio::test]
async fn update_memory_settings_updates_current_thread_memory_mode() -> Result<()> {
let (mut app, _app_event_rx, _op_rx) = Box::pin(make_test_app_with_channels()).await;
let codex_home = tempdir()?;
app.config.codex_home = codex_home.path().to_path_buf().abs();
// Seed the previous setting so this test exercises the thread-mode update path.
app.config.memories.generate_memories = true;
#[test]
fn update_memory_settings_updates_current_thread_memory_mode() -> Result<()> {
const WORKER_THREADS: usize = 1;
const TEST_STACK_SIZE_BYTES: usize = 8 * 1024 * 1024;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(&app.config)).await?;
let started = app_server.start_thread(&app.config).await?;
let thread_id = started.session.thread_id;
app.active_thread_id = Some(thread_id);
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(WORKER_THREADS)
.thread_stack_size(TEST_STACK_SIZE_BYTES)
.enable_all()
.build()?;
Box::pin(app.update_memory_settings_with_app_server(
&mut app_server,
/*use_memories*/ true,
/*generate_memories*/ false,
))
.await;
runtime.block_on(async {
let (mut app, _app_event_rx, _op_rx) = Box::pin(make_test_app_with_channels()).await;
let codex_home = tempdir()?;
app.config.codex_home = codex_home.path().to_path_buf().abs();
// Seed the previous setting so this test exercises the thread-mode update path.
app.config.memories.generate_memories = true;
let state_db = codex_state::StateRuntime::init(
codex_home.path().to_path_buf(),
app.config.model_provider_id.clone(),
)
.await
.expect("state db should initialize");
let memory_mode = state_db
.get_thread_memory_mode(thread_id)
let mut app_server =
Box::pin(crate::start_embedded_app_server_for_picker(&app.config)).await?;
let started = app_server.start_thread(&app.config).await?;
let thread_id = started.session.thread_id;
app.active_thread_id = Some(thread_id);
Box::pin(app.update_memory_settings_with_app_server(
&mut app_server,
/*use_memories*/ true,
/*generate_memories*/ false,
))
.await;
let state_db = codex_state::StateRuntime::init(
codex_home.path().to_path_buf(),
app.config.model_provider_id.clone(),
)
.await
.expect("thread memory mode should be readable");
assert_eq!(memory_mode.as_deref(), Some("disabled"));
.expect("state db should initialize");
let memory_mode = state_db
.get_thread_memory_mode(thread_id)
.await
.expect("thread memory mode should be readable");
assert_eq!(memory_mode.as_deref(), Some("disabled"));
app_server.shutdown().await?;
Ok(())
app_server.shutdown().await?;
Ok(())
})
}
#[tokio::test]
async fn reset_memories_clears_local_memory_directories() -> Result<()> {
let (mut app, _app_event_rx, _op_rx) = Box::pin(make_test_app_with_channels()).await;
let codex_home = tempdir()?;
app.config.codex_home = codex_home.path().to_path_buf().abs();
app.config.sqlite_home = codex_home.path().to_path_buf();
Box::pin(async {
let (mut app, _app_event_rx, _op_rx) = Box::pin(make_test_app_with_channels()).await;
let codex_home = tempdir()?;
app.config.codex_home = codex_home.path().to_path_buf().abs();
app.config.sqlite_home = codex_home.path().to_path_buf();
let memory_root = codex_home.path().join("memories");
let extensions_root = memory_root.join("extensions");
std::fs::create_dir_all(memory_root.join("rollout_summaries"))?;
std::fs::create_dir_all(&extensions_root)?;
std::fs::write(memory_root.join("MEMORY.md"), "stale memory\n")?;
std::fs::write(
memory_root.join("rollout_summaries").join("stale.md"),
"stale summary\n",
)?;
std::fs::write(extensions_root.join("stale.txt"), "stale extension\n")?;
let memory_root = codex_home.path().join("memories");
let extensions_root = memory_root.join("extensions");
std::fs::create_dir_all(memory_root.join("rollout_summaries"))?;
std::fs::create_dir_all(&extensions_root)?;
std::fs::write(memory_root.join("MEMORY.md"), "stale memory\n")?;
std::fs::write(
memory_root.join("rollout_summaries").join("stale.md"),
"stale summary\n",
)?;
std::fs::write(extensions_root.join("stale.txt"), "stale extension\n")?;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(&app.config)).await?;
let mut app_server =
Box::pin(crate::start_embedded_app_server_for_picker(&app.config)).await?;
Box::pin(app.reset_memories_with_app_server(&mut app_server)).await;
Box::pin(app.reset_memories_with_app_server(&mut app_server)).await;
assert_eq!(std::fs::read_dir(&memory_root)?.count(), 0);
assert_eq!(std::fs::read_dir(&memory_root)?.count(), 0);
app_server.shutdown().await?;
Ok(())
app_server.shutdown().await?;
Ok(())
})
.await
}
#[tokio::test]
@@ -1661,15 +1678,17 @@ async fn update_feature_flags_enabling_guardian_selects_auto_review() -> Result<
auto_review.approvals_reviewer
);
assert_eq!(
app.config.permissions.approval_policy.value(),
AskForApproval::from(app.config.permissions.approval_policy.value()),
auto_review.approval_policy
);
assert_eq!(
app.chat_widget
.config_ref()
.permissions
.approval_policy
.value(),
AskForApproval::from(
app.chat_widget
.config_ref()
.permissions
.approval_policy
.value(),
),
auto_review.approval_policy
);
assert_eq!(
@@ -1694,7 +1713,6 @@ async fn update_feature_flags_enabling_guardian_selects_auto_review() -> Result<
cwd: None,
approval_policy: Some(auto_review.approval_policy),
approvals_reviewer: Some(auto_review.approvals_reviewer),
sandbox_policy: None,
permission_profile: Some(auto_review.permission_profile.clone()),
windows_sandbox_level: None,
model: None,
@@ -1750,7 +1768,7 @@ async fn update_feature_flags_disabling_guardian_clears_review_policy_and_restor
app.config
.permissions
.approval_policy
.set(AskForApproval::OnRequest)?;
.set(AskForApproval::OnRequest.to_core())?;
app.config
.permissions
.set_permission_profile(PermissionProfile::workspace_write())?;
@@ -1771,7 +1789,7 @@ async fn update_feature_flags_disabling_guardian_clears_review_policy_and_restor
);
assert_eq!(app.config.approvals_reviewer, ApprovalsReviewer::User);
assert_eq!(
app.config.permissions.approval_policy.value(),
AskForApproval::from(app.config.permissions.approval_policy.value()),
AskForApproval::OnRequest
);
assert_eq!(
@@ -1785,7 +1803,6 @@ async fn update_feature_flags_disabling_guardian_clears_review_policy_and_restor
cwd: None,
approval_policy: None,
approvals_reviewer: Some(ApprovalsReviewer::User),
sandbox_policy: None,
permission_profile: None,
windows_sandbox_level: None,
model: None,
@@ -1848,7 +1865,7 @@ async fn update_feature_flags_enabling_guardian_overrides_explicit_manual_review
auto_review.approvals_reviewer
);
assert_eq!(
app.config.permissions.approval_policy.value(),
AskForApproval::from(app.config.permissions.approval_policy.value()),
auto_review.approval_policy
);
assert_eq!(
@@ -1864,7 +1881,6 @@ async fn update_feature_flags_enabling_guardian_overrides_explicit_manual_review
cwd: None,
approval_policy: Some(auto_review.approval_policy),
approvals_reviewer: Some(auto_review.approvals_reviewer),
sandbox_policy: None,
permission_profile: Some(auto_review.permission_profile.clone()),
windows_sandbox_level: None,
model: None,
@@ -1922,7 +1938,6 @@ async fn update_feature_flags_disabling_guardian_clears_manual_review_policy_wit
cwd: None,
approval_policy: None,
approvals_reviewer: Some(ApprovalsReviewer::User),
sandbox_policy: None,
permission_profile: None,
windows_sandbox_level: None,
model: None,
@@ -1982,7 +1997,6 @@ async fn update_feature_flags_enabling_guardian_in_profile_sets_profile_auto_rev
cwd: None,
approval_policy: Some(auto_review.approval_policy),
approvals_reviewer: Some(auto_review.approvals_reviewer),
sandbox_policy: None,
permission_profile: Some(auto_review.permission_profile.clone()),
windows_sandbox_level: None,
model: None,
@@ -2070,7 +2084,6 @@ guardian_approval = true
cwd: None,
approval_policy: None,
approvals_reviewer: Some(ApprovalsReviewer::User),
sandbox_policy: None,
permission_profile: None,
windows_sandbox_level: None,
model: None,
@@ -2502,35 +2515,37 @@ async fn inactive_thread_exec_approval_preserves_context() {
assert_eq!(
network_approval_context,
Some(NetworkApprovalContext {
Some(AppServerNetworkApprovalContext {
host: "example.com".to_string(),
protocol: NetworkApprovalProtocol::Socks5Tcp,
protocol: AppServerNetworkApprovalProtocol::Socks5Tcp,
})
);
assert_eq!(
additional_permissions,
Some(CoreAdditionalPermissionProfile {
network: Some(NetworkPermissions {
Some(AdditionalPermissionProfile {
network: Some(AdditionalNetworkPermissions {
enabled: Some(true),
}),
file_system: Some(FileSystemPermissions::from_read_write_roots(
Some(vec![test_absolute_path("/tmp/read-only")]),
Some(vec![test_absolute_path("/tmp/write")]),
)),
file_system: Some(AdditionalFileSystemPermissions {
read: Some(vec![test_absolute_path("/tmp/read-only")]),
write: Some(vec![test_absolute_path("/tmp/write")]),
glob_scan_max_depth: None,
entries: None,
}),
})
);
assert_eq!(
available_decisions,
vec![
codex_protocol::protocol::ReviewDecision::Approved,
codex_protocol::protocol::ReviewDecision::ApprovedForSession,
codex_protocol::protocol::ReviewDecision::NetworkPolicyAmendment {
network_policy_amendment: codex_protocol::approvals::NetworkPolicyAmendment {
codex_app_server_protocol::CommandExecutionApprovalDecision::Accept,
codex_app_server_protocol::CommandExecutionApprovalDecision::AcceptForSession,
codex_app_server_protocol::CommandExecutionApprovalDecision::ApplyNetworkPolicyAmendment {
network_policy_amendment: AppServerNetworkPolicyAmendment {
host: "example.com".to_string(),
action: codex_protocol::approvals::NetworkPolicyRuleAction::Allow,
action: AppServerNetworkPolicyRuleAction::Allow,
},
},
codex_protocol::protocol::ReviewDecision::Abort,
codex_app_server_protocol::CommandExecutionApprovalDecision::Cancel,
]
);
}
@@ -2775,35 +2790,14 @@ async fn inactive_thread_started_notification_initializes_replay_session() -> Re
);
let rollout_path = temp_dir.path().join("agent-rollout.jsonl");
let permission_profile = PermissionProfile::workspace_write();
let turn_context = TurnContextItem {
turn_id: None,
trace_id: None,
cwd: test_path_buf("/tmp/agent"),
current_date: None,
timezone: None,
approval_policy: primary_session.approval_policy,
sandbox_policy: permission_profile
.to_legacy_sandbox_policy(test_path_buf("/tmp/agent").as_path())
.expect("workspace profile must be legacy-compatible"),
permission_profile: Some(permission_profile),
network: None,
file_system_sandbox_policy: None,
model: "gpt-agent".to_string(),
personality: None,
collaboration_mode: None,
realtime_active: Some(false),
effort: primary_session.reasoning_effort,
summary: app.config.model_reasoning_summary.unwrap_or_default(),
user_instructions: None,
developer_instructions: None,
final_output_json_schema: None,
truncation_policy: None,
};
let rollout = RolloutLine {
timestamp: "t0".to_string(),
item: RolloutItem::TurnContext(turn_context),
};
let rollout = serde_json::json!({
"timestamp": "t0",
"type": "turn_context",
"payload": {
"cwd": test_path_buf("/tmp/agent"),
"model": "gpt-agent",
},
});
std::fs::write(
&rollout_path,
format!("{}\n", serde_json::to_string(&rollout)?),
@@ -3241,6 +3235,7 @@ async fn side_thread_snapshot_hides_forked_parent_transcript() {
let mut store = ThreadEventStore::new(/*capacity*/ 4);
let session = ThreadSessionState {
forked_from_id: Some(parent_thread_id),
fork_parent_title: None,
..test_thread_session(side_thread_id, test_path_buf("/tmp/side"))
};
let parent_turn = test_turn(
@@ -3303,6 +3298,7 @@ async fn side_thread_snapshot_skips_session_header_preamble() {
let snapshot = ThreadEventSnapshot {
session: Some(ThreadSessionState {
forked_from_id: Some(parent_thread_id),
fork_parent_title: None,
..test_thread_session(side_thread_id, test_path_buf("/tmp/side"))
}),
turns: Vec::new(),
@@ -3390,30 +3386,33 @@ async fn side_discard_selection_keeps_current_side_thread() {
#[tokio::test]
async fn discard_side_thread_removes_agent_navigation_entry() -> Result<()> {
let mut app = make_test_app().await;
let mut app_server =
crate::start_embedded_app_server_for_picker(app.chat_widget.config_ref()).await?;
let mut side_config = app.chat_widget.config_ref().clone();
side_config.ephemeral = true;
let started = app_server.start_thread(&side_config).await?;
let side_thread_id = started.session.thread_id;
app.side_threads
.insert(side_thread_id, SideThreadState::new(ThreadId::new()));
app.agent_navigation.upsert(
side_thread_id,
Some("Side".to_string()),
Some("side".to_string()),
/*is_closed*/ false,
);
Box::pin(async {
let mut app = make_test_app().await;
let mut app_server =
crate::start_embedded_app_server_for_picker(app.chat_widget.config_ref()).await?;
let mut side_config = app.chat_widget.config_ref().clone();
side_config.ephemeral = true;
let started = app_server.start_thread(&side_config).await?;
let side_thread_id = started.session.thread_id;
app.side_threads
.insert(side_thread_id, SideThreadState::new(ThreadId::new()));
app.agent_navigation.upsert(
side_thread_id,
Some("Side".to_string()),
Some("side".to_string()),
/*is_closed*/ false,
);
assert!(
app.discard_side_thread(&mut app_server, side_thread_id)
.await
);
assert!(
app.discard_side_thread(&mut app_server, side_thread_id)
.await
);
assert_eq!(app.agent_navigation.get(&side_thread_id), None);
assert!(!app.side_threads.contains_key(&side_thread_id));
Ok(())
assert_eq!(app.agent_navigation.get(&side_thread_id), None);
assert!(!app.side_threads.contains_key(&side_thread_id));
Ok(())
})
.await
}
#[tokio::test]
@@ -3598,9 +3597,10 @@ async fn render_clear_ui_header_after_long_transcript_for_snapshot() -> String {
)) as Arc<dyn HistoryCell>
};
let make_header = |is_first| -> Arc<dyn HistoryCell> {
let event = SessionConfiguredEvent {
session_id: ThreadId::new(),
let session = ThreadSessionState {
thread_id: ThreadId::new(),
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -3610,17 +3610,17 @@ async fn render_clear_ui_header_after_long_transcript_for_snapshot() -> String {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/tmp/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: Some(ReasoningEffortConfig::High),
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
};
Arc::new(new_session_info(
app.chat_widget.config_ref(),
app.chat_widget.current_model(),
event,
&session,
is_first,
/*tooltip_override*/ None,
/*auth_plan*/ None,
@@ -4248,7 +4248,7 @@ fn test_session_telemetry(config: &Config, model: &str) -> SessionTelemetry {
"test_originator".to_string(),
/*log_user_prompts*/ false,
"test".to_string(),
SessionSource::Cli,
crate::test_support::session_source_cli(),
)
}
@@ -4351,9 +4351,10 @@ async fn backtrack_selection_with_duplicate_history_targets_unique_turn() {
};
let make_header = |is_first| {
let event = SessionConfiguredEvent {
session_id: ThreadId::new(),
let session = ThreadSessionState {
thread_id: ThreadId::new(),
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -4363,17 +4364,17 @@ async fn backtrack_selection_with_duplicate_history_targets_unique_turn() {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/home/user/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
};
Arc::new(new_session_info(
app.chat_widget.config_ref(),
app.chat_widget.current_model(),
event,
&session,
is_first,
/*tooltip_override*/ None,
/*auth_plan*/ None,
@@ -4413,11 +4414,11 @@ async fn backtrack_selection_with_duplicate_history_targets_unique_turn() {
assert_eq!(user_count(&app.transcript_cells), 2);
let base_id = ThreadId::new();
app.chat_widget.handle_codex_event(Event {
id: String::new(),
msg: EventMsg::SessionConfigured(SessionConfiguredEvent {
session_id: base_id,
app.chat_widget
.handle_thread_session(crate::session_state::ThreadSessionState {
thread_id: base_id,
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -4427,14 +4428,13 @@ async fn backtrack_selection_with_duplicate_history_targets_unique_turn() {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/home/user/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
}),
});
});
app.backtrack.base_id = Some(base_id);
app.backtrack.primed = true;
@@ -4507,11 +4507,11 @@ async fn backtrack_resubmit_preserves_data_image_urls_in_user_turn() {
let (mut app, _app_event_rx, mut op_rx) = make_test_app_with_channels().await;
let thread_id = ThreadId::new();
app.chat_widget.handle_codex_event(Event {
id: String::new(),
msg: EventMsg::SessionConfigured(SessionConfiguredEvent {
session_id: thread_id,
app.chat_widget
.handle_thread_session(crate::session_state::ThreadSessionState {
thread_id,
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -4521,14 +4521,13 @@ async fn backtrack_resubmit_preserves_data_image_urls_in_user_turn() {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/home/user/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
}),
});
});
let data_image_url = "data:image/png;base64,abc123".to_string();
app.transcript_cells = vec![Arc::new(UserHistoryCell {
@@ -4564,7 +4563,7 @@ async fn backtrack_resubmit_preserves_data_image_urls_in_user_turn() {
assert!(items.iter().any(|item| {
matches!(
item,
UserInput::Image { image_url } if image_url == &data_image_url
UserInput::Image { url } if url == &data_image_url
)
}));
}
@@ -4749,6 +4748,7 @@ async fn refreshed_snapshot_session_persists_resumed_turns() {
)];
let resumed_session = ThreadSessionState {
cwd: test_path_buf("/tmp/refreshed").abs(),
instruction_source_paths: Vec::new(),
..initial_session.clone()
};
let mut snapshot = ThreadEventSnapshot {
@@ -4873,7 +4873,7 @@ async fn thread_rollback_response_discards_queued_active_thread_events() {
path: None,
cwd: test_path_buf("/tmp/project").abs(),
cli_version: "0.0.0".to_string(),
source: SessionSource::Cli.into(),
source: SessionSource::Cli,
agent_nickname: None,
agent_role: None,
git_info: None,
@@ -4893,48 +4893,49 @@ async fn thread_rollback_response_discards_queued_active_thread_events() {
#[tokio::test]
async fn new_session_requests_shutdown_for_previous_conversation() {
let (mut app, mut app_event_rx, mut op_rx) = Box::pin(make_test_app_with_channels()).await;
Box::pin(async {
let (mut app, mut app_event_rx, mut op_rx) = Box::pin(make_test_app_with_channels()).await;
let thread_id = ThreadId::new();
let event = SessionConfiguredEvent {
session_id: thread_id,
forked_from_id: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
service_tier: None,
approval_policy: AskForApproval::Never,
approvals_reviewer: ApprovalsReviewer::User,
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/home/user/project").abs(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
};
let thread_id = ThreadId::new();
let event = crate::session_state::ThreadSessionState {
thread_id,
forked_from_id: None,
fork_parent_title: None,
thread_name: None,
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
service_tier: None,
approval_policy: AskForApproval::Never,
approvals_reviewer: ApprovalsReviewer::User,
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/home/user/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
};
app.chat_widget.handle_codex_event(Event {
id: String::new(),
msg: EventMsg::SessionConfigured(event),
});
app.chat_widget.handle_thread_session(event);
while app_event_rx.try_recv().is_ok() {}
while op_rx.try_recv().is_ok() {}
while app_event_rx.try_recv().is_ok() {}
while op_rx.try_recv().is_ok() {}
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("embedded app server");
Box::pin(app.shutdown_current_thread(&mut app_server)).await;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("embedded app server");
Box::pin(app.shutdown_current_thread(&mut app_server)).await;
assert!(
op_rx.try_recv().is_err(),
"shutdown should not submit Op::Shutdown"
);
assert!(
op_rx.try_recv().is_err(),
"shutdown should not submit Op::Shutdown"
);
})
.await;
}
#[tokio::test]
@@ -4983,39 +4984,45 @@ async fn shutdown_first_exit_uses_app_server_shutdown_without_submitting_op() {
#[tokio::test]
async fn interrupt_without_active_turn_is_treated_as_handled() {
let mut app = make_test_app().await;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("embedded app server");
let started = app_server
.start_thread(app.chat_widget.config_ref())
Box::pin(async {
let mut app = make_test_app().await;
let mut app_server = Box::pin(crate::start_embedded_app_server_for_picker(
app.chat_widget.config_ref(),
))
.await
.expect("thread/start should succeed");
let thread_id = started.session.thread_id;
app.enqueue_primary_thread_session(started.session, started.turns)
.await
.expect("primary thread should be registered");
let op = AppCommand::interrupt();
let handled =
Box::pin(app.try_submit_active_thread_op_via_app_server(&mut app_server, thread_id, &op))
.expect("embedded app server");
let started = app_server
.start_thread(app.chat_widget.config_ref())
.await
.expect("interrupt submission should not fail");
.expect("thread/start should succeed");
let thread_id = started.session.thread_id;
app.enqueue_primary_thread_session(started.session, started.turns)
.await
.expect("primary thread should be registered");
let op = AppCommand::interrupt();
assert_eq!(handled, true);
let handled = Box::pin(app.try_submit_active_thread_op_via_app_server(
&mut app_server,
thread_id,
&op,
))
.await
.expect("interrupt submission should not fail");
assert_eq!(handled, true);
})
.await;
}
#[tokio::test]
async fn clear_only_ui_reset_preserves_chat_session_state() {
let mut app = make_test_app().await;
let thread_id = ThreadId::new();
app.chat_widget.handle_codex_event(Event {
id: String::new(),
msg: EventMsg::SessionConfigured(SessionConfiguredEvent {
session_id: thread_id,
app.chat_widget
.handle_thread_session(crate::session_state::ThreadSessionState {
thread_id,
forked_from_id: None,
fork_parent_title: None,
thread_name: Some("keep me".to_string()),
model: "gpt-test".to_string(),
model_provider_id: "test-provider".to_string(),
@@ -5025,14 +5032,13 @@ async fn clear_only_ui_reset_preserves_chat_session_state() {
permission_profile: PermissionProfile::read_only(),
active_permission_profile: None,
cwd: test_path_buf("/tmp/project").abs(),
instruction_source_paths: Vec::new(),
reasoning_effort: None,
history_log_id: 0,
history_entry_count: 0,
initial_messages: None,
network_proxy: None,
rollout_path: Some(PathBuf::new()),
}),
});
});
app.chat_widget
.apply_external_edit("draft prompt".to_string());
app.transcript_cells = vec![Arc::new(UserHistoryCell {
+2 -2
View File
@@ -19,7 +19,7 @@ pub(super) struct ThreadEventSnapshot {
pub(super) enum ThreadBufferedEvent {
Notification(ServerNotification),
Request(ServerRequest),
HistoryEntryResponse(GetHistoryEntryResponseEvent),
HistoryEntryResponse(HistoryLookupResponse),
FeedbackSubmission(FeedbackThreadEvent),
}
@@ -318,6 +318,7 @@ mod tests {
use super::*;
use crate::test_support::PathBufExt;
use crate::test_support::test_path_buf;
use codex_app_server_protocol::AskForApproval;
use codex_app_server_protocol::CommandExecutionRequestApprovalParams;
use codex_app_server_protocol::HookCompletedNotification;
use codex_app_server_protocol::HookEventName as AppServerHookEventName;
@@ -334,7 +335,6 @@ mod tests {
use codex_app_server_protocol::TurnStartedNotification;
use codex_config::types::ApprovalsReviewer;
use codex_protocol::models::PermissionProfile;
use codex_protocol::protocol::AskForApproval;
use pretty_assertions::assert_eq;
use std::path::PathBuf;
+50 -81
View File
@@ -5,6 +5,7 @@
//! when the visible thread changes.
use super::*;
use crate::session_resume::read_session_model;
impl App {
pub(super) async fn shutdown_current_thread(&mut self, app_server: &mut AppServerSession) {
@@ -213,24 +214,11 @@ impl App {
let thread_label = Some(self.thread_label(thread_id));
match request {
ServerRequest::CommandExecutionRequestApproval { params, .. } => {
let network_approval_context = params
.network_approval_context
.clone()
.map(network_approval_context_to_core);
let additional_permissions = params.additional_permissions.clone().map(Into::into);
let proposed_execpolicy_amendment = params
.proposed_execpolicy_amendment
.clone()
.map(codex_app_server_protocol::ExecPolicyAmendment::into_core);
let proposed_network_policy_amendments = params
.proposed_network_policy_amendments
.clone()
.map(|amendments| {
amendments
.into_iter()
.map(codex_app_server_protocol::NetworkPolicyAmendment::into_core)
.collect::<Vec<_>>()
});
let network_approval_context = params.network_approval_context.clone();
let additional_permissions = params.additional_permissions.clone();
let proposed_execpolicy_amendment = params.proposed_execpolicy_amendment.clone();
let proposed_network_policy_amendments =
params.proposed_network_policy_amendments.clone();
Some(ThreadInteractiveRequest::Approval(ApprovalRequest::Exec {
thread_id,
thread_label,
@@ -244,23 +232,14 @@ impl App {
.map(split_command_string)
.unwrap_or_default(),
reason: params.reason.clone(),
available_decisions: params
.available_decisions
.clone()
.map(|decisions| {
decisions
.into_iter()
.map(command_execution_decision_to_review_decision)
.collect()
})
.unwrap_or_else(|| {
default_exec_approval_decisions(
network_approval_context.as_ref(),
proposed_execpolicy_amendment.as_ref(),
proposed_network_policy_amendments.as_deref(),
additional_permissions.as_ref(),
)
}),
available_decisions: params.available_decisions.clone().unwrap_or_else(|| {
default_exec_approval_decisions(
network_approval_context.as_ref(),
proposed_execpolicy_amendment.as_ref(),
proposed_network_policy_amendments.as_deref(),
additional_permissions.as_ref(),
)
}),
network_approval_context,
additional_permissions,
}))
@@ -278,14 +257,14 @@ impl App {
changes: self
.thread_file_change_changes(thread_id, &params.turn_id, &params.item_id)
.await
.map(crate::app_server_approval_conversions::file_update_changes_to_core)
.map(crate::app_server_approval_conversions::file_update_changes_to_display)
.unwrap_or_default(),
}),
),
ServerRequest::McpServerElicitationRequest { request_id, params } => {
if let Some(request) = McpServerElicitationFormRequest::from_app_server_request(
thread_id,
app_server_request_id_to_mcp_request_id(request_id),
request_id.clone(),
params.clone(),
) {
Some(ThreadInteractiveRequest::McpServerElicitation(request))
@@ -295,7 +274,7 @@ impl App {
thread_id,
thread_label,
server_name: params.server_name.clone(),
request_id: app_server_request_id_to_mcp_request_id(request_id),
request_id: request_id.clone(),
message: match &params.request {
codex_app_server_protocol::McpServerElicitationRequest::Form {
message,
@@ -457,9 +436,9 @@ impl App {
thread_id: ThreadId,
op: &AppCommand,
) -> Result<bool> {
match op.view() {
AppCommandView::Other(Op::AddToHistory { text }) => {
let text = text.clone();
match op {
AppCommand::AddToHistory { text } => {
let text = text.to_string();
let config = self.chat_widget.config_ref().clone();
tokio::spawn(async move {
if let Err(err) = append_message_history_entry(&text, &thread_id, &config).await
@@ -473,11 +452,11 @@ impl App {
});
Ok(true)
}
AppCommandView::Other(Op::GetHistoryEntryRequest { offset, log_id }) => {
let offset = *offset;
let log_id = *log_id;
AppCommand::GetHistoryEntryRequest { offset, log_id } => {
let config = self.chat_widget.config_ref().clone();
let app_event_tx = self.app_event_tx.clone();
let offset = *offset;
let log_id = *log_id;
tokio::spawn(async move {
let entry_opt = tokio::task::spawn_blocking(move || {
lookup_message_history_entry(log_id, offset, &config)
@@ -490,7 +469,7 @@ impl App {
app_event_tx.send(AppEvent::ThreadHistoryEntryResponse {
thread_id,
event: GetHistoryEntryResponseEvent {
event: HistoryLookupResponse {
offset,
log_id,
entry: entry_opt.map(|entry| {
@@ -515,8 +494,8 @@ impl App {
thread_id: ThreadId,
op: &AppCommand,
) -> Result<bool> {
match op.view() {
AppCommandView::Interrupt => {
match op {
AppCommand::Interrupt => {
if let Some(turn_id) = self.active_turn_id_for_thread(thread_id).await {
app_server.turn_interrupt(thread_id, turn_id).await?;
} else {
@@ -524,7 +503,7 @@ impl App {
}
Ok(true)
}
AppCommandView::UserTurn {
AppCommand::UserTurn {
items,
cwd,
approval_policy,
@@ -618,12 +597,12 @@ impl App {
thread_id,
items.to_vec(),
cwd.clone(),
approval_policy,
*approval_policy,
approvals_reviewer,
permission_profile.clone(),
active_permission_profile,
model.to_string(),
effort,
*effort,
*summary,
*service_tier,
collaboration_mode.clone(),
@@ -634,12 +613,12 @@ impl App {
}
Ok(true)
}
AppCommandView::ListSkills { cwds, force_reload } => {
AppCommand::ListSkills { cwds, force_reload } => {
self.handle_skills_list_result(
app_server
.skills_list(codex_app_server_protocol::SkillsListParams {
cwds: cwds.to_vec(),
force_reload,
cwds: cwds.clone(),
force_reload: *force_reload,
per_cwd_extra_user_roots: None,
})
.await,
@@ -647,76 +626,67 @@ impl App {
);
Ok(true)
}
AppCommandView::Compact => {
AppCommand::Compact => {
app_server.thread_compact_start(thread_id).await?;
Ok(true)
}
AppCommandView::SetThreadName { name } => {
AppCommand::SetThreadName { name } => {
app_server
.thread_set_name(thread_id, name.to_string())
.await?;
Ok(true)
}
AppCommandView::ThreadRollback { num_turns } => {
let response = match app_server.thread_rollback(thread_id, num_turns).await {
AppCommand::ThreadRollback { num_turns } => {
let response = match app_server.thread_rollback(thread_id, *num_turns).await {
Ok(response) => response,
Err(err) => {
self.handle_backtrack_rollback_failed();
return Err(err);
}
};
self.handle_thread_rollback_response(thread_id, num_turns, &response)
self.handle_thread_rollback_response(thread_id, *num_turns, &response)
.await;
Ok(true)
}
AppCommandView::Review { review_request } => {
app_server
.review_start(thread_id, review_request.clone())
.await?;
AppCommand::Review { target } => {
app_server.review_start(thread_id, target.clone()).await?;
Ok(true)
}
AppCommandView::CleanBackgroundTerminals => {
AppCommand::CleanBackgroundTerminals => {
app_server
.thread_background_terminals_clean(thread_id)
.await?;
Ok(true)
}
AppCommandView::RealtimeConversationStart(params) => {
AppCommand::RealtimeConversationStart { transport, voice } => {
app_server
.thread_realtime_start(thread_id, params.clone())
.thread_realtime_start(thread_id, transport.clone(), voice.clone())
.await?;
Ok(true)
}
AppCommandView::RealtimeConversationAudio(params) => {
AppCommand::RealtimeConversationAudio(frame) => {
app_server
.thread_realtime_audio(thread_id, params.clone())
.thread_realtime_audio(thread_id, frame.clone())
.await?;
Ok(true)
}
AppCommandView::RealtimeConversationText(params) => {
app_server
.thread_realtime_text(thread_id, params.clone())
.await?;
Ok(true)
}
AppCommandView::RealtimeConversationClose => {
AppCommand::RealtimeConversationClose => {
app_server.thread_realtime_stop(thread_id).await?;
Ok(true)
}
AppCommandView::RunUserShellCommand { command } => {
AppCommand::RunUserShellCommand { command } => {
app_server
.thread_shell_command(thread_id, command.to_string())
.await?;
Ok(true)
}
AppCommandView::ReloadUserConfig => {
AppCommand::ReloadUserConfig => {
app_server.reload_user_config().await?;
self.refresh_in_memory_config_from_disk().await?;
Ok(true)
}
AppCommandView::OverrideTurnContext { .. } => Ok(true),
AppCommandView::ApproveGuardianDeniedAction { event }
| AppCommandView::Other(Op::ApproveGuardianDeniedAction { event }) => {
AppCommand::OverrideTurnContext { .. } => Ok(true),
AppCommand::ApproveGuardianDeniedAction { event } => {
app_server
.thread_approve_guardian_denied_action(thread_id, event)
.await?;
@@ -1016,7 +986,7 @@ impl App {
pub(super) async fn enqueue_thread_history_entry_response(
&mut self,
thread_id: ThreadId,
event: GetHistoryEntryResponseEvent,
event: HistoryLookupResponse,
) -> Result<()> {
let (sender, store) = {
let channel = self.ensure_thread_channel(thread_id);
@@ -1327,7 +1297,6 @@ impl App {
#[allow(clippy::too_many_arguments)]
pub(super) fn handle_skills_list_response(&mut self, response: SkillsListResponse) {
let response = list_skills_response_to_core(response);
let cwd = self.chat_widget.config_ref().cwd.clone();
let errors = errors_for_cwd(&cwd, &response);
emit_skill_load_warnings(&self.app_event_tx, &errors);
+36 -28
View File
@@ -1,6 +1,7 @@
use super::App;
use crate::app_server_session::ThreadSessionState;
use crate::read_session_model;
use crate::session_resume::read_session_model;
use crate::session_state::ThreadSessionState;
use codex_app_server_protocol::AskForApproval;
use codex_app_server_protocol::Thread;
use codex_protocol::ThreadId;
use codex_protocol::models::ActivePermissionProfile;
@@ -12,7 +13,7 @@ impl App {
return;
};
let approval_policy = self.config.permissions.approval_policy.value();
let approval_policy = AskForApproval::from(self.config.permissions.approval_policy.value());
let approvals_reviewer = self.config.approvals_reviewer;
let permission_profile = self
.chat_widget
@@ -63,7 +64,9 @@ impl App {
model: self.chat_widget.current_model().to_string(),
model_provider_id: self.config.model_provider_id.clone(),
service_tier: self.chat_widget.current_service_tier(),
approval_policy: self.config.permissions.approval_policy.value(),
approval_policy: AskForApproval::from(
self.config.permissions.approval_policy.value(),
),
approvals_reviewer: self.config.approvals_reviewer,
permission_profile: permission_profile.clone(),
active_permission_profile: active_permission_profile.clone(),
@@ -118,15 +121,16 @@ mod tests {
use crate::app::thread_events::ThreadEventChannel;
use crate::test_support::PathBufExt;
use crate::test_support::test_path_buf;
use codex_app_server_protocol::AskForApproval;
use codex_app_server_protocol::FileSystemAccessMode;
use codex_app_server_protocol::FileSystemPath;
use codex_app_server_protocol::FileSystemSandboxEntry;
use codex_app_server_protocol::FileSystemSpecialPath;
use codex_app_server_protocol::PermissionProfile as AppServerPermissionProfile;
use codex_app_server_protocol::PermissionProfileFileSystemPermissions;
use codex_app_server_protocol::PermissionProfileNetworkPermissions;
use codex_config::types::ApprovalsReviewer;
use codex_protocol::models::PermissionProfile;
use codex_protocol::protocol::AskForApproval;
use codex_protocol::protocol::FileSystemAccessMode;
use codex_protocol::protocol::FileSystemPath;
use codex_protocol::protocol::FileSystemSandboxEntry;
use codex_protocol::protocol::FileSystemSandboxPolicy;
use codex_protocol::protocol::FileSystemSpecialPath;
use codex_protocol::protocol::NetworkSandboxPolicy;
use pretty_assertions::assert_eq;
use std::path::PathBuf;
@@ -189,7 +193,7 @@ mod tests {
app.side_threads
.insert(side_thread_id, SideThreadState::new(main_thread_id));
app.config.permissions.approval_policy =
codex_config::Constrained::allow_any(AskForApproval::OnRequest);
codex_config::Constrained::allow_any(AskForApproval::OnRequest.to_core());
app.config.approvals_reviewer = ApprovalsReviewer::AutoReview;
let expected_permission_profile = PermissionProfile::workspace_write();
app.chat_widget.handle_thread_session(main_session.clone());
@@ -243,23 +247,27 @@ mod tests {
let mut app = make_test_app().await;
let thread_id =
ThreadId::from_string("00000000-0000-0000-0000-000000000403").expect("valid thread");
let profile = PermissionProfile::from_runtime_permissions(
&FileSystemSandboxPolicy::restricted(vec![
FileSystemSandboxEntry {
path: FileSystemPath::Special {
value: FileSystemSpecialPath::Root,
let profile: PermissionProfile = AppServerPermissionProfile::Managed {
network: PermissionProfileNetworkPermissions { enabled: false },
file_system: PermissionProfileFileSystemPermissions::Restricted {
entries: vec![
FileSystemSandboxEntry {
path: FileSystemPath::Special {
value: FileSystemSpecialPath::Root,
},
access: FileSystemAccessMode::Read,
},
access: FileSystemAccessMode::Read,
},
FileSystemSandboxEntry {
path: FileSystemPath::GlobPattern {
pattern: "**/.env".to_string(),
FileSystemSandboxEntry {
path: FileSystemPath::GlobPattern {
pattern: "**/.env".to_string(),
},
access: FileSystemAccessMode::None,
},
access: FileSystemAccessMode::None,
},
]),
NetworkSandboxPolicy::Restricted,
);
],
glob_scan_max_depth: None,
},
}
.into();
let session = ThreadSessionState {
permission_profile: profile.clone(),
..test_thread_session(thread_id, test_path_buf("/tmp/main"))
@@ -274,7 +282,7 @@ mod tests {
);
app.chat_widget.handle_thread_session(session.clone());
app.config.permissions.approval_policy =
codex_config::Constrained::allow_any(AskForApproval::OnRequest);
codex_config::Constrained::allow_any(AskForApproval::OnRequest.to_core());
app.sync_active_thread_permission_settings_to_cached_session()
.await;