magi-code 0.96.2

Repository-aware CLI coding agent for terminal work
Documentation
//! One child-owned snapshot, strict startup worker, and bounded cleanup.
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 {
                // Inspect expiry/refreshability without token exchange. The existing runner
                // still rereads and refreshes credentials immediately before provider work.
                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();
        }
        // Never run MCP destruction in the protocol loop or on its Drop stack.
        // The one cleanup worker owns the startup result/manager until closed.
        // Spawn failure drops the closure: suppress synchronous MCP destruction then.
        // These resources deliberately survive until process exit under thread exhaustion.
        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();
            }
            // Absolute process-exit bound: a slow HTTP teardown may be detached here.
        }
    }
}