//! Package upgrade job: runs the host updater, streams its output into the job log, //! then refreshes the inventory. use std::sync::Arc; use async_trait::async_trait; use domain::host::validate_package_name; use domain::ports::{HostUpdater, LineSink}; use domain::DomainError; use serde::{Deserialize, Serialize}; use crate::jobs::{JobHandler, JobLog}; use crate::InventoryService; #[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] pub struct UpgradeParams { /// Empty = upgrade all packages. #[serde(default)] pub packages: Vec, } impl UpgradeParams { pub fn validate(&self) -> Result<(), DomainError> { self.packages .iter() .try_for_each(|p| validate_package_name(p)) } pub fn to_json(&self) -> String { serde_json::to_string(self).unwrap_or_default() } pub fn from_json(s: Option<&str>) -> Result { 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) } } pub struct PackageUpgradeJob { pub updater: Arc, pub inventory: Arc, } /// Bridges the synchronous `LineSink` of the updater to the async job log. pub struct ChannelSink(tokio::sync::mpsc::UnboundedSender); pub fn channel_sink(tx: tokio::sync::mpsc::UnboundedSender) -> ChannelSink { ChannelSink(tx) } 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, 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(()) } }