use std::borrow::Cow; use std::time::Duration; use std::time::Instant; use tokio::process::Child; use tracing::error; use crate::Result; use crate::nix_util::NixOutputStream; use crate::nix_util::nix_output_stream::NixAction; use crate::nix_util::output_stream::OutputStream; use crate::nix_util::tree_iter::ForwardTreeIter; use super::activity_tree::ActivityTreeEntry; use super::activity_tree_stream::ActivityTreeStream; use super::nix_output_stream::ActivityResultMessage; use super::nix_output_stream::NixMessage; pub(crate) struct RunningUpdate { activity_tree: ActivityTreeStream, last_announce: Option, } impl RunningUpdate { pub(crate) fn new() -> Result { Ok(RunningUpdate { activity_tree: ActivityTreeStream::new(), last_announce: None, }) } pub(crate) async fn run_to_completion(&mut self, mut child: Child) -> Result<()> { let output_stream = OutputStream::from_child(&mut child)?; let mut nix_output_stream: NixOutputStream = NixOutputStream::new(output_stream); let exit_status_handle = tokio::spawn(async move { let status = child .wait() .await .expect("nixos-rebuild encountered an error"); status }); while let Some(message) = nix_output_stream.next().await? { self.handle_message(message)?; } let exit_status = exit_status_handle.await?; println!("nix update status was: {}", exit_status); Ok(()) } pub(crate) fn handle_message(&mut self, message: NixMessage) -> Result<()> { self.activity_tree.handle_message(&message)?; let message = match message { NixMessage::ParseFailure(line) => { error!("FAIL PARSE: {line}"); return Ok(()); } NixMessage::Generic(_value, line) => { error!("GENERIC PARSE: {line}"); return Ok(()); } NixMessage::Action(nix_action) => nix_action, }; match message { NixAction::Msg(msg_message) => { // if msg_message.level > 0 && msg_message.level < 5 { // eprintln!("LOG MESSAGE {}: {}", msg_message.level, msg_message.msg); // } } NixAction::Start(activity_start_message) => { // println!("START: {}", serde_json::to_string(&activity_start_message)?); self.print_current_status(); } NixAction::Stop(stop_message) => { // println!("STOP: {}", serde_json::to_string(&stop_message)?); self.print_current_status(); } NixAction::Result(activity_result_message) => { match activity_result_message { ActivityResultMessage::FileLinked(_activity_result_file_linked) => {} ActivityResultMessage::BuildLogLine(_activity_result_build_log_line) => {} ActivityResultMessage::UntrustedPath(_activity_result_untrusted_path) => {} ActivityResultMessage::CorruptedPath(_activity_result_corrupted_path) => {} ActivityResultMessage::SetPhase(_activity_result_set_phase) => {} ActivityResultMessage::Progress(activity_result_progress) => { // if activity_result_progress.expected != 0 { // println!( // "PROGRESS: {}", // serde_json::to_string(&activity_result_progress)? // ); // } self.maybe_print_current_status(); } ActivityResultMessage::SetExpected(activity_result_set_expected) => { // if activity_result_set_expected.expected != 0 { // println!( // "EXPECTED: {}", // serde_json::to_string(&activity_result_set_expected)? // ); // } self.maybe_print_current_status(); } ActivityResultMessage::PostBuildLogLine( _activity_result_post_build_log_line, ) => {} ActivityResultMessage::FetchStatus(_activity_result_fetch_status) => {} }; } }; Ok(()) } fn maybe_print_current_status(&mut self) -> () { let last_announce = match self.last_announce { Some(instant) => instant, None => { // If we haven't announced before, always announce. return self.print_current_status(); } }; let now = Instant::now(); let time_since_last_announce = now.duration_since(last_announce); if time_since_last_announce > Duration::new(5, 0) { return self.print_current_status(); } } fn print_current_status(&mut self) -> () { let nodes = ForwardTreeIter::new( self.activity_tree.get_tree(), None, is_match_predicate, is_transparent_predicate, is_alive_predicate, ); let mut out = Vec::new(); for (depth, _is_match, _is_transparent, node) in nodes { let progress_text = node .get_activity() .get_progress_text() .unwrap_or(Cow::Borrowed("")); let name = node .get_activity() .display_name() .unwrap_or(Cow::Borrowed("null")); out.push(format!("{depth}\t{progress_text}\t{name}")) } if out.is_empty() { println!("No active activities."); } else { println!("\n"); for l in out { println!("{l}\n"); } println!("\n"); } self.last_announce = Some(Instant::now()); } } fn is_match_predicate(entry: &ActivityTreeEntry) -> bool { entry.get_activity().get_progress_text().is_some() } fn is_transparent_predicate(entry: &ActivityTreeEntry) -> bool { return false; } fn is_alive_predicate(entry: &ActivityTreeEntry) -> bool { entry.get_activity().is_active() }