WP-11: package and OS updates from the UI
package_upgrade job streams apt-get output into the job log (streaming CommandRunner, DebianUpdater with dist-upgrade or --only-upgrade), refreshes the inventory afterwards and flags reboot-required. POST /api/system/upgrade (admin), package name validation, selectable package table with confirm dialog and live log on the Updates page. Fake updater for dev. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@ -4,7 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use domain::host::validate_package_name;
|
||||
use domain::ports::HostUpdater;
|
||||
use domain::ports::{HostUpdater, LineSink};
|
||||
use domain::DomainError;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
@ -30,8 +30,13 @@ impl UpgradeParams {
|
||||
}
|
||||
|
||||
pub fn from_json(s: Option<&str>) -> Result<Self, DomainError> {
|
||||
let _ = s;
|
||||
todo!()
|
||||
let params = match s.map(str::trim).filter(|s| !s.is_empty()) {
|
||||
None => Self::default(),
|
||||
Some(s) => serde_json::from_str(s)
|
||||
.map_err(|e| DomainError::Validation(format!("invalid params: {e}")))?,
|
||||
};
|
||||
params.validate()?;
|
||||
Ok(params)
|
||||
}
|
||||
}
|
||||
|
||||
@ -40,9 +45,52 @@ pub struct PackageUpgradeJob {
|
||||
pub inventory: Arc<InventoryService>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl JobHandler for PackageUpgradeJob {
|
||||
async fn run(&self, _params: Option<String>, _log: &dyn JobLog) -> Result<(), String> {
|
||||
todo!()
|
||||
/// Bridges the synchronous `LineSink` of the updater to the async job log.
|
||||
struct ChannelSink(tokio::sync::mpsc::UnboundedSender<String>);
|
||||
|
||||
impl LineSink for ChannelSink {
|
||||
fn line(&self, text: &str) {
|
||||
let _ = self.0.send(text.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl JobHandler for PackageUpgradeJob {
|
||||
async fn run(&self, params: Option<String>, log: &dyn JobLog) -> Result<(), String> {
|
||||
let params = UpgradeParams::from_json(params.as_deref()).map_err(|e| e.to_string())?;
|
||||
if params.packages.is_empty() {
|
||||
log.line("upgrading all packages").await;
|
||||
} else {
|
||||
log.line(&format!("upgrading: {}", params.packages.join(" ")))
|
||||
.await;
|
||||
}
|
||||
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
let sink = ChannelSink(tx);
|
||||
let upgrade = async {
|
||||
let r = self.updater.upgrade(¶ms.packages, &sink).await;
|
||||
drop(sink);
|
||||
r
|
||||
};
|
||||
let drain = async {
|
||||
while let Some(line) = rx.recv().await {
|
||||
log.line(&line).await;
|
||||
}
|
||||
};
|
||||
let (result, _) = tokio::join!(upgrade, drain);
|
||||
result.map_err(|e| e.to_string())?;
|
||||
|
||||
log.line("refreshing inventory").await;
|
||||
let inv = self.inventory.refresh().await.map_err(|e| e.to_string())?;
|
||||
log.line(&format!(
|
||||
"{} packages, {} upgradable",
|
||||
inv.packages.len(),
|
||||
inv.upgradable()
|
||||
))
|
||||
.await;
|
||||
if inv.os.reboot_required {
|
||||
log.line("NOTE: reboot required to complete the update")
|
||||
.await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user