diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3b7ad05696d..fcf6af9382e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -111,7 +111,7 @@ jobs: working-directory: ./quickwit - name: cargo nextest if: always() && steps.modified.outputs.rust_src == 'true' - run: cargo nextest run --features=postgres --retries 1 + run: cargo nextest run --features=postgres --profile ci working-directory: ./quickwit - name: Install python packages if: always() && steps.modified.outputs.rust_src == 'true' diff --git a/.github/workflows/coverage.yml b/.github/workflows/coverage.yml index 3b1be41fa16..dd45a377897 100644 --- a/.github/workflows/coverage.yml +++ b/.github/workflows/coverage.yml @@ -190,9 +190,9 @@ jobs: - name: Generate code coverage run: | cargo llvm-cov clean --workspace - cargo llvm-cov nextest --no-report --test failpoints --features fail/failpoints --retries 4 + cargo llvm-cov nextest --no-report --profile ci --test failpoints --features fail/failpoints # increase stack size for test_all_with_s3_localstack_cli, see quickwit#4963 - RUST_MIN_STACK=67108864 CARGO_BUILD_JOBS=4 cargo llvm-cov nextest --no-report --all-features --retries 4 + RUST_MIN_STACK=67108864 CARGO_BUILD_JOBS=4 cargo llvm-cov nextest --no-report --all-features --profile ci cargo llvm-cov report --lcov --output-path lcov.info working-directory: ./quickwit diff --git a/AGENTS.md b/AGENTS.md index f9d8206fd6b..a92db39e62b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -140,7 +140,7 @@ Run `cargo check` after editing `Cargo.toml` to update `Cargo.lock`. ### Testing - Single crate test: `cargo nextest run -p quickwit-search my_test_name` -- `make test-all` — starts Docker services (LocalStack S3, PostgreSQL, Pub/Sub emulator) and runs the full test suite with `cargo nextest run --all-features --retries 5`. +- `make test-all` — starts Docker services (LocalStack S3, PostgreSQL, Pub/Sub emulator) and runs the full test suite with `cargo nextest run --all-features --profile ci` (the `ci` profile in `quickwit/.config/nextest.toml` retries failing tests twice). - `make test-failpoints` — runs failpoint tests only: `cargo nextest run --test failpoints --features fail/failpoints`. - Docker services: `make docker-compose-up` / `make docker-compose-down` (subset: `DOCKER_SERVICES=kafka,postgres`). - Integration tests are under `rest-api-tests`; compile and run the `quickwit-cli` binary to have an instance available to test against. diff --git a/quickwit/.config/nextest.toml b/quickwit/.config/nextest.toml index 3c8dd80439a..8e88cf9e541 100644 --- a/quickwit/.config/nextest.toml +++ b/quickwit/.config/nextest.toml @@ -6,4 +6,6 @@ slow-timeout = "10s" # of the run (for easy scrollability). failure-output = "immediate-final" # Do not cancel the test run on the first failure. -fail-fast = false \ No newline at end of file +fail-fast = false +slow-timeout = { period = "60s", terminate-after = 10 } +retries = { backoff = "exponential", count = 2, delay = "1s", max-delay = "4s", jitter = true } diff --git a/quickwit/Makefile b/quickwit/Makefile index 5bd42980a24..8b9222736ff 100644 --- a/quickwit/Makefile +++ b/quickwit/Makefile @@ -55,8 +55,8 @@ test-all: QW_S3_FORCE_PATH_STYLE_ACCESS=1 \ QW_TEST_DATABASE_URL=postgres://quickwit-dev:quickwit-dev@localhost:5432/quickwit-metastore-dev \ RUST_MIN_STACK=67108864 \ - cargo nextest run --all-features --retries 5 - cargo nextest run --test failpoints --features fail/failpoints + cargo nextest run --all-features --profile ci + cargo nextest run --profile ci --test failpoints --features fail/failpoints test-failpoints: cargo nextest run --test failpoints --features fail/failpoints diff --git a/quickwit/quickwit-common/src/test_utils.rs b/quickwit/quickwit-common/src/test_utils.rs index 85d50b87a3e..a1854a425eb 100644 --- a/quickwit/quickwit-common/src/test_utils.rs +++ b/quickwit/quickwit-common/src/test_utils.rs @@ -39,36 +39,26 @@ where .await } -/// Tries to connect at most 3 times to `SocketAddr`. +/// Tries to connect to `SocketAddr` for up to 30 seconds. /// If not successful, returns an error. /// This is a convenient function to wait before sending gRPC requests /// to this `SocketAddr`. pub async fn wait_for_server_ready(socket_addr: SocketAddr) -> anyhow::Result<()> { - let mut num_attempts = 0; - let max_num_attempts = 10; let uri = Uri::builder() .scheme("http") .authority(socket_addr.to_string().as_str()) .path_and_query("/") .build()?; - - while num_attempts < max_num_attempts { - tokio::time::sleep(Duration::from_millis(50 * (num_attempts + 1))).await; - let mut http = hyper_util::client::legacy::connect::HttpConnector::new(); - match http.call(uri.clone()).await { - Ok(_) => break, - Err(_) => { - println!( - "Failed to connect to `{}` failed, retrying {}/{}", - socket_addr, - num_attempts + 1, - max_num_attempts - ); - num_attempts += 1; - } - } - } - if num_attempts == max_num_attempts { + let wait_result = wait_until_predicate( + || async { + let mut http = hyper_util::client::legacy::connect::HttpConnector::new(); + http.call(uri.clone()).await.is_ok() + }, + Duration::from_secs(30), + Duration::from_millis(100), + ) + .await; + if wait_result.is_err() { anyhow::bail!("too many attempts to connect to `{}`", socket_addr); } Ok(()) diff --git a/quickwit/quickwit-indexing/src/actors/indexer.rs b/quickwit/quickwit-indexing/src/actors/indexer.rs index 742c94c25c1..6cf97ab8fbf 100644 --- a/quickwit/quickwit-indexing/src/actors/indexer.rs +++ b/quickwit/quickwit-indexing/src/actors/indexer.rs @@ -1475,7 +1475,7 @@ mod tests { #[tokio::test] async fn test_indexer_exceeding_max_num_partitions() { - let universe = Universe::with_accelerated_time(); + let universe = Universe::new(); let pipeline_id = IndexingPipelineId { index_uid: IndexUid::new_with_random_ulid("test-index"), source_id: "test-source".to_string(), diff --git a/quickwit/quickwit-indexing/src/source/kafka_source.rs b/quickwit/quickwit-indexing/src/source/kafka_source.rs index a6769fc169d..1fe54abf4da 100644 --- a/quickwit/quickwit-indexing/src/source/kafka_source.rs +++ b/quickwit/quickwit-indexing/src/source/kafka_source.rs @@ -759,9 +759,11 @@ fn message_payload_to_doc(message: &BorrowedMessage) -> Option { #[cfg(all(test, feature = "kafka-broker-tests"))] mod kafka_broker_tests { use std::num::NonZeroUsize; + use std::sync::Arc; use quickwit_actors::{ActorContext, Universe}; use quickwit_common::rand::append_random_suffix; + use quickwit_common::test_utils::wait_until_predicate; use quickwit_config::{SourceConfig, SourceInputFormat, SourceParams}; use quickwit_metastore::checkpoint::SourceCheckpointDelta; use quickwit_metastore::metastore_for_test; @@ -817,7 +819,45 @@ mod kafka_broker_tests { err_code ) })?; - Ok(()) + wait_for_topic_ready(topic, num_partitions).await + } + + async fn wait_for_topic_ready(topic: &str, num_partitions: i32) -> anyhow::Result<()> { + let consumer = Arc::new(create_base_consumer("quickwit-test-wait-for-topic-ready")); + let is_topic_ready = || { + let consumer = consumer.clone(); + let topic = topic.to_string(); + async move { + spawn_blocking(move || { + let Ok(cluster_metadata) = + consumer.fetch_metadata(Some(&topic), Duration::from_secs(1)) + else { + return false; + }; + cluster_metadata.topics().iter().any(|topic_metadata| { + topic_metadata.name() == topic + && topic_metadata.error().is_none() + && topic_metadata.partitions().len() == num_partitions as usize + && topic_metadata + .partitions() + .iter() + .all(|partition_metadata| { + partition_metadata.error().is_none() + && partition_metadata.leader() >= 0 + }) + }) + }) + .await + .unwrap_or(false) + } + }; + wait_until_predicate( + is_topic_ready, + Duration::from_secs(30), + Duration::from_millis(100), + ) + .await + .with_context(|| format!("topic `{topic}` is not ready after 30 seconds")) } async fn populate_topic( @@ -1280,7 +1320,10 @@ mod kafka_broker_tests { let source_actor = SourceActor::new(source, doc_processor_mailbox); let (_source_mailbox, source_handle) = universe.spawn_builder().spawn(source_actor); let (exit_status, exit_state) = source_handle.join().await; - assert!(exit_status.is_success()); + assert!( + exit_status.is_success(), + "unexpected exit status: {exit_status:?}" + ); let messages: Vec = doc_processor_inbox.drain_for_test_typed(); assert!(messages.is_empty()); @@ -1330,7 +1373,10 @@ mod kafka_broker_tests { let source_actor = SourceActor::new(source, doc_processor_mailbox); let (_source_mailbox, source_handle) = universe.spawn_builder().spawn(source_actor); let (exit_status, exit_state) = source_handle.join().await; - assert!(exit_status.is_success()); + assert!( + exit_status.is_success(), + "unexpected exit status: {exit_status:?}" + ); let messages: Vec = doc_processor_inbox.drain_for_test_typed(); assert!(!messages.is_empty()); @@ -1402,7 +1448,10 @@ mod kafka_broker_tests { let source_actor = SourceActor::new(source, doc_processor_mailbox); let (_source_mailbox, source_handle) = universe.spawn_builder().spawn(source_actor); let (exit_status, exit_state) = source_handle.join().await; - assert!(exit_status.is_success()); + assert!( + exit_status.is_success(), + "unexpected exit status: {exit_status:?}" + ); let messages: Vec = doc_processor_inbox.drain_for_test_typed(); assert!(!messages.is_empty()); @@ -1453,7 +1502,10 @@ mod kafka_broker_tests { let source_actor = SourceActor::new(source, doc_processor_mailbox); let (_source_mailbox, source_handle) = universe.spawn_builder().spawn(source_actor); let (exit_status, exit_state) = source_handle.join().await; - assert!(exit_status.is_success()); + assert!( + exit_status.is_success(), + "unexpected exit status: {exit_status:?}" + ); let messages: Vec = doc_processor_inbox.drain_for_test_typed(); assert!(messages.is_empty()); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs index bb02bef59a8..95633826939 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs @@ -3065,6 +3065,7 @@ mod tests { let state_guard = ingester.state.lock_partially("test").await.unwrap(); let shard = state_guard.shards.get(&queue_id).unwrap(); shard.assert_is_closed(); + drop(state_guard); let fetch_response = timeout(Duration::from_millis(100), fetch_stream.next()) .await diff --git a/quickwit/quickwit-integration-tests/src/test_utils/cluster_sandbox.rs b/quickwit/quickwit-integration-tests/src/test_utils/cluster_sandbox.rs index 7b4980f428d..ced09060bed 100644 --- a/quickwit/quickwit-integration-tests/src/test_utils/cluster_sandbox.rs +++ b/quickwit/quickwit-integration-tests/src/test_utils/cluster_sandbox.rs @@ -500,7 +500,7 @@ impl ClusterSandbox { } } }, - Duration::from_secs(10), + Duration::from_secs(30), Duration::from_millis(100), ) .await?; diff --git a/quickwit/quickwit-integration-tests/src/tests/ingest_v1_tests.rs b/quickwit/quickwit-integration-tests/src/tests/ingest_v1_tests.rs index af28f5b21a1..0b40b423601 100644 --- a/quickwit/quickwit-integration-tests/src/tests/ingest_v1_tests.rs +++ b/quickwit/quickwit-integration-tests/src/tests/ingest_v1_tests.rs @@ -12,10 +12,15 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::time::Duration; + +use quickwit_common::test_utils::wait_until_predicate; use quickwit_config::ConfigFormat; use quickwit_config::service::QuickwitService; use quickwit_metastore::SplitState; +use quickwit_rest_client::error::Error as RestClientError; use quickwit_rest_client::rest_client::CommitType; +use reqwest::StatusCode; use serde_json::json; use crate::ingest_json; @@ -58,14 +63,34 @@ async fn test_ingest_v1_happy_path() { .await .unwrap(); - ingest( - &indexer_client, - index_id, - ingest_json!({"body": "my-doc"}), - CommitType::Auto, + let indexer_client_ref = &indexer_client; + wait_until_predicate( + || async move { + let ingest_result = ingest( + indexer_client_ref, + index_id, + ingest_json!({"body": "my-doc"}), + CommitType::Auto, + ) + .await; + let Err(error) = ingest_result else { + return true; + }; + let status_code_opt = error + .downcast_ref::() + .and_then(RestClientError::status_code); + assert_eq!( + status_code_opt, + Some(StatusCode::NOT_FOUND), + "unexpected ingest error: {error:?}" + ); + false + }, + Duration::from_secs(10), + Duration::from_millis(100), ) .await - .unwrap(); + .expect("the indexer should create the ingest v1 queue of the index"); sandbox .wait_for_splits(index_id, Some(vec![SplitState::Published]), 1) diff --git a/quickwit/quickwit-integration-tests/src/tests/ingest_v2_tests.rs b/quickwit/quickwit-integration-tests/src/tests/ingest_v2_tests.rs index 4036e47bc36..64557eb5733 100644 --- a/quickwit/quickwit-integration-tests/src/tests/ingest_v2_tests.rs +++ b/quickwit/quickwit-integration-tests/src/tests/ingest_v2_tests.rs @@ -29,7 +29,36 @@ use quickwit_serve::{ListSplitsQueryParams, RestIngestResponse, RestParseFailure use serde_json::json; use crate::ingest_json; -use crate::test_utils::{ClusterSandboxBuilder, ingest}; +use crate::test_utils::{ClusterSandbox, ClusterSandboxBuilder, ingest}; + +async fn wait_for_each_query_to_match( + sandbox: &ClusterSandbox, + index_id: &str, + queries: &[&str], + timeout: Duration, +) -> anyhow::Result<()> { + let searcher_client = &sandbox.rest_client(QuickwitService::Searcher); + wait_until_predicate( + || async move { + for query in queries { + let search_request = quickwit_serve::SearchRequestQueryString { + query: query.to_string(), + max_hits: 1, + ..Default::default() + }; + match searcher_client.search(index_id, search_request).await { + Ok(search_response) if search_response.num_hits > 0 => {} + _ => return false, + } + } + true + }, + timeout, + Duration::from_millis(200), + ) + .await + .map_err(|_| anyhow::anyhow!("some queries still match no document after {timeout:?}")) +} /// Ingesting on a freshly re-created index sometimes fails, see #5430 #[tokio::test] @@ -1017,26 +1046,11 @@ async fn test_graceful_shutdown_no_data_loss() { // All 3 documents should eventually be searchable. Documents 1 & 2 were // in-flight on the decommissioning indexer and should have been committed during // the decommission step. Document 3 was ingested to the surviving indexer. - wait_until_predicate( - || async { - match sandbox - .rest_client(QuickwitService::Searcher) - .search( - index_id, - quickwit_serve::SearchRequestQueryString { - query: "*".to_string(), - max_hits: 10, - ..Default::default() - }, - ) - .await - { - Ok(resp) => resp.num_hits == 3, - Err(_) => false, - } - }, + wait_for_each_query_to_match( + &sandbox, + index_id, + &["body:1", "body:2", "body:during"], Duration::from_secs(30), - Duration::from_millis(500), ) .await .expect("expected 3 documents after decommission shutdown, some data may have been lost"); @@ -1143,7 +1157,7 @@ async fn test_retiring_indexer_receives_empty_plan() { Err(_) => false, } }, - Duration::from_secs(1), + Duration::from_secs(10), Duration::from_millis(25), ) .await @@ -1155,18 +1169,18 @@ async fn test_retiring_indexer_receives_empty_plan() { Err(_) => false, } }, - Duration::from_secs(1), + Duration::from_secs(10), Duration::from_millis(25), ) .await .expect("the surviving indexer should also be indexing a shard"); - // Trigger the retiring node's decommission in the background. We only assert that it is told - // to shed its plan; we deliberately do not await full decommission. + // Trigger the retiring node's decommission in the background. We assert that it is told to + // shed its plan while it is still decommissioning, then await the end of the decommission. let shutdown_handle = sandbox .remove_node(&retiring_node_id) .expect("the retiring node should be in the sandbox"); - tokio::spawn(shutdown_handle.shutdown()); + let retiring_shutdown_join_handle = tokio::spawn(shutdown_handle.shutdown()); // The fix: the new plan excludes the retiring node but is still sent to it, so it shuts down // its indexing pipelines (`num_running_pipelines` excludes merge pipelines) while still alive. @@ -1178,11 +1192,21 @@ async fn test_retiring_indexer_receives_empty_plan() { Err(_) => false, } }, - Duration::from_secs(1), + Duration::from_secs(10), Duration::from_millis(25), ) .await .expect("retiring indexer should be sent an empty plan and shut down its indexing pipelines"); + + tokio::time::timeout(Duration::from_secs(30), retiring_shutdown_join_handle) + .await + .expect("decommission of the retiring indexer timed out") + .expect("retiring indexer shutdown task panicked") + .expect("retiring indexer shutdown returned an error"); + tokio::time::timeout(Duration::from_secs(30), sandbox.shutdown()) + .await + .expect("sandbox shutdown timed out") + .unwrap(); } #[tokio::test] @@ -1264,7 +1288,7 @@ async fn test_retiring_indexer_decommissions_gracefully() { Err(_) => false, } }, - Duration::from_secs(1), + Duration::from_secs(10), Duration::from_millis(25), ) .await @@ -1276,7 +1300,7 @@ async fn test_retiring_indexer_decommissions_gracefully() { Err(_) => false, } }, - Duration::from_secs(1), + Duration::from_secs(10), Duration::from_millis(25), ) .await @@ -1291,26 +1315,11 @@ async fn test_retiring_indexer_decommissions_gracefully() { .expect("retiring indexer shutdown returned an error"); // No data lost: all 6 docs remain searchable after the decommission. - wait_until_predicate( - || async { - match sandbox - .rest_client(QuickwitService::Searcher) - .search( - index_id, - quickwit_serve::SearchRequestQueryString { - query: "*".to_string(), - max_hits: 10, - ..Default::default() - }, - ) - .await - { - Ok(resp) => resp.num_hits == 6, - Err(_) => false, - } - }, - Duration::from_secs(3), - Duration::from_millis(200), + wait_for_each_query_to_match( + &sandbox, + index_id, + &["body:0", "body:1", "body:2", "body:3", "body:4", "body:5"], + Duration::from_secs(10), ) .await .expect("all 6 documents should be searchable after decommission"); diff --git a/quickwit/quickwit-integration-tests/src/tests/sqs_tests.rs b/quickwit/quickwit-integration-tests/src/tests/sqs_tests.rs index e9b01834791..68c06f65c81 100644 --- a/quickwit/quickwit-integration-tests/src/tests/sqs_tests.rs +++ b/quickwit/quickwit-integration-tests/src/tests/sqs_tests.rs @@ -23,10 +23,32 @@ use quickwit_common::uri::Uri; use quickwit_config::ConfigFormat; use quickwit_config::service::QuickwitService; use quickwit_indexing::source::sqs_queue::test_helpers as sqs_test_helpers; -use quickwit_metastore::SplitState; use tempfile::NamedTempFile; -use crate::test_utils::ClusterSandboxBuilder; +use crate::test_utils::{ClusterSandbox, ClusterSandboxBuilder}; + +async fn wait_for_num_hits_at_least(sandbox: &ClusterSandbox, index_id: &str, num_hits: u64) { + wait_until_predicate( + || async { + sandbox + .rest_client(QuickwitService::Searcher) + .search( + index_id, + quickwit_serve::SearchRequestQueryString { + query: "".to_string(), + max_hits: 0, + ..Default::default() + }, + ) + .await + .is_ok_and(|search_response| search_response.num_hits >= num_hits) + }, + Duration::from_secs(30), + Duration::from_millis(200), + ) + .await + .unwrap_or_else(|_| panic!("index `{index_id}` should have at least {num_hits} documents")); +} fn create_mock_data_file(num_lines: usize) -> (NamedTempFile, Uri) { let mut temp_file = tempfile::NamedTempFile::new().unwrap(); @@ -103,11 +125,7 @@ async fn test_sqs_with_duplicates() { sqs_test_helpers::send_message(&sqs_client, &queue_url, tmp_mock_data_files[5].1.as_str()) .await; - sandbox - .wait_for_splits(index_id, Some(vec![SplitState::Published]), 1) - .await - .unwrap(); - + wait_for_num_hits_at_least(&sandbox, index_id, 10 * 1000).await; sandbox.assert_hit_count(index_id, "", 10 * 1000).await; // The two duplicates could not be acknowledged when the were received @@ -203,11 +221,7 @@ async fn test_sqs_garbage_collect() { sqs_test_helpers::send_message(&sqs_client, &queue_url, uri.as_str()).await; } - sandbox - .wait_for_splits(index_id, Some(vec![SplitState::Published]), 1) - .await - .unwrap(); - + wait_for_num_hits_at_least(&sandbox, index_id, 10 * 1000).await; sandbox.assert_hit_count(index_id, "", 10 * 1000).await; wait_until_predicate( diff --git a/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs b/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs index 3ae458bdfa8..4abcb5a63ba 100644 --- a/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs +++ b/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs @@ -282,10 +282,13 @@ impl Handler for DeleteTaskPipeline { #[cfg(test)] mod tests { + use std::time::Duration; + use async_trait::async_trait; use quickwit_actors::{Handler, Universe}; use quickwit_common::pubsub::EventBroker; use quickwit_common::temp_dir::TempDirectory; + use quickwit_common::test_utils::wait_until_predicate; use quickwit_indexing::TestSandbox; use quickwit_indexing::actors::MergeSchedulerService; use quickwit_metastore::{ListSplitsRequestExt, MetastoreServiceStreamSplitsExt, SplitState}; @@ -400,9 +403,35 @@ mod tests { let (pipeline_mailbox, pipeline_handler) = universe.spawn_builder().spawn(pipeline); // Ensure that the message sent by initialize method is processed. let _ = pipeline_handler.process_pending_and_observe().await.state; - // Pipeline will first fail and we need to wait a OBSERVE_PIPELINE_INTERVAL * some number - // for the pipeline state to be updated. - universe.sleep(OBSERVE_PIPELINE_INTERVAL * 5).await; + // Pipeline will first fail, then publish the split with the delete applied. We wait for + // that split, then a OBSERVE_PIPELINE_INTERVAL * some number for the pipeline state to be + // updated. + wait_until_predicate( + || { + let metastore = metastore.clone(); + let index_uid = index_uid.clone(); + async move { + let Ok(splits) = metastore + .list_splits(ListSplitsRequest::try_from_index_uid(index_uid).unwrap()) + .await + .unwrap() + .collect_splits() + .await + else { + return false; + }; + splits.iter().any(|split| { + split.split_state == SplitState::Published + && split.split_metadata.delete_opstamp == 1 + }) + } + }, + Duration::from_secs(30), + Duration::from_millis(50), + ) + .await + .expect("the split with the delete applied should be published"); + universe.sleep(OBSERVE_PIPELINE_INTERVAL * 2).await; let pipeline_state = pipeline_handler.process_pending_and_observe().await.state; assert_eq!(pipeline_state.delete_task_planner.metrics.num_errors, 1); assert_eq!(pipeline_state.downloader.metrics.num_errors, 0); diff --git a/quickwit/quickwit-janitor/src/actors/garbage_collector.rs b/quickwit/quickwit-janitor/src/actors/garbage_collector.rs index 239c9914bd4..d40526d96a2 100644 --- a/quickwit/quickwit-janitor/src/actors/garbage_collector.rs +++ b/quickwit/quickwit-janitor/src/actors/garbage_collector.rs @@ -247,7 +247,7 @@ mod tests { use std::path::Path; use std::sync::Arc; - use quickwit_actors::Universe; + use quickwit_actors::{ActorHandle, ObservationType, Universe}; use quickwit_common::ServiceStream; use quickwit_common::shared_consts::split_deletion_grace_period; use quickwit_metastore::{ @@ -264,6 +264,19 @@ mod tests { use super::*; + async fn observe_after_pass( + handle: &ActorHandle, + ) -> GarbageCollectorCounters { + loop { + let observation = handle.process_pending_and_observe().await; + match observation.obs_type { + ObservationType::Alive => return observation.state, + ObservationType::Timeout => continue, + ObservationType::PostMortem => panic!("the garbage collector exited"), + } + } + } + fn hashmap(key: K, value: V) -> HashMap { let mut map = HashMap::new(); map.insert(key, value); @@ -763,7 +776,7 @@ mod tests { let universe = Universe::with_accelerated_time(); let (_mailbox, handle) = universe.spawn_builder().spawn(garbage_collect_actor); - let counters = handle.process_pending_and_observe().await.state; + let counters = observe_after_pass(&handle).await; assert_eq!(counters.num_passes, 1); assert_eq!(counters.num_deleted_files, 14000); assert_eq!(counters.num_deleted_bytes, 20 * 14000); diff --git a/quickwit/quickwit-serve/src/health_check_api/handler.rs b/quickwit/quickwit-serve/src/health_check_api/handler.rs index 2ed474284a1..01be1d28533 100644 --- a/quickwit/quickwit-serve/src/health_check_api/handler.rs +++ b/quickwit/quickwit-serve/src/health_check_api/handler.rs @@ -31,7 +31,7 @@ use crate::rest::recover_fn; use crate::with_arg; const HEALTH_CHECK_ASK_TIMEOUT: Duration = if cfg!(any(test, feature = "testsuite")) { - Duration::from_millis(100) + Duration::from_secs(1) } else { Duration::from_secs(5) };