Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 8 additions & 8 deletions codex-rs/app-server-protocol/src/protocol/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -641,12 +641,12 @@ client_request_definitions! {
serialization: None,
response: v2::ThreadTurnsListResponse,
},
#[experimental("thread/turns/items/list")]
ThreadTurnsItemsList => "thread/turns/items/list" {
params: v2::ThreadTurnsItemsListParams,
#[experimental("thread/items/list")]
ThreadItemsList => "thread/items/list" {
params: v2::ThreadItemsListParams,
// Explicitly concurrent: this primarily reads append-only rollout storage.
serialization: None,
response: v2::ThreadTurnsItemsListResponse,
response: v2::ThreadItemsListResponse,
},
/// Append raw Responses API items to the thread history without starting a user turn.
ThreadInjectItems => "thread/inject_items" {
Expand Down Expand Up @@ -2090,17 +2090,17 @@ mod tests {
};
assert_eq!(thread_turns_list.serialization_scope(), None);

let thread_turns_items_list = ClientRequest::ThreadTurnsItemsList {
let thread_items_list = ClientRequest::ThreadItemsList {
request_id: request_id(),
params: v2::ThreadTurnsItemsListParams {
params: v2::ThreadItemsListParams {
thread_id: "thread-1".to_string(),
turn_id: "turn-1".to_string(),
turn_id: None,
cursor: None,
limit: None,
sort_direction: None,
},
};
assert_eq!(thread_turns_items_list.serialization_scope(), None);
assert_eq!(thread_items_list.serialization_scope(), None);

let mcp_resource_read = ClientRequest::McpResourceRead {
request_id: request_id(),
Expand Down
34 changes: 30 additions & 4 deletions codex-rs/app-server-protocol/src/protocol/v2/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -224,10 +224,10 @@ fn thread_resume_response_round_trips_initial_turns_page() {
}

#[test]
fn thread_turns_items_list_round_trips() {
let params = ThreadTurnsItemsListParams {
fn thread_items_list_round_trips() {
let params = ThreadItemsListParams {
thread_id: "thr_123".to_string(),
turn_id: "turn_456".to_string(),
turn_id: Some("turn_456".to_string()),
cursor: Some("cursor_1".to_string()),
limit: Some(50),
sort_direction: Some(SortDirection::Asc),
Expand All @@ -243,7 +243,7 @@ fn thread_turns_items_list_round_trips() {
"sortDirection": "asc",
})
);
let response = ThreadTurnsItemsListResponse {
let response = ThreadItemsListResponse {
data: vec![ThreadItem::ContextCompaction {
id: "item_1".to_string(),
}],
Expand All @@ -259,6 +259,32 @@ fn thread_turns_items_list_round_trips() {
"backwardsCursor": "cursor_0",
})
);

let params_without_turn = ThreadItemsListParams {
thread_id: "thr_123".to_string(),
turn_id: None,
cursor: None,
limit: None,
sort_direction: None,
};

assert_eq!(
serde_json::to_value(&params_without_turn).expect("serialize params without turn"),
json!({
"threadId": "thr_123",
"turnId": null,
"cursor": null,
"limit": null,
"sortDirection": null,
})
);
assert_eq!(
serde_json::from_value::<ThreadItemsListParams>(json!({
"threadId": "thr_123",
}))
.expect("deserialize params without turn"),
params_without_turn
);
}

#[test]
Expand Down
8 changes: 5 additions & 3 deletions codex-rs/app-server-protocol/src/protocol/v2/thread.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1316,9 +1316,11 @@ pub struct ThreadTurnsListResponse {
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct ThreadTurnsItemsListParams {
pub struct ThreadItemsListParams {
pub thread_id: String,
pub turn_id: String,
/// Optional turn id to filter by. When omitted, returns items across the thread.
#[ts(optional = nullable)]
pub turn_id: Option<String>,
/// Opaque cursor to pass to the next call to continue after the last item.
#[ts(optional = nullable)]
pub cursor: Option<String>,
Expand All @@ -1333,7 +1335,7 @@ pub struct ThreadTurnsItemsListParams {
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct ThreadTurnsItemsListResponse {
pub struct ThreadItemsListResponse {
pub data: Vec<ThreadItem>,
/// Opaque cursor to pass to the next call to continue after the last item.
/// if None, there are no more items to return.
Expand Down
8 changes: 4 additions & 4 deletions codex-rs/app-server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,7 @@ Example with notification opt-out:
- `thread/loaded/list` — list the thread ids currently loaded in memory.
- `thread/read` — read a stored thread by id without resuming it; optionally include turns via `includeTurns`. The returned `thread` includes `status` (`ThreadStatus`), defaulting to `notLoaded` when the thread is not currently loaded.
- `thread/turns/list` — experimental; page through a stored thread’s turn history without resuming it; supports cursor-based pagination with `sortDirection`, `itemsView`, `nextCursor`, and `backwardsCursor`.
- `thread/turns/items/list` — experimental; reserved for paging full items for one turn. The API shape is present, but app-server currently returns an unsupported-method JSON-RPC error.
- `thread/items/list` — experimental; page through persisted thread items without resuming the thread. Pass `turnId` to restrict results to one turn, or omit it to page items across the thread. The active thread store must support item pagination.
- `thread/metadata/update` — patch stored thread metadata in sqlite; currently supports updating persisted `gitInfo` fields and returns the refreshed `thread`.
- `thread/settings/update` — experimental; queue a partial update to a loaded thread’s next-turn settings without starting a turn or adding transcript items. Omitted fields leave settings unchanged; `serviceTier: null` clears the tier; `multiAgentMode` selects `none`, `explicitRequestOnly`, or `proactive` for subsequent turns; `sandboxPolicy` and `permissions` cannot be combined. Returns `{}` when the update is accepted and emits `thread/settings/updated` with the full effective settings only if they actually change. `turn/start` settings overrides emit the same notification when they change the stored settings.
- `thread/memoryMode/set` — experimental; set a thread’s persisted memory eligibility to `"enabled"` or `"disabled"` for either a loaded thread or a stored rollout; returns `{}` on success.
Expand Down Expand Up @@ -514,18 +514,18 @@ Every returned `Turn` includes `itemsView`, which tells clients whether the `ite
} }
```

`thread/turns/items/list` is the planned hydration API for fetching full items for one turn:
`thread/items/list` pages full persisted items across a thread, optionally filtered to one turn:

```json
{ "method": "thread/turns/items/list", "id": 25, "params": {
{ "method": "thread/items/list", "id": 25, "params": {
"threadId": "thr_123",
"turnId": "turn_456",
"limit": 100,
"sortDirection": "asc"
} }
```

This method currently returns JSON-RPC `-32601` with message `thread/turns/items/list is not supported yet`.
Omit `turnId` or pass `null` to page items across the thread. Thread stores that do not implement item pagination return JSON-RPC `-32601` with message `thread/items/list is not supported yet`.

### Example: Update stored thread metadata

Expand Down
4 changes: 2 additions & 2 deletions codex-rs/app-server/src/message_processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1223,8 +1223,8 @@ impl MessageProcessor {
ClientRequest::ThreadTurnsList { params, .. } => {
self.thread_processor.thread_turns_list(params).await
}
ClientRequest::ThreadTurnsItemsList { params, .. } => {
self.thread_processor.thread_turns_items_list(params).await
ClientRequest::ThreadItemsList { params, .. } => {
self.thread_processor.thread_items_list(params).await
}
ClientRequest::ThreadShellCommand { params, .. } => {
self.thread_processor
Expand Down
4 changes: 3 additions & 1 deletion codex-rs/app-server/src/request_processors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -211,6 +211,8 @@ use codex_app_server_protocol::ThreadIncrementElicitationResponse;
use codex_app_server_protocol::ThreadInjectItemsParams;
use codex_app_server_protocol::ThreadInjectItemsResponse;
use codex_app_server_protocol::ThreadItem;
use codex_app_server_protocol::ThreadItemsListParams;
use codex_app_server_protocol::ThreadItemsListResponse;
use codex_app_server_protocol::ThreadListCwdFilter;
use codex_app_server_protocol::ThreadListParams;
use codex_app_server_protocol::ThreadListResponse;
Expand Down Expand Up @@ -256,7 +258,6 @@ use codex_app_server_protocol::ThreadStartParams;
use codex_app_server_protocol::ThreadStartResponse;
use codex_app_server_protocol::ThreadStartedNotification;
use codex_app_server_protocol::ThreadStatus;
use codex_app_server_protocol::ThreadTurnsItemsListParams;
use codex_app_server_protocol::ThreadTurnsListParams;
use codex_app_server_protocol::ThreadTurnsListResponse;
use codex_app_server_protocol::ThreadUnarchiveParams;
Expand Down Expand Up @@ -436,6 +437,7 @@ use codex_state::log_db::LogDbLayer;
use codex_thread_store::ArchiveThreadParams as StoreArchiveThreadParams;
use codex_thread_store::DeleteThreadParams as StoreDeleteThreadParams;
use codex_thread_store::GitInfoPatch as StoreGitInfoPatch;
use codex_thread_store::ListItemsParams as StoreListItemsParams;
use codex_thread_store::ListThreadsParams as StoreListThreadsParams;
use codex_thread_store::LocalThreadStore;
use codex_thread_store::ReadThreadByRolloutPathParams as StoreReadThreadByRolloutPathParams;
Expand Down
74 changes: 69 additions & 5 deletions codex-rs/app-server/src/request_processors/thread_processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -683,13 +683,13 @@ impl ThreadRequestProcessor {
.map(|response| Some(response.into()))
}

pub(crate) async fn thread_turns_items_list(
pub(crate) async fn thread_items_list(
&self,
_params: ThreadTurnsItemsListParams,
params: ThreadItemsListParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
Err(method_not_found(
"thread/turns/items/list is not supported yet",
))
self.thread_items_list_response_inner(params)
.await
.map(|response| Some(response.into()))
}

pub(crate) async fn thread_shell_command(
Expand Down Expand Up @@ -2388,6 +2388,68 @@ impl ThreadRequestProcessor {
)
}

async fn thread_items_list_response_inner(
&self,
params: ThreadItemsListParams,
) -> Result<ThreadItemsListResponse, JSONRPCErrorError> {
let ThreadItemsListParams {
thread_id,
turn_id,
cursor,
limit,
sort_direction,
} = params;
let thread_id = ThreadId::from_string(&thread_id)
.map_err(|err| invalid_request(format!("invalid thread id: {err}")))?;
let page_size = limit
.map(|value| value as usize)
.unwrap_or(THREAD_ITEMS_DEFAULT_LIMIT)
.clamp(1, THREAD_ITEMS_MAX_LIMIT);
let page = self
.thread_store
.list_items(StoreListItemsParams {
thread_id,
turn_id,
include_archived: true,
cursor,
page_size,
sort_direction: match sort_direction.unwrap_or(SortDirection::Asc) {
SortDirection::Asc => StoreSortDirection::Asc,
SortDirection::Desc => StoreSortDirection::Desc,
},
})
.await
.map_err(|err| match err {
ThreadStoreError::InvalidRequest { message } => invalid_request(message),
ThreadStoreError::Unsupported { .. } => {
method_not_found("thread/items/list is not supported yet")
}
ThreadStoreError::ThreadNotFound { thread_id } => {
invalid_request(format!("no rollout found for thread id {thread_id}"))
}
err => internal_error(format!("failed to list thread items: {err}")),
})?;
let data =
page.items
.into_iter()
.map(|item| {
serde_json::from_slice::<ThreadItem>(&item.materialized_thread_item_json)
.map_err(|err| {
internal_error(format!(
"failed to deserialize stored thread item {}: {err}",
item.item_key
))
})
})
.collect::<Result<Vec<_>, _>>()?;

Ok(ThreadItemsListResponse {
data,
next_cursor: page.next_cursor,
backwards_cursor: page.backwards_cursor,
})
}

async fn load_thread_turns_list_history(
&self,
thread_id: ThreadId,
Expand Down Expand Up @@ -3714,6 +3776,8 @@ fn xcode_26_4_mcp_elicitations_auto_deny(

const THREAD_TURNS_DEFAULT_LIMIT: usize = 25;
const THREAD_TURNS_MAX_LIMIT: usize = 100;
const THREAD_ITEMS_DEFAULT_LIMIT: usize = 25;
const THREAD_ITEMS_MAX_LIMIT: usize = 100;

fn thread_backwards_cursor_for_sort_key(
thread: &StoredThread,
Expand Down
10 changes: 5 additions & 5 deletions codex-rs/app-server/tests/common/test_app_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ use codex_app_server_protocol::ThreadCompactStartParams;
use codex_app_server_protocol::ThreadDeleteParams;
use codex_app_server_protocol::ThreadForkParams;
use codex_app_server_protocol::ThreadInjectItemsParams;
use codex_app_server_protocol::ThreadItemsListParams;
use codex_app_server_protocol::ThreadListParams;
use codex_app_server_protocol::ThreadLoadedListParams;
use codex_app_server_protocol::ThreadMemoryModeSetParams;
Expand All @@ -102,7 +103,6 @@ use codex_app_server_protocol::ThreadSetNameParams;
use codex_app_server_protocol::ThreadSettingsUpdateParams;
use codex_app_server_protocol::ThreadShellCommandParams;
use codex_app_server_protocol::ThreadStartParams;
use codex_app_server_protocol::ThreadTurnsItemsListParams;
use codex_app_server_protocol::ThreadTurnsListParams;
use codex_app_server_protocol::ThreadUnarchiveParams;
use codex_app_server_protocol::ThreadUnsubscribeParams;
Expand Down Expand Up @@ -602,13 +602,13 @@ impl TestAppServer {
self.send_request("thread/turns/list", params).await
}

/// Send a `thread/turns/items/list` JSON-RPC request.
pub async fn send_thread_turns_items_list_request(
/// Send a `thread/items/list` JSON-RPC request.
pub async fn send_thread_items_list_request(
&mut self,
params: ThreadTurnsItemsListParams,
params: ThreadItemsListParams,
) -> anyhow::Result<i64> {
let params = Some(serde_json::to_value(params)?);
self.send_request("thread/turns/items/list", params).await
self.send_request("thread/items/list", params).await
}

/// Send a `model/list` JSON-RPC request.
Expand Down
12 changes: 6 additions & 6 deletions codex-rs/app-server/tests/suite/v2/thread_read.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ use codex_app_server_protocol::SortDirection;
use codex_app_server_protocol::ThreadForkParams;
use codex_app_server_protocol::ThreadForkResponse;
use codex_app_server_protocol::ThreadItem;
use codex_app_server_protocol::ThreadItemsListParams;
use codex_app_server_protocol::ThreadListParams;
use codex_app_server_protocol::ThreadListResponse;
use codex_app_server_protocol::ThreadNameUpdatedNotification;
Expand All @@ -32,7 +33,6 @@ use codex_app_server_protocol::ThreadSetNameResponse;
use codex_app_server_protocol::ThreadStartParams;
use codex_app_server_protocol::ThreadStartResponse;
use codex_app_server_protocol::ThreadStatus;
use codex_app_server_protocol::ThreadTurnsItemsListParams;
use codex_app_server_protocol::ThreadTurnsListParams;
use codex_app_server_protocol::ThreadTurnsListResponse;
use codex_app_server_protocol::TurnItemsView;
Expand Down Expand Up @@ -1137,7 +1137,7 @@ async fn thread_turns_list_rejects_unmaterialized_loaded_thread() -> Result<()>
}

#[tokio::test]
async fn thread_turns_items_list_returns_unsupported() -> Result<()> {
async fn thread_items_list_returns_unsupported() -> Result<()> {
let server = create_mock_responses_server_repeating_assistant("Done").await;
let codex_home = TempDir::new()?;
create_config_toml(codex_home.path(), &server.uri())?;
Expand All @@ -1146,9 +1146,9 @@ async fn thread_turns_items_list_returns_unsupported() -> Result<()> {
timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??;

let read_id = mcp
.send_thread_turns_items_list_request(ThreadTurnsItemsListParams {
thread_id: "thr_123".to_string(),
turn_id: "turn_456".to_string(),
.send_thread_items_list_request(ThreadItemsListParams {
thread_id: "00000000-0000-4000-8000-000000000123".to_string(),
turn_id: None,
cursor: None,
limit: None,
sort_direction: None,
Expand All @@ -1163,7 +1163,7 @@ async fn thread_turns_items_list_returns_unsupported() -> Result<()> {
assert_eq!(read_err.error.code, -32601);
assert_eq!(
read_err.error.message,
"thread/turns/items/list is not supported yet"
"thread/items/list is not supported yet"
);

Ok(())
Expand Down
Loading