From 337acfd09ce12285e65b211fa7d4eeb236165ed3 Mon Sep 17 00:00:00 2001 From: mariano Date: Thu, 30 Jul 2026 22:58:15 -0500 Subject: [PATCH] fix: consumer module-info metrics requires + test homing - kpipe-consumer imports io.github.eschizoid.kpipe.metrics.* in four files but its module-info never required the metrics module - it compiled only through producer's 'requires transitive'. Declared directly (transitive: withMetrics(ConsumerMetrics) is public API), so tightening producer's edge later cannot break consumer. - Deleted the consumer-module ConsumerMetricsReporterTest: it shared an exact FQCN with the kpipe-metrics test for the same public class (breaking coverage attribution) and its three Mockito scenarios are subsumed by the metrics module's eight plain ones. - Moved CompositeMessageSinkJCStressTest producer->core jcstress: the SUT is a core class in the core-owned io.github.eschizoid.kpipe.sink package; core's concurrency gate no longer depends on producer's jcstress task. --- .../src/main/java/module-info.java | 1 + .../metrics/ConsumerMetricsReporterTest.java | 151 ------------------ .../CompositeMessageSinkJCStressTest.java | 0 3 files changed, 1 insertion(+), 151 deletions(-) delete mode 100644 lib/kpipe-consumer/src/test/java/io/github/eschizoid/kpipe/metrics/ConsumerMetricsReporterTest.java rename lib/{kpipe-producer => kpipe-core}/src/jcstress/java/io/github/eschizoid/kpipe/sink/CompositeMessageSinkJCStressTest.java (100%) diff --git a/lib/kpipe-consumer/src/main/java/module-info.java b/lib/kpipe-consumer/src/main/java/module-info.java index 7daea416..59fdb701 100644 --- a/lib/kpipe-consumer/src/main/java/module-info.java +++ b/lib/kpipe-consumer/src/main/java/module-info.java @@ -4,6 +4,7 @@ /// `kpipe-format-json`, `kpipe-format-avro`, `kpipe-format-protobuf`. module io.github.eschizoid.kpipe.consumer { requires transitive io.github.eschizoid.kpipe.core; + requires transitive io.github.eschizoid.kpipe.metrics; requires transitive io.github.eschizoid.kpipe.producer; requires jdk.httpserver; requires transitive kafka.clients; diff --git a/lib/kpipe-consumer/src/test/java/io/github/eschizoid/kpipe/metrics/ConsumerMetricsReporterTest.java b/lib/kpipe-consumer/src/test/java/io/github/eschizoid/kpipe/metrics/ConsumerMetricsReporterTest.java deleted file mode 100644 index 45a880f9..00000000 --- a/lib/kpipe-consumer/src/test/java/io/github/eschizoid/kpipe/metrics/ConsumerMetricsReporterTest.java +++ /dev/null @@ -1,151 +0,0 @@ -package io.github.eschizoid.kpipe.metrics; - -import static org.junit.jupiter.api.Assertions.*; -import static org.mockito.Mockito.*; - -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; -import java.util.function.Consumer; -import java.util.function.Supplier; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.ArgumentCaptor; -import org.mockito.Captor; -import org.mockito.Mock; -import org.mockito.junit.jupiter.MockitoExtension; - -@ExtendWith(MockitoExtension.class) -class ConsumerMetricsReporterTest { - - @Mock - private Supplier> metricsSupplier; - - @Mock - private Supplier uptimeSupplier; - - @Mock - private Consumer reporter; - - @Captor - private ArgumentCaptor reportCaptor; - - private ConsumerMetricsReporter metricsReporter; - private Map testMetrics; - - @BeforeEach - void setUp() { - testMetrics = new HashMap<>(); - testMetrics.put("messagesReceived", 100L); - testMetrics.put("messagesProcessed", 95L); - testMetrics.put("processingErrors", 5L); - } - - @Test - void shouldUseProvidedReporter() { - // Arrange - when(metricsSupplier.get()).thenReturn(testMetrics); - when(uptimeSupplier.get()).thenReturn(60000L); - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Act - metricsReporter.reportMetrics(); - - // Assert - verify(reporter).accept(reportCaptor.capture()); - final var report = reportCaptor.getValue(); - assertTrue(report.contains("messages received: 100")); - assertTrue(report.contains("messages processed: 95")); - assertTrue(report.contains("errors: 5")); - assertTrue(report.contains("uptime: 60000")); - } - - @Test - void shouldUseDefaultReporterWhenNull() { - // Arrange - when(metricsSupplier.get()).thenReturn(testMetrics); - when(uptimeSupplier.get()).thenReturn(60000L); - - // Act - metricsReporter = ConsumerMetricsReporter.forConsumer(metricsSupplier, uptimeSupplier); - - // Assert - assertDoesNotThrow(() -> metricsReporter.reportMetrics()); - } - - @Test - void shouldHandleNoMetricsGracefully() { - // Arrange - when(metricsSupplier.get()).thenReturn(Collections.emptyMap()); - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Act - metricsReporter.reportMetrics(); - - // Assert - verifyNoInteractions(reporter); - } - - @Test - void shouldHandleNullMetricsGracefully() { - // Arrange - when(metricsSupplier.get()).thenReturn(null); - - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Act - metricsReporter.reportMetrics(); - - // Assert - verifyNoInteractions(reporter); - } - - @Test - void shouldIncludeBackpressureMetricsInReportWhenPresent() { - // Arrange - testMetrics.put("backpressurePauseCount", 3L); - testMetrics.put("backpressureTimeMs", 1500L); - when(metricsSupplier.get()).thenReturn(testMetrics); - when(uptimeSupplier.get()).thenReturn(60000L); - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Act - metricsReporter.reportMetrics(); - - // Assert - verify(reporter).accept(reportCaptor.capture()); - final var report = reportCaptor.getValue(); - assertTrue(report.contains("backpressure pauses: 3")); - assertTrue(report.contains("backpressure time: 1500 ms")); - } - - @Test - void shouldNotIncludeBackpressureMetricsInReportWhenAbsent() { - // Arrange: testMetrics has no backpressure keys - when(metricsSupplier.get()).thenReturn(testMetrics); - when(uptimeSupplier.get()).thenReturn(60000L); - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Act - metricsReporter.reportMetrics(); - - // Assert - verify(reporter).accept(reportCaptor.capture()); - final var report = reportCaptor.getValue(); - assertFalse(report.contains("backpressure")); - } - - @Test - void shouldHandleExceptionInMetricsSupplier() { - // Arrange - when(metricsSupplier.get()).thenThrow(new RuntimeException("Test exception")); - - // Act - metricsReporter = new ConsumerMetricsReporter(metricsSupplier, uptimeSupplier, reporter); - - // Assert - assertDoesNotThrow(() -> metricsReporter.reportMetrics()); - verifyNoInteractions(reporter); - } -} diff --git a/lib/kpipe-producer/src/jcstress/java/io/github/eschizoid/kpipe/sink/CompositeMessageSinkJCStressTest.java b/lib/kpipe-core/src/jcstress/java/io/github/eschizoid/kpipe/sink/CompositeMessageSinkJCStressTest.java similarity index 100% rename from lib/kpipe-producer/src/jcstress/java/io/github/eschizoid/kpipe/sink/CompositeMessageSinkJCStressTest.java rename to lib/kpipe-core/src/jcstress/java/io/github/eschizoid/kpipe/sink/CompositeMessageSinkJCStressTest.java