use tokio::time::Instant;
use tracing::warn;
use super::request::{CodexRequestIdentity, transcript_matches_child};
use super::{CodexSessionFile, open_jsonl_sessions_by_pid, session as codex_session};
use crate::attribution::pipeline::{AttributionConfig, AttributionState, RequestFacts};
use tapes_capture::peer_pid;
pub trait CodexHookEvidence: std::fmt::Debug + Send + Sync {
fn has_hook_session(&self, session_id: &str) -> bool;
fn has_hook_subagent(&self, parent_session_id: &str, agent_id: &str) -> bool;
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum CodexSelection {
HookExact {
session_id: String,
},
HookSubagentExact {
session_id: String,
parent_session_id: String,
},
ChildTranscriptExact {
session_id: String,
parent_session_id: String,
},
MarkerUnique {
session_id: String,
},
MarkerNewest {
session_id: String,
},
PeerPidUnique {
session_id: String,
},
PeerPidNewest {
session_id: String,
},
RecentUnique {
session_id: String,
},
#[default]
NoMatch,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct CodexSelectionEvidence {
pub child_candidates: Vec<CodexSessionFile>,
pub marker_candidates: Vec<CodexSessionFile>,
pub peer_pid: Option<i32>,
pub peer_candidates: Vec<CodexSessionFile>,
pub recent_candidates: Vec<CodexSessionFile>,
pub selection: CodexSelection,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct CodexSelectionResult {
pub session: Option<CodexSessionFile>,
pub evidence: CodexSelectionEvidence,
}
struct Attempt {
result: CodexSelectionResult,
recent_fallback: Option<CodexSessionFile>,
}
impl Attempt {
fn with_recent_fallback(mut self) -> CodexSelectionResult {
if let Some(session) = self.recent_fallback {
self.result.evidence.selection = CodexSelection::RecentUnique {
session_id: session.session_id.clone(),
};
self.result.session = Some(session);
}
self.result
}
}
pub async fn select(
state: &AttributionState,
config: &AttributionConfig,
facts: RequestFacts<'_>,
identity: &CodexRequestIdentity,
) -> CodexSelectionResult {
let wait_for_exact = !identity.conflicting_metadata
&& identity.thread_id.as_deref().is_some_and(|thread_id| {
facts
.codex_hook_evidence
.is_some_and(|hooks| hooks.has_hook_session(thread_id))
|| identity.is_child_shaped()
});
let deadline = Instant::now() + config.codex_timeout;
let mut retained: Option<CodexSelectionResult> = None;
loop {
let attempt = select_once(state, config, facts, identity);
let exact = matches!(
attempt.result.evidence.selection,
CodexSelection::HookExact { .. }
| CodexSelection::HookSubagentExact { .. }
| CodexSelection::ChildTranscriptExact { .. }
);
if exact || (!wait_for_exact && attempt.result.session.is_some()) {
return attempt.result;
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
if wait_for_exact && attempt.result.session.is_some() {
return attempt.result;
}
if wait_for_exact && facts.codex_marker.is_none() && attempt.recent_fallback.is_some() {
return attempt.with_recent_fallback();
}
if let Some(retained) = retained {
return retained;
}
return if facts.codex_marker.is_none() {
attempt.with_recent_fallback()
} else {
attempt.result
};
}
if wait_for_exact {
if attempt.result.session.is_some() {
retained = Some(attempt.result);
} else if facts.codex_marker.is_none() && attempt.recent_fallback.is_some() {
retained = Some(attempt.with_recent_fallback());
}
}
tokio::time::sleep(config.codex_poll.min(remaining)).await;
}
}
fn rollout_evidence<'a>(
facts: RequestFacts<'a>,
identity: &'a CodexRequestIdentity,
) -> Option<&'a str> {
match facts.codex_identity {
Some(_) => identity.rollout_id(),
None => facts.codex_rollout_id,
}
}
fn select_once(
state: &AttributionState,
config: &AttributionConfig,
facts: RequestFacts<'_>,
identity: &CodexRequestIdentity,
) -> Attempt {
let cutoff = time::OffsetDateTime::now_utc() - config.codex_recent_window;
let snapshot = state.codex_watcher.load_full();
let recent: Vec<&CodexSessionFile> = snapshot
.sessions
.iter()
.filter(|session| is_live_candidate(session, config, cutoff))
.collect();
let rollout_id = rollout_evidence(facts, identity);
let recent_candidates: Vec<CodexSessionFile> = recent.iter().copied().cloned().collect();
let recent_fallback = narrow_by_rollout_id(rollout_id, recent_candidates.clone())
.and_then(|candidates| {
one_live_session_or_refuse(
"recent-session",
"multiple recent sessions and the request named no rollout",
candidates,
)
})
.map(|resolved| resolved.session);
let child_matches: Vec<&CodexSessionFile> = if identity.is_child_shaped() {
recent
.iter()
.copied()
.filter(|session| transcript_matches_child(session, identity))
.collect()
} else {
Vec::new()
};
let child_candidates: Vec<CodexSessionFile> = child_matches.iter().copied().cloned().collect();
let mut evidence = CodexSelectionEvidence {
child_candidates,
recent_candidates,
..CodexSelectionEvidence::default()
};
if !identity.conflicting_metadata
&& let Some(hook_session_id) = identity.thread_id.as_deref()
&& facts
.codex_hook_evidence
.is_some_and(|hooks| hooks.has_hook_session(hook_session_id))
&& let Some(session) = recent
.iter()
.copied()
.find(|session| session.session_id == hook_session_id)
{
evidence.selection = CodexSelection::HookExact {
session_id: session.session_id.clone(),
};
return Attempt {
result: CodexSelectionResult {
session: Some(session.clone()),
evidence,
},
recent_fallback,
};
}
if let ([session], Some(parent_session_id)) = (
child_matches.as_slice(),
identity.parent_thread_id.as_deref(),
) {
let parent_session_id = parent_session_id.to_owned();
let hook_backed = identity.thread_id.as_deref().is_some_and(|agent_id| {
facts
.codex_hook_evidence
.is_some_and(|hooks| hooks.has_hook_subagent(&parent_session_id, agent_id))
});
evidence.selection = if hook_backed {
CodexSelection::HookSubagentExact {
session_id: session.session_id.clone(),
parent_session_id,
}
} else {
CodexSelection::ChildTranscriptExact {
session_id: session.session_id.clone(),
parent_session_id,
}
};
return Attempt {
result: CodexSelectionResult {
session: Some((*session).clone()),
evidence,
},
recent_fallback,
};
}
if let Some(marker) = facts.codex_marker {
let matches: Vec<CodexSessionFile> = recent
.iter()
.copied()
.filter(|session| session.has_model_provider(marker))
.cloned()
.collect();
evidence.marker_candidates = matches.clone();
if let Some(resolved) = marker_match(rollout_id, matches) {
let session_id = resolved.session.session_id.clone();
evidence.selection = if resolved.from_unique_candidate {
CodexSelection::MarkerUnique { session_id }
} else {
CodexSelection::MarkerNewest { session_id }
};
return Attempt {
result: CodexSelectionResult {
session: Some(resolved.session),
evidence,
},
recent_fallback,
};
}
}
let peer_pid = facts
.peer
.and_then(|peer| peer_pid::lookup_owner_once(peer).pid);
evidence.peer_pid = peer_pid;
let matches: Vec<CodexSessionFile> = peer_pid
.into_iter()
.flat_map(open_jsonl_sessions_by_pid)
.filter_map(|path| codex_session::read(&path))
.filter(|session| is_live_candidate(session, config, cutoff))
.collect();
evidence.peer_candidates = matches.clone();
let resolved = peer_match(rollout_id, matches);
evidence.selection = resolved
.as_ref()
.map_or(CodexSelection::NoMatch, |resolved| {
let session_id = resolved.session.session_id.clone();
if resolved.from_unique_candidate {
CodexSelection::PeerPidUnique { session_id }
} else {
CodexSelection::PeerPidNewest { session_id }
}
});
Attempt {
result: CodexSelectionResult {
session: resolved.map(|resolved| resolved.session),
evidence,
},
recent_fallback,
}
}
fn is_live_candidate(
session: &CodexSessionFile,
config: &AttributionConfig,
cutoff: time::OffsetDateTime,
) -> bool {
config.codex_provider.matches_session(session)
&& session.modified_at.is_some_and(|ts| ts >= cutoff)
}
pub(crate) fn narrow_by_rollout_id(
rollout_id: Option<&str>,
candidates: Vec<CodexSessionFile>,
) -> Option<Vec<CodexSessionFile>> {
let Some(rollout_id) = rollout_id else {
return Some(candidates);
};
if candidates.is_empty() {
return Some(candidates);
}
let selected: Vec<CodexSessionFile> = candidates
.iter()
.filter(|session| session.session_id == rollout_id)
.cloned()
.collect();
if selected.is_empty() {
let refs: Vec<&CodexSessionFile> = candidates.iter().collect();
warn!(
rollout_id,
count = candidates.len(),
sample = ?sample_of(&refs),
"codex-session: the request names a rollout none of the live candidates \
are; refusing to guess",
);
return None;
}
Some(selected)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Resolved {
pub(crate) session: CodexSessionFile,
pub(crate) from_unique_candidate: bool,
}
pub(crate) fn one_live_session_or_refuse(
reason: &str,
detail: &str,
candidates: Vec<CodexSessionFile>,
) -> Option<Resolved> {
let from_unique_candidate = candidates.len() == 1;
let mut ids: Vec<&str> = candidates
.iter()
.map(|session| session.session_id.as_str())
.collect();
ids.sort_unstable();
ids.dedup();
if ids.len() > 1 {
let refs: Vec<&CodexSessionFile> = candidates.iter().collect();
warn!(
reason,
count = candidates.len(),
sample = ?sample_of(&refs),
"codex-session: {detail}; refusing to guess",
);
return None;
}
unique_or_newest(reason, candidates).map(|session| Resolved {
session,
from_unique_candidate,
})
}
pub(crate) fn marker_match(
rollout_id: Option<&str>,
candidates: Vec<CodexSessionFile>,
) -> Option<Resolved> {
let candidates = narrow_by_rollout_id(rollout_id, candidates)?;
one_live_session_or_refuse(
"marker",
"marker shared by multiple LIVE sessions and the request named no rollout",
candidates,
)
}
pub(crate) fn peer_match(
rollout_id: Option<&str>,
candidates: Vec<CodexSessionFile>,
) -> Option<Resolved> {
let candidates = narrow_by_rollout_id(rollout_id, candidates)?;
one_live_session_or_refuse(
"peer-open-file",
"one process holds multiple LIVE rollouts open (a sub-thread family) and the \
request named no rollout",
candidates,
)
}
pub(crate) fn unique_or_newest(
reason: &str,
candidates: Vec<CodexSessionFile>,
) -> Option<CodexSessionFile> {
match candidates.as_slice() {
[] => None,
[session] => Some(session.clone()),
_ => {
let refs: Vec<&CodexSessionFile> = candidates.iter().collect();
warn!(
reason,
count = candidates.len(),
sample = ?sample_of(&refs),
"codex-session: one session across multiple rollout files; \
using most recently modified",
);
candidates
.into_iter()
.max_by_key(|session| session.modified_at.unwrap_or(session.timestamp))
}
}
}
fn sample_of(candidates: &[&CodexSessionFile]) -> Vec<String> {
candidates
.iter()
.take(3)
.map(|session| format!("{} ({})", session.session_id, session.path.display()))
.collect()
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::attribution::codex::CodexWatcherSnapshot;
use crate::attribution::pipeline::CodexProviderFilter;
use crate::attribution::{WatcherSnapshot, pipeline::AttributionConfig};
use arc_swap::ArcSwap;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration as StdDuration;
#[derive(Debug, Default)]
struct Hooks {
sessions: Vec<String>,
subagents: Vec<(String, String)>,
}
impl Hooks {
fn with_session(session_id: &str) -> Self {
Self {
sessions: vec![session_id.to_owned()],
subagents: Vec::new(),
}
}
fn with_subagent(parent: &str, agent: &str) -> Self {
Self {
sessions: Vec::new(),
subagents: vec![(parent.to_owned(), agent.to_owned())],
}
}
}
impl CodexHookEvidence for Hooks {
fn has_hook_session(&self, session_id: &str) -> bool {
self.sessions.iter().any(|id| id == session_id)
}
fn has_hook_subagent(&self, parent_session_id: &str, agent_id: &str) -> bool {
self.subagents
.iter()
.any(|(parent, agent)| parent == parent_session_id && agent == agent_id)
}
}
fn config() -> AttributionConfig {
AttributionConfig::new(
CodexProviderFilter::new("paper-openai"),
crate::harness::RegistryUserAgents,
)
}
fn state_with(sessions: Vec<CodexSessionFile>) -> AttributionState {
AttributionState::new(
Arc::new(ArcSwap::from_pointee(WatcherSnapshot::default())),
Arc::new(ArcSwap::from_pointee(CodexWatcherSnapshot { sessions })),
)
}
fn session(id: &str, modified_at: time::OffsetDateTime) -> CodexSessionFile {
session_with_provider(id, modified_at, "paper-openai")
}
fn session_with_provider(
id: &str,
modified_at: time::OffsetDateTime,
provider: &str,
) -> CodexSessionFile {
CodexSessionFile {
session_id: id.to_owned(),
root_session_id: Some(id.to_owned()),
parent_thread_id: None,
subagent_kind: None,
timestamp: modified_at,
modified_at: Some(modified_at),
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(provider.to_owned()),
path: PathBuf::from(format!("/tmp/{id}.jsonl")),
}
}
fn descendant_session(
id: &str,
root: &str,
parent: &str,
kind: &str,
modified_at: time::OffsetDateTime,
) -> CodexSessionFile {
let mut session = session(id, modified_at);
session.root_session_id = Some(root.to_owned());
session.parent_thread_id = Some(parent.to_owned());
session.thread_source = Some("subagent".to_owned());
session.source = Some(format!(r#"{{"subagent":{{"other":"{kind}"}}}}"#));
session.subagent_kind = Some(kind.to_owned());
session
}
fn child_session(
id: &str,
parent: &str,
kind: &str,
modified_at: time::OffsetDateTime,
) -> CodexSessionFile {
descendant_session(id, parent, parent, kind, modified_at)
}
fn descendant_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,
}
}
fn child_identity(parent: &str, child: &str, kind: &str) -> CodexRequestIdentity {
descendant_identity(parent, parent, child, kind)
}
fn thread_identity(thread_id: Option<&str>) -> CodexRequestIdentity {
CodexRequestIdentity {
correlation_id: "correlation-test".to_owned(),
thread_id: thread_id.map(str::to_owned),
..CodexRequestIdentity::default()
}
}
fn facts<'a>(
marker: Option<&'a str>,
identity: &'a CodexRequestIdentity,
hooks: Option<&'a dyn CodexHookEvidence>,
) -> RequestFacts<'a> {
RequestFacts {
codex_marker: marker,
codex_route: true,
codex_identity: Some(identity),
codex_hook_evidence: hooks,
..RequestFacts::default()
}
}
fn once(
state: &AttributionState,
marker: Option<&str>,
identity: &CodexRequestIdentity,
hooks: Option<&dyn CodexHookEvidence>,
) -> CodexSelectionResult {
select_once(state, &config(), facts(marker, identity, hooks), identity).result
}
#[test]
fn hook_evidence_and_a_named_thread_beat_a_newer_sibling() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("hook-session", now - time::Duration::seconds(1)),
session("newest-session", now),
]);
let hooks = Hooks::with_session("hook-session");
let got = once(
&state,
Some("paper-openai"),
&thread_identity(Some("hook-session")),
Some(&hooks),
);
assert_eq!(got.session.unwrap().session_id, "hook-session");
assert_eq!(
got.evidence.selection,
CodexSelection::HookExact {
session_id: "hook-session".to_owned(),
},
);
}
#[test]
fn a_child_shaped_request_joins_its_own_transcript() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("parent", now),
child_session(
"child",
"parent",
"guardian",
now - time::Duration::seconds(1),
),
]);
let identity = child_identity("parent", "child", "guardian");
let got = once(&state, Some("paper-openai"), &identity, None);
assert_eq!(got.session.unwrap().session_id, "child");
assert_eq!(
got.evidence.selection,
CodexSelection::ChildTranscriptExact {
session_id: "child".to_owned(),
parent_session_id: "parent".to_owned(),
},
);
assert_eq!(got.evidence.child_candidates.len(), 1);
}
#[test]
fn a_nested_child_joins_on_its_direct_parent_not_the_root() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("root", now),
child_session(
"parent",
"root",
"explorer",
now - time::Duration::seconds(1),
),
descendant_session(
"child",
"root",
"parent",
"guardian",
now - time::Duration::seconds(2),
),
]);
let identity = descendant_identity("root", "parent", "child", "guardian");
let got = once(&state, Some("paper-openai"), &identity, None);
assert_eq!(got.session.unwrap().session_id, "child");
assert_eq!(
got.evidence.selection,
CodexSelection::ChildTranscriptExact {
session_id: "child".to_owned(),
parent_session_id: "parent".to_owned(),
},
);
}
#[test]
fn lifecycle_evidence_upgrades_a_child_join_to_hook_backed() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("parent", now),
child_session(
"child",
"parent",
"guardian",
now - time::Duration::seconds(1),
),
]);
let hooks = Hooks::with_subagent("parent", "child");
let identity = child_identity("parent", "child", "guardian");
let got = once(&state, Some("paper-openai"), &identity, Some(&hooks));
assert_eq!(
got.evidence.selection,
CodexSelection::HookSubagentExact {
session_id: "child".to_owned(),
parent_session_id: "parent".to_owned(),
},
);
}
#[test]
fn duplicate_child_transcripts_are_not_an_identification() {
let now = time::OffsetDateTime::now_utc();
let child = child_session(
"child",
"parent",
"guardian",
now - time::Duration::seconds(1),
);
let state = state_with(vec![session("parent", now), child.clone(), child]);
let identity = child_identity("parent", "child", "guardian");
let got = once(&state, Some("paper-openai"), &identity, None);
assert_eq!(got.evidence.child_candidates.len(), 2);
assert_eq!(got.session.unwrap().session_id, "child");
assert_eq!(
got.evidence.selection,
CodexSelection::MarkerNewest {
session_id: "child".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn a_self_contradicting_request_gets_no_identity_rungs() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("hook-session", now - time::Duration::seconds(1)),
session("newest-session", now),
]);
let hooks = Hooks::with_session("hook-session");
let identity = CodexRequestIdentity::from_headers(&{
let mut headers = http::HeaderMap::new();
headers.insert("thread-id", "hook-session".parse().unwrap());
headers.insert(
"x-codex-turn-metadata",
r#"{"thread_id":"contradictory-session"}"#.parse().unwrap(),
);
headers
});
assert!(identity.conflicting_metadata);
let config = config();
let got = tokio::time::timeout(
config.codex_timeout + config.codex_poll,
select(
&state,
&config,
facts(Some("paper-openai"), &identity, Some(&hooks)),
&identity,
),
)
.await
.expect("the ladder must still return inside its budget");
assert_eq!(got.session, None);
assert_eq!(got.evidence.selection, CodexSelection::NoMatch);
}
#[test]
fn a_marker_selects_its_own_launch() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session_with_provider("sid-1", now, "paper-openai-one"),
session_with_provider("sid-2", now, "paper-openai-two"),
]);
let got = once(
&state,
Some("paper-openai-two"),
&CodexRequestIdentity::default(),
None,
);
assert_eq!(got.session.unwrap().session_id, "sid-2");
assert_eq!(
got.evidence.selection,
CodexSelection::MarkerUnique {
session_id: "sid-2".to_owned(),
},
);
}
#[test]
fn a_named_thread_narrows_an_ambiguous_marker_without_hook_evidence() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("named-thread", now - time::Duration::seconds(1)),
session("newest-session", now),
]);
let got = once(
&state,
Some("paper-openai"),
&thread_identity(Some("named-thread")),
None,
);
assert_eq!(got.session.unwrap().session_id, "named-thread");
}
#[test]
fn an_ambiguous_marker_with_no_thread_evidence_refuses() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("a", now - time::Duration::seconds(1)),
session("b", now),
]);
let got = once(
&state,
Some("paper-openai"),
&CodexRequestIdentity::default(),
None,
);
assert_eq!(got.session, None);
assert_eq!(got.evidence.selection, CodexSelection::NoMatch);
}
#[test]
fn a_request_naming_an_invisible_rollout_refuses_the_visible_one() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![session("visible", now)]);
let got = once(
&state,
Some("paper-openai"),
&thread_identity(Some("not-on-disk-yet")),
None,
);
assert_eq!(got.session, None);
}
#[test]
fn marker_rotation_of_one_session_still_resolves_to_the_newest_file() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("sid-a", now - time::Duration::minutes(5)),
session("sid-a", now),
]);
let got = once(
&state,
Some("paper-openai"),
&CodexRequestIdentity::default(),
None,
);
assert_eq!(got.session.unwrap().session_id, "sid-a");
assert_eq!(
got.evidence.selection,
CodexSelection::MarkerNewest {
session_id: "sid-a".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn a_child_shaped_request_waits_for_its_transcript_to_appear() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![session("parent", now)]);
let identity = child_identity("parent", "child", "guardian");
let config = config();
let mut task = std::pin::pin!(select(
&state,
&config,
facts(Some("paper-openai"), &identity, None),
&identity,
));
assert!(
tokio::time::timeout(StdDuration::from_millis(1), task.as_mut())
.await
.is_err(),
"a child-shaped request should wait briefly for its exact transcript",
);
state.codex_watcher.store(Arc::new(CodexWatcherSnapshot {
sessions: vec![
session("parent", now),
child_session("child", "parent", "guardian", now),
],
}));
tokio::time::advance(config.codex_poll).await;
let got = tokio::time::timeout(config.codex_timeout, task.as_mut())
.await
.expect("the child transcript must join inside the budget");
assert_eq!(
got.evidence.selection,
CodexSelection::ChildTranscriptExact {
session_id: "child".to_owned(),
parent_session_id: "parent".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn hook_evidence_waits_out_a_ready_heuristic_answer() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![session("fallback-session", now)]);
let hooks = Hooks::with_session("hook-session");
let identity = thread_identity(Some("hook-session"));
let config = config();
let mut task = std::pin::pin!(select(
&state,
&config,
facts(Some("paper-openai"), &identity, Some(&hooks)),
&identity,
));
assert!(
tokio::time::timeout(StdDuration::from_millis(1), task.as_mut())
.await
.is_err(),
"a ready heuristic must not win while exact hook evidence is pending",
);
state.codex_watcher.store(Arc::new(CodexWatcherSnapshot {
sessions: vec![
session("fallback-session", now),
session("hook-session", now),
],
}));
tokio::time::advance(config.codex_poll).await;
let got = tokio::time::timeout(config.codex_timeout, task.as_mut())
.await
.expect("the exact transcript must join inside the budget");
assert_eq!(got.session.unwrap().session_id, "hook-session");
assert_eq!(
got.evidence.selection,
CodexSelection::HookExact {
session_id: "hook-session".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn an_exact_wait_that_times_out_still_returns_within_the_budget() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![
session("older-fallback", now - time::Duration::seconds(1)),
session("newest-fallback", now),
]);
let hooks = Hooks::with_session("missing-hook-session");
let identity = thread_identity(Some("missing-hook-session"));
let config = config();
let mut task = std::pin::pin!(select(
&state,
&config,
facts(Some("paper-openai"), &identity, Some(&hooks)),
&identity,
));
assert!(
tokio::time::timeout(StdDuration::from_millis(1), task.as_mut())
.await
.is_err(),
"hook-backed attribution should wait before giving up",
);
let got = tokio::time::timeout(config.codex_timeout, task.as_mut())
.await
.expect("the ladder must return within the existing budget");
assert_eq!(got.session, None);
}
#[tokio::test(start_paused = true)]
async fn an_exact_wait_retains_a_usable_answer_across_an_empty_scan() {
let now = time::OffsetDateTime::now_utc();
let state = state_with(vec![session("child", now)]);
let identity = child_identity("parent", "child", "guardian");
let config = config();
let mut task = std::pin::pin!(select(
&state,
&config,
facts(Some("paper-openai"), &identity, None),
&identity,
));
assert!(
tokio::time::timeout(StdDuration::from_millis(1), task.as_mut())
.await
.is_err(),
"a child-shaped request waits for its exact transcript",
);
state
.codex_watcher
.store(Arc::new(CodexWatcherSnapshot::default()));
tokio::time::advance(config.codex_poll).await;
tokio::time::advance(config.codex_timeout).await;
let got = task.await;
assert_eq!(got.session.unwrap().session_id, "child");
assert_eq!(
got.evidence.selection,
CodexSelection::MarkerUnique {
session_id: "child".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn an_explicit_marker_never_falls_back_to_a_recent_session() {
let state = state_with(vec![session_with_provider(
"sid-1",
time::OffsetDateTime::now_utc(),
"paper-openai-one",
)]);
let identity = CodexRequestIdentity::default();
let got = select(
&state,
&config(),
facts(Some("paper-openai-missing"), &identity, None),
&identity,
)
.await;
assert_eq!(got.session, None);
}
#[tokio::test(start_paused = true)]
async fn an_unmarked_request_falls_back_to_a_single_recent_session() {
let state = state_with(vec![session("sole", time::OffsetDateTime::now_utc())]);
let identity = CodexRequestIdentity::default();
let got = select(&state, &config(), facts(None, &identity, None), &identity).await;
assert_eq!(got.session.unwrap().session_id, "sole");
assert_eq!(
got.evidence.selection,
CodexSelection::RecentUnique {
session_id: "sole".to_owned(),
},
);
}
#[tokio::test(start_paused = true)]
async fn stale_and_ambiguous_recent_sessions_are_both_refused() {
let stale = time::OffsetDateTime::now_utc() - time::Duration::hours(1);
let identity = CodexRequestIdentity::default();
let got = select(
&state_with(vec![session("stale", stale)]),
&config(),
facts(None, &identity, None),
&identity,
)
.await;
assert_eq!(got.session, None);
let now = time::OffsetDateTime::now_utc();
let got = select(
&state_with(vec![session("a", now), session("b", now)]),
&config(),
facts(None, &identity, None),
&identity,
)
.await;
assert_eq!(got.session, None);
}
fn subagent_family() -> Vec<CodexSessionFile> {
let now = time::OffsetDateTime::now_utc();
vec![
session("sid-parent", now - time::Duration::seconds(30)),
session("sid-child-a", now - time::Duration::seconds(20)),
session("sid-child-b", now - time::Duration::seconds(1)),
]
}
#[test]
fn peer_lane_selects_the_thread_the_request_names_not_the_newest() {
let got = peer_match(Some("sid-parent"), subagent_family())
.expect("the named rollout is right there among the candidates");
assert_eq!(got.session.session_id, "sid-parent");
assert!(got.from_unique_candidate, "narrowing left exactly one");
}
#[test]
fn peer_lane_refuses_a_live_family_with_no_thread_evidence() {
assert!(peer_match(None, subagent_family()).is_none());
}
#[test]
fn peer_lane_refuses_when_the_named_rollout_is_not_among_the_candidates() {
assert!(peer_match(Some("sid-elsewhere"), subagent_family()).is_none());
}
#[test]
fn a_lone_rollout_still_attributes_without_thread_evidence() {
let sole = vec![session("sid-sole", time::OffsetDateTime::now_utc())];
let got = peer_match(None, sole).expect("one live rollout is unambiguous");
assert_eq!(got.session.session_id, "sid-sole");
}
#[test]
fn thread_evidence_still_collapses_rotation_of_the_named_session() {
let now = time::OffsetDateTime::now_utc();
let got = peer_match(
Some("sid-child-a"),
vec![
session("sid-child-a", now - time::Duration::minutes(4)),
session("sid-child-a", now - time::Duration::seconds(2)),
session("sid-parent", now - time::Duration::seconds(1)),
],
)
.expect("rotation of the named session resolves");
assert_eq!(got.session.session_id, "sid-child-a");
assert!(
!got.from_unique_candidate,
"two files of one session is a collapse, not an identification",
);
}
#[test]
fn an_empty_candidate_set_is_a_miss_not_a_refusal() {
assert!(peer_match(Some("sid-parent"), Vec::new()).is_none());
assert!(narrow_by_rollout_id(Some("sid-parent"), Vec::new()).is_some());
}
#[test]
fn the_marker_lane_is_thread_aware_too() {
let got = marker_match(Some("sid-child-b"), subagent_family())
.expect("thread evidence resolves what the marker cannot");
assert_eq!(got.session.session_id, "sid-child-b");
}
#[test]
fn marker_shared_by_distinct_live_sessions_refuses() {
let now = time::OffsetDateTime::now_utc();
assert!(
marker_match(
None,
vec![
session_with_provider("sid-a", now, "tapesctl-openai-x"),
session_with_provider("sid-b", now, "tapesctl-openai-x"),
],
)
.is_none()
);
}
#[test]
fn rotation_ties_break_to_the_most_recently_modified() {
let now = time::OffsetDateTime::now_utc();
let got = unique_or_newest(
"marker",
vec![
session("sid-a", now - time::Duration::minutes(5)),
session("sid-a", now - time::Duration::seconds(5)),
],
);
assert_eq!(got.map(|s| s.session_id), Some("sid-a".to_owned()));
}
}