From 0784cae616139e61c773343c3431d689356ddab9 Mon Sep 17 00:00:00 2001 From: chroniclekevinpowe Date: Tue, 15 Apr 2025 13:59:38 +1000 Subject: [PATCH 1/3] [WIP] Test allowing read past empty messages with MethodReaderQueueEntryReader --- .../reader/queueentryreaders/MethodReaderQueueEntryReader.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/net/openhft/chronicle/queue/internal/reader/queueentryreaders/MethodReaderQueueEntryReader.java b/src/main/java/net/openhft/chronicle/queue/internal/reader/queueentryreaders/MethodReaderQueueEntryReader.java index c1bba13f60..fc419f8a1f 100644 --- a/src/main/java/net/openhft/chronicle/queue/internal/reader/queueentryreaders/MethodReaderQueueEntryReader.java +++ b/src/main/java/net/openhft/chronicle/queue/internal/reader/queueentryreaders/MethodReaderQueueEntryReader.java @@ -61,7 +61,7 @@ public boolean read() { } if (bytes.isEmpty()) { // we read something from the queue but the MR filtered it i.e. did not dispatch - return false; + return true; } messageConsumer.consume(tailer.lastReadIndex(), bytes.toString()); bytes.clear(); From 20a32827f0d45aee119c3de85a6d38f41b6b4346 Mon Sep 17 00:00:00 2001 From: chroniclekevinpowe Date: Tue, 15 Apr 2025 16:32:56 +1000 Subject: [PATCH 2/3] [WIP] Test demonstrating behaviour of ChronicleReader encountering empty message --- .../internal/reader/ChronicleReaderTest.java | 55 +++++++++++++++++-- 1 file changed, 50 insertions(+), 5 deletions(-) diff --git a/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java b/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java index b131a6575c..0a9f02c467 100644 --- a/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java +++ b/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java @@ -91,7 +91,7 @@ private static long getCurrentQueueFileLength(final Path dataDir) throws IOExcep @Before public void before() { - assumeFalse(Jvm.isArm()); + // assumeFalse(Jvm.isArm()); // Reader opens queues in read-only mode if (OS.isWindows()) @@ -296,14 +296,59 @@ public void shouldNotShowIndexForHistoryMessages() { MessageHistory.writeHistory(dc); } } + } - basicReader() - .asMethodReader(SayWhen.class.getName()) - .execute(); +// basicReader() +// .asMethodReader(SayWhen.class.getName()) +// .execute(); +// +// assertTrue(capturedOutput.isEmpty()); +// } - assertTrue(capturedOutput.isEmpty()); + @Test + public void canReadPastEmptyMessage() { + dataDir = getTmpDir().toPath(); + try (final ChronicleQueue queue = SingleChronicleQueueBuilder.binary(dataDir) + .sourceId(1) + .testBlockSize() + .build()) { + final VanillaMethodWriterBuilder methodWriterBuilder = queue.methodWriterBuilder(ChronicleMethodReaderTest.All.class); + final ChronicleMethodReaderTest.All events = methodWriterBuilder.build(); + final ExcerptAppender appender = queue.createAppender(); + + for (int i = 0; i < 3; ) { + ChronicleMethodReaderTest.Method1Type m1 = new ChronicleMethodReaderTest.Method1Type(); + m1.text = "hello"; + m1.value = i; + m1.number = i; + events.method1(m1); + + try (DocumentContext dc = appender.writingDocument()) { + MessageHistory.writeHistory(dc); + } + + i++; + ChronicleMethodReaderTest.Method2Type m2 = new ChronicleMethodReaderTest.Method2Type(); + m2.text = "goodbye"; + m2.value = i; + m2.number = i; + events.method2(m2); + i++; + } + } + + System.out.println("Done with the queue!"); + + ChronicleReader methodReaderForQueue = new ChronicleReader() + .withBasePath(dataDir) + .asMethodReader(ChronicleMethodReaderTest.All.class.getName()) + .inReverseOrder() + .withMessageSink(System.out::println); + + methodReaderForQueue.execute(); } + @Test public void shouldNotIncludeMessageHistoryByDefaultMethodReader() { basicReader(). From 907398cb556c6bc3e026d7c466cce46c81396342 Mon Sep 17 00:00:00 2001 From: Sam Ross Date: Tue, 15 Apr 2025 13:24:45 +0100 Subject: [PATCH 3/3] Adds assertions to canReadPastEmptyMessageInReverseOrder test --- .../internal/reader/ChronicleReaderTest.java | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java b/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java index 0a9f02c467..f9ff935ae8 100644 --- a/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java +++ b/src/test/java/net/openhft/chronicle/queue/internal/reader/ChronicleReaderTest.java @@ -306,7 +306,7 @@ public void shouldNotShowIndexForHistoryMessages() { // } @Test - public void canReadPastEmptyMessage() { + public void canReadPastEmptyMessageInReverseOrder() { dataDir = getTmpDir().toPath(); try (final ChronicleQueue queue = SingleChronicleQueueBuilder.binary(dataDir) .sourceId(1) @@ -343,9 +343,20 @@ public void canReadPastEmptyMessage() { .withBasePath(dataDir) .asMethodReader(ChronicleMethodReaderTest.All.class.getName()) .inReverseOrder() - .withMessageSink(System.out::println); + .withMessageSink(capturedOutput::add); methodReaderForQueue.execute(); + + assertThat(capturedOutput.size(), is(8)); + capturedOutput.poll(); + assertThat(capturedOutput.poll(), containsString("goodbye")); + capturedOutput.poll(); + assertThat(capturedOutput.poll(), containsString("hello")); + capturedOutput.poll(); + assertThat(capturedOutput.poll(), containsString("goodbye")); + capturedOutput.poll(); + assertThat(capturedOutput.poll(), containsString("hello")); + capturedOutput.poll(); }