File size: 6,213 Bytes
52a9af3 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 | use super::TurnInput as PendingTurnInput;
use super::session::Session;
use super::turn_context::TurnContext;
use codex_features::Feature;
use codex_history::CodexHarnessMetadata;
use codex_history::ResponseItemEnvelope;
use codex_protocol::models::ResponseItem;
use codex_protocol::openai_models::ModelInfo;
impl Session {
/// Returns the input if there is no active turn to inject into.
#[expect(
clippy::await_holding_invalid_type,
reason = "active turn checks and turn state updates must remain atomic"
)]
pub(crate) async fn inject_if_running<T: Into<ResponseItemEnvelope>>(
&self,
input: Vec<T>,
) -> Result<(), Vec<T>> {
let mut active = self.active_turn.lock().await;
match active.as_mut() {
Some(active_turn) => {
self.input_queue
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
input
.into_iter()
.map(Into::into)
.map(PendingTurnInput::ResponseItem)
.collect(),
)
.await;
Ok(())
}
None => Err(input),
}
}
/// Injects hook context into the running turn atomically.
#[expect(
clippy::await_holding_invalid_type,
reason = "active turn provenance and turn state updates must remain atomic"
)]
pub(crate) async fn inject_hook_context_if_running(
&self,
input: Vec<ResponseItem>,
) -> Result<(), Vec<ResponseItem>> {
let mut active = self.active_turn.lock().await;
let Some(active_turn) = active.as_mut() else {
return Err(input);
};
if active_turn.task.is_none() {
return Err(input);
}
self.input_queue
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
input
.into_iter()
.map(ResponseItemEnvelope::new)
.map(PendingTurnInput::ResponseItem)
.collect(),
)
.await;
Ok(())
}
/// Preserves trusted client provenance while items wait for an active turn.
#[expect(
clippy::await_holding_invalid_type,
reason = "active turn checks and turn state updates must remain atomic"
)]
pub(crate) async fn inject_client_response_items(
&self,
items: Vec<ResponseItem>,
turn_context: &TurnContext,
) {
let items = items
.into_iter()
.map(|item| self.annotate_client_response_item(item))
.collect::<Vec<_>>();
let mut active = self.active_turn.lock().await;
if let Some(active_turn) = active.as_mut() {
self.input_queue
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
items
.into_iter()
.map(PendingTurnInput::ResponseItem)
.collect(),
)
.await;
return;
}
drop(active);
self.record_annotated_conversation_items(turn_context, turn_context.model_info(), items)
.await;
}
pub(crate) fn annotate_client_response_item(&self, item: ResponseItem) -> ResponseItemEnvelope {
let metadata = (self.enabled(Feature::RetainClientDeveloperMessages)
&& matches!(&item, ResponseItem::Message { role, .. } if role == "developer"))
.then_some(CodexHarnessMetadata {
client_authored: true,
..Default::default()
});
ResponseItemEnvelope { item, metadata }
}
pub(crate) async fn record_annotated_conversation_items(
&self,
turn_context: &TurnContext,
model_info: &ModelInfo,
items: Vec<ResponseItemEnvelope>,
) {
if items.iter().all(|item| item.metadata.is_none()) {
let items = items
.into_iter()
.map(ResponseItemEnvelope::into_item)
.collect::<Vec<_>>();
self.record_conversation_items(turn_context, model_info, &items)
.await;
return;
}
let mut annotated_items = Vec::with_capacity(items.len());
let mut image_preparations = Vec::new();
for envelope in items {
let (prepared_items, prepared_images) = self.prepare_conversation_items_for_history(
turn_context,
model_info,
std::slice::from_ref(&envelope.item),
);
image_preparations.extend(prepared_images);
let mut metadata = envelope.metadata;
annotated_items.extend(prepared_items.into_owned().into_iter().map(|item| {
ResponseItemEnvelope {
item,
metadata: metadata.take(),
}
}));
}
self.record_prepared_conversation_items(
turn_context,
model_info,
annotated_items,
image_preparations,
)
.await;
}
/// Injects items into active work, or records them without starting a turn.
pub(crate) async fn inject_no_new_turn(
&self,
items: Vec<ResponseItem>,
current_turn_context: Option<&TurnContext>,
) {
let Err(items) = self.inject_if_running(items).await else {
return;
};
let default_turn_context;
let turn_context = match current_turn_context {
Some(turn_context) => turn_context,
None => {
default_turn_context = self.new_default_turn().await;
default_turn_context.as_ref()
}
};
self.record_conversation_items(turn_context, turn_context.model_info(), &items)
.await;
}
}
#[cfg(test)]
#[path = "inject_tests.rs"]
mod tests;
|