mirror of
https://github.com/openai/codex.git
synced 2026-04-24 22:54:54 +00:00
54 lines
1.8 KiB
Rust
54 lines
1.8 KiB
Rust
use std::sync::Arc;
|
|
|
|
use codex_core::CodexConversation;
|
|
use codex_core::protocol::Op;
|
|
use tokio::sync::mpsc::UnboundedSender;
|
|
use tokio::sync::mpsc::unbounded_channel;
|
|
|
|
use crate::app_event::AppEvent;
|
|
use crate::app_event_sender::AppEventSender;
|
|
|
|
/// Spawn agent loops for an existing conversation (e.g., a forked conversation).
|
|
/// Sends the provided `SessionConfiguredEvent` immediately, then forwards subsequent
|
|
/// events and accepts Ops for submission.
|
|
pub(crate) fn spawn_agent_from_existing(
|
|
conversation: Arc<CodexConversation>,
|
|
session_configured: codex_core::protocol::SessionConfiguredEvent,
|
|
app_event_tx: AppEventSender,
|
|
) -> UnboundedSender<Op> {
|
|
let conversation_id = session_configured.session_id;
|
|
let (codex_op_tx, mut codex_op_rx) = unbounded_channel::<Op>();
|
|
|
|
let app_event_tx_clone = app_event_tx;
|
|
tokio::spawn(async move {
|
|
// Forward the captured `SessionConfigured` event so it can be rendered in the UI.
|
|
let ev = codex_core::protocol::Event {
|
|
id: "".to_string(),
|
|
msg: codex_core::protocol::EventMsg::SessionConfigured(session_configured),
|
|
};
|
|
app_event_tx_clone.send(AppEvent::CodexEventForConversation {
|
|
conversation_id,
|
|
event: ev,
|
|
});
|
|
|
|
let conversation_clone = conversation.clone();
|
|
tokio::spawn(async move {
|
|
while let Some(op) = codex_op_rx.recv().await {
|
|
let id = conversation_clone.submit(op).await;
|
|
if let Err(e) = id {
|
|
tracing::error!("failed to submit op: {e}");
|
|
}
|
|
}
|
|
});
|
|
|
|
while let Ok(event) = conversation.next_event().await {
|
|
app_event_tx_clone.send(AppEvent::CodexEventForConversation {
|
|
conversation_id,
|
|
event,
|
|
});
|
|
}
|
|
});
|
|
|
|
codex_op_tx
|
|
}
|