54 lines
1.8 KiB
Rust
54 lines
1.8 KiB
Rust
|
|
use std::sync::Arc;
|
||
|
|
|
||
|
|
use codex_core::codex_wrapper::CodexConversation;
|
||
|
|
use codex_core::codex_wrapper::init_codex;
|
||
|
|
use codex_core::config::Config;
|
||
|
|
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 the agent bootstrapper and op forwarding loop, returning the
|
||
|
|
/// `UnboundedSender<Op>` used by the UI to submit operations.
|
||
|
|
pub(crate) fn spawn_agent(config: Config, app_event_tx: AppEventSender) -> UnboundedSender<Op> {
|
||
|
|
let (codex_op_tx, mut codex_op_rx) = unbounded_channel::<Op>();
|
||
|
|
|
||
|
|
let app_event_tx_clone = app_event_tx.clone();
|
||
|
|
tokio::spawn(async move {
|
||
|
|
let CodexConversation {
|
||
|
|
codex,
|
||
|
|
session_configured,
|
||
|
|
..
|
||
|
|
} = match init_codex(config).await {
|
||
|
|
Ok(vals) => vals,
|
||
|
|
Err(e) => {
|
||
|
|
// TODO: surface this error to the user.
|
||
|
|
tracing::error!("failed to initialize codex: {e}");
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
};
|
||
|
|
|
||
|
|
// Forward the captured `SessionInitialized` event that was consumed
|
||
|
|
// inside `init_codex()` so it can be rendered in the UI.
|
||
|
|
app_event_tx_clone.send(AppEvent::CodexEvent(session_configured.clone()));
|
||
|
|
let codex = Arc::new(codex);
|
||
|
|
let codex_clone = codex.clone();
|
||
|
|
tokio::spawn(async move {
|
||
|
|
while let Some(op) = codex_op_rx.recv().await {
|
||
|
|
let id = codex_clone.submit(op).await;
|
||
|
|
if let Err(e) = id {
|
||
|
|
tracing::error!("failed to submit op: {e}");
|
||
|
|
}
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
while let Ok(event) = codex.next_event().await {
|
||
|
|
app_event_tx_clone.send(AppEvent::CodexEvent(event));
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
codex_op_tx
|
||
|
|
}
|