Skip to content
Open
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
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/coverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 3 additions & 1 deletion quickwit/.config/nextest.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
fail-fast = false
slow-timeout = { period = "60s", terminate-after = 10 }
retries = { backoff = "exponential", count = 2, delay = "1s", max-delay = "4s", jitter = true }
4 changes: 2 additions & 2 deletions quickwit/Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
32 changes: 11 additions & 21 deletions quickwit/quickwit-common/src/test_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
Expand Down
2 changes: 1 addition & 1 deletion quickwit/quickwit-indexing/src/actors/indexer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
62 changes: 57 additions & 5 deletions quickwit/quickwit-indexing/src/source/kafka_source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -759,9 +759,11 @@ fn message_payload_to_doc(message: &BorrowedMessage) -> Option<Bytes> {
#[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;
Expand Down Expand Up @@ -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<K, M, J, Q>(
Expand Down Expand Up @@ -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<RawDocBatch> = doc_processor_inbox.drain_for_test_typed();
assert!(messages.is_empty());
Expand Down Expand Up @@ -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<RawDocBatch> = doc_processor_inbox.drain_for_test_typed();
assert!(!messages.is_empty());
Expand Down Expand Up @@ -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<RawDocBatch> = doc_processor_inbox.drain_for_test_typed();
assert!(!messages.is_empty());
Expand Down Expand Up @@ -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<RawDocBatch> = doc_processor_inbox.drain_for_test_typed();
assert!(messages.is_empty());
Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/ingester.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -500,7 +500,7 @@ impl ClusterSandbox {
}
}
},
Duration::from_secs(10),
Duration::from_secs(30),
Duration::from_millis(100),
)
.await?;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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::<RestClientError>()
.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)
Expand Down
Loading
Loading