mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
Persist agent messages as response items (#29829)
## Why Inter-agent messages are recorded in live history as `ResponseItem::AgentMessage`, but rollouts stored `InterAgentCommunication` and rebuilt the response item during resume. This made the rollout differ from the actual Responses history. ## What changed - store the prepared `agent_message` response item directly - keep `trigger_turn` in a small local metadata record for fork truncation - keep reading older `inter_agent_communication` rollout items
This commit is contained in:
@@ -2841,7 +2841,7 @@ impl Session {
|
||||
std::slice::from_ref(&response_item),
|
||||
);
|
||||
let items = items.as_ref();
|
||||
communication.id = items.first().and_then(ResponseItem::id).map(str::to_string);
|
||||
let response_item = items[0].clone();
|
||||
{
|
||||
let mut state = self.state.lock().await;
|
||||
state.record_items(
|
||||
@@ -2849,8 +2849,13 @@ impl Session {
|
||||
turn_context.model_info.truncation_policy.into(),
|
||||
);
|
||||
}
|
||||
self.persist_rollout_items(&[RolloutItem::InterAgentCommunication(communication)])
|
||||
.await;
|
||||
self.persist_rollout_items(&[
|
||||
RolloutItem::InterAgentCommunicationMetadata {
|
||||
trigger_turn: communication.trigger_turn,
|
||||
},
|
||||
RolloutItem::ResponseItem(response_item),
|
||||
])
|
||||
.await;
|
||||
self.send_raw_response_items(turn_context, items).await;
|
||||
}
|
||||
|
||||
|
||||
@@ -264,6 +264,7 @@ impl Session {
|
||||
active_segment.get_or_insert_with(ActiveReplaySegment::default);
|
||||
active_segment.counts_as_user_turn = true;
|
||||
}
|
||||
RolloutItem::InterAgentCommunicationMetadata { .. } => {}
|
||||
RolloutItem::EventMsg(_) | RolloutItem::SessionMeta(_) => {}
|
||||
}
|
||||
|
||||
@@ -320,6 +321,7 @@ impl Session {
|
||||
turn_context.model_info.truncation_policy.into(),
|
||||
);
|
||||
}
|
||||
RolloutItem::InterAgentCommunicationMetadata { .. } => {}
|
||||
RolloutItem::Compacted(compacted) => {
|
||||
if let Some(replacement_history) = &compacted.replacement_history {
|
||||
// This should actually never happen, because the reverse loop above (to build rollout_suffix)
|
||||
|
||||
@@ -1766,6 +1766,29 @@ async fn record_inter_agent_communication_sets_turn_id_in_rollout_and_resume() {
|
||||
else {
|
||||
panic!("expected resumed rollout history");
|
||||
};
|
||||
let persisted_items = resumed
|
||||
.history
|
||||
.iter()
|
||||
.filter(|item| {
|
||||
matches!(
|
||||
item,
|
||||
RolloutItem::ResponseItem(_)
|
||||
| RolloutItem::InterAgentCommunication(_)
|
||||
| RolloutItem::InterAgentCommunicationMetadata { .. }
|
||||
)
|
||||
})
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
let expected_persisted_items = vec![
|
||||
RolloutItem::InterAgentCommunicationMetadata {
|
||||
trigger_turn: false,
|
||||
},
|
||||
RolloutItem::ResponseItem(expected_item.clone()),
|
||||
];
|
||||
assert_eq!(
|
||||
serde_json::to_value(persisted_items).unwrap(),
|
||||
serde_json::to_value(expected_persisted_items).unwrap()
|
||||
);
|
||||
|
||||
let (resumed_session, _resumed_turn_context) = make_session_and_context().await;
|
||||
resumed_session
|
||||
@@ -1818,14 +1841,11 @@ async fn record_inter_agent_communication_preserves_item_id_in_rollout_and_resum
|
||||
else {
|
||||
panic!("expected resumed rollout history");
|
||||
};
|
||||
let persisted_communication = resumed.history.iter().find_map(|item| match item {
|
||||
RolloutItem::InterAgentCommunication(communication) => Some(communication),
|
||||
let persisted_item_id = resumed.history.iter().find_map(|item| match item {
|
||||
RolloutItem::ResponseItem(item @ ResponseItem::AgentMessage { .. }) => item.id(),
|
||||
_ => None,
|
||||
});
|
||||
assert_eq!(
|
||||
persisted_communication.and_then(|communication| communication.id.as_deref()),
|
||||
Some(live_item_id.as_str())
|
||||
);
|
||||
assert_eq!(persisted_item_id, Some(live_item_id.as_str()));
|
||||
|
||||
let (resumed_session, _resumed_turn_context, _rx) =
|
||||
make_session_and_context_with_auth_and_config_and_rx(
|
||||
@@ -2727,6 +2747,7 @@ async fn start_new_context_window_assigns_and_persists_item_ids() {
|
||||
RolloutItem::SessionMeta(_)
|
||||
| RolloutItem::ResponseItem(_)
|
||||
| RolloutItem::InterAgentCommunication(_)
|
||||
| RolloutItem::InterAgentCommunicationMetadata { .. }
|
||||
| RolloutItem::TurnContext(_)
|
||||
| RolloutItem::EventMsg(_) => None,
|
||||
});
|
||||
@@ -2783,6 +2804,7 @@ async fn record_initial_history_assigns_and_persists_id_for_forked_response_item
|
||||
RolloutItem::ResponseItem(response_item) => response_item.id(),
|
||||
RolloutItem::SessionMeta(_)
|
||||
| RolloutItem::InterAgentCommunication(_)
|
||||
| RolloutItem::InterAgentCommunicationMetadata { .. }
|
||||
| RolloutItem::Compacted(_)
|
||||
| RolloutItem::TurnContext(_)
|
||||
| RolloutItem::EventMsg(_) => None,
|
||||
|
||||
Reference in New Issue
Block a user