diff --git a/quickwit/Cargo.lock b/quickwit/Cargo.lock index 05da65b9ec2..df6526f7742 100644 --- a/quickwit/Cargo.lock +++ b/quickwit/Cargo.lock @@ -7538,6 +7538,8 @@ dependencies = [ "thiserror 2.0.18", "time", "tokio", + "tokio-util", + "tower 0.5.3", "tracing", "ulid", "unicode-segmentation", @@ -10775,6 +10777,7 @@ dependencies = [ "futures-core", "futures-io", "futures-sink", + "futures-util", "pin-project-lite", "tokio", ] diff --git a/quickwit/quickwit-indexing/Cargo.toml b/quickwit/quickwit-indexing/Cargo.toml index a0dc08b0ee5..e23257c00f0 100644 --- a/quickwit/quickwit-indexing/Cargo.toml +++ b/quickwit/quickwit-indexing/Cargo.toml @@ -44,6 +44,7 @@ tempfile = { workspace = true } thiserror = { workspace = true } time = { workspace = true } tokio = { workspace = true } +tokio-util = { workspace = true, default-features = false, features = ["rt"] } tracing = { workspace = true } ulid = { workspace = true } utoipa = { workspace = true } @@ -119,6 +120,7 @@ rand = { workspace = true } reqwest = { workspace = true } sqlx = { workspace = true, features = ["runtime-tokio", "postgres"] } tempfile = { workspace = true } +tower = { workspace = true } unicode-segmentation = { workspace = true } quickwit-actors = { workspace = true, features = ["testsuite"] } diff --git a/quickwit/quickwit-indexing/src/actors/uploader.rs b/quickwit/quickwit-indexing/src/actors/uploader.rs index 9ee80f35610..173b2ecd30c 100644 --- a/quickwit/quickwit-indexing/src/actors/uploader.rs +++ b/quickwit/quickwit-indexing/src/actors/uploader.rs @@ -13,14 +13,17 @@ // limitations under the License. use std::collections::HashSet; +use std::future::Future; use std::iter::FromIterator; use std::mem; +use std::panic::{AssertUnwindSafe, resume_unwind}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, OnceLock}; use anyhow::{Context, bail}; use async_trait::async_trait; use fail::fail_point; +use futures::FutureExt; use itertools::Itertools; use quickwit_actors::{Actor, ActorContext, ActorExitStatus, Handler, Mailbox, QueueCapacity}; use quickwit_common::pubsub::EventBroker; @@ -38,6 +41,7 @@ use quickwit_storage::{SplitPayload, SplitPayloadBuilder}; use serde::Serialize; use tokio::sync::oneshot::Sender; use tokio::sync::{Semaphore, SemaphorePermit, oneshot}; +use tokio_util::task::TaskTracker; use tracing::{Instrument, Span, debug, error, info, instrument, warn}; use crate::actors::Publisher; @@ -171,6 +175,7 @@ pub struct Uploader { max_concurrent_split_uploads: usize, counters: UploaderCounters, event_broker: EventBroker, + upload_tasks: TaskTracker, } impl Uploader { @@ -195,6 +200,7 @@ impl Uploader { max_concurrent_split_uploads, counters: Default::default(), event_broker, + upload_tasks: TaskTracker::new(), } } async fn acquire_semaphore( @@ -257,6 +263,17 @@ impl Actor for Uploader { fn name(&self) -> String { format!("{:?}", self.uploader_type) } + + async fn finalize( + &mut self, + _exit_status: &ActorExitStatus, + ctx: &ActorContext, + ) -> anyhow::Result<()> { + // Joining an uploader must also quiesce its storage and metastore writes. + self.upload_tasks.close(); + ctx.protect_future(self.upload_tasks.wait()).await; + Ok(()) + } } #[async_trait] @@ -270,6 +287,21 @@ impl Handler for Uploader { &mut self, batch: PackagedSplitBatch, ctx: &ActorContext, + ) -> Result<(), ActorExitStatus> { + drain_uploads_on_panic( + self.upload_tasks.clone(), + ctx, + self.handle_batch(batch, ctx), + ) + .await + } +} + +impl Uploader { + async fn handle_batch( + &mut self, + batch: PackagedSplitBatch, + ctx: &ActorContext, ) -> Result<(), ActorExitStatus> { fail_point!("uploader:before"); let split_update_sender = self @@ -301,8 +333,10 @@ impl Handler for Uploader { let retention_policy = self.retention_policy.clone(); debug!(split_ids=?split_ids, "start-stage-and-store-splits"); let event_broker = self.event_broker.clone(); + let upload_task = self.upload_tasks.token(); spawn_named_task( async move { + let _upload_task = upload_task; fail_point!("uploader:intask:before"); let mut split_metadata_list = Vec::with_capacity(batch.splits.len()); @@ -510,6 +544,21 @@ impl Handler for Uploader { &mut self, empty_split: EmptySplit, ctx: &ActorContext, + ) -> Result<(), ActorExitStatus> { + drain_uploads_on_panic( + self.upload_tasks.clone(), + ctx, + self.handle_empty_split(empty_split, ctx), + ) + .await + } +} + +impl Uploader { + async fn handle_empty_split( + &mut self, + empty_split: EmptySplit, + ctx: &ActorContext, ) -> Result<(), ActorExitStatus> { let split_update_sender = self .split_update_mailbox @@ -530,6 +579,22 @@ impl Handler for Uploader { } } +async fn drain_uploads_on_panic( + upload_tasks: TaskTracker, + ctx: &ActorContext, + work: impl Future, +) -> T { + match AssertUnwindSafe(work).catch_unwind().await { + Ok(result) => result, + Err(panic) => { + // Actor finalization is skipped on panic, but uploads must still finish before join. + upload_tasks.close(); + ctx.protect_future(upload_tasks.wait()).await; + resume_unwind(panic) + } + } +} + fn make_publish_operation( index_uid: IndexUid, packaged_splits_and_metadatas: Vec<(PackagedSplit, SplitMetadata)>, @@ -581,6 +646,10 @@ async fn upload_split( Ok(()) } +#[cfg(test)] +#[path = "uploader_lifecycle_tests.rs"] +mod lifecycle_tests; + #[cfg(test)] mod tests { use std::path::PathBuf; diff --git a/quickwit/quickwit-indexing/src/actors/uploader_lifecycle_tests.rs b/quickwit/quickwit-indexing/src/actors/uploader_lifecycle_tests.rs new file mode 100644 index 00000000000..a58430f46eb --- /dev/null +++ b/quickwit/quickwit-indexing/src/actors/uploader_lifecycle_tests.rs @@ -0,0 +1,156 @@ +// Copyright 2021-Present Datadog, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::time::Duration; + +use quickwit_actors::Universe; +use quickwit_common::temp_dir::TempDirectory; +use quickwit_common::tower::BoxService; +use quickwit_proto::metastore::{EmptyResponse, MetastoreError, MockMetastoreService}; +use quickwit_proto::types::{DocMappingUid, NodeId}; +use quickwit_storage::RamStorage; +use tower::ServiceExt; + +use super::*; +use crate::merge_policy::NopMergePolicy; +use crate::models::SplitAttrs; + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_uploader_shutdown_waits_for_in_flight_upload() { + for kill in [false, true] { + assert_uploader_waits_for_upload(kill, false).await; + } +} + +#[tokio::test] +async fn test_uploader_panic_waits_for_in_flight_upload() { + assert_uploader_waits_for_upload(false, true).await; +} + +async fn assert_uploader_waits_for_upload(kill: bool, panic: bool) { + let universe = Universe::new(); + let (publisher, _inbox) = universe.create_test_mailbox::(); + let entered = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let entered_clone = entered.clone(); + let release_clone = release.clone(); + let hold_stage = tower::layer::layer_fn( + move |service: BoxService| { + let entered = entered_clone.clone(); + let release = release_clone.clone(); + tower::service_fn(move |request| { + let service = service.clone(); + let entered = entered.clone(); + let release = release.clone(); + async move { + entered.notify_one(); + release.notified().await; + service.oneshot(request).await + } + }) + }, + ); + let mut metastore = MockMetastoreService::new(); + metastore + .expect_stage_splits() + .times(1) + .returning(|_| Ok(EmptyResponse {})); + let metastore = MetastoreServiceClient::tower() + .stack_stage_splits_layer(hold_stage) + .build_from_mock(metastore); + let storage = RamStorage::default(); + let uploader = Uploader::new( + UploaderType::IndexUploader, + metastore, + Arc::new(NopMergePolicy), + None, + IndexingSplitStore::create_without_local_store_for_test(Arc::new(storage.clone())), + publisher.into(), + 4, + EventBroker::default(), + ); + let (mailbox, handle) = universe.spawn_builder().spawn(uploader); + let batch = PackagedSplitBatch::new( + vec![PackagedSplit { + split_attrs: SplitAttrs { + node_id: NodeId::from_str("test-node"), + index_uid: IndexUid::for_test("test-index", 0), + source_id: "test-source".to_string(), + doc_mapping_uid: DocMappingUid::default(), + partition_id: 0, + time_range: None, + uncompressed_docs_size_in_bytes: 1, + num_docs: 1, + replaced_split_ids: Vec::new(), + split_id: "held-upload".into(), + delete_opstamp: 0, + num_merge_ops: 0, + }, + serialized_split_fields: Vec::new(), + split_scratch_directory: TempDirectory::for_test(), + tags: Default::default(), + hotcache_bytes: Vec::new(), + split_files: Vec::new(), + }], + None, + PublishLock::default(), + None, + Span::none(), + ); + mailbox.ask(batch).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), entered.notified()) + .await + .unwrap(); + let shutdown = async move { + if panic { + // A malformed batch panics after dispatching the first, held upload. + mailbox + .send_message(PackagedSplitBatch { + splits: Vec::new(), + checkpoint_delta_opt: None, + publish_lock: PublishLock::default(), + merge_task_opt: None, + batch_parent_span: Span::none(), + }) + .await + .unwrap(); + handle.join().await + } else if kill { + handle.kill().await + } else { + handle.quit().await + } + }; + tokio::pin!(shutdown); + let early_result = tokio::time::timeout(Duration::from_millis(100), &mut shutdown).await; + release.notify_one(); + let returned_early = early_result.is_ok(); + if !returned_early { + let (status, _) = tokio::time::timeout(Duration::from_secs(5), &mut shutdown) + .await + .unwrap(); + if panic { + assert!(matches!(status, ActorExitStatus::Panicked), "{status:?}"); + } + } + universe.quit().await; + assert!( + !returned_early, + "shutdown returned while staging/uploading could still write" + ); + assert_eq!( + storage.list_files().await, + vec![std::path::PathBuf::from("held-upload.split")] + ); +}