From 7e5a5a3c5ae93fd6dc8b4d8a4fc871d87876f36e Mon Sep 17 00:00:00 2001 From: vignesh_A Date: Fri, 10 Jul 2026 01:53:00 +0530 Subject: [PATCH] test: wait for async refresh events in testEventsAreEmitted Fixes flaky LocalIcebergCatalogRelationalTest.testEventsAreEmitted where AFTER_REFRESH_TABLE was asserted before Vert.x event-bus delivery completed. - Clear the test listener after setup so only property-commit refresh events are observed - Await BEFORE/AFTER refresh events via Awaitility (same async path as production PolarisServiceBusEventDispatcher) - Mirror the wait for view refresh events - Add TestPolarisEventListener.hasEvent for await conditions Fixes #5019 --- .../iceberg/AbstractLocalIcebergCatalogTest.java | 15 +++++++++++++++ .../AbstractLocalIcebergCatalogViewTest.java | 13 +++++++++++++ .../listeners/TestPolarisEventListener.java | 5 +++++ 3 files changed, 33 insertions(+) diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java index 735fcc85207..c69c2bf900a 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java @@ -21,6 +21,7 @@ import static java.nio.charset.StandardCharsets.UTF_8; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Fail.fail; +import static org.awaitility.Awaitility.await; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyList; @@ -45,6 +46,7 @@ import java.lang.reflect.Method; import java.nio.file.Path; import java.time.Clock; +import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.Comparator; @@ -3087,12 +3089,25 @@ public void testEventsAreEmitted() { catalog.createNamespace(TestData.NAMESPACE); Table table = catalog.buildTable(TestData.TABLE, TestData.SCHEMA).create(); + // Clear after setup so we only observe refresh events from the property commits below. + // Events are published asynchronously via the Vert.x event bus + listener executor + // (see PolarisServiceBusEventDispatcher / PolarisEventListeners), so the test must wait + // for delivery rather than asserting immediately after commit. + testPolarisEventListener.clear(); + String key = "foo"; String valOld = "bar1"; String valNew = "bar2"; table.updateProperties().set(key, valOld).commit(); table.updateProperties().set(key, valNew).commit(); + await() + .atMost(Duration.ofSeconds(10)) + .until( + () -> + testPolarisEventListener.hasEvent(PolarisEventType.BEFORE_REFRESH_TABLE) + && testPolarisEventListener.hasEvent(PolarisEventType.AFTER_REFRESH_TABLE)); + PolarisEvent beforeRefreshEvent = testPolarisEventListener.getLatest(PolarisEventType.BEFORE_REFRESH_TABLE); Assertions.assertThat( diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogViewTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogViewTest.java index f18678b4adf..b708062c6c0 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogViewTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogViewTest.java @@ -18,12 +18,15 @@ */ package org.apache.polaris.service.catalog.iceberg; +import static org.awaitility.Awaitility.await; + import com.google.common.collect.ImmutableMap; import io.quarkus.test.junit.QuarkusMock; import io.smallrye.common.annotation.Identifier; import jakarta.inject.Inject; import java.io.IOException; import java.lang.reflect.Method; +import java.time.Duration; import java.util.List; import java.util.Map; import java.util.Set; @@ -245,12 +248,22 @@ public void testEventsAreEmitted() { .withQuery("a", "b") .create(); + // Clear after setup; refresh events are delivered asynchronously via Vert.x. + testPolarisEventListener.clear(); + String key = "foo"; String valOld = "bar1"; String valNew = "bar2"; view.updateProperties().set(key, valOld).commit(); view.updateProperties().set(key, valNew).commit(); + await() + .atMost(Duration.ofSeconds(10)) + .until( + () -> + testPolarisEventListener.hasEvent(PolarisEventType.BEFORE_REFRESH_VIEW) + && testPolarisEventListener.hasEvent(PolarisEventType.AFTER_REFRESH_VIEW)); + PolarisEvent beforeRefreshEvent = testPolarisEventListener.getLatest(PolarisEventType.BEFORE_REFRESH_VIEW); Assertions.assertThat( diff --git a/runtime/service/src/testFixtures/java/org/apache/polaris/service/events/listeners/TestPolarisEventListener.java b/runtime/service/src/testFixtures/java/org/apache/polaris/service/events/listeners/TestPolarisEventListener.java index e23a7a92646..2264689cb98 100644 --- a/runtime/service/src/testFixtures/java/org/apache/polaris/service/events/listeners/TestPolarisEventListener.java +++ b/runtime/service/src/testFixtures/java/org/apache/polaris/service/events/listeners/TestPolarisEventListener.java @@ -36,6 +36,11 @@ public void clear() { latestEvents.clear(); } + /** Returns whether an event of the given type has been recorded. */ + public boolean hasEvent(PolarisEventType type) { + return latestEvents.containsKey(type); + } + public PolarisEvent getLatest(PolarisEventType type) { var latest = latestEvents.get(type); if (latest == null) {