182 lines
6.4 KiB
Rust
182 lines
6.4 KiB
Rust
|
|
use std::borrow::Cow;
|
||
|
|
use std::time::Duration;
|
||
|
|
use std::time::Instant;
|
||
|
|
|
||
|
|
use tokio::process::Child;
|
||
|
|
use tracing::error;
|
||
|
|
use tracing::info;
|
||
|
|
|
||
|
|
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 RunningVerify {
|
||
|
|
activity_tree: ActivityTreeStream,
|
||
|
|
last_announce: Option<Instant>,
|
||
|
|
}
|
||
|
|
|
||
|
|
impl RunningVerify {
|
||
|
|
pub(crate) fn new() -> Result<Self> {
|
||
|
|
Ok(RunningVerify {
|
||
|
|
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<OutputStream> =
|
||
|
|
NixOutputStream::new(output_stream);
|
||
|
|
|
||
|
|
info!("Verifying nix store.");
|
||
|
|
|
||
|
|
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?;
|
||
|
|
info!("nix store verify 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()
|
||
|
|
}
|
||
|
|
|
||
|
|
pub(crate) fn is_transparent_predicate(entry: &ActivityTreeEntry) -> bool {
|
||
|
|
return false;
|
||
|
|
}
|
||
|
|
|
||
|
|
fn is_alive_predicate(entry: &ActivityTreeEntry) -> bool {
|
||
|
|
entry.get_activity().is_active()
|
||
|
|
}
|