Handled suggestions from PR.

This commit is contained in:
Ember
2026-02-11 10:13:15 -08:00
committed by Cooper Quintin
parent bc7dcc97c6
commit 24e79aad9d
7 changed files with 106 additions and 146 deletions
+78 -78
View File
@@ -29,11 +29,13 @@ use crate::qmdl_store::{RecordingStore, RecordingStoreError};
use crate::server::ServerState;
use crate::stats::DiskStats;
const SPACE_CHECK_INTERVAL_CONTAINERS: usize = 100;
const DISK_CHECK_INTERVAL: usize = 100;
pub enum DiagDeviceCtrlMessage {
StopRecording,
StartRecording,
StartRecording {
response_tx: Option<oneshot::Sender<Result<(), String>>>,
},
DeleteEntry {
name: String,
response_tx: oneshot::Sender<Result<(), RecordingStoreError>>,
@@ -53,7 +55,8 @@ pub struct DiagTask {
min_space_to_continue_mb: u64,
state: DiagState,
max_type_seen: EventType,
container_count: usize,
messages_since_space_check: usize,
low_space_warned: bool,
}
enum DiagState {
@@ -82,48 +85,27 @@ impl DiagTask {
min_space_to_continue_mb,
state: DiagState::Stopped,
max_type_seen: EventType::Informational,
container_count: 0,
messages_since_space_check: 0,
low_space_warned: false,
}
}
/// Start recording
async fn start(&mut self, qmdl_store: &mut RecordingStore) {
/// Start recording, returning an error if disk space is too low.
async fn start(&mut self, qmdl_store: &mut RecordingStore) -> Result<(), String> {
self.max_type_seen = EventType::Informational;
self.container_count = 0;
self.messages_since_space_check = 0;
self.low_space_warned = false;
let min_space_bytes = self.min_space_to_start_mb * 1024 * 1024;
match DiskStats::new(qmdl_store.path.to_str().unwrap()) {
Ok(disk_stats) if disk_stats.available_bytes.unwrap_or(0) < min_space_bytes => {
let available_mb = disk_stats.available_bytes.unwrap_or(0) / 1024 / 1024;
error!(
"Insufficient disk space to start recording: {}MB available, {}MB required",
let msg = format!(
"Insufficient disk space: {}MB available, {}MB required",
available_mb, self.min_space_to_start_mb
);
if let Err(e) = self
.notification_channel
.send(Notification::new(
NotificationType::Warning,
format!(
"Cannot start recording: only {}MB free (need {}MB minimum)",
available_mb, self.min_space_to_start_mb
),
None,
))
.await
{
warn!("Failed to send notification: {e}");
}
if let Err(e) = self
.ui_update_sender
.send(display::DisplayState::Paused)
.await
{
warn!("couldn't send ui update message: {e}");
}
return;
error!("{msg}");
return Err(msg);
}
Ok(disk_stats) => {
let available_mb = disk_stats.available_bytes.unwrap_or(0) / 1024 / 1024;
@@ -140,8 +122,9 @@ impl DiagTask {
let (qmdl_file, analysis_file) = match qmdl_store.new_entry().await {
Ok(files) => files,
Err(e) => {
error!("failed creating QMDL file entry: {e}");
return;
let msg = format!("failed creating QMDL file entry: {e}");
error!("{msg}");
return Err(msg);
}
};
self.stop_current_recording().await;
@@ -150,8 +133,9 @@ impl DiagTask {
{
Ok(writer) => Box::new(writer),
Err(e) => {
error!("failed to create analysis writer: {e}");
return;
let msg = format!("failed to create analysis writer: {e}");
error!("{msg}");
return Err(msg);
}
};
self.state = DiagState::Recording {
@@ -165,11 +149,17 @@ impl DiagTask {
{
warn!("couldn't send ui update message: {e}");
}
Ok(())
}
/// Stop recording
async fn stop(&mut self, qmdl_store: &mut RecordingStore) {
/// Stop recording, optionally annotating the entry with a reason.
async fn stop(&mut self, qmdl_store: &mut RecordingStore, reason: Option<String>) {
self.stop_current_recording().await;
if let Some(reason) = reason
&& let Err(e) = qmdl_store.set_current_stop_reason(reason).await
{
warn!("couldn't set stop reason: {e}");
}
if let Some((_, entry)) = qmdl_store.get_current_entry()
&& let Err(e) = self
.analysis_sender
@@ -198,7 +188,7 @@ impl DiagTask {
name: &str,
) -> Result<(), RecordingStoreError> {
if qmdl_store.is_current_entry(name) {
self.stop(qmdl_store).await;
self.stop(qmdl_store, None).await;
}
let res = qmdl_store.delete_entry(name).await;
if let Err(e) = res.as_ref() {
@@ -211,7 +201,7 @@ impl DiagTask {
&mut self,
qmdl_store: &mut RecordingStore,
) -> Result<(), RecordingStoreError> {
self.stop(qmdl_store).await;
self.stop(qmdl_store, None).await;
let res = qmdl_store.delete_all_entries().await;
if let Err(e) = res.as_ref() {
error!("Error deleting QMDL entries {e}");
@@ -249,41 +239,36 @@ impl DiagTask {
analysis_writer,
} = &mut self.state
{
self.container_count += 1;
self.messages_since_space_check += 1;
if self
.container_count
.is_multiple_of(SPACE_CHECK_INTERVAL_CONTAINERS)
{
let min_continue_bytes = self.min_space_to_continue_mb * 1024 * 1024;
let min_start_bytes = self.min_space_to_start_mb * 1024 * 1024;
if self.messages_since_space_check >= DISK_CHECK_INTERVAL {
self.messages_since_space_check = 0;
let critical_bytes = self.min_space_to_continue_mb * 1024 * 1024;
let warning_bytes = self.min_space_to_start_mb * 1024 * 1024;
match DiskStats::new(qmdl_store.path.to_str().unwrap()) {
Ok(disk_stats)
if disk_stats.available_bytes.unwrap_or(0) < min_continue_bytes =>
{
Ok(disk_stats) if disk_stats.available_bytes.unwrap_or(0) < critical_bytes => {
let available_mb = disk_stats.available_bytes.unwrap_or(0) / 1024 / 1024;
error!(
"Disk space critically low ({}MB), stopping recording",
let reason = format!(
"Disk space critically low ({}MB free), recording stopped automatically",
available_mb
);
error!("{reason}");
self.notification_channel.send(Notification::new(
NotificationType::Warning,
format!(
"Disk space critically low ({}MB), recording stopped automatically",
available_mb
),
None,
)).await.ok();
self.notification_channel
.send(Notification::new(
NotificationType::Warning,
reason.clone(),
None,
))
.await
.ok();
self.stop(qmdl_store).await;
self.stop(qmdl_store, Some(reason)).await;
return;
}
Ok(disk_stats) if disk_stats.available_bytes.unwrap_or(0) < min_start_bytes => {
if self
.container_count
.is_multiple_of(SPACE_CHECK_INTERVAL_CONTAINERS * 10)
{
Ok(disk_stats) if disk_stats.available_bytes.unwrap_or(0) < warning_bytes => {
if !self.low_space_warned {
self.low_space_warned = true;
let available_mb =
disk_stats.available_bytes.unwrap_or(0) / 1024 / 1024;
warn!("Disk space low: {}MB remaining", available_mb);
@@ -305,8 +290,9 @@ impl DiagTask {
}
if let Err(e) = qmdl_writer.write_container(&container).await {
error!("failed to write to QMDL (disk full?): {e}");
self.stop(qmdl_store).await;
let reason = format!("failed to write to QMDL (disk full?): {e}");
error!("{reason}");
self.stop(qmdl_store, Some(reason)).await;
return;
}
debug!(
@@ -320,8 +306,9 @@ impl DiagTask {
.update_entry_qmdl_size(index, qmdl_writer.total_written)
.await
{
error!("failed to update manifest (disk full?): {e}");
self.stop(qmdl_store).await;
let reason = format!("failed to update manifest (disk full?): {e}");
error!("{reason}");
self.stop(qmdl_store, Some(reason)).await;
return;
}
debug!("done!");
@@ -380,20 +367,23 @@ pub fn run_diag_read_thread(
let mut diag_stream = pin!(dev.as_stream().into_stream());
let mut diag_task = DiagTask::new(ui_update_sender, analysis_sender, analyzer_config, notification_channel, min_space_to_start_mb, min_space_to_continue_mb);
qmdl_file_tx
.send(DiagDeviceCtrlMessage::StartRecording)
.send(DiagDeviceCtrlMessage::StartRecording { response_tx: None })
.await
.unwrap();
loop {
tokio::select! {
msg = qmdl_file_rx.recv() => {
match msg {
Some(DiagDeviceCtrlMessage::StartRecording) => {
Some(DiagDeviceCtrlMessage::StartRecording { response_tx }) => {
let mut qmdl_store = qmdl_store_lock.write().await;
diag_task.start(qmdl_store.deref_mut()).await;
let result = diag_task.start(qmdl_store.deref_mut()).await;
if let Some(tx) = response_tx {
tx.send(result).ok();
}
},
Some(DiagDeviceCtrlMessage::StopRecording) => {
let mut qmdl_store = qmdl_store_lock.write().await;
diag_task.stop(qmdl_store.deref_mut()).await;
diag_task.stop(qmdl_store.deref_mut(), None).await;
},
// None means all the Senders have been dropped, so it's
// time to go
@@ -443,9 +433,12 @@ pub async fn start_recording(
return Err((StatusCode::FORBIDDEN, "server is in debug mode".to_string()));
}
let (response_tx, response_rx) = oneshot::channel();
state
.diag_device_ctrl_sender
.send(DiagDeviceCtrlMessage::StartRecording)
.send(DiagDeviceCtrlMessage::StartRecording {
response_tx: Some(response_tx),
})
.await
.map_err(|e| {
(
@@ -454,7 +447,14 @@ pub async fn start_recording(
)
})?;
Ok((StatusCode::ACCEPTED, "ok".to_string()))
match response_rx.await {
Ok(Ok(())) => Ok((StatusCode::ACCEPTED, "ok".to_string())),
Ok(Err(reason)) => Err((StatusCode::INSUFFICIENT_STORAGE, reason)),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("failed to receive start recording response: {e}"),
)),
}
}
/// Stop recording API for web thread
+3 -2
View File
@@ -81,8 +81,9 @@ pub fn run_key_input_thread(
{
error!("Failed to send StopRecording: {e}");
}
if let Err(e) =
diag_tx.send(DiagDeviceCtrlMessage::StartRecording).await
if let Err(e) = diag_tx
.send(DiagDeviceCtrlMessage::StartRecording { response_tx: None })
.await
{
error!("Failed to send StartRecording: {e}");
}
+15
View File
@@ -54,6 +54,8 @@ pub struct ManifestEntry {
pub rayhunter_version: Option<String>,
pub system_os: Option<String>,
pub arch: Option<String>,
#[serde(default)]
pub stop_reason: Option<String>,
}
impl ManifestEntry {
@@ -68,6 +70,7 @@ impl ManifestEntry {
rayhunter_version: Some(metadata.rayhunter_version),
system_os: Some(metadata.system_os),
arch: Some(metadata.arch),
stop_reason: None,
}
}
@@ -197,6 +200,7 @@ impl RecordingStore {
rayhunter_version: None,
system_os: None,
arch: None,
stop_reason: None,
});
}
@@ -342,6 +346,17 @@ impl RecordingStore {
Some((entry_index, &self.manifest.entries[entry_index]))
}
pub async fn set_current_stop_reason(
&mut self,
reason: String,
) -> Result<(), RecordingStoreError> {
if let Some(idx) = self.current_entry {
self.manifest.entries[idx].stop_reason = Some(reason);
self.write_manifest().await?;
}
Ok(())
}
pub fn is_current_entry(&self, name: &str) -> bool {
match self.current_entry {
Some(idx) => match self.manifest.entries.get(idx) {