This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new da75706e3 test(connectors): add random source liveness helper (#3377)
da75706e3 is described below

commit da75706e337f35c89e5bccd90a3e384042918104
Author: Arun Singh <[email protected]>
AuthorDate: Fri Jun 5 14:18:05 2026 +0530

    test(connectors): add random source liveness helper (#3377)
---
 core/integration/tests/connectors/mod.rs           |  1 +
 .../tests/connectors/random/random_source.rs       | 97 +---------------------
 .../tests/connectors/random_source_liveness.rs     | 66 +++++++++++++++
 3 files changed, 69 insertions(+), 95 deletions(-)

diff --git a/core/integration/tests/connectors/mod.rs 
b/core/integration/tests/connectors/mod.rs
index bb4bcc69f..1016c6e51 100644
--- a/core/integration/tests/connectors/mod.rs
+++ b/core/integration/tests/connectors/mod.rs
@@ -30,6 +30,7 @@ mod mongodb;
 mod postgres;
 mod quickwit;
 mod random;
+mod random_source_liveness;
 mod runtime;
 mod stdout;
 
diff --git a/core/integration/tests/connectors/random/random_source.rs 
b/core/integration/tests/connectors/random/random_source.rs
index 5efb6fd08..5a05d7f53 100644
--- a/core/integration/tests/connectors/random/random_source.rs
+++ b/core/integration/tests/connectors/random/random_source.rs
@@ -17,107 +17,14 @@
  * under the License.
  */
 
-use iggy_common::MessageClient;
-use iggy_common::{Consumer, Identifier, PollingStrategy};
+use crate::connectors::random_source_liveness;
 use integration::harness::seeds;
 use integration::iggy_harness;
-use std::time::Duration;
-use tokio::time::sleep;
 
 #[iggy_harness(
     server(connectors_runtime(config_path = 
"tests/connectors/random/source.toml")),
     seed = seeds::connector_stream
 )]
 async fn random_source_produces_messages(harness: &TestHarness) {
-    sleep(Duration::from_secs(1)).await;
-
-    let client = harness.root_client().await.unwrap();
-    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
-    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
-    let consumer_id: Identifier = "test_consumer".try_into().unwrap();
-
-    let messages = client
-        .poll_messages(
-            &stream_id,
-            &topic_id,
-            None,
-            &Consumer::new(consumer_id),
-            &PollingStrategy::next(),
-            10,
-            true,
-        )
-        .await
-        .expect("Failed to poll messages");
-
-    assert!(
-        !messages.messages.is_empty(),
-        "No messages received from random source"
-    );
-    assert!(
-        messages.current_offset > 0,
-        "Current offset should be greater than 0"
-    );
-}
-
-#[iggy_harness(
-    server(connectors_runtime(config_path = 
"tests/connectors/random/source.toml")),
-    seed = seeds::connector_stream
-)]
-async fn state_persists_across_connector_restart(harness: &mut TestHarness) {
-    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
-    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
-    let consumer_id: Identifier = "state_test_consumer".try_into().unwrap();
-
-    sleep(Duration::from_secs(1)).await;
-
-    let client = harness.root_client().await.unwrap();
-    let offset_before = {
-        let messages = client
-            .poll_messages(
-                &stream_id,
-                &topic_id,
-                None,
-                &Consumer::new(consumer_id.clone()),
-                &PollingStrategy::next(),
-                100,
-                true,
-            )
-            .await
-            .expect("Failed to poll messages before restart");
-        assert!(
-            messages.current_offset > 0,
-            "Should have messages before restart"
-        );
-        messages.current_offset
-    };
-
-    harness
-        .server_mut()
-        .stop_dependents()
-        .expect("Failed to stop connectors");
-    harness
-        .server_mut()
-        .start_dependents()
-        .await
-        .expect("Failed to restart connectors");
-    sleep(Duration::from_secs(1)).await;
-
-    let offset_after = client
-        .poll_messages(
-            &stream_id,
-            &topic_id,
-            None,
-            &Consumer::new(consumer_id),
-            &PollingStrategy::next(),
-            100,
-            true,
-        )
-        .await
-        .expect("Failed to poll messages after restart")
-        .current_offset;
-
-    assert!(
-        offset_after > offset_before,
-        "After restart, offset {offset_after} should be greater than before 
{offset_before}"
-    );
+    random_source_liveness::assert_produces_messages(harness).await;
 }
diff --git a/core/integration/tests/connectors/random_source_liveness.rs 
b/core/integration/tests/connectors/random_source_liveness.rs
new file mode 100644
index 000000000..eb52eacb7
--- /dev/null
+++ b/core/integration/tests/connectors/random_source_liveness.rs
@@ -0,0 +1,66 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you 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 iggy_common::{Consumer, Identifier, MessageClient, PollingStrategy};
+use integration::harness::{TestHarness, seeds};
+use tokio::time::{sleep, timeout};
+
+const CONSUMER_NAME: &str = "random_source_liveness_consumer";
+const POLL_BATCH: u32 = 100;
+const RETRY_INTERVAL: Duration = Duration::from_millis(100);
+const POLL_TIMEOUT: Duration = Duration::from_secs(5);
+
+pub(crate) async fn assert_produces_messages(harness: &TestHarness) {
+    let client = harness.root_client().await.expect("root client");
+    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
+    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
+    let consumer_id: Identifier = CONSUMER_NAME.try_into().unwrap();
+
+    let poll = async {
+        loop {
+            if let Ok(polled) = client
+                .poll_messages(
+                    &stream_id,
+                    &topic_id,
+                    None,
+                    &Consumer::new(consumer_id.clone()),
+                    &PollingStrategy::next(),
+                    POLL_BATCH,
+                    true,
+                )
+                .await
+            {
+                if !polled.messages.is_empty() {
+                    return;
+                }
+            }
+
+            sleep(RETRY_INTERVAL).await;
+        }
+    };
+
+    timeout(POLL_TIMEOUT, poll).await.unwrap_or_else(|_| {
+        panic!(
+            "random source liveness timed out after {:?} waiting for messages",
+            POLL_TIMEOUT
+        )
+    })
+}

Reply via email to