Skip to content
//! Background worker: owns the tokio runtime and the current connection.
//!
//! egui runs on the main thread and must never block, so every request is
//! handed to this worker and every result comes back as a [`BackendEvent`].

use std::sync::mpsc::{Receiver, Sender};

use tokio::sync::mpsc::{UnboundedSender, unbounded_channel};
use twice_core::{
    AppSummary, Host, InfrastructureSpec, JobOutcome, JobView, LogLine, UpdateImageRequest,
};

use crate::client::{ApiClient, Lifecycle, StreamItem};
use crate::profiles::ServerProfile;

/// A request from the UI thread.
#[derive(Clone, Debug)]
pub(crate) enum UiCommand {
    /// Point the worker at a different server and refresh.
    Connect(Box<ServerProfile>),
    /// Re-read the application list.
    RefreshApps,
    /// Start or stop an application.
    Lifecycle {
        /// Target application.
        host: Host,
        /// Which operation.
        op: Lifecycle,
    },
    /// Roll an application onto a new image.
    UpdateImage {
        /// Target application.
        host: Host,
        /// Image and flags.
        request: Box<UpdateImageRequest>,
    },
    /// Apply a declarative application specification.
    ApplyInfrastructure(Box<InfrastructureSpec>),
}

/// A result destined for the UI thread.
#[derive(Clone, Debug)]
pub(crate) enum BackendEvent {
    /// Fresh application inventory.
    Apps(Vec<AppSummary>),
    /// A job was accepted by the server.
    JobStarted(Box<JobView>),
    /// A line of live job output.
    JobLine(LogLine),
    /// The job reached a terminal state.
    JobFinished(JobOutcome),
    /// Something went wrong; the message is operator-facing.
    Failed(String),
}

/// Handle used by the UI to talk to the worker.
#[derive(Clone, Debug)]
pub(crate) struct Backend {
    commands: UnboundedSender<UiCommand>,
}

impl Backend {
    /// Starts the worker thread and returns it with its event receiver.
    pub(crate) fn spawn(ctx: egui::Context) -> (Self, Receiver<BackendEvent>) {
        let (commands, mut inbox) = unbounded_channel::<UiCommand>();
        let (events, results) = std::sync::mpsc::channel::<BackendEvent>();

        std::thread::spawn(move || {
            let runtime = match tokio::runtime::Builder::new_current_thread()
                .enable_all()
                .build()
            {
                Ok(runtime) => runtime,
                Err(err) => {
                    let _ = events.send(BackendEvent::Failed(format!("runtime: {err}")));
                    return;
                }
            };
            runtime.block_on(async move {
                let mut client: Option<ApiClient> = None;
                while let Some(command) = inbox.recv().await {
                    handle(command, &mut client, &events, &ctx).await;
                }
            });
        });

        (Self { commands }, results)
    }

    /// Queues a command; a dead worker is reported through the event channel.
    pub(crate) fn send(&self, command: UiCommand) {
        let _ = self.commands.send(command);
    }
}

async fn handle(
    command: UiCommand,
    client: &mut Option<ApiClient>,
    events: &Sender<BackendEvent>,
    ctx: &egui::Context,
) {
    let emit = |event: BackendEvent| {
        let _ = events.send(event);
        ctx.request_repaint();
    };

    match command {
        UiCommand::Connect(profile) => match ApiClient::new(&profile) {
            Ok(fresh) => {
                *client = Some(fresh);
                refresh(client.as_ref(), &emit).await;
            }
            Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
        },
        UiCommand::RefreshApps => refresh(client.as_ref(), &emit).await,
        UiCommand::Lifecycle { host, op } => {
            let Some(api) = client.clone() else {
                emit(BackendEvent::Failed("no server selected".to_owned()));
                return;
            };
            match api.lifecycle(&host, op).await {
                Ok(view) => follow(api, view, events.clone(), ctx.clone()),
                Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
            }
        }
        UiCommand::UpdateImage { host, request } => {
            let Some(api) = client.clone() else {
                emit(BackendEvent::Failed("no server selected".to_owned()));
                return;
            };
            match api.update_image(&host, &request).await {
                Ok(view) => follow(api, view, events.clone(), ctx.clone()),
                Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
            }
        }
        UiCommand::ApplyInfrastructure(spec) => {
            let Some(api) = client.clone() else {
                emit(BackendEvent::Failed("no server selected".to_owned()));
                return;
            };
            match api.apply_infrastructure(&spec).await {
                Ok(view) => follow(api, view, events.clone(), ctx.clone()),
                Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
            }
        }
    }
}

async fn refresh(client: Option<&ApiClient>, emit: &(impl Fn(BackendEvent) + Sync)) {
    let Some(api) = client else {
        emit(BackendEvent::Failed("no server selected".to_owned()));
        return;
    };
    match api.apps().await {
        Ok(apps) => emit(BackendEvent::Apps(apps)),
        Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
    }
}

/// Announces a newly accepted job, then streams its output until it ends.
fn follow(api: ApiClient, view: JobView, events: Sender<BackendEvent>, ctx: egui::Context) {
    let id = view.id;
    let _ = events.send(BackendEvent::JobStarted(Box::new(view)));
    ctx.request_repaint();

    tokio::spawn(async move {
        let emit = |event: BackendEvent| {
            let _ = events.send(event);
            ctx.request_repaint();
        };
        let sink = |item: StreamItem| match item {
            StreamItem::Line(line) => emit(BackendEvent::JobLine(line)),
            StreamItem::Finished(outcome) => emit(BackendEvent::JobFinished(outcome)),
        };
        if let Err(err) = api.stream_job(id, &sink).await {
            emit(BackendEvent::Failed(format!("{err:#}")));
        }
        // The inventory changes as a result of every job.
        match api.apps().await {
            Ok(apps) => emit(BackendEvent::Apps(apps)),
            Err(err) => emit(BackendEvent::Failed(format!("{err:#}"))),
        }
    });
}