diff --git a/quickwit/quickwit-actors/src/actor_context.rs b/quickwit/quickwit-actors/src/actor_context.rs index 9a55a6f91eb..6308dd9ca6f 100644 --- a/quickwit/quickwit-actors/src/actor_context.rs +++ b/quickwit/quickwit-actors/src/actor_context.rs @@ -23,7 +23,7 @@ use std::time::Duration; use quickwit_common::{KillSwitch, Progress, ProtectedZoneGuard}; use quickwit_metrics::Counter; use tokio::sync::{oneshot, watch}; -use tracing::{debug, error}; +use tracing::debug; #[cfg(any(test, feature = "testsuite"))] use crate::Universe; @@ -209,12 +209,17 @@ impl ActorContext { obs_state } - pub(crate) fn exit(&self, exit_status: &ActorExitStatus) { - self.actor_state.exit(exit_status.is_success()); + pub(crate) fn exit(&self, exit_status: &ActorExitStatus, fault_opt: Option) { + // The fault has to be recorded before the failed state becomes observable: a supervisor + // that sees the failure first would terminate the pipeline and kill this very switch, + // and the fault would then be dropped as a mere consequence of that kill. if should_activate_kill_switch(exit_status) { - error!(actor=%self.actor_instance_id(), exit_status=?exit_status, "exit activating-kill-switch"); - self.kill_switch().kill(); + match fault_opt { + Some(fault) => self.kill_switch().kill_with_fault(fault), + None => self.kill_switch().kill(), + } } + self.actor_state.exit(exit_status.is_success()); } /// Posts a message in an actor's mailbox. diff --git a/quickwit/quickwit-actors/src/actor_handle.rs b/quickwit/quickwit-actors/src/actor_handle.rs index 7de27a736e6..33b61ab1722 100644 --- a/quickwit/quickwit-actors/src/actor_handle.rs +++ b/quickwit/quickwit-actors/src/actor_handle.rs @@ -89,7 +89,6 @@ impl Supervisable for ActorHandle { return Health::Success; } if actor_state == ActorState::Failure { - error!(actor = self.name(), "actor-exit-without-success"); return Health::FailureOrUnhealthy; } if !check_for_progress @@ -100,7 +99,12 @@ impl Supervisable for ActorHandle { { Health::Healthy } else { - error!(actor = self.name(), "actor-timeout"); + self.actor_context + .kill_switch() + .kill_with_fault(anyhow::anyhow!( + "{} stopped reporting progress", + self.name() + )); Health::FailureOrUnhealthy } } diff --git a/quickwit/quickwit-actors/src/spawn_builder.rs b/quickwit/quickwit-actors/src/spawn_builder.rs index 922cfc4d71d..e8308693497 100644 --- a/quickwit/quickwit-actors/src/spawn_builder.rs +++ b/quickwit/quickwit-actors/src/spawn_builder.rs @@ -19,7 +19,7 @@ use anyhow::Context; use quickwit_metrics::Counter; use sync_wrapper::SyncWrapper; use tokio::sync::watch; -use tracing::{debug, error, info}; +use tracing::{debug, error}; use crate::envelope::Envelope; use crate::mailbox::{Inbox, create_mailbox}; @@ -400,23 +400,24 @@ async fn actor_loop( actor_env.process_messages().await }; - let actor_id = actor_env.ctx.actor_instance_id(); - match after_process_exit_status { - ActorExitStatus::Success - | ActorExitStatus::Quit - | ActorExitStatus::DownstreamClosed - | ActorExitStatus::Killed => { - info!(actor_id, phase = ?exit_phase, exit_status = ?after_process_exit_status, "actor-exit"); - } - ActorExitStatus::Failure(_) | ActorExitStatus::Panicked => { - error!(actor_id, phase = ?exit_phase, exit_status = ?after_process_exit_status, "actor-exit"); - } - }; + let actor_name = actor_env.actor.get_mut().name(); // TODO the no advance time guard for finalize has a race condition. Ideally we would // like to have the guard before we drop the last envelope. let final_exit_status = actor_env.finalize(after_process_exit_status).await; + let fault_opt: Option = match &final_exit_status { + ActorExitStatus::Failure(cause) => Some(anyhow::anyhow!( + "{actor_name} failed while {exit_phase:?}: {cause:#}" + )), + ActorExitStatus::Panicked => Some(anyhow::anyhow!( + "{actor_name} panicked while {exit_phase:?}" + )), + ActorExitStatus::Success + | ActorExitStatus::Quit + | ActorExitStatus::DownstreamClosed + | ActorExitStatus::Killed => None, + }; // The last observation is collected on `ActorExecutionEnv::Drop`. - actor_env.ctx.exit(&final_exit_status); + actor_env.ctx.exit(&final_exit_status, fault_opt); final_exit_status } diff --git a/quickwit/quickwit-actors/src/supervisor.rs b/quickwit/quickwit-actors/src/supervisor.rs index ae26d7933fa..9d3ae906234 100644 --- a/quickwit/quickwit-actors/src/supervisor.rs +++ b/quickwit/quickwit-actors/src/supervisor.rs @@ -14,7 +14,6 @@ use async_trait::async_trait; use serde::Serialize; -use tracing::{info, warn}; use crate::mailbox::Inbox; use crate::{Actor, ActorContext, ActorExitStatus, ActorHandle, Handler, Health, Supervisable}; @@ -140,14 +139,12 @@ impl Supervisor { return Err(ActorExitStatus::Success); } } - warn!("unhealthy-actor"); - // The actor is failing we need to restart it. + // The actor is failing, we need to restart it. let actor_handle = self.handle_opt.take().unwrap(); let actor_mailbox = actor_handle.mailbox().clone(); let (actor_exit_status, _last_state) = if !actor_handle.state().is_exit() { // The actor is probably frozen. // Let's kill it. - warn!("killing"); actor_handle.kill().await } else { actor_handle.join().await @@ -172,7 +169,6 @@ impl Supervisor { self.metrics.num_panics += 1; } } - info!("respawning-actor"); let (_, actor_handle) = ctx .spawn_actor() .set_mailboxes(actor_mailbox, self.inbox.clone()) @@ -203,7 +199,6 @@ mod tests { use std::time::Duration; use async_trait::async_trait; - use tracing::info; use crate::supervisor::SupervisorMetrics; use crate::tests::{Ping, PingReceiverActor}; @@ -239,7 +234,6 @@ mod tests { _exit_status: &ActorExitStatus, _ctx: &ActorContext, ) -> anyhow::Result<()> { - info!("finalize-failing-actor"); Ok(()) } } diff --git a/quickwit/quickwit-actors/src/tests.rs b/quickwit/quickwit-actors/src/tests.rs index 0842e03d6e5..f1d039baf16 100644 --- a/quickwit/quickwit-actors/src/tests.rs +++ b/quickwit/quickwit-actors/src/tests.rs @@ -24,7 +24,7 @@ use serde::Serialize; use crate::observation::ObservationType; use crate::{ Actor, ActorContext, ActorExitStatus, ActorHandle, ActorState, Command, Handler, Health, - Mailbox, Observation, Supervisable, Universe, + KillSwitch, Mailbox, Observation, Supervisable, Universe, }; // An actor that receives ping messages. @@ -220,6 +220,9 @@ struct DoNothing; #[derive(Clone, Debug)] struct Block; +#[derive(Clone, Debug)] +struct Fail; + impl Actor for BuggyActor { type ObservableState = (); @@ -259,6 +262,49 @@ impl Handler for BuggyActor { } } +#[async_trait] +impl Handler for BuggyActor { + type Reply = (); + + async fn handle( + &mut self, + _message: Fail, + _ctx: &ActorContext, + ) -> Result<(), ActorExitStatus> { + Err(ActorExitStatus::from(anyhow::anyhow!("handler blew up"))) + } +} + +#[tokio::test] +async fn test_failing_actor_records_root_cause_fault() { + let universe = Universe::with_accelerated_time(); + let kill_switch = KillSwitch::default(); + let (failing_mailbox, failing_handle) = universe + .spawn_builder() + .set_kill_switch(kill_switch.clone()) + .spawn(BuggyActor); + let (sibling_mailbox, sibling_handle) = universe + .spawn_builder() + .set_kill_switch(kill_switch.clone()) + .spawn(BuggyActor); + sibling_mailbox.send_message(Block).await.unwrap(); + failing_mailbox.send_message(Fail).await.unwrap(); + + let (exit_status, _) = failing_handle.join().await; + assert!(matches!(exit_status, ActorExitStatus::Failure(_))); + + // The sibling is taken down by the shared kill switch. Being a casualty rather than the + // root cause, it must not overwrite the fault. + let (sibling_exit_status, _) = sibling_handle.join().await; + assert!(matches!(sibling_exit_status, ActorExitStatus::Killed)); + + let fault = kill_switch.fault().expect("fault should be recorded"); + let cause = format!("{fault:#}"); + assert!(cause.contains("BuggyActor"), "{cause}"); + assert!(cause.contains("handling"), "{cause}"); + assert!(cause.contains("handler blew up"), "{cause}"); +} + #[tokio::test] async fn test_timeouting_actor() { let universe = Universe::with_accelerated_time(); diff --git a/quickwit/quickwit-common/src/kill_switch.rs b/quickwit/quickwit-common/src/kill_switch.rs index 15c2154c8d9..5231c989603 100644 --- a/quickwit/quickwit-common/src/kill_switch.rs +++ b/quickwit/quickwit-common/src/kill_switch.rs @@ -13,9 +13,9 @@ // limitations under the License. use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::{Arc, Mutex, Weak}; +use std::sync::{Arc, Mutex, OnceLock, Weak}; -use tracing::debug; +use tracing::{debug, error}; #[derive(Clone, Default)] pub struct KillSwitch { @@ -24,6 +24,7 @@ pub struct KillSwitch { struct Inner { alive: AtomicBool, + fault: OnceLock>, children: Mutex>>, } @@ -31,6 +32,7 @@ impl Default for Inner { fn default() -> Self { Self { alive: AtomicBool::new(true), + fault: OnceLock::new(), children: Mutex::default(), } } @@ -60,6 +62,20 @@ impl KillSwitch { self.inner.kill(); } + pub fn kill_with_fault(&self, fault: anyhow::Error) { + if self.is_alive() { + let fault = Arc::new(fault); + if self.inner.fault.set(fault.clone()).is_ok() { + error!(cause = %format!("{fault:#}"), "actor-fault"); + } + } + self.kill(); + } + + pub fn fault(&self) -> Option> { + self.inner.fault.get().cloned() + } + // Creates a child killswitch. // // If the parent kill switch is dead to begin with, the child will be dead too. @@ -133,6 +149,34 @@ mod tests { assert!(grandchild_kill_switch.is_dead()); } + #[test] + fn test_kill_switch_fault() { + let kill_switch = KillSwitch::default(); + assert!(kill_switch.fault().is_none()); + + kill_switch.kill_with_fault(anyhow::anyhow!("indexer blew up")); + assert!(kill_switch.is_dead()); + let fault = kill_switch.fault().expect("fault should be recorded"); + assert_eq!(fault.to_string(), "indexer blew up"); + + // The first fault is the root cause: actors dying because of it must not overwrite it. + kill_switch.kill_with_fault(anyhow::anyhow!("publisher noticed and gave up")); + let fault = kill_switch.fault().expect("fault should be recorded"); + assert_eq!(fault.to_string(), "indexer blew up"); + } + + #[test] + fn test_kill_switch_without_fault_records_nothing() { + let kill_switch = KillSwitch::default(); + kill_switch.kill(); + assert!(kill_switch.is_dead()); + assert!(kill_switch.fault().is_none()); + + // An error hit while already dying is a consequence of the kill, not a root cause. + kill_switch.kill_with_fault(anyhow::anyhow!("directory kill switch was activated")); + assert!(kill_switch.fault().is_none()); + } + #[test] fn test_kill_switch_to_quoque_me_fili() { let kill_switch = KillSwitch::default(); diff --git a/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs b/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs index f8183f995e2..ceee1ed0f3f 100644 --- a/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs @@ -187,13 +187,13 @@ impl IndexingPipeline { } if !failure_or_unhealthy_actors.is_empty() { - error!( + debug!( pipeline_id=?self.params.pipeline_id, generation=self.generation(), healthy_actors=?healthy_actors, failed_or_unhealthy_actors=?failure_or_unhealthy_actors, success_actors=?success_actors, - "Indexing pipeline failure." + "indexing pipeline failure" ); return Health::FailureOrUnhealthy; } diff --git a/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs b/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs index 093687285d8..6bd8d77ffbf 100644 --- a/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs @@ -205,7 +205,7 @@ impl MergePipeline { } } if !failure_or_unhealthy_actors.is_empty() { - error!( + debug!( index_uid=%self.params.pipeline_id.index_uid, source_id=%self.params.pipeline_id.source_id, generation=self.generation(), diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs index a76c3ee2e03..af39b9c5746 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs @@ -231,7 +231,7 @@ impl ParquetMergePipeline { } } if !failure_or_unhealthy_actors.is_empty() { - error!( + debug!( generation = self.generation(), healthy_actors = ?healthy_actors, failed_or_unhealthy_actors = ?failure_or_unhealthy_actors, diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs index 620d4720f4c..61ead13d45d 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs @@ -269,11 +269,11 @@ impl Handler for ParquetUploader { let stage_result = stage_splits(metastore.clone(), index_uid.clone(), &splits).await; - if let Err(e) = stage_result { - warn!(error = %e, "failed to stage splits"); + if let Err(error) = stage_result { // Discard sequencer position on error let _ = tx.send(SequencerCommand::Discard); - kill_switch.kill(); + kill_switch + .kill_with_fault(error.context("ParquetUploader failed to stage splits")); return; } @@ -293,17 +293,17 @@ impl Handler for ParquetUploader { let local_path = output_dir.join(&parquet_file); let file_content = match tokio::fs::read(&local_path).await { Ok(content) => content, - Err(e) => { - warn!( - error = %e, - local_path = %local_path.display(), - split_id = %split.split_id_str(), - parquet_file = %parquet_file, - "failed to read local parquet file" - ); + Err(error) => { // Discard sequencer position on error let _ = tx.send(SequencerCommand::Discard); - kill_switch.kill(); + kill_switch.kill_with_fault(anyhow::Error::from(error).context( + format!( + "ParquetUploader failed to read local parquet file {} for \ + split {}", + local_path.display(), + split.split_id_str() + ), + )); return; } }; @@ -312,16 +312,14 @@ impl Handler for ParquetUploader { let payload: Box = Box::new(file_content); // Upload to S3 using the filename directly (matches logs pipeline) - if let Err(e) = split_store.put(Path::new(&parquet_file), payload).await { - warn!( - error = %e, - split_id = %split.split_id_str(), - parquet_file = %parquet_file, - "failed to upload parquet file" - ); + if let Err(error) = split_store.put(Path::new(&parquet_file), payload).await { // Discard sequencer position on error let _ = tx.send(SequencerCommand::Discard); - kill_switch.kill(); + kill_switch.kill_with_fault(anyhow::Error::from(error).context(format!( + "ParquetUploader failed to upload parquet file {} for split {}", + parquet_file, + split.split_id_str() + ))); return; } diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs index d0ed12c2c7a..4056fbeaaa5 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs @@ -209,7 +209,7 @@ impl MetricsPipeline { } if !failure_or_unhealthy_actors.is_empty() { - error!( + debug!( pipeline_id=?self.params.pipeline_id, generation=self.generation(), healthy_actors=?healthy_actors, diff --git a/quickwit/quickwit-indexing/src/actors/uploader.rs b/quickwit/quickwit-indexing/src/actors/uploader.rs index ac1dfeeb5af..f3b9f85e126 100644 --- a/quickwit/quickwit-indexing/src/actors/uploader.rs +++ b/quickwit/quickwit-indexing/src/actors/uploader.rs @@ -376,8 +376,10 @@ impl Handler for Uploader { .await; if let Err(cause) = upload_result { - warn!(cause=?cause, split_id=packaged_split.split_id_str(), "Failed to upload split. Killing!"); - kill_switch.kill(); + kill_switch.kill_with_fault(cause.context(format!( + "Uploader failed to upload split {}", + packaged_split.split_id_str() + ))); return; }