Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions quickwit/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions quickwit/quickwit-indexing/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down Expand Up @@ -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"] }
Expand Down
69 changes: 69 additions & 0 deletions quickwit/quickwit-indexing/src/actors/uploader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -171,6 +175,7 @@ pub struct Uploader {
max_concurrent_split_uploads: usize,
counters: UploaderCounters,
event_broker: EventBroker,
upload_tasks: TaskTracker,
}

impl Uploader {
Expand All @@ -195,6 +200,7 @@ impl Uploader {
max_concurrent_split_uploads,
counters: Default::default(),
event_broker,
upload_tasks: TaskTracker::new(),
}
}
async fn acquire_semaphore(
Expand Down Expand Up @@ -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<Self>,
) -> 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]
Expand All @@ -270,6 +287,21 @@ impl Handler<PackagedSplitBatch> for Uploader {
&mut self,
batch: PackagedSplitBatch,
ctx: &ActorContext<Self>,
) -> 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<Self>,
) -> Result<(), ActorExitStatus> {
fail_point!("uploader:before");
let split_update_sender = self
Expand Down Expand Up @@ -301,8 +333,10 @@ impl Handler<PackagedSplitBatch> 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());
Expand Down Expand Up @@ -510,6 +544,21 @@ impl Handler<EmptySplit> for Uploader {
&mut self,
empty_split: EmptySplit,
ctx: &ActorContext<Self>,
) -> 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<Self>,
) -> Result<(), ActorExitStatus> {
let split_update_sender = self
.split_update_mailbox
Expand All @@ -530,6 +579,22 @@ impl Handler<EmptySplit> for Uploader {
}
}

async fn drain_uploads_on_panic<T>(
upload_tasks: TaskTracker,
ctx: &ActorContext<Uploader>,
work: impl Future<Output = T>,
) -> 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)>,
Expand Down Expand Up @@ -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;
Expand Down
156 changes: 156 additions & 0 deletions quickwit/quickwit-indexing/src/actors/uploader_lifecycle_tests.rs
Original file line number Diff line number Diff line change
@@ -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::<Publisher>();
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<StageSplitsRequest, EmptyResponse, MetastoreError>| {
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")]
);
}
Loading