use http::HeaderMap;
use serde::Deserialize;
use super::CodexSessionFile;
use tapes_capture::envelope::{
CODEX_PARENT_THREAD_ID_HEADER, CODEX_SESSION_ID_HEADER, CODEX_THREAD_ID_HEADER,
CODEX_TURN_METADATA_HEADER, OPENAI_SUBAGENT_HEADER, TapesAttribution, X_TAPES_METADATA_RAW_CAP,
};
pub const REQUEST_CORRELATION_METADATA_KEY: &str = "paperProxyRequestId";
const THREAD_SOURCE_SUBAGENT: &str = "subagent";
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CodexRequestIdentity {
pub correlation_id: String,
pub session_id: Option<String>,
pub thread_id: Option<String>,
pub parent_thread_id: Option<String>,
pub turn_id: Option<String>,
pub subagent_kind: Option<String>,
pub conflicting_metadata: bool,
}
#[derive(Debug, Deserialize)]
struct CodexTurnMetadata {
session_id: Option<String>,
thread_id: Option<String>,
parent_thread_id: Option<String>,
turn_id: Option<String>,
subagent_kind: Option<String>,
}
impl CodexRequestIdentity {
#[must_use]
pub fn from_headers(headers: &HeaderMap) -> Self {
let metadata = header_string(headers, CODEX_TURN_METADATA_HEADER)
.and_then(|raw| serde_json::from_str::<CodexTurnMetadata>(&raw).ok());
let header_session_id = header_string(headers, CODEX_SESSION_ID_HEADER);
let header_thread_id = header_string(headers, CODEX_THREAD_ID_HEADER);
let header_parent_thread_id = header_string(headers, CODEX_PARENT_THREAD_ID_HEADER);
let header_subagent_kind = header_string(headers, OPENAI_SUBAGENT_HEADER)
.map(|kind| canonical_subagent_kind(&kind).to_owned());
let metadata_subagent_kind = metadata
.as_ref()
.and_then(|value| value.subagent_kind.as_deref())
.map(|kind| canonical_subagent_kind(kind).to_owned());
let conflicting_metadata = [
(
header_session_id.as_ref(),
metadata
.as_ref()
.and_then(|value| value.session_id.as_ref()),
),
(
header_thread_id.as_ref(),
metadata.as_ref().and_then(|value| value.thread_id.as_ref()),
),
(
header_parent_thread_id.as_ref(),
metadata
.as_ref()
.and_then(|value| value.parent_thread_id.as_ref()),
),
(
header_subagent_kind.as_ref(),
metadata_subagent_kind.as_ref(),
),
]
.into_iter()
.any(|(header, metadata)| matches!((header, metadata), (Some(a), Some(b)) if a != b));
Self {
correlation_id: String::new(),
session_id: header_session_id
.or_else(|| metadata.as_ref().and_then(|value| value.session_id.clone())),
thread_id: header_thread_id
.or_else(|| metadata.as_ref().and_then(|value| value.thread_id.clone())),
parent_thread_id: header_parent_thread_id.or_else(|| {
metadata
.as_ref()
.and_then(|value| value.parent_thread_id.clone())
}),
turn_id: metadata.as_ref().and_then(|value| value.turn_id.clone()),
subagent_kind: metadata_subagent_kind.or(header_subagent_kind),
conflicting_metadata,
}
}
#[must_use]
pub fn with_correlation_id(mut self, correlation_id: impl Into<String>) -> Self {
self.correlation_id = correlation_id.into();
self
}
#[must_use]
pub fn is_child_shaped(&self) -> bool {
!self.conflicting_metadata
&& matches!(
(
self.session_id.as_deref(),
self.thread_id.as_deref(),
self.parent_thread_id.as_deref(),
),
(Some(root), Some(thread), Some(parent)) if thread != root && thread != parent
)
}
#[must_use]
pub fn rollout_id(&self) -> Option<&str> {
if self.conflicting_metadata {
return None;
}
self.thread_id
.as_deref()
.or(self.session_id.as_deref())
.filter(|value| !value.is_empty())
}
}
#[must_use]
pub fn canonical_subagent_kind(kind: &str) -> &str {
match kind {
"collab_spawn" => "thread_spawn",
kind => kind,
}
}
fn header_string(headers: &HeaderMap, name: &str) -> Option<String> {
headers
.get(name)
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_owned)
}
#[must_use]
pub fn transcript_matches_child(
session: &CodexSessionFile,
identity: &CodexRequestIdentity,
) -> bool {
let (Some(root), Some(thread), Some(parent)) = (
identity.session_id.as_deref(),
identity.thread_id.as_deref(),
identity.parent_thread_id.as_deref(),
) else {
return false;
};
session.session_id == thread
&& session.root_session_id.as_deref() == Some(root)
&& session.parent_thread_id.as_deref() == Some(parent)
&& session.thread_source.as_deref() == Some(THREAD_SOURCE_SUBAGENT)
&& match (
identity.subagent_kind.as_deref(),
session.subagent_kind.as_deref(),
) {
(Some(request), Some(transcript)) => request == transcript,
_ => true,
}
}
#[must_use]
pub fn codex_envelope(
session: &CodexSessionFile,
identity: &CodexRequestIdentity,
) -> TapesAttribution {
if let Some(envelope) = child_envelope(identity, Some(session)) {
return envelope;
}
let mut metadata = request_metadata(identity);
for (key, value) in [
("originator", session.originator.as_deref()),
("source", session.source.as_deref()),
("threadSource", session.thread_source.as_deref()),
("modelProvider", session.model_provider.as_deref()),
] {
if let Some(value) = value {
try_insert_metadata(&mut metadata, key, value);
}
}
try_insert_metadata(
&mut metadata,
"transcriptPath",
&session.path.display().to_string(),
);
let parent_sid = if session.thread_source.as_deref() == Some(THREAD_SOURCE_SUBAGENT) {
None
} else {
session.parent_thread_id.as_deref()
};
TapesAttribution::codex_session_with_parent(
&session.session_id,
parent_sid,
session.cwd.as_deref(),
session.cli_version.as_deref(),
metadata,
)
}
#[must_use]
pub fn request_envelope(identity: &CodexRequestIdentity) -> TapesAttribution {
child_envelope(identity, None)
.unwrap_or_else(|| TapesAttribution::codex_with_metadata(request_metadata(identity)))
}
#[must_use]
pub fn child_envelope(
identity: &CodexRequestIdentity,
session: Option<&CodexSessionFile>,
) -> Option<TapesAttribution> {
if !identity.is_child_shaped() {
return None;
}
let root = identity.session_id.as_deref()?;
let mut metadata = serde_json::Map::new();
insert_correlation_id(&mut metadata, identity);
try_insert_metadata(&mut metadata, "codexSessionId", root);
if let Some(provider) = session.and_then(|session| session.model_provider.as_deref()) {
try_insert_metadata(&mut metadata, "modelProvider", provider);
}
Some(TapesAttribution::codex_session_with_parent(
root,
None,
session.and_then(|session| session.cwd.as_deref()),
session.and_then(|session| session.cli_version.as_deref()),
metadata,
))
}
#[must_use]
pub fn envelope_session_id<'a>(
identity: &'a CodexRequestIdentity,
session: Option<&'a CodexSessionFile>,
) -> Option<&'a str> {
if identity.is_child_shaped() {
return identity.session_id.as_deref();
}
session.map(|session| session.session_id.as_str())
}
#[must_use]
pub fn request_metadata(
identity: &CodexRequestIdentity,
) -> serde_json::Map<String, serde_json::Value> {
let mut metadata = serde_json::Map::new();
insert_correlation_id(&mut metadata, identity);
for (key, value) in [
("codexSessionId", identity.session_id.as_ref()),
("codexThreadId", identity.thread_id.as_ref()),
("codexParentThreadId", identity.parent_thread_id.as_ref()),
("codexTurnId", identity.turn_id.as_ref()),
("codexSubagentKind", identity.subagent_kind.as_ref()),
] {
let Some(value) = value else {
continue;
};
try_insert_metadata(&mut metadata, key, value);
}
metadata
}
fn insert_correlation_id(
metadata: &mut serde_json::Map<String, serde_json::Value>,
identity: &CodexRequestIdentity,
) {
if identity.correlation_id.is_empty() {
return;
}
metadata.insert(
REQUEST_CORRELATION_METADATA_KEY.to_owned(),
serde_json::Value::String(identity.correlation_id.clone()),
);
}
pub fn try_insert_metadata(
metadata: &mut serde_json::Map<String, serde_json::Value>,
key: &str,
value: &str,
) {
if value.len() >= X_TAPES_METADATA_RAW_CAP {
return;
}
metadata.insert(key.to_owned(), serde_json::Value::String(value.to_owned()));
let fits = serde_json::to_vec(&metadata).is_ok_and(|raw| raw.len() <= X_TAPES_METADATA_RAW_CAP);
if !fits {
metadata.remove(key);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use std::path::PathBuf;
fn headers(pairs: &[(&str, &str)]) -> HeaderMap {
let mut headers = HeaderMap::new();
for (name, value) in pairs {
headers.insert(
http::HeaderName::from_bytes(name.as_bytes()).unwrap(),
value.parse().unwrap(),
);
}
headers
}
fn root_session(id: &str) -> CodexSessionFile {
CodexSessionFile {
session_id: id.to_owned(),
root_session_id: Some(id.to_owned()),
parent_thread_id: None,
subagent_kind: None,
timestamp: time::OffsetDateTime::now_utc(),
modified_at: Some(time::OffsetDateTime::now_utc()),
cwd: Some("/tmp/work".to_owned()),
originator: Some("codex-tui".to_owned()),
cli_version: Some("0.139.0".to_owned()),
source: Some("cli".to_owned()),
thread_source: Some("user".to_owned()),
model_provider: Some("paper-openai".to_owned()),
path: PathBuf::from(format!("/tmp/{id}.jsonl")),
}
}
fn child_session(id: &str, root: &str, parent: &str, kind: &str) -> CodexSessionFile {
let mut session = root_session(id);
session.root_session_id = Some(root.to_owned());
session.parent_thread_id = Some(parent.to_owned());
session.thread_source = Some("subagent".to_owned());
session.subagent_kind = Some(kind.to_owned());
session
}
fn child_identity(root: &str, parent: &str, child: &str, kind: &str) -> CodexRequestIdentity {
CodexRequestIdentity {
correlation_id: "correlation-child".to_owned(),
session_id: Some(root.to_owned()),
thread_id: Some(child.to_owned()),
parent_thread_id: Some(parent.to_owned()),
turn_id: Some("turn-child".to_owned()),
subagent_kind: Some(kind.to_owned()),
conflicting_metadata: false,
}
}
#[test]
fn parses_simple_headers_and_allowlisted_turn_metadata() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "parent"),
("thread-id", "child"),
("x-codex-parent-thread-id", "parent"),
("x-openai-subagent", "guardian"),
(
"x-codex-turn-metadata",
r#"{"turn_id":"turn-child","session_id":"parent","thread_id":"child","parent_thread_id":"parent","subagent_kind":"guardian","prompt":"must not survive"}"#,
),
]));
assert_eq!(identity.session_id.as_deref(), Some("parent"));
assert_eq!(identity.thread_id.as_deref(), Some("child"));
assert_eq!(identity.parent_thread_id.as_deref(), Some("parent"));
assert_eq!(identity.turn_id.as_deref(), Some("turn-child"));
assert_eq!(identity.subagent_kind.as_deref(), Some("guardian"));
assert!(identity.is_child_shaped());
assert!(!identity.conflicting_metadata);
}
#[test]
fn nested_child_allows_direct_parent_to_differ_from_root() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "root"),
("thread-id", "child"),
("x-codex-parent-thread-id", "parent"),
(
"x-codex-turn-metadata",
r#"{"session_id":"root","thread_id":"child","parent_thread_id":"parent"}"#,
),
]));
assert_eq!(identity.session_id.as_deref(), Some("root"));
assert_eq!(identity.thread_id.as_deref(), Some("child"));
assert_eq!(identity.parent_thread_id.as_deref(), Some("parent"));
assert!(identity.is_child_shaped());
}
#[test]
fn canonicalizes_collab_spawn_to_the_structured_thread_spawn() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "root"),
("thread-id", "child"),
("x-codex-parent-thread-id", "parent"),
("x-openai-subagent", "collab_spawn"),
(
"x-codex-turn-metadata",
r#"{"session_id":"root","thread_id":"child","parent_thread_id":"parent","subagent_kind":"thread_spawn"}"#,
),
]));
assert_eq!(identity.subagent_kind.as_deref(), Some("thread_spawn"));
assert!(!identity.conflicting_metadata);
assert!(identity.is_child_shaped());
}
#[test]
fn genuinely_different_subagent_kinds_still_conflict() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "root"),
("thread-id", "child"),
("x-codex-parent-thread-id", "parent"),
("x-openai-subagent", "guardian"),
(
"x-codex-turn-metadata",
r#"{"session_id":"root","thread_id":"child","parent_thread_id":"parent","subagent_kind":"thread_spawn"}"#,
),
]));
assert_eq!(identity.subagent_kind.as_deref(), Some("thread_spawn"));
assert!(identity.conflicting_metadata);
assert!(!identity.is_child_shaped());
}
#[test]
fn contradictory_sources_disable_the_child_shape() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "parent"),
("thread-id", "child"),
("x-codex-parent-thread-id", "parent"),
(
"x-codex-turn-metadata",
r#"{"session_id":"other-parent","thread_id":"child","parent_thread_id":"parent"}"#,
),
]));
assert!(identity.conflicting_metadata);
assert!(!identity.is_child_shaped());
}
#[test]
fn a_contradicted_request_names_no_rollout_to_narrow_by() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("thread-id", "claimed"),
("x-codex-turn-metadata", r#"{"thread_id":"contradictory"}"#),
]));
assert!(identity.conflicting_metadata);
assert_eq!(identity.rollout_id(), None);
}
#[test]
fn rollout_id_prefers_the_thread_over_the_root() {
let identity = CodexRequestIdentity::from_headers(&headers(&[
("session-id", "root"),
("thread-id", "child"),
]));
assert_eq!(identity.rollout_id(), Some("child"));
let root_turn = CodexRequestIdentity::from_headers(&headers(&[("session-id", "root")]));
assert_eq!(root_turn.rollout_id(), Some("root"));
}
#[test]
fn a_matching_subagent_transcript_joins_the_request() {
let identity = child_identity("root", "root", "child", "guardian");
assert!(transcript_matches_child(
&child_session("child", "root", "root", "guardian"),
&identity,
));
}
#[test]
fn a_transcript_that_is_not_a_subagent_never_joins() {
let identity = child_identity("root", "root", "child", "guardian");
let mut session = child_session("child", "root", "root", "guardian");
session.thread_source = Some("user".to_owned());
assert!(!transcript_matches_child(&session, &identity));
}
#[test]
fn a_missing_kind_on_either_side_is_not_a_contradiction() {
let mut identity = child_identity("root", "root", "child", "guardian");
identity.subagent_kind = None;
assert!(transcript_matches_child(
&child_session("child", "root", "root", "guardian"),
&identity,
));
}
#[test]
fn a_child_request_rekeys_to_the_root_and_drops_per_child_metadata() {
let session = child_session("child", "parent", "parent", "guardian");
let identity = child_identity("parent", "parent", "child", "guardian");
let envelope = codex_envelope(&session, &identity);
assert_eq!(envelope.session_id.as_deref(), Some("parent"));
assert_eq!(envelope.parent_sid, None);
assert_eq!(
envelope.metadata[REQUEST_CORRELATION_METADATA_KEY],
"correlation-child",
);
assert_eq!(envelope.metadata["codexSessionId"], "parent");
assert_eq!(envelope.metadata["modelProvider"], "paper-openai");
for key in [
"codexThreadId",
"codexParentThreadId",
"codexTurnId",
"codexSubagentKind",
"originator",
"source",
"threadSource",
"transcriptPath",
] {
assert!(
!envelope.metadata.contains_key(key),
"unexpected per-child {key}",
);
}
}
#[test]
fn a_child_request_with_no_rollout_still_names_its_root() {
let identity = child_identity("parent", "parent", "child", "guardian");
let envelope = request_envelope(&identity);
assert_eq!(envelope.session_id.as_deref(), Some("parent"));
assert_eq!(envelope.cwd, None);
assert_eq!(envelope.metadata["codexSessionId"], "parent");
assert!(!envelope.metadata.contains_key("modelProvider"));
}
#[test]
fn a_subagent_rollout_never_emits_a_parent_session_header() {
let session = child_session("child", "root", "parent", "guardian");
let envelope = codex_envelope(&session, &CodexRequestIdentity::default());
assert_eq!(envelope.session_id.as_deref(), Some("child"));
assert_eq!(envelope.parent_sid, None);
}
#[test]
fn a_resumed_rollouts_parent_is_still_lineage() {
let mut session = root_session("resumed");
session.parent_thread_id = Some("original".to_owned());
let envelope = codex_envelope(&session, &CodexRequestIdentity::default());
assert_eq!(envelope.parent_sid.as_deref(), Some("original"));
}
#[test]
fn an_absent_correlation_id_is_omitted_rather_than_emitted_blank() {
let metadata = request_metadata(&CodexRequestIdentity::default());
assert!(!metadata.contains_key(REQUEST_CORRELATION_METADATA_KEY));
}
#[test]
fn oversized_identifiers_drop_individually_and_spare_the_correlation_id() {
let oversized = "x".repeat(X_TAPES_METADATA_RAW_CAP + 1);
let identity = CodexRequestIdentity {
correlation_id: "correlation-bounded".to_owned(),
session_id: Some(oversized.clone()),
thread_id: Some("safe-thread".to_owned()),
parent_thread_id: Some(oversized.clone()),
turn_id: Some(oversized.clone()),
subagent_kind: Some(oversized),
conflicting_metadata: false,
};
let metadata = request_metadata(&identity);
assert_eq!(
metadata[REQUEST_CORRELATION_METADATA_KEY],
"correlation-bounded",
);
assert_eq!(metadata["codexThreadId"], "safe-thread");
for key in [
"codexSessionId",
"codexParentThreadId",
"codexTurnId",
"codexSubagentKind",
] {
assert!(
!metadata.contains_key(key),
"{key} should have been dropped"
);
}
}
#[test]
fn oversized_rollout_enrichment_cannot_evict_the_correlation_id() {
let oversized = "x".repeat(X_TAPES_METADATA_RAW_CAP + 1);
let mut session = root_session("root");
session.originator = Some("safe-originator".to_owned());
session.source = Some(oversized.clone());
session.path = PathBuf::from(oversized);
let identity = CodexRequestIdentity {
correlation_id: "correlation-root".to_owned(),
session_id: Some("root".to_owned()),
thread_id: Some("root".to_owned()),
..CodexRequestIdentity::default()
};
let envelope = codex_envelope(&session, &identity);
assert_eq!(
envelope.metadata[REQUEST_CORRELATION_METADATA_KEY],
"correlation-root",
);
assert_eq!(envelope.metadata["codexThreadId"], "root");
assert_eq!(envelope.metadata["originator"], "safe-originator");
assert!(!envelope.metadata.contains_key("source"));
assert!(!envelope.metadata.contains_key("transcriptPath"));
}
}