magi-code 0.80.1

Repository-aware CLI coding agent for terminal work
Documentation
//! One bounded source capture slot; publication stays on the coordinator.
use super::*;

pub(super) struct PendingSourceCapture {
    connection: String,
    request: Request,
    pub(super) deadline: Instant,
    worker: std::thread::JoinHandle<Result<super::super::application::PreparedSource, Code>>,
}
impl PendingSourceCapture {
    pub(super) fn matches_request(&self, connection: &str, request: &Request) -> bool {
        self.connection == connection && self.request.request_id == request.request_id
    }
    pub(super) fn session(&self) -> Option<&str> {
        self.request.session_id.as_deref()
    }
    pub(super) fn join(self) {
        let _ = self.worker.join();
    }
}
impl Coordinator {
    pub(super) fn prepare_source_registration(
        &mut self,
        connection: &str,
        request: &Request,
        now: Instant,
    ) -> Option<Result<Value, Code>> {
        if !self.connections[connection].application_profile {
            return Some(Err(Code::UnsupportedCapability));
        }
        if let Err(code) = self.validate_grant(connection, request) {
            return Some(Err(code));
        }
        if self.source_capture.is_some() {
            return Some(Err(Code::ConfigurationBusy));
        }
        if let Err(code) = self
            .application
            .lock()
            .unwrap_or_else(|e| e.into_inner())
            .registry
            .reserve_source_capture(
                request.session_id.as_deref().expect("validated session"),
                &request.payload,
            )
        {
            return Some(Err(code));
        }
        let runtime = Arc::clone(&self.runtime);
        let payload = request.payload.clone();
        let worker = match std::thread::Builder::new()
            .name("magi-resource-capture".into())
            .spawn(move || super::super::application::prepare_source(&payload, &runtime))
        {
            Ok(worker) => worker,
            Err(_) => {
                self.release_source_reservation();
                return Some(Err(Code::InternalError));
            }
        };
        self.source_capture = Some(PendingSourceCapture {
            connection: connection.into(),
            request: request.clone(),
            deadline: now + Duration::from_secs(30),
            worker,
        });
        None
    }
    pub(super) fn finish_source_capture(&mut self, now: Instant) {
        if !self
            .source_capture
            .as_ref()
            .is_some_and(|pending| pending.worker.is_finished())
        {
            return;
        }
        let pending = self.source_capture.take().expect("finished capture");
        let prepared = pending.worker.join().unwrap_or(Err(Code::InternalError));
        let result = (|| {
            if now >= pending.deadline {
                return Err(Code::RequestTimeout);
            }
            self.validate_connection(&pending.connection, &pending.request)?;
            self.validate_grant(&pending.connection, &pending.request)?;
            if self.operations.lookup(
                &self.instance,
                &self.instance,
                pending.request.operation_id.as_deref().unwrap_or_default(),
            )["state"]
                != "in_progress"
            {
                return Err(Code::OperationAlreadyKnown);
            }
            let attribution = json!({"instance_id":self.instance,"connection_id":pending.connection,
                "operation_id":pending.request.operation_id,"grant_generation":pending.request.control.as_ref().expect("validated grant").generation});
            self.application
                .lock()
                .unwrap_or_else(|e| e.into_inner())
                .registry
                .publish_source(
                    pending
                        .request
                        .session_id
                        .as_deref()
                        .expect("validated session"),
                    prepared?,
                    &self.runtime,
                    attribution,
                )
        })();
        self.release_source_reservation();
        self.settle_immediate(&pending.request, &result, now);
        self.queue(
            &pending.connection,
            wire::response(
                &self.instance,
                &pending.connection,
                &pending.request,
                result,
            ),
            false,
        );
    }
    fn release_source_reservation(&self) {
        self.application
            .lock()
            .unwrap_or_else(|e| e.into_inner())
            .registry
            .release_source_capture();
    }
}