diff --git a/src/main/java/net/openhft/chronicle/queue/ExcerptAppender.java b/src/main/java/net/openhft/chronicle/queue/ExcerptAppender.java
index b4523321c0..dfdc13d8fe 100644
--- a/src/main/java/net/openhft/chronicle/queue/ExcerptAppender.java
+++ b/src/main/java/net/openhft/chronicle/queue/ExcerptAppender.java
@@ -76,6 +76,22 @@ default void writeBytes(@NotNull Bytes> bytes) {
default void pretouch() {
}
+ /**
+ * {@inheritDoc}
+ *
+ * For Queue, an output context is a roll cycle. Each appender's listener runs under the queue
+ * write lock before that appender's first data document in a cycle. It is not called for
+ * metadata, explicit-index writes or append-locked queues. A restarted appender can therefore
+ * write context again in an existing cycle. Context listeners are not supported with double
+ * buffering, and the caller retains ownership of the listener.
+ */
+ @NotNull
+ @Override
+ default ExcerptAppender contextListener(@NotNull Class writerType,
+ @NotNull MarshallableOut.ContextListener super T> listener) {
+ throw new UnsupportedOperationException();
+ }
+
/**
* Creates and returns a new writer proxy for the given interface {@code tclass} and the given {@code additional }
* interfaces.
diff --git a/src/main/java/net/openhft/chronicle/queue/impl/single/ContextListenerState.java b/src/main/java/net/openhft/chronicle/queue/impl/single/ContextListenerState.java
new file mode 100644
index 0000000000..967932f87c
--- /dev/null
+++ b/src/main/java/net/openhft/chronicle/queue/impl/single/ContextListenerState.java
@@ -0,0 +1,227 @@
+/*
+ * Copyright 2013-2025 chronicle.software; SPDX-License-Identifier: Apache-2.0
+ */
+package net.openhft.chronicle.queue.impl.single;
+
+import net.openhft.chronicle.bytes.MethodWriterBuilder;
+import net.openhft.chronicle.wire.BinaryMethodWriterInvocationHandler;
+import net.openhft.chronicle.wire.DocumentContext;
+import net.openhft.chronicle.wire.DocumentContextHolder;
+import net.openhft.chronicle.wire.MarshallableOut;
+import net.openhft.chronicle.wire.VanillaMethodWriterBuilder;
+import org.jetbrains.annotations.NotNull;
+import org.jetbrains.annotations.Nullable;
+
+import static java.util.Objects.requireNonNull;
+
+/** Queue configuration, appender state and locked output used by context listeners. */
+final class ContextListenerState extends DocumentContextHolder implements MarshallableOut {
+ static final ContextListenerState UNSET = new ContextListenerState(false);
+ static final ContextListenerState NONE = new ContextListenerState(true);
+
+ @Nullable
+ private final StoreAppender appender;
+ @Nullable
+ private final StoreAppender.StoreAppenderContext context;
+ @Nullable
+ private final Class> writerType;
+ @Nullable
+ private final MarshallableOut.ContextListener> listener;
+ @Nullable
+ private final Object methodWriter;
+ private boolean started;
+ private boolean notifying;
+ private int lastContextCount = -1;
+ private int nesting;
+
+ private ContextListenerState(boolean started) {
+ this.appender = null;
+ this.context = null;
+ this.writerType = null;
+ this.listener = null;
+ this.methodWriter = null;
+ this.started = started;
+ }
+
+ ContextListenerState(@Nullable Class> writerType,
+ @Nullable MarshallableOut.ContextListener> listener) {
+ this.appender = null;
+ this.context = null;
+ this.writerType = writerType;
+ this.listener = listener;
+ this.methodWriter = null;
+ }
+
+ private ContextListenerState(@NotNull StoreAppender appender,
+ @NotNull StoreAppender.StoreAppenderContext context,
+ @NotNull Class> writerType,
+ @NotNull MarshallableOut.ContextListener> listener) {
+ this.appender = requireNonNull(appender, "appender");
+ this.context = requireNonNull(context, "context");
+ this.writerType = null;
+ this.listener = requireNonNull(listener, "listener");
+ documentContext(context);
+ this.methodWriter = methodWriter(requireNonNull(writerType, "writerType"));
+ }
+
+ ContextListenerState forAppender(@NotNull StoreAppender appender,
+ @NotNull StoreAppender.StoreAppenderContext context) {
+ return listener == null
+ ? UNSET
+ : forAppender(appender, context, writerType, listener);
+ }
+
+ static ContextListenerState forAppender(@NotNull StoreAppender appender,
+ @NotNull StoreAppender.StoreAppenderContext context,
+ @NotNull Class> writerType,
+ @NotNull MarshallableOut.ContextListener> listener) {
+ return new ContextListenerState(
+ appender, context, writerType, listener);
+ }
+
+ boolean started() {
+ return started;
+ }
+
+ void onWriteAttempt() {
+ if (listener == null)
+ return;
+ if (notifying)
+ throw new IllegalStateException("Cannot write to the appender from within a ContextListener; " +
+ "write through the supplied method writer instead");
+ started = true;
+ }
+
+ boolean beforeDocument(boolean metaData) {
+ if (listener == null || metaData)
+ return false;
+ return notifyIfNeeded();
+ }
+
+ boolean beforeRawDocument() {
+ if (listener == null)
+ return false;
+ appender.resetPositionForContextListener();
+ return notifyIfNeeded();
+ }
+
+ private boolean notifyIfNeeded() {
+ StoreAppender appender = this.appender;
+ int contextCount = appender.cycle();
+ if (contextCount <= lastContextCount)
+ return false;
+ lastContextCount = contextCount;
+
+ notifying = true;
+ try {
+ try {
+ notifyListener();
+ } finally {
+ rollbackIfNotComplete();
+ }
+ return true;
+ } finally {
+ nesting = 0;
+ notifying = false;
+ }
+ }
+
+ @SuppressWarnings({"rawtypes", "unchecked"})
+ private void notifyListener() {
+ ((MarshallableOut.ContextListener) listener).onNewContext(methodWriter);
+ }
+
+ @Override
+ public DocumentContext writingDocument(boolean metaData) {
+ return notifying
+ ? acquireWritingDocument(metaData)
+ : appender.writingDocument(metaData);
+ }
+
+ @Override
+ public DocumentContext acquireWritingDocument(boolean metaData) {
+ if (!notifying)
+ return appender.acquireWritingDocument(metaData);
+
+ StoreAppender.StoreAppenderContext context = this.context;
+ if (nesting > 0 && context.wire() != null && context.isOpen()) {
+ if (!context.chainedElement()) {
+ assert metaData == context.isMetaData();
+ nesting++;
+ }
+ return this;
+ }
+
+ appender.openContextForContextListener(metaData);
+ nesting = 1;
+ return this;
+ }
+
+ @NotNull
+ @Override
+ public MethodWriterBuilder methodWriterBuilder(boolean metaData, @NotNull Class type) {
+ VanillaMethodWriterBuilder builder = new VanillaMethodWriterBuilder<>(type,
+ appender.queue().wireType(),
+ () -> new BinaryMethodWriterInvocationHandler(type, metaData,
+ () -> ContextListenerState.this));
+ builder.marshallableOut(this);
+ builder.metaData(metaData);
+ return builder;
+ }
+
+ @Override
+ public void rollbackIfNotComplete() {
+ if (!notifying) {
+ appender.rollbackIfNotComplete();
+ return;
+ }
+ StoreAppender.StoreAppenderContext context = this.context;
+ if (nesting == 0 || !context.isOpen())
+ return;
+ context.chainedElement(false);
+ context.rollbackOnClose();
+ nesting = 1;
+ close();
+ }
+
+ @Override
+ public boolean writingIsComplete() {
+ return notifying
+ ? context.writingIsComplete()
+ : appender.writingIsComplete();
+ }
+
+ @Override
+ public void rollbackOnClose() {
+ requireCallback();
+ context.rollbackOnClose();
+ }
+
+ @Override
+ public void close() {
+ requireCallback();
+ if (nesting == 0)
+ throw new IllegalStateException("No ContextListener document is open");
+ StoreAppender.StoreAppenderContext context = this.context;
+ if (context.chainedElement())
+ return;
+ if (nesting > 1) {
+ nesting--;
+ return;
+ }
+ nesting = 0;
+ appender.closeContextForContextListener();
+ }
+
+ @Override
+ public void reset() {
+ requireCallback();
+ context.reset();
+ nesting = 0;
+ }
+
+ private void requireCallback() {
+ if (!notifying)
+ throw new IllegalStateException("ContextListener document is not active");
+ }
+}
diff --git a/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueue.java b/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueue.java
index 37f29ffafe..434a493500 100644
--- a/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueue.java
+++ b/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueue.java
@@ -128,6 +128,8 @@ public class SingleChronicleQueue extends AbstractCloseable implements RollingCh
@NotNull
private final RollCycle rollCycle;
final AppenderListener appenderListener;
+ @NotNull
+ private final ContextListenerState contextListenerState;
protected int sourceId;
private int cycleFileRenamed = -1;
@NotNull
@@ -186,6 +188,7 @@ protected SingleChronicleQueue(@NotNull final SingleChronicleQueueBuilder builde
}
readOnly = builder.readOnly();
appenderListener = builder.appenderListener();
+ contextListenerState = builder.contextListenerState();
// ReadonlyTableStore is the no-metadata fallback. A SingleTableStore can also be
// read-only, but still contains the persisted cycle listing that a read-only queue
@@ -638,6 +641,12 @@ protected StoreFileListener storeFileListener() {
return storeFileListener;
}
+ @NotNull
+ ContextListenerState newContextListenerState(
+ StoreAppender appender, StoreAppender.StoreAppenderContext context) {
+ return contextListenerState.forAppender(appender, context);
+ }
+
// used by enterprise CQ
WireStoreSupplier storeSupplier() {
return storeSupplier;
@@ -1614,7 +1623,7 @@ private StoreSupplier() {
/**
* Acquires a {@link SingleChronicleQueueStore} for the specified cycle.
- * If the store doesn't exist and the strategy is {@link CreateStrategy.CREATE}, it will create a new store.
+ * If the store doesn't exist and the strategy is {@code CreateStrategy.CREATE}, it will create a new store.
*
* @param cycle the cycle to acquire the store for
* @param createStrategy the strategy for creating or reading the store
diff --git a/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueueBuilder.java b/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueueBuilder.java
index f9ea37ed7c..f1c59c46d4 100644
--- a/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueueBuilder.java
+++ b/src/main/java/net/openhft/chronicle/queue/impl/single/SingleChronicleQueueBuilder.java
@@ -139,6 +139,10 @@ public class SingleChronicleQueueBuilder extends SelfDescribingMarshallable impl
private Function createAppenderConditionCreator;
private long forceDirectoryListingRefreshIntervalMs = 60_000;
private AppenderListener appenderListener;
+ @Nullable
+ private Class> contextListenerWriterType;
+ @Nullable
+ private MarshallableOut.ContextListener> contextListener;
private SyncMode syncMode;
protected SingleChronicleQueueBuilder() {
@@ -370,6 +374,7 @@ public WireStoreFactory storeFactory() {
*/
@NotNull
public SingleChronicleQueue build() {
+ validateContextListenerCompatibility();
preBuild();
SingleChronicleQueue chronicleQueue;
@@ -386,6 +391,17 @@ public SingleChronicleQueue build() {
return chronicleQueue;
}
+ private void validateContextListenerCompatibility() {
+ if (contextListener == null)
+ return;
+ if (key != null || encodingSupplier != null)
+ throw new UnsupportedOperationException("contextListener is not supported on encoded or encrypted Enterprise queues");
+ if (writeBufferMode == BufferMode.Asynchronous)
+ throw new UnsupportedOperationException("contextListener is not supported on asynchronous Enterprise write buffers");
+ if (doubleBuffer)
+ throw new UnsupportedOperationException("contextListener is not supported with double buffering");
+ }
+
/**
* Performs post-build tasks such as setting the appender condition.
* The condition is added after the queue is constructed to avoid circular dependencies
@@ -1617,6 +1633,24 @@ public AppenderListener appenderListener() {
return appenderListener;
}
+ /**
+ * Sets the default context listener for appenders created by this queue.
+ * The caller retains ownership of the listener.
+ *
+ * @see ExcerptAppender#contextListener(Class, MarshallableOut.ContextListener)
+ */
+ public SingleChronicleQueueBuilder contextListener(@NotNull Class writerType,
+ @NotNull MarshallableOut.ContextListener super T> listener) {
+ contextListenerWriterType = requireNonNull(writerType, "writerType");
+ contextListener = requireNonNull(listener, "listener");
+ return this;
+ }
+
+ @NotNull
+ ContextListenerState contextListenerState() {
+ return new ContextListenerState(contextListenerWriterType, contextListener);
+ }
+
/**
* Sets the synchronization mode for the queue's memory-mapped files.
*
diff --git a/src/main/java/net/openhft/chronicle/queue/impl/single/StoreAppender.java b/src/main/java/net/openhft/chronicle/queue/impl/single/StoreAppender.java
index fa518341bc..620de46adf 100644
--- a/src/main/java/net/openhft/chronicle/queue/impl/single/StoreAppender.java
+++ b/src/main/java/net/openhft/chronicle/queue/impl/single/StoreAppender.java
@@ -30,6 +30,7 @@
import java.io.IOException;
import java.io.StreamCorruptedException;
import java.nio.BufferOverflowException;
+import java.util.Objects;
import java.util.concurrent.TimeUnit;
import static net.openhft.chronicle.queue.impl.single.SingleChronicleQueue.WARN_SLOW_APPENDER_MS;
@@ -75,6 +76,8 @@ class StoreAppender extends AbstractCloseable
private Pretoucher pretoucher = null;
private MicroToucher microtoucher = null;
private Wire bufferWire = null;
+ @NotNull
+ private ContextListenerState contextListenerState;
private int count = 0;
/**
@@ -94,6 +97,7 @@ class StoreAppender extends AbstractCloseable
this.writeLock = queue.writeLock();
this.appendLock = queue.appendLock();
this.context = new StoreAppenderContext();
+ this.contextListenerState = queue.newContextListenerState(this, context);
this.finalizer = Jvm.isResourceTracing() ? new Finalizer() : null;
try {
@@ -243,6 +247,7 @@ public void writeBytes(@NotNull final WriteBytesMarshallable marshallable) {
*/
@Override
protected void performClose() {
+ contextListenerState = ContextListenerState.NONE;
releaseBytesFor(wireForIndex);
releaseBytesFor(wire);
releaseBytesFor(bufferWire);
@@ -511,6 +516,21 @@ private boolean checkPositionOfHeader(final Bytes> bytes) {
return isReadyData(header) || isReadyMetaData(header) || isNotComplete(header);
}
+ @NotNull
+ @Override
+ public ExcerptAppender contextListener(@NotNull Class writerType,
+ @NotNull MarshallableOut.ContextListener super T> listener) {
+ throwExceptionIfClosed();
+ Objects.requireNonNull(writerType, "writerType");
+ Objects.requireNonNull(listener, "listener");
+ if (queue.doubleBuffer)
+ throw new UnsupportedOperationException("contextListener is not supported with double buffering");
+ if (contextListenerState.started())
+ throw new IllegalStateException("Cannot change contextListener after this appender has written");
+ contextListenerState = ContextListenerState.forAppender(this, context, writerType, listener);
+ return this;
+ }
+
@NotNull
@Override
// throws UnrecoverableTimeoutException
@@ -532,15 +552,27 @@ public DocumentContext writingDocument(final boolean metaData) {
throwExceptionIfClosed();
// we allow the sink process to write metaData
checkAppendLock(metaData);
+ ContextListenerState listenerState = startContextListenerWriteAttempt();
count++;
try {
- return prepareAndReturnWriteContext(metaData);
- } catch (RuntimeException e) {
+ return prepareAndReturnWriteContext(metaData, listenerState);
+ } catch (Throwable e) {
+ // Throwable, not just RuntimeException: an Error from a context listener must also
+ // restore count, or the next write takes the count>1 fast path and is handed a stale,
+ // never-opened context - permanently wedging the appender.
count--;
- throw e;
+ throw Jvm.rethrow(e);
}
}
+ private ContextListenerState startContextListenerWriteAttempt() {
+ ContextListenerState state = contextListenerState;
+ if (state == ContextListenerState.UNSET)
+ contextListenerState = state = ContextListenerState.NONE;
+ state.onWriteAttempt();
+ return state;
+ }
+
/**
* Prepares and returns the {@link StoreAppenderContext} for writing data. This method checks if
* the context needs to be reopened, locks the writeLock, handles double buffering if enabled,
@@ -549,7 +581,8 @@ public DocumentContext writingDocument(final boolean metaData) {
* @param metaData indicates if the write context is for metadata
* @return the prepared {@link StoreAppenderContext} ready for writing
*/
- private StoreAppender.StoreAppenderContext prepareAndReturnWriteContext(boolean metaData) {
+ private StoreAppender.StoreAppenderContext prepareAndReturnWriteContext(
+ boolean metaData, ContextListenerState listenerState) {
if (count > 1) {
assert metaData == context.metaData;
return context;
@@ -570,18 +603,22 @@ private StoreAppender.StoreAppenderContext prepareAndReturnWriteContext(boolean
if (this.cycle != cycle)
rollCycleTo(cycle);
- long safeLength = queue.overlapSize();
resetPosition();
+ if (listenerState.beforeDocument(metaData))
+ resetPosition();
assert !QueueSystemProperties.CHECK_INDEX || checkWritePositionHeaderNumber();
- // sets the writeLimit based on the safeLength
- openContext(metaData, safeLength);
+ // sets the writeLimit based on the overlap size
+ openContext(metaData, queue.overlapSize());
// Move readPosition to the start of the context. i.e. readRemaining() == 0
wire.bytes().readPosition(wire.bytes().writePosition());
- } catch (RuntimeException e) {
+ } catch (Throwable e) {
+ // Catch Throwable, not just RuntimeException: a context listener (or a corrupt-index
+ // AssertionError) can throw an Error, and leaking the cross-process write lock would
+ // stall every appender in every process until the lock times out.
writeLock.unlock();
- throw e;
+ throw Jvm.rethrow(e);
}
}
@@ -726,6 +763,42 @@ private void openContext(final boolean metaData, final long safeLength) {
context.metaData(metaData);
}
+ /**
+ * Opens a document for a context listener while this appender's write lock is already held.
+ *
+ * This is package-private solely for {@link ContextListenerState}. Context
+ * listeners cannot use the regular appender entry points because those paths attempt to acquire
+ * the non-reentrant write lock again.
+ *
+ * @param metaData whether the listener document contains metadata
+ */
+ void openContextForContextListener(boolean metaData) {
+ resetPosition();
+ openContext(metaData, queue.overlapSize());
+ }
+
+ void resetPositionForContextListener() {
+ resetPosition();
+ }
+
+ /**
+ * Closes the document written by a context listener without releasing the appender's write
+ * lock, which remains owned by the outer application write.
+ *
+ * The temporary count prevents any nesting state from the application write from suppressing
+ * the listener document's commit. The original count is restored for the application document
+ * that follows.
+ */
+ void closeContextForContextListener() {
+ int savedCount = count;
+ try {
+ count = 1;
+ context.close(false);
+ } finally {
+ count = savedCount;
+ }
+ }
+
/**
* Checks if the current header number matches the expected sequence in the queue.
* Throws an {@link AssertionError} if there is a mismatch.
@@ -778,6 +851,7 @@ public int sourceId() {
public void writeBytes(@NotNull final BytesStore, ?> bytes) {
throwExceptionIfClosed();
checkAppendLock();
+ ContextListenerState listenerState = startContextListenerWriteAttempt();
writeLock.lock();
try {
int cycle = queue.cycle();
@@ -787,7 +861,10 @@ public void writeBytes(@NotNull final BytesStore, ?> bytes) {
if (this.cycle != cycle)
rollCycleTo(cycle);
- this.positionOfHeader = writeHeader(wire, (int) queue.overlapSize()); // writeHeader sets wire.byte().writePosition
+ if (listenerState.beforeRawDocument())
+ resetPosition();
+
+ this.positionOfHeader = writeHeader(wire, queue.overlapSize()); // writeHeader sets wire.byte().writePosition
assert isInsideHeader(wire);
beforeAppend(wire, wire.headerNumber() + 1);
@@ -827,6 +904,7 @@ private boolean isInsideHeader(Wire wire) {
public void writeBytes(final long index, @NotNull final BytesStore, ?> bytes) {
throwExceptionIfClosed();
checkAppendLock();
+ startContextListenerWriteAttempt();
writeLock.lock();
try {
writeBytesInternal(index, bytes);
@@ -885,9 +963,8 @@ protected void writeBytesInternal(final long index, @NotNull final BytesStore,
private void writeBytesInternal(@NotNull final BytesStore, ?> bytes, boolean metadata) {
assert writeLock.locked();
try {
- int safeLength = (int) queue.overlapSize();
assert count == 0 : "count=" + count;
- openContext(metadata, safeLength);
+ openContext(metadata, queue.overlapSize());
try {
final Bytes> bytes0 = context.wire().bytes();
@@ -1433,6 +1510,16 @@ private void doRollback() {
}
}
+ @Override
+ public int contextCount() {
+ // Reject on any double-buffered queue, not just when this write happened to hit lock
+ // contention: otherwise the same code works or throws depending on runtime contention.
+ // Progressive contextCount usage and double buffering are an unsupported combination.
+ if (queue.doubleBuffer)
+ throw new IndexNotAvailableException("Context count is unavailable when double buffering because the target cycle is selected when the buffer is flushed");
+ return isClosed ? -1 : StoreAppender.this.cycle();
+ }
+
/**
* Returns the index of the current context. If the context is using double buffering, an
* {@link IndexNotAvailableException} will be thrown as the index is not available in this case.
diff --git a/src/main/java/net/openhft/chronicle/queue/impl/single/StoreTailer.java b/src/main/java/net/openhft/chronicle/queue/impl/single/StoreTailer.java
index b18903b1cd..4ccb4fb4b2 100644
--- a/src/main/java/net/openhft/chronicle/queue/impl/single/StoreTailer.java
+++ b/src/main/java/net/openhft/chronicle/queue/impl/single/StoreTailer.java
@@ -1814,6 +1814,11 @@ public int sourceId() {
return StoreTailer.this.sourceId();
}
+ @Override
+ public int contextCount() {
+ return queue.rollCycle().toCycle(index());
+ }
+
/**
* Closes the context, and if necessary, increments the index after reading a document.
*/
diff --git a/src/test/java/net/openhft/chronicle/queue/impl/single/ContextListenerCoreTest.java b/src/test/java/net/openhft/chronicle/queue/impl/single/ContextListenerCoreTest.java
new file mode 100644
index 0000000000..2659f3af66
--- /dev/null
+++ b/src/test/java/net/openhft/chronicle/queue/impl/single/ContextListenerCoreTest.java
@@ -0,0 +1,588 @@
+/*
+ * Copyright 2013-2025 chronicle.software; SPDX-License-Identifier: Apache-2.0
+ */
+package net.openhft.chronicle.queue.impl.single;
+
+import net.openhft.chronicle.bytes.Bytes;
+import net.openhft.chronicle.core.time.SetTimeProvider;
+import net.openhft.chronicle.core.time.SystemTimeProvider;
+import net.openhft.chronicle.queue.ChronicleQueue;
+import net.openhft.chronicle.queue.ExcerptAppender;
+import net.openhft.chronicle.queue.ExcerptTailer;
+import net.openhft.chronicle.queue.QueueTestCommon;
+import net.openhft.chronicle.wire.DocumentContext;
+import net.openhft.chronicle.wire.DocumentWritten;
+import net.openhft.chronicle.wire.MarshallableOut;
+import net.openhft.chronicle.wire.ProgressiveContext;
+import net.openhft.chronicle.wire.SelfDescribingMarshallable;
+import net.openhft.chronicle.wire.Wire;
+import net.openhft.chronicle.wire.WireType;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.io.File;
+import java.io.StringWriter;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static net.openhft.chronicle.queue.rollcycles.TestRollCycles.TEST_SECONDLY;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+public class ContextListenerCoreTest extends QueueTestCommon {
+ private final SetTimeProvider timeProvider = new SetTimeProvider();
+
+ @Before
+ public void useDeterministicSystemTimeProvider() {
+ timeProvider.currentTimeNanos(1_000_000_000L);
+ SystemTimeProvider.CLOCK = timeProvider;
+ }
+
+ @After
+ public void resetSystemTimeProvider() {
+ SystemTimeProvider.CLOCK = SystemTimeProvider.INSTANCE;
+ }
+
+ @Test
+ public void writesContextBeforeDataAndAllowsRetainingWriter() {
+ File path = getTmpDir();
+ AtomicInteger callbacks = new AtomicInteger();
+ AtomicReference retainedWriter = new AtomicReference<>();
+
+ try (ChronicleQueue queue = builder(path)
+ .contextListener(Events.class, writer -> {
+ callbacks.incrementAndGet();
+ retainedWriter.set(writer);
+ writer.context(new ServiceContext("queue"));
+ })
+ .build()) {
+ Events events = queue.methodWriter(Events.class);
+ events.message(new Message("one"));
+ retainedWriter.get().context(new ServiceContext("retained"));
+ timeProvider.advanceMillis(TEST_SECONDLY.lengthInMillis());
+ events.message(new Message("two"));
+ }
+
+ assertEquals(2, callbacks.get());
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: queue\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: {\n" +
+ " text: one\n" +
+ "}\n" +
+ "# index: 100000002\n" +
+ "context: {\n" +
+ " name: retained\n" +
+ "}\n" +
+ "# index: 200000000\n" +
+ "context: {\n" +
+ " name: queue\n" +
+ "}\n" +
+ "# index: 200000001\n" +
+ "message: {\n" +
+ " text: two\n" +
+ "}\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test
+ public void writesContextBeforeHeldDataDocument() {
+ File path = getTmpDir();
+
+ try (SingleChronicleQueue queue = builder(path).build()) {
+ ExcerptAppender appender = queue.acquireAppender();
+ appender.contextListener(Events.class,
+ writer -> writer.context(new ServiceContext("checkpoint")));
+ ProgressiveEvents events = appender.methodWriter(ProgressiveEvents.class);
+
+ writeMessages(events, "one", "two");
+ timeProvider.advanceMillis(TEST_SECONDLY.lengthInMillis());
+ writeMessages(events, "three");
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: checkpoint\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: {\n" +
+ " text: one\n" +
+ "}\n" +
+ "message: {\n" +
+ " text: two\n" +
+ "}\n" +
+ "# index: 200000000\n" +
+ "context: {\n" +
+ " name: checkpoint\n" +
+ "}\n" +
+ "# index: 200000001\n" +
+ "message: {\n" +
+ " text: three\n" +
+ "}\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test
+ public void writesProgressiveContextInsideHeldDocumentWhenCycleAdvances() {
+ File path = getTmpDir();
+ ServiceContext context = new ServiceContext("progressive");
+ int firstCycle;
+ int sameCycle;
+ int secondCycle;
+
+ try (ChronicleQueue queue = builder(path).build()) {
+ ProgressiveEvents events = queue.createAppender().methodWriter(ProgressiveEvents.class);
+
+ firstCycle = writeProgressively(queue, events, context, "one");
+ sameCycle = writeProgressively(queue, events, context, "two");
+ timeProvider.advanceMillis(TEST_SECONDLY.lengthInMillis());
+ secondCycle = writeProgressively(queue, events, context, "three");
+ }
+
+ assertEquals(firstCycle, sameCycle);
+ assertTrue(secondCycle > firstCycle);
+ assertEquals(secondCycle, context.lastContextCount);
+ assertContextCounts(path, firstCycle, sameCycle, secondCycle);
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: progressive\n" +
+ "}\n" +
+ "message: {\n" +
+ " text: one\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: {\n" +
+ " text: two\n" +
+ "}\n" +
+ "# index: 200000000\n" +
+ "context: {\n" +
+ " name: progressive\n" +
+ "}\n" +
+ "message: {\n" +
+ " text: three\n" +
+ "}\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test
+ public void indexedWriteDoesNotNotifyButStartsAppender() {
+ AtomicInteger callbacks = new AtomicInteger();
+
+ try (ChronicleQueue queue = builder(getTmpDir())
+ .contextListener(Events.class, writer -> callbacks.incrementAndGet())
+ .build()) {
+ StoreAppender appender = (StoreAppender) queue.createAppender();
+ long index = queue.rollCycle().toIndex(((SingleChronicleQueue) queue).cycle(), 0);
+ Bytes> payload = binaryPayload("indexed");
+ try {
+ appender.writeBytes(index, payload);
+ } finally {
+ payload.releaseLast();
+ }
+
+ assertEquals(0, callbacks.get());
+ assertThrows(IllegalStateException.class,
+ () -> appender.contextListener(Events.class, writer -> { }));
+ }
+ }
+
+ @Test
+ public void notifiesEachAppenderWithItsOwnContextOncePerRoll() {
+ File path = getTmpDir();
+ AtomicInteger firstCallbacks = new AtomicInteger();
+ AtomicInteger secondCallbacks = new AtomicInteger();
+
+ try (ChronicleQueue queue = builder(path).build()) {
+ StoreAppender first = (StoreAppender) queue.createAppender()
+ .contextListener(Events.class, writer -> {
+ firstCallbacks.incrementAndGet();
+ writer.context(new ServiceContext("first"));
+ });
+ StoreAppender second = (StoreAppender) queue.createAppender()
+ .contextListener(Events.class, writer -> {
+ secondCallbacks.incrementAndGet();
+ writer.context(new ServiceContext("second"));
+ });
+
+ first.writeMessage("message", "one");
+ first.writeMessage("message", "two");
+ assertEquals(1, firstCallbacks.get());
+ assertEquals(0, secondCallbacks.get());
+ second.writeMessage("message", "three");
+ second.writeMessage("message", "four");
+ assertEquals(1, firstCallbacks.get());
+ assertEquals(1, secondCallbacks.get());
+
+ timeProvider.advanceMillis(TEST_SECONDLY.lengthInMillis());
+ second.writeMessage("message", "five");
+ second.writeMessage("message", "six");
+ assertEquals(1, firstCallbacks.get());
+ assertEquals(2, secondCallbacks.get());
+ first.writeMessage("message", "seven");
+ first.writeMessage("message", "eight");
+ assertEquals(2, firstCallbacks.get());
+ assertEquals(2, secondCallbacks.get());
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: first\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: one\n" +
+ "# index: 100000002\n" +
+ "message: two\n" +
+ "# index: 100000003\n" +
+ "context: {\n" +
+ " name: second\n" +
+ "}\n" +
+ "# index: 100000004\n" +
+ "message: three\n" +
+ "# index: 100000005\n" +
+ "message: four\n" +
+ "# index: 200000000\n" +
+ "context: {\n" +
+ " name: second\n" +
+ "}\n" +
+ "# index: 200000001\n" +
+ "message: five\n" +
+ "# index: 200000002\n" +
+ "message: six\n" +
+ "# index: 200000003\n" +
+ "context: {\n" +
+ " name: first\n" +
+ "}\n" +
+ "# index: 200000004\n" +
+ "message: seven\n" +
+ "# index: 200000005\n" +
+ "message: eight\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test
+ public void listenerFailureAfterWritingContextIsNotRetriedInTheSameRoll() {
+ File path = getTmpDir();
+ AtomicInteger callbacks = new AtomicInteger();
+
+ try (ChronicleQueue queue = builder(path)
+ .contextListener(Events.class, writer -> {
+ callbacks.incrementAndGet();
+ writer.context(new ServiceContext("partial"));
+ throw new IllegalStateException("boom");
+ })
+ .build()) {
+ StoreAppender appender = (StoreAppender) queue.createAppender();
+
+ assertThrows(IllegalStateException.class,
+ () -> appender.writeMessage("message", "first"));
+ appender.writeMessage("message", "second");
+ assertEquals(1, callbacks.get());
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: partial\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: second\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test(timeout = 5_000)
+ public void listenerErrorRollsBackHeldDocumentAndLeavesAppenderUsable() {
+ File path = getTmpDir();
+ AtomicInteger callbacks = new AtomicInteger();
+
+ try (ChronicleQueue queue = builder(path)
+ .contextListener(ProgressiveEvents.class, writer -> {
+ callbacks.incrementAndGet();
+ DocumentContext document = writer.writingDocument();
+ writer.context(new ServiceContext("partial"));
+ assertTrue(document.isNotComplete());
+ throw new AssertionError("boom");
+ })
+ .build()) {
+ StoreAppender appender = (StoreAppender) queue.createAppender();
+
+ assertThrows(AssertionError.class,
+ () -> appender.writeMessage("message", "first"));
+ appender.writeMessage("message", "second");
+ assertEquals(1, callbacks.get());
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "message: second\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test(timeout = 5_000)
+ public void appenderWriteFromListenerFailsFast() {
+ File path = getTmpDir();
+ AtomicReference appenderRef = new AtomicReference<>();
+
+ try (ChronicleQueue queue = builder(path)
+ .contextListener(Events.class,
+ writer -> appenderRef.get().writeMessage("message", "reentrant"))
+ .build()) {
+ StoreAppender appender = (StoreAppender) queue.createAppender();
+ appenderRef.set(appender);
+
+ IllegalStateException exception = assertThrows(IllegalStateException.class,
+ () -> appender.writeMessage("message", "first"));
+ assertTrue(exception.getMessage().contains("supplied method writer"));
+ appender.writeMessage("message", "second");
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "message: second\n" +
+ "# no more messages at 8000000000000000\n", dump(path));
+ }
+
+ @Test
+ public void listenerCanHoldOneDocumentWhileWritingContext() {
+ AtomicBoolean outerDocumentRemainedOpen = new AtomicBoolean();
+ AtomicInteger contextCount = new AtomicInteger();
+
+ try (ChronicleQueue queue = builder(getTmpDir())
+ .contextListener(ProgressiveEvents.class, writer -> {
+ try (DocumentContext document = writer.writingDocument()) {
+ contextCount.set(document.contextCount());
+ writer.context(new ServiceContext("queue"));
+ outerDocumentRemainedOpen.set(document.isNotComplete());
+ }
+ })
+ .build()) {
+ int expectedCycle = ((SingleChronicleQueue) queue).cycle();
+ queue.createAppender().writeMessage("message", "one");
+
+ assertEquals(expectedCycle, contextCount.get());
+ assertTrue(outerDocumentRemainedOpen.get());
+ }
+ }
+
+ @Test
+ public void listenerClosePreservesNestingForDocumentWrite() {
+ try (ChronicleQueue queue = builder(getTmpDir())
+ .timeoutMS(100)
+ .contextListener(Events.class,
+ writer -> writer.context(new ServiceContext("queue")))
+ .build()) {
+ assertNestedWriteReusesContext((StoreAppender) queue.createAppender());
+ }
+ }
+
+ @Test
+ public void listenerClosePreservesNestingAfterRawWrite() {
+ try (ChronicleQueue queue = builder(getTmpDir())
+ .timeoutMS(100)
+ .contextListener(Events.class,
+ writer -> writer.context(new ServiceContext("queue")))
+ .build()) {
+ StoreAppender appender = (StoreAppender) queue.createAppender();
+ Bytes> payload = binaryPayload("raw");
+ try {
+ appender.writeBytes(payload);
+ } finally {
+ payload.releaseLast();
+ }
+
+ assertNestedWriteReusesContext(appender);
+ }
+ }
+
+ @Test
+ public void listenerRemainsCallerOwned() {
+ CloseableListener listener = new CloseableListener();
+
+ try (ChronicleQueue queue = builder(getTmpDir())
+ .contextListener(Events.class, listener)
+ .build()) {
+ queue.methodWriter(Events.class).message(new Message("one"));
+ }
+
+ assertFalse(listener.closed);
+ }
+
+ @Test
+ public void rejectsDoubleBuffering() {
+ assertThrows(UnsupportedOperationException.class, () -> builder(getTmpDir())
+ .doubleBuffer(true)
+ .contextListener(Events.class, writer -> { })
+ .build());
+
+ try (ChronicleQueue queue = builder(getTmpDir()).doubleBuffer(true).build()) {
+ assertThrows(UnsupportedOperationException.class,
+ () -> queue.createAppender().contextListener(Events.class, writer -> { }));
+ }
+ }
+
+ @Test
+ public void encodesContextWithQueueWireType() {
+ File path = getTmpDir();
+
+ try (ChronicleQueue queue = builder(path, WireType.BINARY_LIGHT)
+ .contextListener(Events.class,
+ writer -> writer.context(new ServiceContext("queue")))
+ .build()) {
+ queue.createAppender().writeMessage("message", "one");
+ }
+
+ assertEquals("" +
+ "# firstIndex: 100000000\n" +
+ "# index: 100000000\n" +
+ "context: {\n" +
+ " name: queue\n" +
+ "}\n" +
+ "# index: 100000001\n" +
+ "message: one\n" +
+ "# no more messages at 8000000000000000\n", dump(path, WireType.BINARY_LIGHT));
+ }
+
+ private static SingleChronicleQueueBuilder builder(File path) {
+ return builder(path, WireType.BINARY);
+ }
+
+ private static SingleChronicleQueueBuilder builder(File path, WireType wireType) {
+ return SingleChronicleQueueBuilder.builder(path, wireType)
+ .testBlockSize()
+ .rollCycle(TEST_SECONDLY);
+ }
+
+ private static void writeMessages(ProgressiveEvents events, String... messages) {
+ try (DocumentContext document = events.writingDocument()) {
+ assertTrue(document.isOpen());
+ for (String message : messages)
+ events.message(new Message(message));
+ }
+ }
+
+ private static int writeProgressively(ChronicleQueue queue,
+ ProgressiveEvents events,
+ ServiceContext context,
+ String message) {
+ try (DocumentContext document = events.writingDocument()) {
+ int contextCount = document.contextCount();
+ assertEquals(queue.rollCycle().toCycle(document.index()), contextCount);
+ if (context.needsResending(contextCount))
+ events.context(context);
+ events.message(new Message(message));
+ return contextCount;
+ }
+ }
+
+ private static void assertContextCounts(File path, int... expected) {
+ try (ChronicleQueue queue = builder(path).build();
+ ExcerptTailer tailer = queue.createTailer()) {
+ for (int contextCount : expected) {
+ try (DocumentContext document = tailer.readingDocument()) {
+ assertTrue(document.isPresent());
+ assertEquals(contextCount, document.contextCount());
+ assertEquals(queue.rollCycle().toCycle(document.index()), document.contextCount());
+ }
+ }
+ try (DocumentContext document = tailer.readingDocument()) {
+ assertFalse(document.isPresent());
+ }
+ }
+ }
+
+ private static Bytes> binaryPayload(String value) {
+ Bytes> payload = Bytes.allocateElasticOnHeap();
+ Wire wire = WireType.BINARY.apply(payload);
+ wire.write("message").text(value);
+ return payload;
+ }
+
+ private static void assertNestedWriteReusesContext(StoreAppender appender) {
+ try (DocumentContext outer = appender.writingDocument()) {
+ outer.wire().write("outer").text("one");
+ try (DocumentContext nested = appender.writingDocument()) {
+ assertSame(outer, nested);
+ nested.wire().write("nested").text("two");
+ }
+ assertTrue(outer.isNotComplete());
+ }
+ }
+
+ private static String dump(File path) {
+ return dump(path, WireType.BINARY);
+ }
+
+ private static String dump(File path, WireType wireType) {
+ try (ChronicleQueue queue = builder(path, wireType).build()) {
+ StringWriter writer = new StringWriter();
+ queue.dump(writer, 0, Long.MAX_VALUE);
+ return writer.toString();
+ }
+ }
+
+ interface Events {
+ void context(ServiceContext context);
+
+ void message(Message message);
+ }
+
+ interface ProgressiveEvents extends Events, DocumentWritten {
+ }
+
+ static final class ServiceContext extends SelfDescribingMarshallable implements ProgressiveContext {
+ private final String name;
+ private transient int lastContextCount = -1;
+
+ ServiceContext(String name) {
+ this.name = name;
+ }
+
+ @Override
+ public boolean needsResending(int contextCount) {
+ if (contextCount <= lastContextCount)
+ return false;
+ lastContextCount = contextCount;
+ return true;
+ }
+ }
+
+ static final class Message extends SelfDescribingMarshallable {
+ private final String text;
+
+ Message(String text) {
+ this.text = text;
+ }
+ }
+
+ private static final class CloseableListener
+ implements MarshallableOut.ContextListener, AutoCloseable {
+ private boolean closed;
+
+ @Override
+ public void onNewContext(Events writer) {
+ writer.context(new ServiceContext("queue"));
+ }
+
+ @Override
+ public void close() {
+ closed = true;
+ }
+ }
+}