//! 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:#}"))),
}
});
}