Skip to content

请求追踪与 Analytics

沿 App Server 请求追踪父 trace 的选择、span 的异步所有权、Analytics 的请求关联与事实汇合,解释字段来源、耗时边界和投递缺口。

基于rust-v0.150.0
CodexRustAppServerTracingAnalytics

请求追踪与 Analytics ​

一次 turn/start 能在 trace 中找到,分析事件却没有出现;请求处理函数已经返回,诊断里的在途计数仍未归零; flush().await 已经结束,后端还是查不到事件。这些现象不能用一句“异步日志有延迟”解释。 App Server 中至少存在三种不同的完成边界:请求 span 的引用释放、Analytics 的事实汇合、事件的最终投递。

本文面向熟悉 Rust async、Arc 和 Result、准备追踪一次真实请求的读者。 MessageProcessor读循环 解释请求进入处理器的路径, AppServer错误码体系 区分 RPC 返回与运行中失败, TurnTiming与Metadata 解释 Core 提供的时间与轮次元数据。 这里把这些入口接到 tracing 和 Analytics 的消费者,不展开全部遥测 exporter、工具统计或后台 worker 实现。

tracing span 表示一段带属性和父级关系的工作范围;OpenTelemetry(OTel) 为它提供可传播的 Trace/Span ID 及导出能力;Analytics Fact 是进程内交给归约器的事实,归约器将多条事实组合成可投递的分析事件。 它们有不同的数据结构与筛选条件。读完应能从 RPC 的 trace 字段找到 Core 子链,从请求 ID 找到分析关联, 并判断一份缺失数据是在采集、归约、认证筛选、排队还是发送阶段消失。

1. 观测坐标 ​

1.1 请求属性 ​

源码文件:codex-rs/app-server/src/message_processor.rs

相关函数/类型:MessageProcessor::process_request(L609–L619,局部节选)

rust
// 作者注:请求 ID 先加入连接范围,再构造 span 与请求上下文;尚未执行 typed 业务分派。
// ...
let request_id = ConnectionRequestId {
    connection_id,
    request_id: request.id.clone(),
};
let request_span =
    crate::app_server_tracing::request_span(&request, transport, connection_id, &session);
let request_trace = request.trace.as_ref().map(|trace| W3cTraceContext {
    traceparent: trace.traceparent.clone(),
    tracestate: trace.tracestate.clone(),
});
let request_context = RequestContext::new(request_id.clone(), request_span, request_trace);
// ...

process_request 在 typed 业务解码之前创建上下文,因此一个最终被拒绝的请求也可能已经创建 span。 ConnectionRequestId 保留连接与请求两个部分,避免两个客户端都使用 ID 7 时混为一条请求。 后面的 Analytics reducer 也采用这个组合键,但最终分析事件主要通过 Thread/Turn ID 关联。

源码文件:codex-rs/app-server/src/app_server_tracing.rs

相关函数/类型:app_server_request_span_template / record_client_info(L94–L123,摘录)

rust
// 作者注:只预声明列出的属性,turn.id 留待接纳输入后回填;这里没有通用 RPC 成功或耗时字段。
fn app_server_request_span_template(
    method: &str,
    transport: &'static str,
    request_id: &impl std::fmt::Display,
    connection_id: ConnectionId,
) -> Span {
    info_span!(
        "app_server.request",
        otel.kind = "server",
        otel.name = method,
        rpc.system = "jsonrpc",
        rpc.method = method,
        rpc.transport = transport,
        rpc.request_id = %request_id,
        app_server.connection_id = %connection_id,
        app_server.api_version = "v2",
        app_server.client_name = field::Empty,
        app_server.client_version = field::Empty,
        turn.id = field::Empty,
    )
}

fn record_client_info(span: &Span, client_name: Option<&str>, client_version: Option<&str>) {
    if let Some(client_name) = client_name {
        span.record("app_server.client_name", client_name);
    }
    if let Some(client_version) = client_version {
        span.record("app_server.client_version", client_version);
    }
}

tracing 内部名称固定为 app_server.request,OTel 导出名称由 otel.name 使用实际 method。 turn.id 先预声明为 field::Empty,等 Core 接纳输入后才回填。 这个模板没有 rpc.duration_ms、结果枚举或统一的错误状态字段,不能从一篇通用 tracing 教程推定它们一定存在。 app_server.api_version = v2 也是这里写入的标签,不是一次协议协商的结果。

标识或字段产生位置用途与边界
rpc.request_id客户端 RPC ID 的 Display便于读日志;整数 7 和字符串 7 的显示值可能相同
app_server.connection_idtransport 分配与请求 ID 一起定位本次 RPC
Trace IDOTel 上下文关联同一条跨任务父链
Span ID每段工作自己的 span下游生成自己的 ID,并引用上一段为父级
turn.idCore 接纳后回填将 RPC 观测与工作轮次连接
product_client_idAnalytics 的连接/线程来源产品归因字段,与 Trace ID 不是一回事

内部 map 的 RequestId 枚举仍区分整数与字符串,不能把格式化属性直接当成无损主键。 同理,一次 RPC 可能只是接纳某个 Turn 的输入,而一个 Turn 可以包含多次模型采样;两者的耗时并不相等。

1.2 身份来源 ​

源码文件:codex-rs/app-server/src/app_server_tracing.rs

相关函数/类型:initialize_client_info / client_name(L144–L170,摘录)

rust
// 作者注:initialize 的客户端身份先从参数提取;其他请求使用连接已有状态。
fn client_name<'a>(
    initialize_client_info: Option<&'a InitializeParams>,
    session: &'a ConnectionSessionState,
) -> Option<&'a str> {
    if let Some(params) = initialize_client_info {
        return Some(params.client_info.name.as_str());
    }
    session.app_server_client_name()
}

fn client_version<'a>(
    initialize_client_info: Option<&'a InitializeParams>,
    session: &'a ConnectionSessionState,
) -> Option<&'a str> {
    if let Some(params) = initialize_client_info {
        return Some(params.client_info.version.as_str());
    }
    session.client_version()
}

fn initialize_client_info(request: &JSONRPCRequest) -> Option<InitializeParams> {
    if request.method != "initialize" {
        return None;
    }
    let params = request.params.clone()?;
    serde_json::from_value(params).ok()
}

initialize 请求尚没有已建立的连接身份,所以 helper 尝试从初始化参数取 name/version。 其他请求使用连接中的已有身份。这个提取发生在实际 initialize handler 验证之前, 所以 span 中有 client_name 不代表握手已经成功,更不能把它当成经过认证的用户身份。

span 模板只选取这些属性,没有在此序列化完整 params。 这是一处字段选择,不是全局脱敏器:其他日志和 Fact 有各自的数据路径,后文会继续区分。

2. 父链选择 ​

2.1 两种入口 ​

源码文件:codex-rs/protocol/src/protocol.rs

相关函数/类型:W3cTraceContext(L202–L210,摘录)

rust
// 作者注:载体只是两个可选字符串,Rust 类型本身不保证 W3C 内容有效。
#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct W3cTraceContext {
    #[serde(default, skip_serializing_if = "Option::is_none")]
    #[ts(optional)]
    pub traceparent: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    #[ts(optional)]
    pub tracestate: Option<String>,
}

W3C carrier 只是两个可选字符串,serde 成功并不说明内容合法。 例如测试中的 00-…-…-01 包含 32 个十六进制字符的 Trace ID 和 16 个十六进制字符的 Span ID, 但实际有效性由 TraceContextPropagator 判断,不由这个结构体检查。

源码文件:codex-rs/app-server/src/app_server_tracing.rs

相关函数/类型:request_span(L24–L55,摘录)

rust
// 作者注:显式 carrier 只有带 traceparent 时才进入父链选择,tracestate 单独出现不满足条件。
pub(crate) fn request_span(
    request: &JSONRPCRequest,
    transport: &AppServerTransport,
    connection_id: ConnectionId,
    session: &ConnectionSessionState,
) -> Span {
    let initialize_client_info = initialize_client_info(request);
    let method = request.method.as_str();
    let span = app_server_request_span_template(
        method,
        transport_name(transport),
        &request.id,
        connection_id,
    );

    record_client_info(
        &span,
        client_name(initialize_client_info.as_ref(), session),
        client_version(initialize_client_info.as_ref(), session),
    );

    let parent_trace = request.trace.as_ref().and_then(|trace| {
        trace.traceparent.as_ref()?;
        Some(W3cTraceContext {
            traceparent: trace.traceparent.clone(),
            tracestate: trace.tracestate.clone(),
        })
    });
    attach_parent_context(&span, method, &request.id, parent_trace.as_ref());

    span
}

request_span 只把含 traceparent 的请求 carrier 交给父链选择函数。 若请求只给 tracestate,and_then 返回 None,随后走缺省路径。 如果有一个内容无效的 traceparent 字符串,则它仍进入显式 carrier 分支,再由解析器拒绝。 这两种输入在 JSON 层看起来都“有 trace”,控制流却不同。

源码文件:codex-rs/app-server/src/app_server_tracing.rs

相关函数/类型:typed_request_span(L62–L83,摘录)

rust
// 作者注:typed 入口共享 span 形状,transport 为 in-process,父链选择不从 JSON request.trace 读取。
pub(crate) fn typed_request_span(
    request: &ClientRequest,
    connection_id: ConnectionId,
    session: &ConnectionSessionState,
) -> Span {
    let method = request.method_name();
    let span = app_server_request_span_template(method, "in-process", request.id(), connection_id);

    let client_info = initialize_client_info_from_typed_request(request);
    record_client_info(
        &span,
        client_info
            .map(|(client_name, _)| client_name)
            .or(session.app_server_client_name()),
        client_info
            .map(|(_, client_version)| client_version)
            .or(session.client_version()),
    );

    attach_parent_context(&span, method, request.id(), /*parent_trace*/ None);
    span
}

嵌入式 typed 请求沿用相同属性模板,rpc.transport 为 in-process。 这里没有 JSON envelope 的 request.trace 参数,调用 attach_parent_context 时传 None。 不能因为字段形状相同,就假定两种入口使用完全相同的显式父链来源。

2.2 显式与缺省 ​

源码文件:codex-rs/app-server/src/app_server_tracing.rs

相关函数/类型:attach_parent_context(L125–L142,摘录)

rust
// 作者注:有显式 carrier 的分支即使解析失败,也不会转入环境 fallback 分支。
fn attach_parent_context(
    span: &Span,
    method: &str,
    request_id: &impl std::fmt::Display,
    parent_trace: Option<&W3cTraceContext>,
) {
    if let Some(trace) = parent_trace {
        if !set_parent_from_w3c_trace_context(span, trace) {
            tracing::warn!(
                rpc_method = method,
                rpc_request_id = %request_id,
                "ignoring invalid inbound request trace carrier"
            );
        }
    } else if let Some(context) = traceparent_context_from_env() {
        set_parent_from_context(span, context);
    }
}

这段 if/else if 给出了真正的优先级:显式 carrier 存在时只尝试它;解析失败记录 warning, 不会再使用环境变量。没有显式 carrier 时才查询环境缓存。 没有设置新父级的路径也不能一律命名为“新根 trace”,因为 span 创建时还可能已经继承当前执行范围的上下文。

下图从已经创建的 span 出发,只描述这个 helper 是否尝试设置新父级。

  • 显式无效 carrier 与缺少 carrier 不等价。
  • tracestate 单独存在在 socket helper 中属于缺少 traceparent。
  • typed 入口从缺省分支进入;是否最终有可导出的父链,还受 subscriber/tracer 是否存在影响。

源码文件:codex-rs/otel/src/trace_context.rs

相关函数/类型:set_parent_from_w3c_trace_context / context_from_trace_headers(L106–L121、L129–L145,摘录)

rust
// 作者注:bool 反映 carrier 能否解析;set_parent 的返回结果在内部被忽略,仍需验证实际导出的父链。
pub fn context_from_w3c_trace_context(trace: &W3cTraceContext) -> Option<Context> {
    context_from_trace_headers(trace.traceparent.as_deref(), trace.tracestate.as_deref())
}

pub fn set_parent_from_w3c_trace_context(span: &Span, trace: &W3cTraceContext) -> bool {
    if let Some(context) = context_from_w3c_trace_context(trace) {
        set_parent_from_context(span, context);
        true
    } else {
        false
    }
}

pub fn set_parent_from_context(span: &Span, context: Context) {
    let _ = span.set_parent(context);
}

// ...

pub(crate) fn context_from_trace_headers(
    traceparent: Option<&str>,
    tracestate: Option<&str>,
) -> Option<Context> {
    let traceparent = traceparent?;
    let mut headers = HashMap::new();
    headers.insert("traceparent".to_string(), traceparent.to_string());
    if let Some(tracestate) = tracestate {
        headers.insert("tracestate".to_string(), tracestate.to_string());
    }

    let context = TraceContextPropagator::new().extract(&headers);
    if !context.span().span_context().is_valid() {
        return None;
    }
    Some(context)
}

解析器提取后检查 SpanContext::is_valid。set_parent_from_w3c_trace_context 返回 true 表示 carrier 解析出了有效 Context;内部 span.set_parent 的返回值被忽略。 因此这个布尔值不应被写成“远端追踪平台已经收到正确父子关系”的证明。 实际父链还要查看导出的 SpanData,后面的测试正是这样验证。

源码文件:codex-rs/otel/src/trace_context.rs

相关函数/类型:traceparent_context_from_env / load_traceparent_context(L20–L22、L123–L127、L147–L161,摘录)

rust
// 作者注:OnceLock 连 None 结果也缓存,首次读取之后改环境变量不会重新解析。
const TRACEPARENT_ENV_VAR: &str = "TRACEPARENT";
const TRACESTATE_ENV_VAR: &str = "TRACESTATE";
static TRACEPARENT_CONTEXT: OnceLock<Option<Context>> = OnceLock::new();

// ...

pub fn traceparent_context_from_env() -> Option<Context> {
    TRACEPARENT_CONTEXT
        .get_or_init(load_traceparent_context)
        .clone()
}

// ...

fn load_traceparent_context() -> Option<Context> {
    let traceparent = env::var(TRACEPARENT_ENV_VAR).ok()?;
    let tracestate = env::var(TRACESTATE_ENV_VAR).ok();

    match context_from_trace_headers(Some(&traceparent), tracestate.as_deref()) {
        Some(context) => {
            debug!("TRACEPARENT detected; continuing trace from parent context");
            Some(context)
        }
        None => {
            warn!("TRACEPARENT is set but invalid; ignoring trace context");
            None
        }
    }
}

OnceLock<Option<Context>> 缓存首次查询的结果,包含 None。 它不是每次 RPC 都重新读环境,也不是一段必定在进程刚启动时执行的 eager 初始化。 一旦首次查询发生,后来再修改环境变量不会让这个函数重新解析。

2.3 下游载体 ​

源码文件:codex-rs/otel/src/trace_context.rs

相关函数/类型:span_w3c_trace_context(L33–L50,摘录)

rust
// 作者注:下游 carrier 优先来自当前请求 span,它有自己的 span ID,并可能合并配置的 tracestate。
pub fn span_w3c_trace_context(span: &Span) -> Option<W3cTraceContext> {
    let context = span.context();
    if !context.span().span_context().is_valid() {
        return None;
    }

    let mut headers = HashMap::new();
    TraceContextPropagator::new().inject_context(&context, &mut headers);
    let tracestate = headers.remove("tracestate");
    let configured_tracestate_guard = tracestate_entries()
        .read()
        .unwrap_or_else(std::sync::PoisonError::into_inner);

    Some(W3cTraceContext {
        traceparent: headers.remove("traceparent"),
        tracestate: merge_tracestate_entries(tracestate.as_deref(), &configured_tracestate_guard),
    })
}

向下游发出 carrier 时,优先从请求自身的有效 span 生成 Trace/Span ID。 下游应成为当前 RPC span 的子链,而不是始终挂在最初客户端 span 下面。 如果当前 span 没有有效 OTel Context,这个函数返回 None;请求上下文还有单独的原始 carrier fallback。

源码文件:codex-rs/otel/src/trace_context.rs

相关函数/类型:merge_tracestate_entries(L167–L196,摘录)

rust
// 作者注:配置成员按确定顺序 upsert;非法原 tracestate 退为空,非法合并记录 warning。
fn merge_tracestate_entries(
    tracestate: Option<&str>,
    configured_entries: &BTreeMap<String, BTreeMap<String, String>>,
) -> Option<String> {
    let mut trace_state = tracestate
        .and_then(|tracestate| match TraceState::from_str(tracestate) {
            Ok(trace_state) => Some(trace_state),
            Err(err) => {
                warn!("ignoring invalid tracestate while propagating trace context: {err}");
                None
            }
        })
        .unwrap_or_default();

    // TraceState::insert places members at the front. Reverse iteration keeps
    // deterministic map order while upserting fields inside configured members.
    for (key, fields) in configured_entries.iter().rev() {
        let value = merge_tracestate_member_fields(trace_state.get(key), fields);
        trace_state = match trace_state.insert(key.clone(), value) {
            Ok(trace_state) => trace_state,
            Err(err) => {
                warn!("ignoring configured tracestate while propagating trace context: {err}");
                break;
            }
        };
    }

    let tracestate = trace_state.header();
    (!tracestate.is_empty()).then_some(tracestate)
}

tracestate 不再只是原样镜像。当前实现读取进程级配置成员,对原有值解析并按确定顺序合并, 通过逆序遍历配合 TraceState::insert 保持成员顺序。 非法原值可退为空状态,配置合并失败记录 warning 并停止继续插入。 这里既不会给无效 span 凭空生成 Trace ID,也不能被概括为“所有原始 tracestate 字节永远不变”。

源码文件:codex-rs/otel/src/trace_context.rs

相关函数/类型:parses_valid_w3c_trace_context / missing_traceparent_returns_none(L351–L368、L377–L385,摘录)

rust
// 作者注:合法 fixture 验证远端 Trace/Span ID;只有 tracestate 的 fixture 不产生有效上下文。
fn parses_valid_w3c_trace_context() {
    let trace_id = "00000000000000000000000000000001";
    let span_id = "0000000000000002";
    let context = context_from_w3c_trace_context(&W3cTraceContext {
        traceparent: Some(format!("00-{trace_id}-{span_id}-01")),
        tracestate: None,
    })
    .expect("trace context");

    let span = context.span();
    let span_context = span.span_context();
    assert_eq!(
        span_context.trace_id(),
        TraceId::from_hex(trace_id).unwrap()
    );
    assert_eq!(span_context.span_id(), SpanId::from_hex(span_id).unwrap());
    assert!(span_context.is_remote());
}

// ...

#[test]
fn missing_traceparent_returns_none() {
    assert!(
        context_from_w3c_trace_context(&W3cTraceContext {
            traceparent: None,
            tracestate: Some("vendor=value".to_string()),
        })
        .is_none()
    );

有效 fixture 断言 Trace ID、Span ID 和 remote 标记;缺少 traceparent 的 fixture 返回 None。 另一个 invalid_traceparent_returns_none 验证无效字符串。 这些测试证明 carrier 提取,不覆盖所有 App Server 入口和 exporter;真正的跨队列父链要继续跟到消费者。

3. 引用与异步 ​

3.1 上下文所有者 ​

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:ConnectionRequestId / RequestContext(L44–L89,摘录)

rust
// 作者注:Context 的 clone 共享一个 GaugeGuard;span 方法只克隆 Span,不复制诊断 guard。
static IN_FLIGHT_REQUESTS: Gauge = Gauge::new("app.requests.in_flight");
static PENDING_SERVER_REQUESTS: Gauge = Gauge::new("app.server_requests.pending");

/// Stable identifier for a client request scoped to a transport connection.
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
pub(crate) struct ConnectionRequestId {
    pub(crate) connection_id: ConnectionId,
    pub(crate) request_id: RequestId,
}

/// Trace data we keep for an incoming request until we send its final
/// response or error.
#[derive(Clone)]
pub(crate) struct RequestContext {
    request_id: ConnectionRequestId,
    span: Span,
    parent_trace: Option<W3cTraceContext>,
    _diagnostics_guard: Arc<GaugeGuard>,
}

impl RequestContext {
    pub(crate) fn new(
        request_id: ConnectionRequestId,
        span: Span,
        parent_trace: Option<W3cTraceContext>,
    ) -> Self {
        Self {
            request_id,
            span,
            parent_trace,
            _diagnostics_guard: Arc::new(IN_FLIGHT_REQUESTS.track()),
        }
    }

    pub(crate) fn request_trace(&self) -> Option<W3cTraceContext> {
        span_w3c_trace_context(&self.span).or_else(|| self.parent_trace.clone())
    }

    pub(crate) fn span(&self) -> Span {
        self.span.clone()
    }

    fn record_turn_id(&self, turn_id: &str) {
        self.span.record("turn.id", turn_id);
    }
}

RequestContext 保存的是 span、原始父 carrier 和一个共享的诊断 guard。 克隆 Context 会增加 Arc<GaugeGuard> 的引用计数,不会重新调用 track()。 span() 则只克隆 span,所以 span 生命周期、guard 生命周期和 map 中是否还有条目可以在某些时刻不同。

request_trace() 优先序列化当前请求 span,失败才返回最初保存的 carrier。 后者没有在 getter 中重新校验,不能因为该方法返回 Some 就宣称本地已经生成了有效 trace。 这条 fallback 允许没有本地有效 span 时继续传递已有载体,实际接收端仍会解析。

下图画出引用关系。队列 future 不靠“当前线程刚好还在 span 里”来维持关联,而显式持有 context/span。

观察值实际测量对象
request_contexts.len()仍在等待响应处理的 map 项
app.requests.in_flight尚未释放共享 guard 的请求上下文
span 生命周期span 及其相关引用所维持的工作范围

这些量都不是“客户端已经收到多少成功结果”的计数。

源码文件:codex-rs/diagnostics/src/lib.rs

相关函数/类型:Gauge::track / GaugeGuard(L31–L44、L56–L65,摘录)

rust
// 作者注:创建 guard 加一,最后释放时减一;饱和减法避免下溢。
/// Decrements the gauge without allowing an underflow.
pub fn decrement(&self) {
    let _ = self
        .value
        .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |value| {
            Some(value.saturating_sub(1))
        });
}

/// Increments the gauge for the lifetime of the returned guard.
pub fn track(&'static self) -> GaugeGuard {
    self.increment();
    GaugeGuard { gauge: self }
}

// ...

/// Decrements its gauge when the measured object is dropped.
pub struct GaugeGuard {
    gauge: &'static Gauge,
}

impl Drop for GaugeGuard {
    fn drop(&mut self) {
        self.gauge.decrement();
    }
}

guard 创建时加一,释放时以饱和减法减一,原子操作使用 relaxed ordering。 它提供进程诊断计数,不构成业务状态的同步屏障。 如果最后一份 Context 仍被异步任务捕获,map 项已经取走也不必立即归零。

源码文件:codex-rs/app-server/src/request_processors/diagnostics.rs

相关函数/类型:read_server_diagnostics(L5–L23,摘录)

rust
// 作者注:消费者读进程诊断快照;这一 gauge 不等同于自动发送的 OTel metric。
pub(crate) fn read_server_diagnostics() -> ServerDiagnosticsResponse {
    let diagnostics = codex_diagnostics::snapshot();

    ServerDiagnosticsResponse {
        process: ServerDiagnosticsProcess {
            id: diagnostics.process.id,
            resident_memory_bytes: diagnostics.process.resident_memory_bytes,
            physical_footprint_bytes: diagnostics.process.physical_footprint_bytes,
        },
        gauges: diagnostics
            .gauges
            .into_iter()
            .map(|gauge| ServerDiagnosticsGauge {
                name: gauge.name.to_string(),
                value: gauge.value,
            })
            .collect(),
    }
}

读取方把 codex_diagnostics::snapshot() 转成 ServerDiagnosticsResponse。 因此要排查在途数量,应该查看这个诊断面与上下文所有者;不能先假设该 gauge 已作为 OTel metric 上传到后端。

3.2 排队与持有 ​

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:register_request_context / connection_closed / take_request_context(L232–L245、L268–L274,摘录)

rust
// 作者注:登记、按连接清理与终止响应取走 map 项是不同路径;同键覆盖会记录 warning。
pub(crate) async fn register_request_context(&self, request_context: RequestContext) {
    let mut request_contexts = self.request_contexts.lock().await;
    if request_contexts
        .insert(request_context.request_id.clone(), request_context)
        .is_some()
    {
        warn!("replaced unresolved request context");
    }
}

pub(crate) async fn connection_closed(&self, connection_id: ConnectionId) {
    let mut request_contexts = self.request_contexts.lock().await;
    request_contexts.retain(|request_id, _| request_id.connection_id != connection_id);
}

// ...

async fn take_request_context(
    &self,
    request_id: &ConnectionRequestId,
) -> Option<RequestContext> {
    let mut request_contexts = self.request_contexts.lock().await;
    request_contexts.remove(request_id)
}

登记保留 Context,最终响应取走一项,断连按 connection ID 批量清理。 同键重复登记会替换旧项并记录 warning;旧 Context 若仍被其他任务持有,其 guard 和 span 不一定当场释放。 这也不是一个阻止业务重复执行的幂等表。

源码文件:codex-rs/app-server/src/message_processor.rs

相关函数/类型:run_request_with_context / dispatch_initialized_client_request(L711–L722、L912–L947,局部节选)

rust
// 作者注:请求分派 future 和排队后的业务 future 都使用同一 span;派发返回不会清掉等待响应的 context。
// ...
async fn run_request_with_context<F>(
    outgoing: Arc<OutgoingMessageSender>,
    request_context: RequestContext,
    request_fut: F,
) where
    F: Future<Output = ()>,
{
    outgoing
        .register_request_context(request_context.clone())
        .await;
    request_fut.instrument(request_context.span()).await;
}

// ...

let serialization_scope = codex_request.serialization_scope();
let error_request_id = connection_request_id.clone();
let rpc_gate = Arc::clone(&session.rpc_gate);
let processor = Arc::clone(self);
let span = request_context.span();
let request = QueuedInitializedRequest::new(
    rpc_gate,
    async move {
        let processor_for_request = Arc::clone(&processor);
        let result = processor_for_request
            .handle_initialized_client_request(
                connection_request_id,
                codex_request,
                request_context,
                session,
                event_stream_ready,
            )
            .await;
        if let Err(error) = result {
            processor.outgoing.send_error(error_request_id, error).await;
        }
    }
    .instrument(span),
);

if let Some(scope) = serialization_scope {
    let (key, access) = RequestSerializationQueueKey::from_scope(connection_id, scope);
    self.request_serialization_queues
        .enqueue(key, access, request)
        .await;
} else {
    tokio::spawn(async move {
        request.run().await;
    });
}
Ok(())
// ...

初次请求 future 与排队后的业务 future 都被显式 instrument。 有串行化 scope 时,业务 future 先进入对应队列;没有时再由 tokio::spawn 运行。 dispatch_initialized_client_request 返回,只说明分派结束,不能据此结束需要跨这段等待的观测上下文。

这里使用 .instrument(span),让 future 在被 poll 时进入对应 span。 不能用一个跨 .await 的普通 entered guard 代替这种绑定,并默认线程切换后上下文仍然正确。 串行化 key 与公平性可另见请求串行化与顺序,本节关注的是引用与父链。

3.3 结束位置 ​

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:send_response_as_inner / send_error / send_outgoing_message_to_connection(L570–L582、L656–L664、L701–L722,局部节选)

rust
// 作者注:取走 context 后仍借它 instrument 入队;结束于实际引用释放,不是客户端接收确认。
// ...
let response = Box::new(response);
let request_context = self.take_request_context(&request_id).await;
let outgoing_message = OutgoingMessage::Response(OutgoingResponse {
    id: request_id.request_id,
    result: response,
});
self.send_outgoing_message_to_connection(
    request_context,
    connection_id,
    outgoing_message,
    "response",
)
.await;

// ...

pub(crate) async fn send_error(
    &self,
    request_id: ConnectionRequestId,
    error: impl Into<JSONRPCErrorError>,
) {
    let request_context = self.take_request_context(&request_id).await;
    self.send_error_inner(request_context, request_id, error.into())
        .await;
}

// ...

async fn send_outgoing_message_to_connection(
    &self,
    request_context: Option<RequestContext>,
    connection_id: ConnectionId,
    message: OutgoingMessage,
    message_kind: &'static str,
) {
    let send_fut = self.sender.send(OutgoingEnvelope::ToConnection {
        connection_id,
        message,
        write_complete_tx: None,
    });
    let send_result = if let Some(request_context) = request_context {
        send_fut.instrument(request_context.span()).await
    } else {
        send_fut.await
    };

    if let Err(err) = send_result {
        warn!("failed to send {message_kind} to client: {err:?}");
    }
}
// ...

response/error 先从 map 取 Context,再用它 instrument 出站队列的发送。 这里等待的是全局出站队列接纳,没有携带 socket 写完或客户端确认。 请求 span 的时间包含哪些异步等待,取决于它及相关句柄实际保留到哪里;不能把 handler 返回、map remove、 span close 和客户端收到响应压缩成同一个时间点。

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:send_response_clears_registered_request_context(L1107–L1135,摘录)

rust
// 作者注:此测试验证 request_contexts 清理;没有把 map 为空断言成 exporter 已交付。
async fn send_response_clears_registered_request_context() {
    let (tx, _rx) = mpsc::channel::<OutgoingEnvelope>(4);
    let outgoing =
        OutgoingMessageSender::new(tx, codex_analytics::AnalyticsEventsClient::disabled());
    let request_id = ConnectionRequestId {
        connection_id: ConnectionId(42),
        request_id: RequestId::Integer(7),
    };

    outgoing
        .register_request_context(RequestContext::new(
            request_id.clone(),
            tracing::info_span!("app_server.request", rpc.method = "thread/start"),
            /*parent_trace*/ None,
        ))
        .await;
    assert_eq!(outgoing.request_context_count().await, 1);

    outgoing
        .send_response(
            request_id,
            ClientResponsePayload::ThreadArchive(
                codex_app_server_protocol::ThreadArchiveResponse {},
            ),
        )
        .await;

    assert_eq!(outgoing.request_context_count().await, 0);
}

测试先登记一项,再调用 send_response,断言 map 从 1 变成 0。 它没有启动网络 writer,因而不能证明 wire 交付或 exporter 交付。 connection_closed_clears_registered_request_contexts 另用两个连接验证只清理指定连接,不误清仍存活的请求。

4. 跨队列父链 ​

4.1 Core 接力 ​

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:request_trace_context / record_request_turn_id(L247–L266,摘录)

rust
// 作者注:Core carrier 从当前注册项取得;Turn ID 回填到同一 span,便于连接两条观测链。
pub(crate) async fn request_trace_context(
    &self,
    request_id: &ConnectionRequestId,
) -> Option<W3cTraceContext> {
    let request_contexts = self.request_contexts.lock().await;
    request_contexts
        .get(request_id)
        .and_then(RequestContext::request_trace)
}

pub(crate) async fn record_request_turn_id(
    &self,
    request_id: &ConnectionRequestId,
    turn_id: &str,
) {
    let request_contexts = self.request_contexts.lock().await;
    if let Some(request_context) = request_contexts.get(request_id) {
        request_context.record_turn_id(turn_id);
    }
}

请求 ID 是向当前 Context 查询 carrier 的入口。 record_request_turn_id 在同一 span 上回填 turn.id,把 RPC 与实际被 Core 接纳的 Turn 关联起来。 它不创建新的 Turn,也不会给本来不存在的 map 项补建 span。

源码文件:codex-rs/app-server/src/request_processors/turn_processor.rs

相关函数/类型:request_trace_context / turn_start_inner(L418–L434、L534–L550、L580–L582,局部节选)

rust
// 作者注:普通 Op 与 TurnInput 都显式取请求 trace;接纳后再记录实际 Turn ID。
// ...
async fn request_trace_context(
    &self,
    request_id: &ConnectionRequestId,
) -> Option<codex_protocol::protocol::W3cTraceContext> {
    self.outgoing.request_trace_context(request_id).await
}

async fn submit_core_op(
    &self,
    request_id: &ConnectionRequestId,
    thread: &CodexThread,
    op: Op,
) -> CodexResult<String> {
    thread
        .submit_with_trace(op, self.request_trace_context(request_id).await)
        .await
}

// ...

let submission = thread
    .start_or_steer_turn(
        TurnInputRequest::new(TurnInput::UserInput {
            content: mapped_items,
            client_id: client_user_message_id,
        })
        .with_thread_settings(thread_settings)
        .on_start(TurnStartOptions {
            final_output_json_schema: params.output_schema,
            ..Default::default()
        })
        .with_additional_context(additional_context)
        .with_responses_metadata(params.responsesapi_client_metadata)
        .with_trace(self.request_trace_context(&request_id).await),
    )
    .await
    .map_err(|err| {

// ...

self.outgoing
    .record_request_turn_id(&request_id, &turn_id)
    .await;
// ...

普通 Op 通过 submit_with_trace 传 carrier;Turn 输入则在 builder 中显式 .with_trace(...)。 Core 返回实际 Turn ID 后,App Server 再回填 span。 一次 turn/start 也可能 steer 现有工作,所以不能预先把 RPC request ID 当作将来一定出现的 Turn ID。

源码文件:codex-rs/core/src/session/mod.rs

相关函数/类型:SessionIo::submit_with_id / submit_turn_input(L843–L879,摘录)

rust
// 作者注:入队前把 trace 从请求取到 Submission;只在整个 carrier 缺省时补当前 span。
pub(crate) async fn submit_with_id(&self, mut sub: Submission) -> CodexResult<()> {
    if sub.trace.is_none() {
        sub.trace = current_span_w3c_trace_context();
    }
    self.tx_sub
        .send(sub)
        .await
        .map_err(|_| CodexErr::InternalAgentDied)?;
    Ok(())
}

/// Submits an ordered turn-input call and waits only for Core's routing decision.
///
/// Once queued, dropping the waiter does not retract the call. If the
/// session loop exits before replying, the caller gets `InternalAgentDied`.
pub(crate) async fn submit_turn_input(
    &self,
    mut request: TurnInputRequest,
    mode: TurnInputMode,
) -> CodexResult<TurnInputSubmission> {
    let id = new_submission_id();
    let (reply_tx, reply_rx) = oneshot::channel();
    let trace = request.trace.take();
    self.submit_with_id(Submission {
        id,
        op: Op::TurnInput {
            request: Box::new(request),
            mode,
            reply: reply_tx,
        },
        trace,
        parent_turn_id: None,
        root_turn_id: None,
    })
    .await?;
    reply_rx.await.unwrap_or(Err(CodexErr::InternalAgentDied))
}

SessionIo 将 carrier 放在 Submission 上跨 channel 传递。 只在 sub.trace 整体缺省时尝试从当前 span 补充,不会覆盖调用方已经提供的 carrier。 oneshot 等待的是 Core 输入路由结果;队列接纳之后,丢弃等待者也不会撤销已提交的工作。

源码文件:codex-rs/core/src/session/handlers.rs

相关函数/类型:submission_dispatch_span(L735–L763,摘录)

rust
// 作者注:Core 在接收端创建自己的 dispatch span,再以 carrier 设置父级;不是复用 RPC span ID。
pub(super) fn submission_dispatch_span(sub: &Submission) -> tracing::Span {
    let op_name = sub.op.kind();
    let span_name = format!("op.dispatch.{op_name}");
    let dispatch_span = match &sub.op {
        Op::RealtimeConversationAudio(_) => {
            debug_span!(
                "submission_dispatch",
                otel.name = span_name.as_str(),
                submission.id = sub.id.as_str(),
                codex.op = op_name
            )
        }
        _ => info_span!(
            "submission_dispatch",
            otel.name = span_name.as_str(),
            submission.id = sub.id.as_str(),
            codex.op = op_name
        ),
    };
    if let Some(trace) = sub.trace.as_ref()
        && !set_parent_from_w3c_trace_context(&dispatch_span, trace)
    {
        warn!(
            submission.id = sub.id.as_str(),
            "ignoring invalid submission trace carrier"
        );
    }
    dispatch_span
}

Core 接收端创建 submission_dispatch,OTel 名称为 op.dispatch.<操作名>,并解析传入 carrier。 submission_loop 随后将实际分派 future instrument 到这个 span。 传递的是父链信息,不是把一份跨线程 entered guard 塞进 channel,也不是全链路重用同一个 Span ID。

下面沿 turn/start 画出三个追踪范围。图中只要求父子关联成立,不把所有 Core 工作都画成 RPC 返回前完成。

  • 客户端 Trace ID 延续,RPC 与 Core 分派各有 Span ID。
  • turn.id 是接纳后的桥接字段,不能用它代替 trace 父子关系。
  • Core task 与 RPC 的生命周期不同,比较耗时前应先明确正在看哪个 span。

4.2 导出断言 ​

源码文件:codex-rs/app-server/src/message_processor_tracing_tests.rs

相关函数/类型:turn_start_jsonrpc_span_parents_core_turn_spans(L649–L653、L699–L714,局部节选)

rust
// 作者注:fixture 的远端父级已给定;测试检查 RPC 的真实 parent、回填 turn.id 和 Core 子孙关系。
// ...
let RemoteTrace {
    trace_id: remote_trace_id,
    parent_span_id: remote_parent_span_id,
    context: remote_trace,
} = RemoteTrace::new("00000000000000000000000000000077", "0000000000000088");

// ...

let server_request_span =
    find_rpc_span_with_trace(&spans, SpanKind::Server, "turn/start", remote_trace_id);
let core_turn_span =
    find_span_with_trace(&spans, remote_trace_id, "codex.op=turn_input", |span| {
        span_attr(span, "codex.op") == Some("turn_input")
    });

assert_eq!(server_request_span.parent_span_id, remote_parent_span_id);
assert!(server_request_span.parent_span_is_remote);
assert_eq!(server_request_span.span_context.trace_id(), remote_trace_id);
assert_eq!(
    span_attr(server_request_span, "turn.id"),
    Some(turn_start_response.turn.id.as_str())
);
assert_span_descends_from(&spans, core_turn_span, server_request_span);
harness.shutdown().await;
// ...

测试先建立 thread,再为 turn/start 提供合成的远端 carrier。 内存 OTel exporter 中的 RPC span 必须保留给定父 Span ID、remote 标记和 Trace ID, 同时含有实际 Turn ID;Core 的 codex.op=turn_input span 必须是该 RPC span 的子孙。 这是从真实处理器到 Core 的反向验证,不只是断言 carrier 字符串被复制了一次。

另一项 thread/start 测试同时覆盖缺少显式 carrier 与带 carrier 两种输入,并寻找至少两层内部子孙。 这些测试使用内存 exporter,验证父链与属性;不证明某个远端观测平台已经保存记录。

4.3 观察层安装 ​

源码文件:codex-rs/app-server/src/lib.rs

相关函数/类型:run_main_with_transport_options(L667–L697,局部节选)

rust
// 作者注:EnvFilter 附在 stderr fmt 层,feedback、数据库和 OTel 各有自己的层;stderr 不占用协议 stdout。
// ...
// Install a simple subscriber so `tracing` output is visible. Users can
// control the log level with `RUST_LOG` and switch to JSON logs with
// `LOG_FORMAT=json`.
let stderr_fmt: StderrLogLayer = match log_format_from_env() {
    LogFormat::Json => tracing_subscriber::fmt::layer()
        .json()
        .with_writer(std::io::stderr)
        .with_span_events(tracing_subscriber::fmt::format::FmtSpan::FULL)
        .with_filter(EnvFilter::from_default_env())
        .boxed(),
    LogFormat::Default => tracing_subscriber::fmt::layer()
        .with_writer(std::io::stderr)
        .with_span_events(tracing_subscriber::fmt::format::FmtSpan::FULL)
        .with_filter(EnvFilter::from_default_env())
        .boxed(),
};

let feedback_layer = feedback.logger_layer();
let feedback_metadata_layer = feedback.metadata_layer();
let log_db = state_db.clone().map(log_db::start);
let log_db_layer = log_db
    .clone()
    .map(|layer| layer.with_filter(log_db::default_filter()));
let (otel_layers, otel_logger_reload_handle) = otel_reloader::layers(otel.as_ref());
let _ = tracing_subscriber::registry()
    .with(stderr_fmt)
    .with(feedback_layer)
    .with(feedback_metadata_layer)
    .with(log_db_layer)
    .with(otel_layers)
    .try_init();
// ...

standalone 入口尝试把 stderr、反馈日志、SQLite 日志与 OTel 层组合进 subscriber。 RUST_LOG 的 EnvFilter 直接附着在 stderr fmt 层,LOG_FORMAT=json 选择其输出格式; 不能因此认为所有其他层都用同一个过滤器。协议 stdout 仍与这些 stderr 输出分开。

源码文件:codex-rs/app-server/src/otel_reloader.rs

相关函数/类型:layers(L22–L42,摘录)

rust
// 作者注:只有 provider 中存在 tracer 时才加入可重载 tracing 层;本地日志可见不自动等于 trace 可导出。
pub(crate) fn layers<S>(provider: Option<&OtelProvider>) -> OtelReloadLayers<S>
where
    S: Subscriber + for<'span> LookupSpan<'span> + Send + Sync + 'static,
{
    let logger_export_layer: OtelExportLayer<S> = provider
        .and_then(OtelProvider::logger_export_layer)
        .map(Layer::boxed);
    let (logger_layer, logger_handle) = reload::Layer::new(logger_export_layer);

    let mut layers: Vec<Box<dyn Layer<S> + Send + Sync + 'static>> = vec![
        logger_layer
            .with_filter(tracing_subscriber::filter::filter_fn(
                OtelProvider::log_export_filter,
            ))
            .boxed(),
    ];
    if provider.is_some_and(|provider| provider.tracer.is_some()) {
        layers.push(OtelProvider::reloadable_tracing_layer(OTEL_SERVICE_NAME).boxed());
    }
    (layers, logger_handle)
}

只有 provider 中存在 tracer 时才加 tracing 导出层。 所以“本地日志看得见”与“span 有可传播的 OTel Context”不是同一条件。 嵌入式使用者还需要核对宿主的 subscriber;本篇的 parent helper 并不会自行安装一个 exporter。

5. 采集选择 ​

5.1 独立客户端 ​

源码文件:codex-rs/app-server/src/analytics_utils.rs

相关函数/类型:analytics_events_client_from_config(L7–L16,摘录)

rust
// 作者注:目的地址来自 chatgpt_base_url,开关传入独立的 AnalyticsEventsClient。
pub(crate) fn analytics_events_client_from_config(
    auth_manager: Arc<AuthManager>,
    config: &Config,
) -> AnalyticsEventsClient {
    AnalyticsEventsClient::new(
        auth_manager,
        config.chatgpt_base_url.trim_end_matches('/').to_string(),
        config.analytics_enabled,
    )
}

Analytics 使用 chatgpt_base_url 和 analytics_enabled,不是把 request span 自动转换成 HTTP 分析事件。 standalone 初始化时把同一 AnalyticsEventsClient 克隆给 outgoing、MessageProcessor 和 Core 相关依赖, 使不同生产者的事实进入同一个归约器。 这个共享关系可在 lib.rs 的装配和 MessageProcessor::new 的 set_analytics_events_client 处继续跟读。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:AnalyticsEventsClient / AnalyticsEventsQueueMessage(L68–L89、L228–L242,摘录)

rust
// 作者注:启用时才持有队列;常规 Fact 与 Flush 共用有界 channel,容量按消息数计算。
const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256;
const ANALYTICS_EVENTS_TIMEOUT: Duration = Duration::from_secs(10);
// Covers two sequential POSTs plus queue/barrier scheduling; additional queued sends remain best-effort.
const ANALYTICS_EVENTS_FLUSH_TIMEOUT: Duration = Duration::from_secs(25);
const ANALYTICS_EVENT_DEDUPE_MAX_KEYS: usize = 4096;

pub(crate) enum AnalyticsEventsQueueMessage {
    Fact(Box<AnalyticsFact>),
    Flush(oneshot::Sender<()>),
}

#[derive(Clone)]
pub(crate) struct AnalyticsEventsQueue {
    pub(crate) sender: mpsc::Sender<AnalyticsEventsQueueMessage>,
    pub(crate) app_used_emitted_keys: Arc<Mutex<HashSet<(String, String)>>>,
    pub(crate) plugin_used_emitted_keys: Arc<Mutex<HashSet<(String, String)>>>,
}

#[derive(Clone)]
pub struct AnalyticsEventsClient {
    queue: Option<AnalyticsEventsQueue>,
}

// ...

pub fn new(
    auth_manager: Arc<AuthManager>,
    base_url: String,
    analytics_enabled: Option<bool>,
) -> Self {
    let destination = AnalyticsEventsDestination::from_base_url(base_url);
    Self {
        queue: (analytics_enabled != Some(false))
            .then(|| AnalyticsEventsQueue::new(Arc::clone(&auth_manager), destination)),
    }
}

pub fn disabled() -> Self {
    Self { queue: None }
}

只有 Some(false) 禁用队列;None 不等于关闭。 禁用时 record_fact 不入队,但不能据此声称所有调用方都没有克隆、计数或序列化开销。 256 是队列消息条数,并非事件正文的总字节预算;一个 Fact 可以暂时持有较大的 typed request clone。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:AnalyticsEventsQueue::new / AnalyticsEventsClient::record_fact(L147–L183、L584–L588,摘录)

rust
// 作者注:一个消费者独占 reducer;每条 Fact 处理完立即发送产生的 events,没有固定时间窗的批量计时器。
fn new(auth_manager: Arc<AuthManager>, destination: AnalyticsEventsDestination) -> Self {
    let (sender, mut receiver) = mpsc::channel(ANALYTICS_EVENTS_QUEUE_SIZE);
    tokio::spawn(async move {
        let mut reducer = AnalyticsReducer::default();
        while let Some(input) = receiver.recv().await {
            let input = match input {
                AnalyticsEventsQueueMessage::Fact(input) => *input,
                AnalyticsEventsQueueMessage::Flush(done_tx) => {
                    let mut events = Vec::new();
                    reducer.flush(&mut events);
                    send_track_events(&auth_manager, &destination, events).await;
                    let _ = done_tx.send(());
                    continue;
                }
            };
            let mut events = Vec::new();
            reducer.ingest(input, &mut events).await;
            send_track_events(&auth_manager, &destination, events).await;
        }
    });
    Self {
        sender,
        app_used_emitted_keys: Arc::new(Mutex::new(HashSet::new())),
        plugin_used_emitted_keys: Arc::new(Mutex::new(HashSet::new())),
    }
}

fn try_send(&self, input: AnalyticsFact) {
    if self
        .sender
        .try_send(AnalyticsEventsQueueMessage::Fact(Box::new(input)))
        .is_err()
    {
        //TODO: add a metric for this
        tracing::warn!("dropping analytics events: queue is full");
    }
}

// ...

pub(crate) fn record_fact(&self, input: AnalyticsFact) {
    if let Some(queue) = self.queue.as_ref() {
        queue.try_send(input);
    }
}

消费者独占 AnalyticsReducer,顺序读取 Fact,归约出零个或多个事件,再立即尝试发送。 这里没有“每隔固定秒数发送一批”的计时器。处理一条 Fact 期间的归约 await 和 HTTP await 都会推迟下一次 recv, 即使前端生产者使用 try_send,队列仍可能填满。

try_send(...).is_err() 同时覆盖 Full 和 Closed,warning 却统一写 queue is full。 排查时不能仅凭这句日志断言一定是容量不足。 正常 Fact 不等待队列空位,后面的 Flush 消息则使用不同发送方式。

5.2 方法筛选 ​

源码文件:codex-rs/app-server/src/request_processors/initialize_processor.rs

相关函数/类型:track_initialized_request(L183–L191,摘录)

rust
// 作者注:这层在握手与实验 gate 通过后调用,早期请求拒绝不经过普通请求 Fact。
pub(crate) fn track_initialized_request(
    &self,
    connection_id: ConnectionId,
    request_id: RequestId,
    request: &ClientRequest,
) {
    self.analytics_events_client
        .track_request(connection_id.0, request_id, request);
}

dispatch_initialized_client_request 先通过初始化和实验能力检查,再调用此方法,之后才按 scope 排队。 因此某次请求可以已有 request span,却还没有经过 Analytics 的请求记录入口。 即使进入这个入口,也要继续经过下面的 method 筛选。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:AnalyticsEventsClient::track_request(L397–L426,摘录)

rust
// 作者注:只接受 Turn start/steer 与非空 ID 的 interrupt;不是全部 RPC 的通用完成日志。
pub fn track_request(
    &self,
    connection_id: u64,
    request_id: RequestId,
    request: &ClientRequest,
) {
    if let ClientRequest::TurnInterrupt { params, .. } = request {
        if params.turn_id.is_empty() {
            return;
        }
        self.record_fact(AnalyticsFact::ExplicitClientInterruptRequest {
            connection_id,
            request_id,
            turn_id: params.turn_id.clone(),
            requested_at_ms: now_unix_millis(),
        });
        return;
    }
    if !matches!(
        request,
        ClientRequest::TurnStart { .. } | ClientRequest::TurnSteer { .. }
    ) {
        return;
    }
    self.record_fact(AnalyticsFact::ClientRequest {
        connection_id,
        request_id,
        request: Box::new(request.clone()),
    });
}

普通请求 Fact 只有 TurnStart 和 TurnSteer;非空 ID 的 TurnInterrupt 另建带时间的专用 Fact。 其余方法直接返回。initialize 的身份缓存由专门的初始化 hook 提供,不能通过这个列表判断它完全没有 Analytics 作用。 这套事件是产品行为分析,不是每个 RPC 都必有一条的访问日志。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:track_response_inner(L614–L641,摘录)

rust
// 作者注:只为六类响应建 Fact,并先检查响应可序列化;这不证明响应已发到客户端。
fn track_response_inner(
    &self,
    connection_id: u64,
    request_id: RequestId,
    response: &ClientResponsePayload,
    thread_originator: Option<String>,
) {
    if !matches!(
        response,
        ClientResponsePayload::ThreadStart(_)
            | ClientResponsePayload::ThreadResume(_)
            | ClientResponsePayload::ThreadFork(_)
            | ClientResponsePayload::TurnStart(_)
            | ClientResponsePayload::TurnSteer(_)
            | ClientResponsePayload::TurnInterrupt(_)
    ) {
        return;
    }
    if serde_json::to_writer(std::io::sink(), response).is_err() {
        return;
    }
    self.record_fact(AnalyticsFact::ClientResponse {
        connection_id,
        request_id,
        response: Box::new(response.clone()),
        thread_originator,
    });
}

响应筛选又是另一张表:Thread start/resume/fork 与 Turn start/steer/interrupt。 它先用 sink 做可序列化检查,再克隆响应放入 Fact。 这避免明显无法编码的 typed response 被当作正常分析响应,却仍不能证明后续出站队列、socket 或客户端接收成功。

源码文件:codex-rs/analytics/src/client_tests.rs

相关函数/类型:track_request_only_enqueues_analytics_relevant_requests(L763–L773、L793–L806,局部节选)

rust
// 作者注:start/steer 入队,archive 和空 turn ID 的 interrupt 不入队;不要用 RPC 总数直接对比 Fact 数。
// ...
for (request_id, request) in [
    (RequestId::Integer(1), sample_turn_start_request()),
    (RequestId::Integer(2), sample_turn_steer_request()),
] {
    client.track_request(/*connection_id*/ 7, request_id, &request);
    assert!(matches!(
        receiver.try_recv(),
        Ok(AnalyticsEventsQueueMessage::Fact(input))
            if matches!(*input, AnalyticsFact::ClientRequest { .. })
    ));
}

// ...

let ignored_request = sample_thread_archive_request();
client.track_request(
    /*connection_id*/ 7,
    RequestId::Integer(3),
    &ignored_request,
);
assert!(matches!(receiver.try_recv(), Err(TryRecvError::Empty)));

client.track_request(
    /*connection_id*/ 7,
    RequestId::Integer(4),
    &sample_turn_interrupt_request(""),
);
assert!(matches!(receiver.try_recv(), Err(TryRecvError::Empty)));
// ...

fixture 用可直接读取的队列观察 Fact:start/steer 入队,archive 与空 turn ID 的 interrupt 不入队。 这是对采集选择的测试,不是对完整 reducer 或远端分析后端的测试。 对应响应测试遍历六种响应,并用 Unix 非 UTF-8 路径验证不可序列化响应被忽略。

5.3 发送前记录 ​

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:send_response_as_inner(L550–L569,局部节选)

rust
// 作者注:成功响应在真正出站前触发 Analytics;有线程 originator 时额外传给 reducer。
// ...
let connection_id = request_id.connection_id;
let request_id_for_analytics = request_id.request_id.clone();
match thread_originator {
    Some(thread_originator) => {
        self.analytics_events_client
            .track_response_with_thread_originator(
                connection_id.0,
                request_id_for_analytics,
                &response,
                thread_originator,
            );
    }
    None => {
        self.analytics_events_client.track_response(
            connection_id.0,
            request_id_for_analytics,
            &response,
        );
    }
}
// ...

成功响应在真正出站之前记录 Analytics。 有线程专属 originator 时,额外把它传入 reducer,避免共享连接上的不同来源线程都被归因成同一个产品客户端。 所以“分析里看到了一个 response Fact”和“用户收到了该 response”之间仍隔着后续发送流程。

源码文件:codex-rs/app-server/src/outgoing_message.rs

相关函数/类型:ThreadScopedOutgoingMessageSender::send_server_notification(L169–L179,摘录)

rust
// 作者注:即使没有订阅连接也先 track notification;分析事件与客户端看到通知并非同一条件。
pub(crate) async fn send_server_notification(&self, notification: ServerNotification) {
    self.outgoing
        .analytics_events_client
        .track_notification(&notification);
    if self.connection_ids.is_empty() {
        return;
    }
    self.outgoing
        .send_server_notification_to_connections(self.connection_ids.as_slice(), notification)
        .await;
}

ThreadScoped sender 即使没有订阅连接,也先调用 track_notification。 该方法自己只接受部分线程、Turn、item 和 review 生命周期通知,再进入 Fact 队列;并不是所有通知都复制。 UI 没收到某条通知,不能直接推出 Analytics 没有相应输入;反方向也不成立。

6. 请求归约 ​

6.1 三种关联键 ​

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:AnalyticsReducer / ConnectionState / ThreadAnalyticsState(L195–L231,摘录)

rust
// 作者注:请求键含连接 ID 与带类型的 RequestId;Thread、Turn 是另两层关联,不用 Trace ID 当主键。
#[derive(Default)]
pub(crate) struct AnalyticsReducer {
    requests: HashMap<(u64, RequestId), RequestState>,
    turns: HashMap<String, TurnState>,
    connections: HashMap<u64, ConnectionState>,
    threads: HashMap<String, ThreadAnalyticsState>,
    tool_items_started_at_ms: HashMap<ToolItemKey, u64>,
    tool_response_states: HashMap<(String, String), ToolResponseState>,
    code_mode_cells: HashMap<String, HashMap<String, CodeModeCellState>>,
    pending_reviews: HashMap<RequestId, PendingReviewState>,
    item_review_summaries: HashMap<ToolItemKey, ItemReviewSummary>,
}

struct ConnectionState {
    app_server_client: CodexAppServerClientMetadata,
    runtime: CodexRuntimeMetadata,
}

#[derive(Default)]
struct ThreadAnalyticsState {
    connection_id: Option<u64>,
    metadata: Option<ThreadMetadataState>,
    originator: Option<String>,
}

impl ThreadAnalyticsState {
    fn app_server_client(
        &self,
        connection_state: &ConnectionState,
    ) -> CodexAppServerClientMetadata {
        let mut app_server_client = connection_state.app_server_client.clone();
        if let Some(originator) = self.originator.as_ref() {
            app_server_client.product_client_id.clone_from(originator);
        }
        app_server_client
    }
}

连接缓存提供客户端与 runtime 身份,线程缓存提供 session/source/血缘和 originator,Turn 状态收集本轮的配置、 计数和结束资料。请求 map 则用 (connection_id, RequestId) 把一次调用的输入与结果配对。 这里没有以 Trace ID 为主键的自动关联表。

ThreadAnalyticsState::app_server_client 克隆连接身份后,可以只覆盖 product_client_id。 client_name/client_version 仍来自连接初始化;这三个字段不应被当作永远相等的别名。 thread_originator_overrides_shared_connection_across_thread_events 用共享连接上的两个线程验证了这种区别。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:RequestState / PendingTurnStartState / PendingTurnSteerState / PendingTurnInterruptState(L384–L405,摘录)

rust
// 作者注:pending 只保存所需标识与计数;steer 的 created_at 在 reducer 处理输入时以秒取样。
enum RequestState {
    TurnStart(PendingTurnStartState),
    TurnSteer(PendingTurnSteerState),
    ExplicitClientInterrupt(PendingTurnInterruptState),
}

struct PendingTurnStartState {
    thread_id: String,
    num_input_images: usize,
}

struct PendingTurnSteerState {
    thread_id: String,
    expected_turn_id: String,
    num_input_images: usize,
    created_at: u64,
}

struct PendingTurnInterruptState {
    turn_id: String,
    requested_at_ms: u64,
}

三种 pending 状态分别保留当前问题真正需要的字段。 特别是 steer 的 created_at 使用 reducer 处理时的 Unix 秒;interrupt 的 requested_at_ms 则由采集函数在入队前取样。 即使名字都包含“请求时间”,它们的单位与采样位置也不同。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:ingest_request(L1119–L1148,摘录)

rust
// 作者注:从 raw ClientRequest 提取有限字段;请求文本未作为 pending 字段保存。
fn ingest_request(
    &mut self,
    connection_id: u64,
    request_id: RequestId,
    request: ClientRequest,
) {
    match request {
        ClientRequest::TurnStart { params, .. } => {
            self.requests.insert(
                (connection_id, request_id),
                RequestState::TurnStart(PendingTurnStartState {
                    thread_id: params.thread_id,
                    num_input_images: num_input_images(&params.input),
                }),
            );
        }
        ClientRequest::TurnSteer { params, .. } => {
            self.requests.insert(
                (connection_id, request_id),
                RequestState::TurnSteer(PendingTurnSteerState {
                    thread_id: params.thread_id,
                    expected_turn_id: params.expected_turn_id,
                    num_input_images: num_input_images(&params.input),
                    created_at: now_unix_seconds(),
                }),
            );
        }
        _ => {}
    }
}

在进入 reducer 之前,AnalyticsFact::ClientRequest 持有完整 typed request 的 clone。 到这里才投影成 thread ID、预期 Turn ID 和图像数等有限字段。 因此不能将“pending 状态没存 prompt”扩大为“原始 Fact 从未持有用户输入”,也不应直接把 Fact 的 Debug 输出当成安全分析格式。

6.2 归因资料 ​

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:ingest_initialize(L1041–L1064,摘录)

rust
// 作者注:初始化只建立连接身份和 runtime,不直接追加分析事件。
fn ingest_initialize(
    &mut self,
    connection_id: u64,
    params: InitializeParams,
    product_client_id: String,
    runtime: CodexRuntimeMetadata,
    rpc_transport: AppServerRpcTransport,
) {
    self.connections.insert(
        connection_id,
        ConnectionState {
            app_server_client: CodexAppServerClientMetadata {
                product_client_id,
                client_name: Some(params.client_info.name),
                client_version: Some(params.client_info.version),
                rpc_transport,
                experimental_api_enabled: params
                    .capabilities
                    .map(|capabilities| capabilities.experimental_api),
            },
            runtime,
        },
    );
}

初始化建立连接身份和运行时信息,没有在此往输出事件数组追加一条通用 RPC 事件。 客户端声明的实验能力也作为可选值保存,不应替代实际 gate 的运行结果。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:emit_thread_initialized(L2003–L2033,局部节选)

rust
// 作者注:先要求连接身份存在,再保存线程来源、originator 与连接归属;缺失初始化时不会自动补偿。
// ...
fn emit_thread_initialized(
    &mut self,
    connection_id: u64,
    thread: codex_app_server_protocol::Thread,
    model: String,
    initialization_mode: ThreadInitializationMode,
    thread_originator: Option<String>,
    out: &mut Vec<TrackEventRequest>,
) {
    let session_source: SessionSource = thread.source.into();
    let session_id = thread.session_id;
    let thread_id = thread.id;
    let parent_thread_id = thread.parent_thread_id;
    let forked_from_thread_id = thread.forked_from_id;
    let Some(connection_state) = self.connections.get(&connection_id) else {
        return;
    };
    let thread_metadata = ThreadMetadataState::from_thread_metadata(
        session_id.clone(),
        &session_source,
        thread.thread_source.map(Into::into),
        parent_thread_id,
        initialization_mode,
    );
    let thread_state = self.threads.entry(thread_id.clone()).or_default();
    if let Some(originator) = thread_originator {
        thread_state.originator = Some(originator);
    }
    thread_state.connection_id = Some(connection_id);
    thread_state.metadata = Some(thread_metadata.clone());
    let app_server_client = thread_state.app_server_client(connection_state);
// ...

Thread 响应先要求对应的连接初始化资料存在,随后才缓存线程数据,后续生成 codex_thread_initialized。 如果第一步缺少连接资料就直接返回,后来单独收到 initialize 不会自动重放此前这份响应。 测试名中的“once initialized”应理解为具备初始化资料之后能产生事件,不是全局 exactly-once 去重承诺。

源码文件:codex-rs/app-server/src/lib.rs

相关函数/类型:analytics_rpc_transport(L1377–L1384,摘录)

rust
// 作者注:Analytics 的 UDS/off 被归入 Websocket,与 request span 的精细 transport_name 不一致。
fn analytics_rpc_transport(transport: &AppServerTransport) -> AppServerRpcTransport {
    match transport {
        AppServerTransport::Stdio => AppServerRpcTransport::Stdio,
        AppServerTransport::UnixSocket { .. }
        | AppServerTransport::WebSocket { .. }
        | AppServerTransport::Off => AppServerRpcTransport::Websocket,
    }
}

这里还有一个容易在报表里误判的维度差异:Analytics 把 Unix Socket 和 off 都归到 Websocket。 request span 的 transport_name 则分别使用 unix_socket、off;typed span 是 in-process, Analytics 的 InProcess 经 snake_case 序列化为 in_process。 比较两条管道的 transport 分组时,必须先统一它们的分类,不能直接按字符串联表。

来源stdioUnix SocketWebSocket嵌入式
request spanstdiounix_socketwebsocketin-process
Analyticsstdiowebsocketwebsocketin_process

表格描述传入 AppServerTransport 的映射。lib.rs 调用 process_request 时传入服务器的 transport 配置, helper 不会按 ConnectionOrigin 自动重分类。同进程有额外接入通道时,这个标签不能替代真实连接来源。

6.3 成功与拒绝 ​

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:ingest_response / ingest_turn_steer_response(L1456–L1471、L2139–L2162,局部节选)

rust
// 作者注:响应按连接与 request ID 消费 pending;start 绑定 Turn,steer 另发接纳事件并增加已有 Turn 的计数。
// ...
ClientResponse::TurnStart {
    request_id,
    response,
} => {
    let turn_id = response.turn.id;
    let Some(RequestState::TurnStart(pending_request)) =
        self.requests.remove(&(connection_id, request_id))
    else {
        return;
    };
    let turn_state = self.turns.entry(turn_id.clone()).or_default();
    turn_state.connection_id = Some(connection_id);
    turn_state.thread_id = Some(pending_request.thread_id);
    turn_state.num_input_images = Some(pending_request.num_input_images);
    self.maybe_emit_turn_event(&turn_id, out).await;
}

// ...

fn ingest_turn_steer_response(
    &mut self,
    connection_id: u64,
    request_id: RequestId,
    response: TurnSteerResponse,
    out: &mut Vec<TrackEventRequest>,
) {
    let Some(RequestState::TurnSteer(pending_request)) =
        self.requests.remove(&(connection_id, request_id))
    else {
        return;
    };
    if let Some(turn_state) = self.turns.get_mut(&response.turn_id) {
        turn_state.steer_count += 1;
    }
    self.emit_turn_steer_event(
        connection_id,
        pending_request,
        Some(response.turn_id),
        TurnSteerResult::Accepted,
        /*rejection_reason*/ None,
        out,
    );
}
// ...

start 响应消费 pending request,把本次连接与实际 Turn ID 绑定,继续尝试发出轮次事件。 steer 响应也先消费 pending,然后给已有 Turn 增加 steer 计数,并发出单独的接纳事件。 没有匹配 pending 时直接返回,不能仅凭一条孤立的成功响应重建原请求的全部归因。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:AnalyticsReducer::ingest / ingest_error_response / ingest_request_error_response(L558–L565、L1698–L1730,局部节选)

rust
// 作者注:raw error 明确忽略;只有显式 error_type 参与拒绝归类,失败 start/interrupt 不在此发通用事件。
// ...
AnalyticsFact::ErrorResponse {
    connection_id,
    request_id,
    error: _,
    error_type,
} => {
    self.ingest_error_response(connection_id, request_id, error_type, out);
}

// ...

fn ingest_error_response(
    &mut self,
    connection_id: u64,
    request_id: RequestId,
    error_type: Option<AnalyticsJsonRpcError>,
    out: &mut Vec<TrackEventRequest>,
) {
    let Some(request) = self.requests.remove(&(connection_id, request_id)) else {
        return;
    };
    self.ingest_request_error_response(connection_id, request, error_type, out);
}

fn ingest_request_error_response(
    &mut self,
    connection_id: u64,
    request: RequestState,
    error_type: Option<AnalyticsJsonRpcError>,
    out: &mut Vec<TrackEventRequest>,
) {
    match request {
        RequestState::TurnStart(_) => {}
        RequestState::TurnSteer(pending_request) => {
            self.ingest_turn_steer_error_response(
                connection_id,
                pending_request,
                error_type,
                out,
            );
        }
        RequestState::ExplicitClientInterrupt(_) => {}
    }
}
// ...

ErrorResponse.error 在归约入口被明确忽略;归类依据是调用点显式传来的 error_type。 失败的 start/interrupt 在这里只消费 pending,不生成一条“所有失败请求通用事件”。 失败的 steer 才继续形成拒绝事件。

OutgoingMessageSender::send_error 负责协议发送和 context 清理,并没有自动调用这条 Analytics 记录链。 Turn processor 的相关失败分支另行调用 track_error_response。 增加一个新错误返回点时,必须检查该 hook;不能将 RPC 错误数量当成完整的拒绝分析样本数。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:emit_turn_steer_event(L2164–L2203,摘录)

rust
// 作者注:使用请求连接的 runtime/client 身份,再以线程 originator 覆盖 product_client_id;缺少上下文会跳过。
fn emit_turn_steer_event(
    &mut self,
    connection_id: u64,
    pending_request: PendingTurnSteerState,
    accepted_turn_id: Option<String>,
    result: TurnSteerResult,
    rejection_reason: Option<TurnSteerRejectionReason>,
    out: &mut Vec<TrackEventRequest>,
) {
    let Some(connection_state) = self.connections.get(&connection_id) else {
        return;
    };
    let drop_site = AnalyticsDropSite::turn_steer(&pending_request.thread_id);
    let Some(thread_state) = self.threads.get(drop_site.thread_id) else {
        warn_missing_analytics_context(&drop_site, MissingAnalyticsContext::ThreadMetadata);
        return;
    };
    let Some(thread_metadata) = thread_state.metadata.as_ref() else {
        warn_missing_analytics_context(&drop_site, MissingAnalyticsContext::ThreadMetadata);
        return;
    };
    out.push(TrackEventRequest::TurnSteer(CodexTurnSteerEventRequest {
        event_type: "codex_turn_steer_event",
        event_params: CodexTurnSteerEventParams {
            thread_id: pending_request.thread_id,
            session_id: thread_metadata.session_id.clone(),
            expected_turn_id: Some(pending_request.expected_turn_id),
            accepted_turn_id,
            app_server_client: thread_state.app_server_client(connection_state),
            runtime: connection_state.runtime.clone(),
            thread_source: thread_metadata.thread_source.clone(),
            subagent_source: thread_metadata.subagent_source.clone(),
            parent_thread_id: thread_metadata.parent_thread_id.clone(),
            num_input_images: pending_request.num_input_images,
            result,
            rejection_reason,
            created_at: pending_request.created_at,
        },
    }));
}

steer 事件使用发起这次请求的连接取得 runtime/client 信息,再叠加线程的来源元数据。 接纳与拒绝在 accepted_turn_id/result/rejection_reason 上不同,但并不依赖对人类错误文本做字符串识别。 缺少必要上下文时这次输出被跳过,已经消费掉的 pending 不会因此自动复原。

源码文件:codex-rs/analytics/src/analytics_client_tests.rs

相关函数/类型:accepted_turn_steer_emits_expected_event(L4811–L4841,局部节选)

rust
// 作者注:在已送入匹配请求/响应与初始化资料的 fixture 中验证实际输出字段和 accepted 结果。
// ...
assert_eq!(out.len(), 1);
let payload = serde_json::to_value(&out[0]).expect("serialize turn steer event");
assert_eq!(payload["event_type"], json!("codex_turn_steer_event"));
assert_eq!(payload["event_params"]["thread_id"], json!("thread-2"));
assert_eq!(
    payload["event_params"]["session_id"],
    json!("session-thread-2")
);
assert_eq!(payload["event_params"]["expected_turn_id"], json!("turn-2"));
assert_eq!(payload["event_params"]["accepted_turn_id"], json!("turn-2"));
assert_eq!(payload["event_params"]["num_input_images"], json!(1));
assert_eq!(payload["event_params"]["result"], json!("accepted"));
assert_eq!(payload["event_params"]["rejection_reason"], json!(null));
assert!(
    payload["event_params"]["created_at"]
        .as_u64()
        .expect("created_at")
        > 0
);
assert_eq!(
    payload["event_params"]["app_server_client"]["product_client_id"],
    json!("codex-tui")
);
assert_eq!(
    payload["event_params"]["runtime"]["codex_rs_version"],
    json!("0.1.0")
);
assert_eq!(payload["event_params"]["thread_source"], json!("user"));
assert_eq!(payload["event_params"]["subagent_source"], json!(null));
assert_eq!(payload["event_params"]["parent_thread_id"], json!(null));
assert!(payload["event_params"].get("product_client_id").is_none());
// ...

测试先输入匹配的 request/response 和身份资料,再检查最终序列化字段:实际 Turn、预期 Turn、图像数、accepted 状态和归因。 拒绝测试则使用另一条请求连接验证 runtime 来源,并检查 accepted_turn_id 为 null。 这类断言比仅确认“事件数组非空”更能发现连接串用或字段来源错误。

7. 轮次汇合 ​

7.1 多源事实 ​

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:CompletedTurnState / TurnState(L407–L433,摘录)

rust
// 作者注:completed、profile、配置等来自独立 Fact;发出 Turn 事件后仍可能为后台工具保留状态。
#[derive(Clone)]
struct CompletedTurnState {
    status: Option<TurnStatus>,
    turn_error: Option<CodexErrorInfo>,
    completed_at: u64,
    duration_ms: Option<u64>,
}

#[derive(Default)]
struct TurnState {
    connection_id: Option<u64>,
    thread_id: Option<String>,
    num_input_images: Option<usize>,
    image_preparations: Vec<ImagePreparationMetadata>,
    resolved_config: Option<TurnResolvedConfigFact>,
    started_at: Option<u64>,
    token_usage: Option<TokenUsage>,
    profile: Option<TurnProfile>,
    completed: Option<CompletedTurnState>,
    explicit_client_interrupt_requested_at_ms: Option<u64>,
    codex_error: Option<TurnCodexError>,
    latest_diff: Option<String>,
    steer_count: usize,
    tool_counts: TurnToolCounts,
    resource_skill_invocations: HashSet<String>,
    turn_event_emitted: bool,
}

TurnState 不是另一个正在驱动模型的运行时 Turn,它是分析消费者用来汇合事实的状态。 配置、profile、token usage、开始通知、结束通知来自不同生产点,字段先后到达时各自更新。 CompletedTurnState 已经收到,也不一定具备发出完整 Turn 事件的条件。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:ingest_turn_profile(L1177–L1186,摘录)

rust
// 作者注:后到的 profile 同样触发一次完成条件检查,不要求通知固定先后顺序。
async fn ingest_turn_profile(
    &mut self,
    input: TurnProfileFact,
    out: &mut Vec<TrackEventRequest>,
) {
    let TurnProfileFact { turn_id, profile } = input;
    let turn_state = self.turns.entry(turn_id.clone()).or_default();
    turn_state.profile = Some(profile);
    self.maybe_emit_turn_event(&turn_id, out).await;
}

profile 到达也调用 maybe_emit_turn_event;配置、token usage、start 响应和完成通知有各自的触发位置。 这让某些事实可以晚于完成通知到达,而不用将整个分析链绑定在单一回调顺序上。 但不是所有缓存更新都会扫描所有待完成 Turn,不能假定只要最后缺的资料出现就必定立即补发。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:ingest_notification(L1923–L1929、L1935–L1956,局部节选)

rust
// 作者注:开始时间可缺省;结束资料来自完成通知,duration 直接转入而非重新计算。
// ...
ServerNotification::TurnStarted(notification) => {
    let turn_state = self.turns.entry(notification.turn.id).or_default();
    turn_state.started_at = notification
        .turn
        .started_at
        .and_then(|started_at| u64::try_from(started_at).ok());
}

// ...

ServerNotification::TurnCompleted(notification) => {
    self.flush_pending_tool_events(&notification.thread_id, &notification.turn.id, out);
    let turn_state = self.turns.entry(notification.turn.id.clone()).or_default();
    turn_state.completed = Some(CompletedTurnState {
        status: analytics_turn_status(notification.turn.status),
        turn_error: notification
            .turn
            .error
            .and_then(|error| error.codex_error_info),
        completed_at: notification
            .turn
            .completed_at
            .and_then(|completed_at| u64::try_from(completed_at).ok())
            .unwrap_or_default(),
        duration_ms: notification
            .turn
            .duration_ms
            .and_then(|duration_ms| u64::try_from(duration_ms).ok()),
    });
    let turn_id = notification.turn.id;
    self.maybe_emit_turn_event(&turn_id, out).await;
}
// ...

开始时间来自通知中的墙钟秒,结束时间与 duration 也沿用通知值;这里没有用两个时间戳重新相减。 无法转成无符号数的时间被丢弃或退为默认值,duration_ms 可为 None。 这套资料与 request span 的开始/结束不是同一时间区间。

7.2 发出条件 ​

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:maybe_emit_turn_event(L2272–L2312,局部节选)

rust
// 作者注:五个数据条件与连接/线程上下文分别检查;started_at 和 token_usage 不在五个必需字段中。
// ...
async fn maybe_emit_turn_event(&mut self, turn_id: &str, out: &mut Vec<TrackEventRequest>) {
    let Some(turn_state) = self.turns.get(turn_id) else {
        return;
    };
    if turn_state.turn_event_emitted {
        return;
    }
    if turn_state.thread_id.is_none()
        || turn_state.num_input_images.is_none()
        || turn_state.resolved_config.is_none()
        || turn_state.profile.is_none()
        || turn_state.completed.is_none()
    {
        return;
    }
    let Some(thread_id) = turn_state.thread_id.as_ref() else {
        return;
    };
    let drop_site = AnalyticsDropSite::turn(thread_id, turn_id);
    let connection_id = turn_state
        .connection_id
        .or_else(|| self.thread_connection_id(drop_site.thread_id));
    let Some(connection_id) = connection_id else {
        warn_missing_analytics_context(&drop_site, MissingAnalyticsContext::ThreadConnection);
        return;
    };
    let Some(connection_state) = self.connections.get(&connection_id) else {
        warn_missing_analytics_context(
            &drop_site,
            MissingAnalyticsContext::Connection { connection_id },
        );
        return;
    };
    let Some(thread_state) = self.threads.get(drop_site.thread_id) else {
        warn_missing_analytics_context(&drop_site, MissingAnalyticsContext::ThreadMetadata);
        return;
    };
    let Some(thread_metadata) = thread_state.metadata.as_ref() else {
        warn_missing_analytics_context(&drop_site, MissingAnalyticsContext::ThreadMetadata);
        return;
    };
// ...

先检查五个必需字段:thread ID、图像数、有效配置、profile、完成资料。 再检查连接归属、初始化资料和线程元数据。 started_at 和 token_usage 不在五个必需项中,因此缺少它们不一定阻止事件。 缺少上下文时的 warning 应与保留的 Turn 状态一起理解,日志里的 dropping 不自动表示所有状态都已删除。

下图按实际 guard 组织流程。它是多个 optional 字段的汇合,不是假设有一份源码中不存在的状态枚举。

  • 五个字段齐全只是第一道条件,连接/线程上下文还需存在。
  • 进入输出数组与 HTTP 交付之间还隔着认证筛选和发送。
  • 没有满足 guard 时通常保留状态,只有后来再次走到检查函数才会尝试发出。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:maybe_emit_turn_event / has_pending_tool_items_for_turn(L2313–L2343,局部节选)

rust
// 作者注:先形成事件,再按在途工具决定标记已发出或移除 Turn;不是跨任意重复输入的永久去重。
// ...
let turn_event = TrackEventRequest::TurnEvent(Box::new(CodexTurnEventRequest {
        event_type: "codex_turn_event",
        event_params: codex_turn_event_params(
            thread_state.app_server_client(connection_state),
            connection_state.runtime.clone(),
            turn_id.to_string(),
            turn_state,
            thread_metadata,
        ),
    }));
    let accepted_line_event = accepted_line_event_input(turn_id, turn_state);

    out.push(turn_event);
    if let Some((mut input, cwd)) = accepted_line_event {
        input.repo_hash = accepted_line_repo_hash_for_cwd(cwd.as_path()).await;
        out.extend(accepted_line_fingerprint_event_requests(input));
    }
    if self.has_pending_tool_items_for_turn(turn_id) {
        if let Some(turn_state) = self.turns.get_mut(turn_id) {
            turn_state.turn_event_emitted = true;
        }
    } else {
        self.turns.remove(turn_id);
    }
}

fn has_pending_tool_items_for_turn(&self, turn_id: &str) -> bool {
    self.tool_items_started_at_ms
        .keys()
        .any(|key| key.turn_id == turn_id)
}
// ...

归约器先追加 Turn 事件,还可能异步补充代码改动聚合,再检查是否仍有在途工具。 有在途工具时保留状态并标记 turn_event_emitted,让后到工具结果继续使用本轮上下文;否则移除。 这个标记不是跨进程或跨任意重复 Fact 的永久去重表。

归约器由一个 worker 独占,所以这段 await 期间不会另有同一 reducer 的 ingest 并发修改状态; 代价是后面的 Fact 暂时不能被处理,队列压力可能增加。 completed_background_tool_item_emits_after_turn_event 验证 Turn 事件之后仍可输出后台工具结果, 不能把 Turn 事件发出当成所有分析状态立刻被销毁。

7.3 可选时间 ​

源码文件:codex-rs/analytics/src/analytics_client_tests.rs

相关函数/类型:turn_completed_without_started_notification_emits_null_started_at(L6088–L6119,局部节选)

rust
// 作者注:fixture 缺开始通知与 token usage,仍可输出有 duration 的事件,相关字段保持 null。
// ...
ingest_turn_prerequisites(
    &mut reducer,
    &mut out,
    /*include_initialize*/ true,
    /*include_resolved_config*/ true,
    /*include_started*/ false,
    /*include_token_usage*/ false,
)
.await;
reducer
    .ingest(
        AnalyticsFact::Notification(Box::new(sample_turn_completed_notification(
            "thread-2",
            "turn-2",
            AppServerTurnStatus::Completed,
            /*codex_error_info*/ None,
        ))),
        &mut out,
    )
    .await;

let payload = serde_json::to_value(&out[0]).expect("serialize turn event");
assert_eq!(payload["event_params"]["started_at"], json!(null));
assert_eq!(payload["event_params"]["duration_ms"], json!(1234));
assert_eq!(payload["event_params"]["input_tokens"], json!(null));
assert_eq!(payload["event_params"]["cached_input_tokens"], json!(null));
assert_eq!(payload["event_params"]["output_tokens"], json!(null));
assert_eq!(
    payload["event_params"]["reasoning_output_tokens"],
    json!(null)
);
assert_eq!(payload["event_params"]["total_tokens"], json!(null));
// ...

fixture 故意省去 started notification 和 token usage,仍断言输出中 duration_ms = 1234, 同时开始时间与 token 字段为 null。这证明“事件已经存在”不等于“所有观测字段齐全”。 另一项 turn_does_not_emit_without_required_prerequisites 分别缺少初始化资料与配置,断言没有输出。 两组测试共同界定了可选数据和真正的发出前置条件。

8. 字段与时间 ​

8.1 输出投影 ​

源码文件:codex-rs/analytics/src/events.rs

相关函数/类型:CodexTurnEventParams(L1008–L1072,摘录)

rust
// 作者注:完整字段分为身份/配置、状态/计数与时间;未包含 raw prompt、RPC error message 或 trace ID。
#[derive(Serialize)]
pub(crate) struct CodexTurnEventParams {
    pub(crate) thread_id: String,
    pub(crate) session_id: String,
    pub(crate) turn_id: String,
    pub(crate) root_turn_id: Option<String>,
    // TODO(rhan-oai): Populate once queued/default submission type is plumbed from
    // the turn/start callsites instead of always being reported as None.
    pub(crate) submission_type: Option<TurnSubmissionType>,
    pub(crate) app_server_client: CodexAppServerClientMetadata,
    pub(crate) runtime: CodexRuntimeMetadata,
    pub(crate) ephemeral: bool,
    pub(crate) thread_source: Option<ThreadSource>,
    pub(crate) initialization_mode: ThreadInitializationMode,
    pub(crate) subagent_source: Option<String>,
    pub(crate) parent_thread_id: Option<String>,
    pub(crate) model: Option<String>,
    pub(crate) model_provider: String,
    pub(crate) sandbox_policy: Option<&'static str>,
    pub(crate) reasoning_effort: Option<String>,
    pub(crate) reasoning_summary: Option<String>,
    pub(crate) service_tier: String,
    pub(crate) approval_policy: String,
    pub(crate) approvals_reviewer: String,
    pub(crate) sandbox_network_access: bool,
    pub(crate) collaboration_mode: Option<&'static str>,
    pub(crate) personality: Option<String>,
    pub(crate) workspace_kind: Option<String>,
    pub(crate) num_input_images: usize,
    pub(crate) image_preparations: Vec<ImagePreparationMetadata>,
    pub(crate) is_first_turn: bool,
    // 作者注:此状态来自 Turn completion,与 RPC response 是否成功不同。
    pub(crate) status: Option<TurnStatus>,
    /// Client wall-clock time for the first non-startup turn/interrupt request
    /// that later received a successful response.
    pub(crate) explicit_client_interrupt_requested_at_ms: Option<u64>,
    pub(crate) turn_error: Option<CodexErrorInfo>,
    pub(crate) codex_error_kind: Option<CodexErrKind>,
    pub(crate) codex_error_http_status_code: Option<u16>,
    pub(crate) steer_count: Option<usize>,
    pub(crate) total_tool_call_count: Option<usize>,
    pub(crate) shell_command_count: Option<usize>,
    pub(crate) file_change_count: Option<usize>,
    pub(crate) mcp_tool_call_count: Option<usize>,
    pub(crate) dynamic_tool_call_count: Option<usize>,
    pub(crate) subagent_tool_call_count: Option<usize>,
    pub(crate) web_search_count: Option<usize>,
    pub(crate) image_generation_count: Option<usize>,
    pub(crate) input_tokens: Option<i64>,
    pub(crate) cached_input_tokens: Option<i64>,
    pub(crate) cache_write_input_tokens: Option<i64>,
    pub(crate) output_tokens: Option<i64>,
    pub(crate) reasoning_output_tokens: Option<i64>,
    pub(crate) total_tokens: Option<i64>,
    // 作者注:phase 毫秒沿用 Core 的 TurnProfile,不能当成 RPC span 的时间桶。
    pub(crate) before_first_sampling_ms: u64,
    pub(crate) sampling_ms: u64,
    pub(crate) compaction_ms: u64,
    pub(crate) between_sampling_overhead_ms: u64,
    pub(crate) tool_blocking_ms: u64,
    pub(crate) after_last_sampling_ms: u64,
    pub(crate) sampling_request_count: u32,
    pub(crate) sampling_retry_count: u32,
    pub(crate) duration_ms: Option<u64>,
    pub(crate) started_at: Option<u64>,
    pub(crate) completed_at: Option<u64>,
}

这份完整结构把身份、模型与策略、调用计数、错误分类和时间并列输出。 它没有 raw prompt、RPC error message 或 Trace ID 字段,但含有产品和线程标识;不能把有限字段投影称为通用匿名化。 其他 Analytics 事件有各自结构,本篇结论只覆盖这里跟读的请求/Turn 路径。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:codex_turn_event_params(L3403–L3443、L3471–L3483,局部节选)

rust
// 作者注:root_turn_id 在发出时读取,配置转为分析字段;duration 与时间戳沿用各自的生产点。
// ...
let token_usage = turn_state.token_usage.clone();
let codex_error = turn_state.codex_error.as_ref();
CodexTurnEventParams {
    thread_id,
    session_id: thread_metadata.session_id.clone(),
    turn_id,
    root_turn_id: turn_metadata.root_turn_id(),
    app_server_client,
    runtime,
    submission_type,
    ephemeral,
    thread_source: thread_metadata.thread_source.clone(),
    initialization_mode: thread_metadata.initialization_mode,
    subagent_source: thread_metadata.subagent_source.clone(),
    parent_thread_id: thread_metadata.parent_thread_id.clone(),
    model: Some(model),
    model_provider,
    sandbox_policy: Some(sandbox_policy_mode(
        &permission_profile,
        permission_profile_cwd.as_path(),
    )),
    reasoning_effort: reasoning_effort.map(|value| value.to_string()),
    reasoning_summary: reasoning_summary_mode(reasoning_summary),
    service_tier: service_tier
        .map(|value| value.to_string())
        .unwrap_or_else(|| "default".to_string()),
    approval_policy: approval_policy.to_string(),
    approvals_reviewer: approvals_reviewer.to_string(),
    sandbox_network_access,
    collaboration_mode: Some(collaboration_mode_mode(collaboration_mode)),
    personality: personality_mode(personality),
    workspace_kind,
    num_input_images,
    image_preparations: turn_state.image_preparations.clone(),
    is_first_turn,
    status: completed.status,
    explicit_client_interrupt_requested_at_ms: turn_state
        .explicit_client_interrupt_requested_at_ms,
    turn_error: completed.turn_error,
    codex_error_kind: codex_error.map(|error| error.kind),
    codex_error_http_status_code: codex_error.and_then(|error| error.http_status_code),

// ...

before_first_sampling_ms,
        sampling_ms,
        compaction_ms,
        between_sampling_overhead_ms,
        tool_blocking_ms,
        after_last_sampling_ms,
        sampling_request_count,
        sampling_retry_count,
        duration_ms: completed.duration_ms,
        started_at,
        completed_at: Some(completed.completed_at),
    }
}
// ...

root_turn_id 在发出时通过 metadata trait 读取,模型/策略来自 resolved config,计数来自 TurnState, 错误语义来自完成资料和独立的 Core 错误 Fact。 phase 毫秒沿用 Core 的 TurnProfile;duration 沿用完成事件,而非 reducer 自己的处理耗时。 构造函数前面再次要求五个数据项齐全,是发出 guard 与序列化之间的不变量。

数据取样或来源不能拿它代表什么
request span 的范围span 创建、引用与异步执行用户点击到收到响应的完整网络耗时
steer created_atreducer 消费请求时的 Unix 秒精确到毫秒的 RPC 到达时间
interrupt requested_at_mstrack_request 入队前取墙钟毫秒,成功响应后纳入 Turn用户操作必然中断成功的时刻
started_at/completed_atTurn 通知中的墙钟秒直接相减得到准确 duration_ms
duration_ms、phase 毫秒Core 单调时钟计算后传入Analytics 自身排队或 HTTP 投递延迟
HTTP timeout单次事件发送请求的配置整个队列清空的最大时间

interrupt 还保留最早一次被成功接纳的请求时间,拒绝的 interrupt 不应给 Turn 打上这个标记。 相关测试分别覆盖拒绝、接纳和重复请求取最早值;不能把一次网络请求的时间覆盖成用户工作的全部生命周期。

8.2 信息边界 ​

字段筛选发生在不同位置:span 只记录模板属性;reducer 将请求降为有限 pending 字段; ErrorResponse 的原始 error 在归约时被忽略;最终 DTO 再选择要序列化的字段。 这条链没有一个可以替所有日志、所有工具事件做脱敏的统一函数。

例如 OutgoingMessageSender 的客户端错误回调只记录 code,源码注释说明 message/data 可能带凭据; 而 Core 的某些 debug 日志会打印完整 Submission。排查需要读取具体日志点,不能因为分析 DTO 未含正文就开启所有 debug 输出 并假定没有用户内容。这里讨论的是已存在的数据流边界,不是新增一套日志过滤保证。

9. 投递与屏障 ​

9.1 认证与批次 ​

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:send_track_events(L713–L737,摘录)

rust
// 作者注:有队列、有事件、具备适合的认证是三个条件;认证过滤发生在 capture 或 HTTP 发送之前。
async fn send_track_events(
    auth_manager: &AuthManager,
    destination: &AnalyticsEventsDestination,
    mut events: Vec<TrackEventRequest>,
) {
    if events.is_empty() {
        return;
    }

    let Some(auth) = auth_manager.auth().await else {
        return;
    };
    if auth.is_api_key_auth() {
        events.retain(TrackEventRequest::can_send_with_api_key_auth);
    } else if !auth.uses_codex_backend() {
        return;
    }
    if events.is_empty() {
        return;
    }

    for events in track_event_request_batches(events) {
        send_track_events_request(&auth, destination, events).await;
    }
}

有队列不等于有可发送事件,也不等于当前认证允许发送。 没有 auth 直接返回;API key 先筛选事件;其他认证还需使用 Codex backend。 这些判断发生在 capture 或 HTTP 之前,所以打开 debug capture 也不会绕过上面的认证筛选。

源码文件:codex-rs/analytics/src/events.rs

相关函数/类型:TrackEventRequest::can_send_with_api_key_auth(L156–L171,摘录)

rust
// 作者注:API key 仅允许符合各类 plugin_id 条件的事件;普通 Turn 事件不在清单中。
impl TrackEventRequest {
    pub(crate) fn should_send_in_isolated_request(&self) -> bool {
        matches!(self, Self::AcceptedLineFingerprints(_))
    }

    pub(crate) fn can_send_with_api_key_auth(&self) -> bool {
        match self {
            Self::PluginUsed(event) => event.event_params.plugin.plugin_id.is_some(),
            Self::SkillInvocation(event) => event.event_params.plugin_id.is_some(),
            Self::McpToolCall(event) => event.event_params.plugin_id.is_some(),
            Self::ArtifactOperation(event) => !event.event_params.plugin_id.is_empty(),
            Self::PluginMeasurement(event) => !event.event_params.plugin_id.is_empty(),
            _ => false,
        }
    }
}

API key 只保留满足各自 plugin ID 条件的五类事件,普通 Turn 事件不在清单中。 不同分支检查 is_some 或非空字符串,不能把所有条件简写成一个并不存在的统一 predicate。 api_key_auth_sends_only_plugin_events_to_codex_backend 用 capture 文件验证最终保留的五类事件及其插件归属, 不把 mock 数据写作线上账户实验。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:track_event_request_batches(L739–L760,摘录)

rust
// 作者注:保留输入顺序,仅把 AcceptedLineFingerprints 单独分批,不按定时或固定条数分批。
fn track_event_request_batches(events: Vec<TrackEventRequest>) -> Vec<Vec<TrackEventRequest>> {
    let mut batches = Vec::new();
    let mut current_batch = Vec::new();

    for event in events {
        if event.should_send_in_isolated_request() {
            if !current_batch.is_empty() {
                batches.push(current_batch);
                current_batch = Vec::new();
            }
            batches.push(vec![event]);
        } else {
            current_batch.push(event);
        }
    }

    if !current_batch.is_empty() {
        batches.push(current_batch);
    }

    batches
}

分批只隔离 AcceptedLineFingerprints,其余连续事件组成当前批次,并保持输入相对顺序。 这不是按固定条数或固定时间窗积累事件。 “一个 Fact 生成多个事件”与“一个 HTTP 请求装多条事件”是相邻的两步,不能混成一个后台定时任务。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:send_track_events_request(L762–L803,摘录)

rust
// 作者注:单次 HTTP 请求设 10 秒超时;失败记录 warning,不在此把事件重新放回队列。
async fn send_track_events_request(
    auth: &CodexAuth,
    destination: &AnalyticsEventsDestination,
    events: Vec<TrackEventRequest>,
) {
    if events.is_empty() {
        return;
    }

    let payload = TrackEventsRequest { events };

    #[cfg(debug_assertions)]
    if capture_track_events_request(destination, &payload) {
        return;
    }

    let url = match destination {
        AnalyticsEventsDestination::Http { url } => url,
        #[cfg(debug_assertions)]
        AnalyticsEventsDestination::CaptureFile { .. } => return,
    };
    let response = create_client()
        .post(url)
        .timeout(ANALYTICS_EVENTS_TIMEOUT)
        .headers(codex_model_provider::auth_provider_from_auth(auth).to_auth_headers())
        .header("Content-Type", "application/json")
        .json(&payload)
        .send()
        .await;

    match response {
        Ok(response) if response.status().is_success() => {}
        Ok(response) => {
            let status = response.status();
            let body = response.text().await.unwrap_or_default();
            tracing::warn!("events failed with status {status}: {body}");
        }
        Err(err) => {
            tracing::warn!("failed to send events request: {err}");
        }
    }
}

最终 payload 是 TrackEventsRequest { events },发向 Analytics 地址,单次请求设置 10 秒 timeout。 HTTP 非成功状态和请求错误都只在此记录 warning,没有把失败事件重新放回 Analytics 队列。 即使这里得到 2xx,也仅说明 HTTP 层返回成功,不能代替后端存储或后续查询可见性的验证。

9.2 捕获与屏障 ​

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:capture_track_events_request(L805–L821,摘录)

rust
// 作者注:返回 true 表示 capture 已接管这次投递,即使写失败也不会回退到网络。
#[cfg(debug_assertions)]
fn capture_track_events_request(
    destination: &AnalyticsEventsDestination,
    payload: &TrackEventsRequest,
) -> bool {
    let AnalyticsEventsDestination::CaptureFile { path } = destination else {
        return false;
    };

    if let Err(err) = crate::analytics_capture::append_payload(path, payload) {
        tracing::error!(
            path = %path.display(),
            "failed to capture analytics events; network delivery remains disabled: {err}"
        );
    }
    true
}

CaptureFile 仅在 debug 构建可用,选择后写最终的序列化 payload。 true 表示 capture 已接管投递,写文件失败也返回 true,因此不会回退网络发送。 capture_write_failure_still_consumes_delivery 故意使用缺少父目录的路径验证这个边界。 错误恢复动作由这一分支决定,不能把函数返回 bool 一律当成“保存成功”。

源码文件:codex-rs/analytics/src/client.rs

相关函数/类型:AnalyticsEventsClient::flush(L244–L265,摘录)

rust
// 作者注:25 秒覆盖等待入队和等待 done;done 由消费者发出,不携带后端保存成功状态。
pub async fn flush(&self) {
    let Some(queue) = self.queue.as_ref() else {
        return;
    };
    let (done_tx, done_rx) = oneshot::channel();
    let flushed = tokio::time::timeout(ANALYTICS_EVENTS_FLUSH_TIMEOUT, async {
        if queue
            .sender
            .send(AnalyticsEventsQueueMessage::Flush(done_tx))
            .await
            .is_err()
        {
            return false;
        }
        done_rx.await.is_ok()
    })
    .await;

    if !matches!(flushed, Ok(true)) {
        tracing::warn!("timed out or failed while flushing analytics events");
    }
}

Flush 使用等待式 send,与普通 Fact 的 try_send 不同。 25 秒预算同时覆盖等待队列空位和等待 done;若 barrier 已入队,调用方超时也不自动撤回队列里的消息或停止 worker。 done 只表示 worker 走完这条 barrier 的处理,不携带每个 HTTP/capture 操作的成功状态。

源码文件:codex-rs/analytics/src/reducer.rs

相关函数/类型:AnalyticsReducer::flush(L1035–L1039,摘录)

rust
// 作者注:只排空暂存工具事件,不为未满足条件的 Turn 伪造配置、profile 或结束状态。
pub(crate) fn flush(&mut self, out: &mut Vec<TrackEventRequest>) {
    for state in self.tool_response_states.values_mut() {
        out.extend(state.pending_tool_events.drain(..));
    }
}

reducer 的 flush 仅排出尚未输出的工具事件。 它不会补造缺失的 Turn 配置、profile 或完成资料,也不会把所有 pending request 强制转换成事件。 因此“显式 flush 后没有某条 Turn 事件”,还需要回到发出条件检查。

下面把入队、处理和网络结果分开。图中的先后只对成功进入同一队列、位于 barrier 前的 Fact 成立。

  • 普通 Fact 可能在入队前就被忽略或丢弃,barrier 不能恢复它。
  • 身份不允许、capture 失败、HTTP 失败,都可能走到 barrier 的 done。
  • 先检查处理是否完成,再检查投递是否成功,最后检查后端是否可见。

源码文件:codex-rs/analytics/src/client_tests.rs

相关函数/类型:flush_waits_for_preceding_fact_delivery(L859–L881,摘录)

rust
// 作者注:控制 receiver 先取 Fact 再取 Flush,只有显式回复 done 后 flush future 才完成。
async fn flush_waits_for_preceding_fact_delivery() {
    let (client, mut receiver) = client_with_receiver();
    client.track_request(
        /*connection_id*/ 7,
        RequestId::Integer(1),
        &sample_turn_start_request(),
    );

    let flush = tokio::spawn(async move { client.flush().await });
    assert!(matches!(
        receiver.recv().await,
        Some(AnalyticsEventsQueueMessage::Fact(input))
            if matches!(*input, AnalyticsFact::ClientRequest { .. })
    ));
    let done_tx = match receiver.recv().await {
        Some(AnalyticsEventsQueueMessage::Flush(done_tx)) => done_tx,
        _ => panic!("expected analytics flush barrier"),
    };
    tokio::time::sleep(std::time::Duration::from_millis(25)).await;
    assert!(!flush.is_finished());
    done_tx.send(()).expect("flush receiver should remain open");
    flush.await.expect("flush task should complete");
}

测试自己控制接收端,先读 Fact 再读 Flush,并在 25 ms 后确认 flush 仍在等待。 发送 done 后才完成。这验证 channel/barrier 顺序,不是网络成功测试,也不是 25 秒超时的真实时间基准。

9.3 入口收尾 ​

源码文件:codex-rs/app-server/src/in_process.rs

相关函数/类型:start_uninitialized(L762–L776,局部节选)

rust
// 作者注:嵌入式运行时先限时回收处理器和出站路由,再 flush,最后回复 shutdown ack。
// ...
if let Err(_elapsed) = timeout(SHUTDOWN_TIMEOUT, &mut processor_handle).await {
    processor_handle.abort();
    let _ = processor_handle.await;
}
let _ = outbound_shutdown_tx.send(());
if let Err(_elapsed) = timeout(SHUTDOWN_TIMEOUT, &mut outbound_handle).await {
    outbound_handle.abort();
    let _ = outbound_handle.await;
}

analytics_events_flush_client.flush().await;

if let Some(done_tx) = shutdown_ack {
    let _ = done_tx.send(());
}
// ...

嵌入式 runtime 在处理器和出站路由的限时回收之后显式 flush,再回复 shutdown ack。 in_process_shutdown_waits_for_analytics_flush_budget 使用暂停的 Tokio 时间,验证客户端给这些步骤留下完整预算。 它不测试真实 HTTP 后端耗时。

standalone socket 路径的 lib.rs 收尾调用连接清理、后台任务 drain 和线程 shutdown, 不是上面这段显式 Analytics flush 代码;不能把嵌入式的保证推广为所有退出路径都等待同一屏障。 进程被强制终止时,也不能依赖队列或 span 的正常析构替代持久化。

10. 故障取证 ​

遇到缺失观测,先把问题落到某个确切状态和消费者上。

现象首先检查
有本地日志但没有可关联 tracesubscriber 的 OTel 层、span Context 有效性、父 carrier 分支
父 Trace ID 没有延续显式 carrier 是否无效、是否只有 tracestate、环境缓存是否早已为 None
handler 返回后在途计数未归零map 项是否取走、Context/guard 是否仍被异步任务持有
RPC 已拒绝但分析里无拒绝事件方法选择、调用点是否 track_error_response、pending 和线程元数据是否齐全
Turn 已完成但分析事件缺失五个必需字段、connection/thread 资料、是否再次触发 maybe_emit
API key 场景只看到插件事件can_send_with_api_key_auth 的变体筛选
flush 已返回但后端无记录认证过滤、capture 接管、HTTP warning、barrier 完成与投递结果的区别
trace 与报表的 transport 对不上两套 transport 分类和 in-process 拼写差异

在 Codex 仓库根目录运行下面的测试,可分别观察父链、carrier 解析和 Analytics 汇合/屏障。

bash
just test --locked -p codex-app-server --lib \
  -E 'test(message_processor_tracing_tests::) | test(send_response_clears_registered_request_context)'
just test --locked -p codex-otel --lib \
  -E 'test(trace_context::tests::)'
just test --locked -p codex-analytics --lib \
  -E 'test(track_request_only_) | test(turn_does_not_emit_without_) | test(turn_completed_without_started_) | test(flush_waits_for_)'

第一组检查真实内存 exporter 的父子关系,第二组只检查 carrier 和当前 span,第三组控制 Fact 与 barrier 的输入顺序。 需要观察整个 App Server 到 mock 分析地址的输出,可以继续运行 turn_start_tracks_thread_originator_in_analytics 与 turn_profile_tracks_blocking_tool_and_follow_up_sampling:后者延迟回答一个工具请求,检查正的 ToolBlocking、 两次采样和零次采样重试,而不是把 RPC span 时间当作工具阻塞时间。

还可以沿源码推演两个改动:去掉排队 future 的 .instrument(span) 后,哪些子 span 失去当前范围、哪些 Core 子链仍能依靠显式 carrier 关联; 为 Turn 事件再增加一个必需字段时,哪些生产点、guard 和缺字段测试必须同时修改。 能回答这两题,就能继续从观测结果反推真正的持有者与消费时刻。

更细的错误返回与副作用边界见AppServer错误码体系, Core 各阶段的计时算法见TurnTiming与Metadata。