ThreadList与分页
thread/list 看起来像一次数据库查询,实际却是一个跨层投影:App Server 先解析双层 Option 和关系过滤器,再把默认 provider、来源类型、cwd、project、归档状态和 search term 翻译成 StoreListThreadsParams;ThreadStore 返回的每一页还可能被过滤掉,于是处理器必须继续取下一页,直到得到足够多的结果或 cursor 耗尽。最后,持久化摘要被转成协议 Thread,并用内存中已加载 Thread 的状态覆盖 NotLoaded。
本文面向熟悉 Rust enum、Option、async 和分页查询的读者,承接ThreadStart处理流程、AppServer错误码体系和源码仓库目录地图。本文只讲 thread/list,不展开 thread/search、turn/item 分页和 ThreadStore 的具体 SQLite SQL。
读完后应能解释:为什么 limit=0 仍会返回一条结果,为什么空 sourceKinds 默认只看交互线程,为什么 cwd 过滤会触发多次 store 请求,以及列表中的一个正在运行的 sub-agent 如何显示为 Active 而不是磁盘里的旧状态。
1. 列表响应的四层
StoredThread 是 ThreadStore 的持久化摘要,Thread 是协议对象;二者之间没有直接 From 转换,而是要补 fallback provider、cwd、source 和时间字段。enrich_loaded_threads 再把运行时状态叠加到这个快照上,所以列表既不是纯数据库结果,也不是把所有 Thread 都加载进内存。
列表请求只有一个 RPC 出口;watch manager 只负责补运行状态,不负责决定持久化分页顺序。
类型图对应 thread_list_response_inner 中的四次投影:参数解构、过滤器装配、存储摘要转换、最终响应包装。
2. 入口与协议字段
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数:MessageProcessor::process_request 的 ClientRequest::ThreadList 分支。
ClientRequest::ThreadList { params, .. } => {
self.thread_processor.thread_list(params).await
}源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:ThreadRequestProcessor::thread_list。
pub(crate) async fn thread_list(
&self,
params: ThreadListParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
self.thread_list_response_inner(params)
.await
.map(|response| Some(response.into()))
}与 thread/start 不同,列表请求在处理器返回值中直接带回 ClientResponsePayload,没有后台 task,也没有通知。列表 response 的完成边界就是 thread_list_response_inner 返回。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关类型:ThreadListParams、ThreadListResponse。
pub struct ThreadListParams {
pub cursor: Option<String>,
pub limit: Option<u32>,
pub sort_key: Option<ThreadSortKey>,
pub sort_direction: Option<SortDirection>,
pub model_providers: Option<Vec<String>>,
pub source_kinds: Option<Vec<ThreadSourceKind>>,
pub archived: Option<bool>,
pub section_id: Option<Option<String>>,
pub project_id: Option<Option<String>>,
pub cwd: Option<ThreadListCwdFilter>,
pub use_state_db_only: bool,
pub search_term: Option<String>,
pub parent_thread_id: Option<String>,
pub ancestor_thread_id: Option<String>,
}
pub struct ThreadListResponse {
pub data: Vec<Thread>,
pub next_cursor: Option<String>,
pub backwards_cursor: Option<String>,
}section_id 和 project_id 使用 Option<Option<String>>:字段省略表示“不限制”,显式 null 表示“没有 section/project”,字符串才表示精确匹配。next_cursor 用于沿当前排序方向继续,backwards_cursor 则是反向同步的锚点,不能把它们都当作“下一页 token”。
3. 校验与排序
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:thread_list_response_inner。
let ThreadListParams {
cursor,
limit,
sort_key,
sort_direction,
model_providers,
source_kinds,
archived,
section_id,
project_id,
cwd,
use_state_db_only,
search_term,
parent_thread_id,
ancestor_thread_id,
} = params;
if project_id.is_some() && !self.thread_store.supports_projects() {
return Err(unsupported_thread_store_operation("projects"));
}
if let Some(Some(project_id)) = project_id.as_ref() {
if project_id.is_empty() {
return Err(invalid_params("projectId must not be empty"));
}
match self.thread_store.read_project(project_id.clone()).await {
Ok(Some(_)) => {}
Ok(None) => return Err(invalid_params(format!("project not found: {project_id}"))),
Err(ThreadStoreError::Unsupported { operation }) => {
return Err(unsupported_thread_store_operation(operation));
}
Err(err) => return Err(internal_error(format!("failed to read project: {err}"))),
}
}project 过滤先验证存储能力,再验证项目存在性;因此“没有这个 project”和“当前 ThreadStore 不支持 project”是两个不同错误。空字符串也在访问 Store 前被拒绝。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:thread_list_response_inner 的关系过滤与 cursor 解析。
let cwd_filters = normalize_thread_list_cwd_filters(cwd)?;
let relation_filter = match (parent_thread_id, ancestor_thread_id) {
(Some(_), Some(_)) => {
return Err(invalid_request(
"parentThreadId and ancestorThreadId are mutually exclusive",
));
}
(Some(parent_thread_id), None) => Some(StoreThreadRelationFilter::DirectChildrenOf(
ThreadId::from_string(&parent_thread_id)
.map_err(|err| invalid_request(format!("invalid parent thread id: {err}")))?,
)),
(None, Some(ancestor_thread_id)) => Some(StoreThreadRelationFilter::DescendantsOf(
ThreadId::from_string(&ancestor_thread_id)
.map_err(|err| invalid_request(format!("invalid ancestor thread id: {err}")))?,
)),
(None, None) => None,
};
let requested_page_size = limit
.map(|value| value as usize)
.unwrap_or(THREAD_LIST_DEFAULT_LIMIT)
.clamp(1, THREAD_LIST_MAX_LIMIT);
let store_sort_key = match sort_key.unwrap_or(ThreadSortKey::CreatedAt) {
ThreadSortKey::CreatedAt => StoreThreadSortKey::CreatedAt,
ThreadSortKey::UpdatedAt => StoreThreadSortKey::UpdatedAt,
ThreadSortKey::RecencyAt => StoreThreadSortKey::RecencyAt,
ThreadSortKey::SectionPosition => StoreThreadSortKey::SectionPosition,
};
let sort_direction = sort_direction.unwrap_or(match store_sort_key {
StoreThreadSortKey::SectionPosition => SortDirection::Asc,
StoreThreadSortKey::CreatedAt
| StoreThreadSortKey::UpdatedAt
| StoreThreadSortKey::RecencyAt => SortDirection::Desc,
});默认按 created_at 降序;section position 默认升序。limit 被限制到 1~100,不能通过 0 请求空页,也不能让单次调用绕过服务端上限。父线程和祖先线程关系过滤互斥:前者只返回直接子节点,后者返回任意深度后代且排除祖先自身。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:normalize_thread_list_cwd_filters。
fn normalize_thread_list_cwd_filters(
cwd: Option<ThreadListCwdFilter>,
) -> Result<Option<Vec<PathBuf>>, JSONRPCErrorError> {
let Some(cwd) = cwd else { return Ok(None); };
let cwds = match cwd {
ThreadListCwdFilter::One(cwd) => vec![cwd],
ThreadListCwdFilter::Many(cwds) => cwds,
};
let mut normalized_cwds = Vec::with_capacity(cwds.len());
for cwd in cwds {
let cwd = AbsolutePathBuf::relative_to_current_dir(cwd.as_str())
.map(AbsolutePathBuf::into_path_buf)
.map_err(|err| {
invalid_params(format!("invalid thread/list cwd filter `{cwd}`: {err}"))
})?;
normalized_cwds.push(cwd);
}
Ok(Some(normalized_cwds))
}协议允许一个字符串或字符串数组,但 Store 收到的永远是绝对 PathBuf。相对 cwd 按服务进程当前目录解析;路径格式错误在查询前变成 invalid_params。真正匹配时还会调用 paths_match_after_normalization,所以不能只比较原始字符串。
4. 来源与默认值
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:list_threads_common。
let model_provider_filter = match model_providers {
Some(providers) => {
if providers.is_empty() { None } else { Some(providers) }
}
None if relation_filter.is_some() => None,
None => Some(vec![self.config.model_provider_id.clone()]),
};
let (allowed_sources_vec, source_kind_filter) =
if relation_filter.is_some() && source_kinds.is_none() {
(Vec::new(), None)
} else {
compute_source_filters(source_kinds)
};省略 modelProviders 时,普通列表只读当前配置 provider;显式空数组表示所有 provider。关系查询省略 provider 时不再偷偷限制为当前 provider,否则父线程与子线程使用不同 provider 时会被截断。来源过滤被拆成两部分:可下推给 Store 的 allowed_sources,以及处理器对 sub-agent 变体执行的 source_kind_filter。
5. 多页补取与过滤
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:list_threads_common。
let mut cursor_obj = cursor;
let mut last_cursor = cursor_obj.clone();
let mut remaining = requested_page_size;
let mut items = Vec::with_capacity(requested_page_size);
let mut next_cursor = None;
while remaining > 0 {
let page_size = remaining.min(THREAD_LIST_MAX_LIMIT);
let page = self.thread_store.list_threads(StoreListThreadsParams {
page_size,
cursor: cursor_obj.clone(),
sort_key,
sort_direction: store_sort_direction,
allowed_sources: allowed_sources.to_vec(),
model_providers: model_provider_filter.clone(),
cwd_filters: cwd_filters.clone(),
archived,
section: section_id.clone(),
project_id: project_id.clone(),
search_term: search_term.clone(),
use_state_db_only,
relation_filter,
}).await.map_err(thread_store_list_error)?;
let mut filtered = Vec::with_capacity(page.items.len());
for it in page.items {
let source = with_thread_spawn_agent_metadata(
it.source.clone(), it.agent_nickname.clone(), it.agent_role.clone(),
);
if source_kind_filter.as_ref().is_none_or(|filter| source_kind_matches(&source, filter))
&& cwd_filters.as_ref().is_none_or(|expected_cwds| {
expected_cwds.iter().any(|expected_cwd| {
path_utils::paths_match_after_normalization(&it.cwd, expected_cwd)
})
})
{
filtered.push(it);
if filtered.len() >= remaining { break; }
}
}
items.extend(filtered);
remaining = requested_page_size.saturating_sub(items.len());
next_cursor = page.next_cursor;
if remaining == 0 { break; }
let Some(cursor_val) = next_cursor.clone() else { break; };
if last_cursor.as_ref() == Some(&cursor_val) {
next_cursor = None;
break;
}
last_cursor = Some(cursor_val.clone());
cursor_obj = Some(cursor_val);
}分页的关键不是“取一页再截断”,而是“以 Store cursor 为边界持续补取”。当 Store 返回的页面因来源或 cwd 被过滤掉,处理器继续使用 page.next_cursor,直到凑满 limit。last_cursor 防止异常 Store 重复返回同一 cursor 导致无限循环;如果没有更多可读页面,response 的 next_cursor 为 None。
这个循环是列表实现最容易遗漏的部分:Store page 的数量和对客户端可见的数量不是同一个数量。
6. 反向游标
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:thread_backwards_cursor_for_sort_key。
fn thread_backwards_cursor_for_sort_key(
thread: &StoredThread,
sort_key: StoreThreadSortKey,
sort_direction: SortDirection,
) -> Option<String> {
if sort_key == StoreThreadSortKey::SectionPosition {
let position = match sort_direction {
SortDirection::Asc => thread.section_position?.checked_add(1)?,
SortDirection::Desc => thread.section_position?.checked_sub(1)?,
};
return Some(format!("{position}|{}", thread.thread_id));
}
let timestamp = match sort_key {
StoreThreadSortKey::CreatedAt => thread.created_at,
StoreThreadSortKey::UpdatedAt => thread.updated_at,
StoreThreadSortKey::RecencyAt => thread.recency_at,
StoreThreadSortKey::SectionPosition => unreachable!(),
};
let timestamp = match sort_direction {
SortDirection::Asc => timestamp.checked_add_signed(ChronoDuration::milliseconds(1))?,
SortDirection::Desc => timestamp.checked_sub_signed(ChronoDuration::milliseconds(1))?,
};
Some(timestamp.to_rfc3339_opts(SecondsFormat::Millis, true))
}时间排序用毫秒偏移构造反向锚点,使切换方向时包含当前页边界;section 排序则用 position|thread_id。这两个格式是 Store cursor 的实现细节,客户端只能原样保存和回传,不能自行解析后拼接。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数:thread_list_response_inner 的 response 组装。
let backwards_cursor = stored_threads.first().and_then(|thread| {
thread_backwards_cursor_for_sort_key(thread, store_sort_key, sort_direction)
});
let fallback_provider = self.config.model_provider_id.clone();
let mut data = Vec::with_capacity(stored_threads.len());
for stored_thread in stored_threads {
let (thread, _) = thread_from_stored_thread(
stored_thread,
fallback_provider.as_str(),
&self.config.cwd,
);
data.push(thread);
}反向 cursor 必须在消费 StoredThread 前从第一页首项生成。之后 thread_from_stored_thread 把摘要转换成协议对象,缺失 provider 时使用当前配置 provider,缺失或非法 cwd 时使用 server cwd fallback。排序在 Store 层完成,转换过程不能重新排序。
7. 运行时状态
源码文件:codex-rs/app-server/src/request_processors/thread_enrichment.rs
相关函数:enrich_loaded_threads。
pub(super) async fn enrich_loaded_threads<T>(
thread_manager: &ThreadManager,
thread_watch_manager: &ThreadWatchManager,
threads: &mut [T],
mut as_thread: impl FnMut(&mut T) -> &mut Thread,
) {
let statuses = thread_watch_manager
.loaded_statuses_for_threads(
threads.iter_mut().map(&mut as_thread).map(|thread| thread.id.clone()),
)
.await;
futures::future::join_all(threads.iter_mut().map(as_thread).map(|thread| {
let statuses = &statuses;
async move {
let watched_status = statuses.get(&thread.id);
if let Some(status) = watched_status {
thread.status = status.clone();
}
if !matches!(
&thread.source,
SessionSource::SubAgent(SubAgentSource::ThreadSpawn { .. })
) || matches!(watched_status, Some(ThreadStatus::NotLoaded)) {
return;
}
let Ok(thread_id) = ThreadId::from_string(&thread.id) else { return; };
let Ok(loaded_thread) = thread_manager.get_thread(thread_id).await else { return; };
match loaded_thread.agent_status().await {
AgentStatus::Running => {
if watched_status.is_none() {
thread.status = resolve_thread_status(ThreadStatus::Idle, true);
}
}
AgentStatus::Errored(_) => thread.status = ThreadStatus::SystemError,
AgentStatus::Shutdown | AgentStatus::NotFound => {
thread.status = ThreadStatus::NotLoaded;
return;
}
AgentStatus::PendingInit | AgentStatus::Interrupted | AgentStatus::Completed(_) => {
if watched_status.is_none() { thread.status = ThreadStatus::Idle; }
}
}
let config_snapshot = loaded_thread.config_snapshot().await;
thread.can_accept_direct_input = Some(can_accept_direct_input(
loaded_thread.multi_agent_version(),
&config_snapshot.session_source,
));
}
})).await;
}普通历史 Thread 只接受 watch manager 的已有状态;只有 ThreadSpawn sub-agent 才会回到 ThreadManager 读取实时 agent_status。这避免为了渲染普通列表而加载每个 Thread,同时让正在运行但尚未注册 watch 的 sub-agent 显示 active。Shutdown/NotFound 会降回 NotLoaded;Errored 显示 SystemError;成功读取配置快照后才计算 can_accept_direct_input。
8. 失败与状态库
use_state_db_only 控制的是 Store 是否扫描 JSONL rollout 修复摘要。false 允许从 rollout 修复 cwd、preview、git 和时间字段;true 只读 SQLite 中已有值。两者都会经过同一套协议转换和运行时 enrichment,因此 state DB-only 不是另一种 response schema,而是数据来源选择。
常见失败可以按层定位:非法 cursor 或关系 id 在参数层返回 -32600;项目不存在在 project 校验层返回 invalid params;Store I/O 错误经过 thread_store_list_error 变成内部错误;Store 返回重复 cursor 时安全停止并把 next_cursor 置空。列表为空不是错误,空数据和最后一页都以 data=[]、next_cursor=null 表示。
9. 测试证据
源码文件:codex-rs/app-server/tests/suite/v2/thread_list.rs
测试:thread_list_pagination_next_cursor_none_on_last_page。
let ThreadListResponse { data: data1, next_cursor: cursor1, .. } =
list_threads(&mut mcp, None, Some(2), Some(vec!["mock_provider".to_string()]), None, None)
.await?;
assert_eq!(data1.len(), 2);
let cursor1 = cursor1.expect("expected nextCursor on first page");
let ThreadListResponse { data: data2, next_cursor: cursor2, .. } =
list_threads(&mut mcp, Some(cursor1), Some(2), Some(vec!["mock_provider".to_string()]), None, None)
.await?;
assert!(data2.len() <= 2);
assert_eq!(cursor2, None);测试创建三个 rollout,以 limit=2 请求两页,断言第一页有 cursor、第二页没有 cursor,并检查 preview、provider、时间、cwd、source 和初始 NotLoaded。它证明 Store cursor 到协议响应的分页闭环,不证明正在运行 Thread 的 enrichment。
测试:thread_list_state_db_only_returns_sqlite_without_jsonl_repair 先让 rollout 修复 SQLite,再手工把 SQLite cwd 改成 stale path。use_state_db_only=true 必须返回 stale path;false 则重新扫描 rollout 并因 cwd 过滤不再匹配而返回空列表。这个测试证明开关确实改变数据来源,不能把 state DB-only 解释成“更快但等价”。
测试:thread_list_relation_filters_reject_invalid_requests 传入非法 parent id,读取同 request id 的 error response,并断言 code 为 -32600;它证明关系过滤的 ID 校验发生在 Store 查询前。
当第一页被二次过滤时,处理器不会让客户端看到一个过小的“最终页”;它会继续查询,再做状态 enrichment。
可执行验证:
cargo test -p codex-app-server --test all thread_list_pagination_next_cursor_none_on_last_page
cargo test -p codex-app-server --test all thread_list_state_db_only_returns_sqlite_without_jsonl_repair
cargo test -p codex-app-server --test all thread_list_relation_filters_reject_invalid_requests10. 学习检查
复述一次带 cwd 和 source filter 的请求:ThreadListParams 如何变成 ThreadListFilters,哪些条件下由 Store 下推,哪些条件由处理器二次过滤,过滤掉一页后 cursor 如何推进,最终又在哪一步覆盖已加载 Thread 的状态?
故障练习:当客户端请求 limit=2 却收到一条结果和 next_cursor=null,分别从“确实只有一条匹配项”“第二页被过滤为空”“Store 重复 cursor 被保护停止”三个方向设计只读日志或测试断言。不要只看 response 数量,要同时检查 Store page、过滤器和 cursor 演进。
下一篇ThreadRead历史加载会继续处理列表之后的单 Thread 读取:它将解释 includeTurns、legacy rollout 和分页 history 如何改变同一个 Thread 快照的内容。
