magi-code 0.96.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::super::*;

/// A single owned outcome, published before its coalescible app wake-up.
/// Keep the handle until it finishes; receiving an outcome must never join it.
pub(super) struct CompletionWorker<T> {
    handle: JoinHandle<()>,
    receiver: Receiver<Result<T, String>>,
    consumed: bool,
}

impl<T: Send + 'static> CompletionWorker<T> {
    pub(super) fn spawn(
        name: &str,
        wake: Sender<TuiEvent>,
        task: impl FnOnce() -> Result<T, String> + Send + 'static,
    ) -> std::io::Result<Self> {
        let (sender, receiver) = crossbeam_channel::bounded(1);
        let handle = thread::Builder::new().name(name.into()).spawn(move || {
            let result = panic::catch_unwind(panic::AssertUnwindSafe(task))
                .unwrap_or_else(|_| Err("completion worker panicked".into()));
            // One producer and one result: this cannot wait for queue capacity.
            if sender.send(result).is_ok() {
                // Full means the app is already awake. Disconnection cannot
                // discard the owned outcome in the dedicated mailbox.
                let _ = wake.try_send(TuiEvent::WorkerOutcomeReady);
            }
        })?;
        Ok(Self {
            handle,
            receiver,
            consumed: false,
        })
    }

    pub(super) fn take_result(&mut self) -> Option<Result<T, String>> {
        if self.consumed {
            return None;
        }
        match self.receiver.try_recv() {
            Ok(result) => {
                self.consumed = true;
                Some(result)
            }
            Err(crossbeam_channel::TryRecvError::Disconnected) => {
                self.consumed = true;
                Some(Err("completion worker exited without an outcome".into()))
            }
            Err(crossbeam_channel::TryRecvError::Empty) => None,
        }
    }

    pub(super) fn is_finished(&self) -> bool {
        self.handle.is_finished()
    }

    pub(super) fn ready_to_reap(&self) -> bool {
        self.consumed && self.is_finished()
    }

    pub(super) fn join(self) -> Result<Option<T>, String> {
        let Self {
            handle,
            receiver,
            consumed,
        } = self;
        collect_completion_outcome(handle.join(), &receiver, consumed)
    }

    /// Returns `None` and detaches the thread when it does not finish in time.
    pub(super) fn join_with_timeout(self, timeout: Duration) -> Option<Result<Option<T>, String>> {
        let Self {
            handle,
            receiver,
            consumed,
        } = self;
        let joined = crate::thread_join::join_with_timeout(handle, timeout).ok()?;
        Some(collect_completion_outcome(joined, &receiver, consumed))
    }
}

fn collect_completion_outcome<T>(
    joined: thread::Result<()>,
    receiver: &Receiver<Result<T, String>>,
    consumed: bool,
) -> Result<Option<T>, String> {
    joined.map_err(|_| "completion worker panicked".to_string())?;
    if consumed {
        Ok(None)
    } else {
        receiver
            .recv()
            .map_err(|_| "completion worker exited without an outcome".to_string())?
            .map(Some)
    }
}