diff --git a/daemon/src/diag.rs b/daemon/src/diag.rs index 7ebca83..3c43395 100644 --- a/daemon/src/diag.rs +++ b/daemon/src/diag.rs @@ -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>>, + }, DeleteEntry { name: String, response_tx: oneshot::Sender>, @@ -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) { 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 diff --git a/daemon/src/key_input.rs b/daemon/src/key_input.rs index 4865fb4..7e75fbb 100644 --- a/daemon/src/key_input.rs +++ b/daemon/src/key_input.rs @@ -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}"); } diff --git a/daemon/src/qmdl_store.rs b/daemon/src/qmdl_store.rs index fd3785c..81fb539 100644 --- a/daemon/src/qmdl_store.rs +++ b/daemon/src/qmdl_store.rs @@ -54,6 +54,8 @@ pub struct ManifestEntry { pub rayhunter_version: Option, pub system_os: Option, pub arch: Option, + #[serde(default)] + pub stop_reason: Option, } 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) { diff --git a/daemon/web/src/lib/components/DiskSpaceAlert.svelte b/daemon/web/src/lib/components/DiskSpaceAlert.svelte deleted file mode 100644 index 4a79449..0000000 --- a/daemon/web/src/lib/components/DiskSpaceAlert.svelte +++ /dev/null @@ -1,64 +0,0 @@ - - -{#if warning_message} -
- - - Low Disk Space - -

{warning_message}

-
-{/if} diff --git a/daemon/web/src/lib/components/ManifestCard.svelte b/daemon/web/src/lib/components/ManifestCard.svelte index 71b2764..05b5b34 100644 --- a/daemon/web/src/lib/components/ManifestCard.svelte +++ b/daemon/web/src/lib/components/ManifestCard.svelte @@ -81,6 +81,11 @@ 'N/A'} + {#if entry.stop_reason} +
+ {entry.stop_reason} +
+ {/if}
diff --git a/daemon/web/src/lib/manifest.svelte.ts b/daemon/web/src/lib/manifest.svelte.ts index 2feedb5..ec5d1cb 100644 --- a/daemon/web/src/lib/manifest.svelte.ts +++ b/daemon/web/src/lib/manifest.svelte.ts @@ -11,6 +11,7 @@ interface JsonManifestEntry { start_time: string; last_message_time: string; qmdl_size_bytes: number; + stop_reason: string | null; } export class Manifest { @@ -57,6 +58,7 @@ export class ManifestEntry { public analysis_size_bytes = $state(0); public analysis_status: AnalysisStatus | undefined = $state(undefined); public analysis_report: AnalysisReport | string | undefined = $state(undefined); + public stop_reason: string | undefined = $state(undefined); constructor(json: JsonManifestEntry) { this.name = json.name; @@ -65,6 +67,9 @@ export class ManifestEntry { if (json.last_message_time) { this.last_message_time = new Date(json.last_message_time); } + if (json.stop_reason) { + this.stop_reason = json.stop_reason; + } } get_readable_qmdl_size(): string { diff --git a/daemon/web/src/routes/+page.svelte b/daemon/web/src/routes/+page.svelte index 2d529d5..5abb235 100644 --- a/daemon/web/src/routes/+page.svelte +++ b/daemon/web/src/routes/+page.svelte @@ -11,7 +11,6 @@ import ConfigForm from '$lib/components/ConfigForm.svelte'; import ActionErrors from '$lib/components/ActionErrors.svelte'; import ClockDriftAlert from '$lib/components/ClockDriftAlert.svelte'; - import DiskSpaceAlert from '$lib/components/DiskSpaceAlert.svelte'; import LogView from '$lib/components/LogView.svelte'; let manager: AnalysisManager = new AnalysisManager(); @@ -211,7 +210,6 @@ {#if loaded} -
{#if current_entry}