use crate::{
cancellation::AgentCancellation,
config::{CliConfigOverrides, McPaths},
providers::ProviderSelection,
runtime_context::{LaunchContext, StartupPolicy},
};
use std::{
path::PathBuf,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
thread::JoinHandle,
time::{Duration, Instant},
};
pub(super) struct RuntimeContext {
pub(super) launch: LaunchContext,
pub(super) policy: StartupPolicy,
pub(super) attachment: Option<super::sessions::Attachment>,
}
pub(super) struct StartupResult {
pub(super) context_id: Option<String>,
pub(super) selection: Option<ProviderSelection>,
pub(super) result: Result<RuntimeContext, &'static str>,
}
impl StartupResult {
fn load(cwd: PathBuf, cancellation: &AgentCancellation) -> Self {
let mut loaded = Self {
context_id: None,
selection: None,
result: Err("configuration_invalid"),
};
let result = (|| {
cancellation.check().map_err(|_| "startup_cancelled")?;
let paths = McPaths::resolve().map_err(|_| "configuration_invalid")?;
let launch = LaunchContext::load(cwd, paths, CliConfigOverrides::default())
.map_err(|_| "configuration_invalid")?;
let canonical_root = launch
.config
.paths
.root
.canonicalize()
.map_err(|_| "configuration_invalid")?;
use sha2::{Digest, Sha256};
let identity = serde_json::to_vec(&(&launch.cwd, &canonical_root))
.map_err(|_| "configuration_invalid")?;
loaded.context_id = Some(crate::hex::lower_hex(Sha256::digest(identity)));
loaded.selection = Some(ProviderSelection::from_config_without_auth(&launch.config));
let policy = StartupPolicy::load(
&launch.settings,
&launch.discovered_skills,
&launch.config.paths,
cancellation,
)
.map_err(|error| error.code)?;
cancellation.check().map_err(|_| "startup_cancelled")?;
let selection = loaded.selection.as_ref().expect("selection captured above");
crate::providers::supported_custom_provider(&launch.config, selection)
.map_err(|_| "provider_unavailable")?;
let ready = if launch.config.provider_id() == crate::providers::OPENAI_CODEX_PROVIDER {
let store = crate::auth::read_auth_store(&launch.config.paths)
.map_err(|_| "provider_unavailable")?;
crate::auth::classify_provider_auth_record(
launch.config.provider_id(),
store.auth().providers.get(launch.config.provider_id()),
chrono::Utc::now().timestamp(),
)
.is_ready()
} else {
launch.config.auth_state().is_ready()
};
if !ready {
return Err("provider_unavailable");
}
Ok(RuntimeContext {
launch,
policy,
attachment: None,
})
})();
loaded.result = result;
loaded
}
}
impl RuntimeContext {
pub(super) fn run_options<'a, 'sink>(
&'a self,
prompt: &'a str,
session: Option<&'a crate::sessions::Session>,
output_sink: Option<&'sink mut dyn crate::agent::AgentOutputSink>,
cancel: Arc<AtomicBool>,
) -> crate::agent::runner::ProviderRunOptions<'a, 'sink> {
crate::agent::runner::ProviderRunOptions {
settings: Some(self.launch.settings.clone()),
prompt,
session,
cwd: &self.launch.cwd,
output_sink,
selected_primary_agent: self
.policy
.selected_primary_agent
.as_ref()
.and_then(|id| self.policy.primary_agent_discovery.profiles.get(id))
.cloned(),
cancellation: Some(cancel),
subagent_controls: None,
session_title_notifier: None,
herdr_reporter: None,
invocation_mode: crate::output::InvocationMode::LocalAgent,
disabled_tools: Some(Arc::new(std::sync::Mutex::new(
self.policy.disabled_tools.clone(),
))),
disabled_subagent_profiles: Some(Arc::new(std::sync::Mutex::new(
self.policy.disabled_subagent_profiles.clone(),
))),
subagent_profile_discovery: Some(self.policy.subagent_profile_discovery.clone()),
mcp: self.policy.mcp.clone(),
}
}
}
pub(super) struct StartupWorker {
cancel: Arc<AtomicBool>,
worker: Option<JoinHandle<StartupResult>>,
}
impl StartupWorker {
pub(super) fn start(cwd: PathBuf) -> std::io::Result<Self> {
let cancel = Arc::new(AtomicBool::new(false));
let cancellation = AgentCancellation::new(Arc::clone(&cancel));
let worker = std::thread::Builder::new()
.name("local-agent-startup".into())
.spawn(move || StartupResult::load(cwd, &cancellation))?;
Ok(Self {
cancel,
worker: Some(worker),
})
}
pub(super) fn poll(&mut self) -> Option<StartupResult> {
if !self.worker.as_ref()?.is_finished() {
return None;
}
Some(self.worker.take()?.join().unwrap_or(StartupResult {
context_id: None,
selection: None,
result: Err("startup_failed"),
}))
}
pub(super) fn cancel(&self) {
self.cancel.store(true, Ordering::SeqCst);
}
pub(super) fn cleanup(
mut self,
runtime: Option<RuntimeContext>,
session_worker: Option<super::sessions::SessionWorker>,
turn_worker: Option<super::turns::TurnWorker>,
deadline: Instant,
) {
self.cancel();
if let Some(worker) = &session_worker {
worker.cancel();
}
if let Some(worker) = &turn_worker {
worker.cancel();
}
let cleanup =
std::mem::ManuallyDrop::new((self.worker.take(), runtime, session_worker, turn_worker));
if let Ok(worker) = std::thread::Builder::new()
.name("local-agent-cleanup".into())
.spawn(move || {
let (startup, runtime, session_worker, turn_worker) =
std::mem::ManuallyDrop::into_inner(cleanup);
if let Some(startup) = startup {
drop(startup.join());
}
if let Some(worker) = session_worker {
worker.finish();
}
if let Some(worker) = turn_worker {
worker.finish();
}
drop(runtime);
})
{
while !worker.is_finished() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(10));
}
if worker.is_finished() {
let _ = worker.join();
}
}
}
}