Repository navigation
Add PublishListener hook for observing every publish #27
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
03051c3
8314779
bd2e44e
4675fd8
e920f74
2005a01
39dd80c
b88a66c
b5ce6f6
e70b880
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
|
|
@@ -5,8 +5,56 @@ All notable changes to this project will be documented in this file. | |||||
| The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), | ||||||
| and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). | ||||||
|
|
||||||
|
|
||||||
| ## [Unreleased] | ||||||
|
|
||||||
| ### Added | ||||||
|
|
||||||
| - `PublishListener`: a pluggable hook notified synchronously on every | ||||||
| `Orchestrator.publish` that reaches the bus, before subscriber dispatch. | ||||||
| Intended for diagnostics -- recording, metrics, tracing -- which previously | ||||||
| had no way to observe the bus without reimplementing `publish`. | ||||||
|
|
||||||
| Registered with `addPublishListener` / `removePublishListener`, declared as | ||||||
| `default` methods so existing `Orchestrator` implementors and test doubles keep | ||||||
| compiling. Listeners run on the publishing thread, must not block, and are | ||||||
| called in registration order. A listener that throws is caught and logged, so | ||||||
| diagnostics can never break the bus -- including when it throws an `Error` | ||||||
| such as `AssertionError` or `NoClassDefFoundError`. The one exception is | ||||||
| `OutOfMemoryError`, and only that: heap exhaustion is the one condition where | ||||||
| the recovery path needs memory too, since logging the error allocates. Every | ||||||
| other `VirtualMachineError` is contained. `StackOverflowError` is routinely | ||||||
| recoverable (an unbounded listener recursion unwinds that listener's frames | ||||||
| and leaves the stack whole, with the heap untouched). `InternalError` and | ||||||
| `UnknownError` are both documented as serious VM failures, but no subclass of | ||||||
| either marks the fatal instance, so `publish` cannot tell a fatal one from a | ||||||
| benign one and does not guess -- the hierarchy's silence is not read as | ||||||
| evidence that they are harmless. What it does settle is the family: neither is | ||||||
| an `OutOfMemoryError`, so no instance of either reaches the one rethrow. That | ||||||
| is read across every module in the boot layer rather than `java.base` alone, | ||||||
| because `catch (OutOfMemoryError)` matches subclasses from any module -- | ||||||
| `UnknownError` has no subclass there at all, and the sole `InternalError` | ||||||
| subclass is `java.util.zip.ZipError`, named as a fact about the family rather | ||||||
| than a hazard: the JDK documents it as no longer used and superseded by | ||||||
| `ZipException`, so a corrupt archive raises something else today. The bus is | ||||||
| not compromised in any of these cases, so the fault stays contained. With no | ||||||
| listeners registered, `publish` costs a single volatile read. | ||||||
|
|
||||||
| `addPublishListener` throws `UnsupportedOperationException` on an | ||||||
| implementation that does not support listeners, rather than accepting the | ||||||
| registration and quietly recording nothing; `removePublishListener` is always | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. P3: Prompt for AI agents
Suggested change
|
||||||
| a safe no-op. Implementors that can support listeners must override both. | ||||||
|
|
||||||
| Listeners are notified before the topic's type is validated, so a publish | ||||||
| rejected for a type mismatch is still reported. The two cases where a call | ||||||
| never reaches the hook are documented rather than reported: a publish to a | ||||||
| closed orchestrator returns early, and a publish of a `null` value throws | ||||||
| `IllegalArgumentException` before any listener runs -- in both, nothing was | ||||||
| published, and in the `null` case there is no value to hand a listener, which | ||||||
| is why `onPublish` documents its value as never null. One publish iterates a | ||||||
| snapshot of the listener list taken when it starts, so a listener | ||||||
| unregistered part-way through still sees that publish but not the next one. | ||||||
|
|
||||||
| ### Changed | ||||||
|
|
||||||
| - **`Orchestrator.hardware()` now returns a shared instance.** The | ||||||
|
|
||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -76,6 +76,11 @@ public interface Orchestrator extends AutoCloseable { | |
| * value's runtime type, so callers never have to pre-register topics. The publish | ||
| * itself is non-blocking even from the hardware thread. | ||
| * | ||
| * <p>A publish on a {@linkplain #isClosed() closed} orchestrator is ignored: it | ||
| * logs a warning and returns without dispatching. Neither that case nor a | ||
| * {@code null} value notifies any registered {@link PublishListener}; a | ||
| * type mismatch does. {@link PublishListener} states which is which. | ||
| * | ||
| * @param name the topic name | ||
| * @param value the value to publish (must not be null) | ||
| * @param <T> the value type | ||
|
|
@@ -270,6 +275,57 @@ public interface Orchestrator extends AutoCloseable { | |
| @Override | ||
| void close(); | ||
|
|
||
| // ---- diagnostics ----------------------------------------------------- | ||
|
|
||
| /** | ||
| * Registers a listener notified on every {@link #publish} that reaches the | ||
| * bus, before subscriber dispatch. See {@link PublishListener} for the | ||
| * threading contract. | ||
| * | ||
| * <p>Listeners run on the publishing thread and must not block. Registering | ||
| * none leaves publish with a single volatile read, so this is cheap to leave | ||
| * enabled permanently. | ||
| * | ||
| * <p>Not every {@code publish} call reaches a listener. A publish to a | ||
| * closed orchestrator returns before the hook, and a publish of a | ||
| * {@code null} value throws {@code IllegalArgumentException} before it, so | ||
| * neither is observable. A publish rejected for a <i>type mismatch</i> is | ||
| * observed, because that one happens after the bus has started acting: | ||
| * listeners are notified, and then the caller receives the exception. That | ||
| * is deliberate -- a type mismatch is a fault worth being able to observe, | ||
| * and the caller still receives the exception. | ||
| * | ||
| * <p>The default implementation throws rather than silently doing nothing. | ||
| * An implementor that cannot support listeners should fail at registration, | ||
| * where the mistake is visible, rather than leave the caller believing a | ||
| * recorder or a metric is attached when nothing is being captured. An | ||
| * implementor that can must override both this and | ||
| * {@link #removePublishListener(PublishListener)}: the default removal is a | ||
| * silent no-op, so a wrapper that forwarded only registration would let a | ||
| * caller detach a listener that is still attached. | ||
| * | ||
| * @param listener the listener; ignored if null | ||
| * @throws UnsupportedOperationException if this implementation cannot | ||
| * support listeners | ||
| */ | ||
| default void addPublishListener(PublishListener listener) { | ||
| if (listener == null) { | ||
| return; | ||
| } | ||
| throw new UnsupportedOperationException("publish listeners are not supported by this orchestrator"); | ||
|
cubic-dev-ai[bot] marked this conversation as resolved.
|
||
| } | ||
|
Comment on lines
+311
to
+316
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win The default Existing 🤖 Prompt for AI Agents |
||
|
|
||
| /** | ||
| * Removes a previously registered listener. Does nothing if it was not | ||
| * registered, or if this implementation does not support listeners, so | ||
| * cleanup is always safe to call. | ||
| * | ||
| * @param listener the listener to remove; ignored if null | ||
| */ | ||
| default void removePublishListener(PublishListener listener) { | ||
| // no-op | ||
| } | ||
|
|
||
| // ---- factories ------------------------------------------------------- | ||
|
|
||
| /** | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -8,6 +8,7 @@ | |
| import java.util.Optional; | ||
| import java.util.concurrent.CompletableFuture; | ||
| import java.util.concurrent.ConcurrentHashMap; | ||
| import java.util.concurrent.CopyOnWriteArrayList; | ||
| import java.util.concurrent.Executors; | ||
| import java.util.concurrent.LinkedBlockingQueue; | ||
| import java.util.concurrent.ScheduledExecutorService; | ||
|
|
@@ -202,6 +203,38 @@ public <T> Optional<Topic<T>> findTopic(String topicName, Class<T> type) { | |
| return Optional.of((Topic<T>) t); | ||
| } | ||
|
|
||
| // ---- publish listeners ------------------------------------------------ | ||
|
|
||
| // Copy-on-write: iteration happens on the publishing thread and must be | ||
| // lock-free, and registration is rare. A volatile read of the field is what | ||
| // keeps publish cheap when nobody is listening. | ||
| private final java.util.List<PublishListener> publishListeners = new CopyOnWriteArrayList<>(); | ||
|
|
||
| @Override | ||
| public void addPublishListener(PublishListener listener) { | ||
| if (listener != null) publishListeners.add(listener); | ||
| } | ||
|
|
||
| @Override | ||
| public void removePublishListener(PublishListener listener) { | ||
| if (listener != null) publishListeners.remove(listener); | ||
| } | ||
|
|
||
| /** | ||
| * Package-private: how many listeners are currently registered. | ||
| * | ||
| * <p>Exists for one test. Concurrent remove/add of the same listener can | ||
| * leave duplicates behind -- {@link java.util.concurrent.CopyOnWriteArrayList} | ||
| * removes only the first equal element -- and a duplicate is invisible to | ||
| * every other assertion available from outside, because the extra | ||
| * registration is a no-op that still has to be iterated and logged. | ||
| * Exposing the size is what lets the churn test enforce the bound it | ||
| * exists to test rather than merely claim it. | ||
| */ | ||
| int publishListenerCount() { | ||
| return publishListeners.size(); | ||
| } | ||
|
|
||
| // ---- publish --------------------------------------------------------- | ||
|
|
||
| @Override | ||
|
|
@@ -215,6 +248,101 @@ public <T> void publish(String topicName, T value) { | |
| throw new IllegalArgumentException("publish value cannot be null"); | ||
| } | ||
|
|
||
| // One timestamp for every listener, taken before the type check so a | ||
| // listener never observes a later instant than the publish itself. | ||
| // Guarded so the common case -- nobody listening -- is one read. | ||
| // | ||
| // Two guards above run *before* this block, so their failures are not | ||
| // observable through a listener, and the two IllegalArgumentExceptions | ||
| // this method can throw are not equivalent: | ||
| // | ||
| // closed orchestrator -> warn, return, listeners not notified | ||
| // null value -> throw, listeners not notified | ||
| // type mismatch -> listeners notified, THEN throw | ||
| // | ||
| // The split is argument validation versus a fault in an otherwise real | ||
| // publish. A closed bus or a null value means the call never became a | ||
| // publish: no topic was resolved, no latest value recorded, nothing | ||
| // dispatched -- and there is no value to hand a listener, which is why | ||
| // onPublish documents its value as never null. A type mismatch happens | ||
| // after the bus has started acting, and is exactly the kind of fault a | ||
| // recording should be able to show; the caller still sees the throw. | ||
| java.util.List<PublishListener> listeners = publishListeners; | ||
| if (!listeners.isEmpty()) { | ||
|
kody-ai[bot] marked this conversation as resolved.
|
||
| long now = System.nanoTime(); | ||
| // Iterate, never index. CopyOnWriteArrayList's size() and get(i) | ||
| // each read the current array independently, so a listener that | ||
| // unregistered a later one mid-publish left the cached size() | ||
| // stale and get(i) threw IndexOutOfBoundsException. That call sits | ||
| // inside the try below, so the bus did not break -- but the loop | ||
| // aborted, every remaining listener was silently skipped for that | ||
| // publish, and the log blamed a listener for "throwing" when the | ||
| // list was merely shorter than expected. The iterator is backed by | ||
| // a single stable snapshot, so one publish always notifies exactly | ||
| // the listeners registered when it started. | ||
| for (PublishListener listener : listeners) { | ||
| try { | ||
| listener.onPublish(topicName, value, now); | ||
| } catch (OutOfMemoryError heapGone) { | ||
| // The one deliberate exception to "a listener can never | ||
| // break the bus", and it is exactly one class. Heap | ||
| // exhaustion is the only condition where there is no | ||
| // publish left worth protecting, because the recovery | ||
| // path itself needs memory: log.error builds a message | ||
| // and fills in a stack trace, so stepping to the next | ||
| // listener would fail a second time in a worse place. | ||
| // Let it out and let the JVM deal with it. | ||
| throw heapGone; | ||
| } catch (Throwable t) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🔵 Trivial | 💤 Low value
The catch logs the error, so this is not silent. Consider catching 🤖 Prompt for AI AgentsSource: Learnings |
||
| // Everything else is contained, including the sibling | ||
| // VirtualMachineErrors. Catching Throwable rather than | ||
| // Exception is the point: a diagnostics module on a | ||
| // robot fails with AssertionError or NoClassDefFoundError | ||
| // at least as often as it fails with a RuntimeException, | ||
| // and none of those may take the bus down. | ||
| // | ||
| // StackOverflowError in particular is the common one and | ||
| // is contained on purpose. A listener that recurses | ||
| // without a bound blows the stack, the JVM unwinds that | ||
| // listener's frames, and the stack is whole again by the | ||
| // time we are here -- the heap was never touched. The | ||
| // listener's bug is its own; failing the publish, and | ||
| // with it every subscriber, would be the hook causing | ||
| // the outage it exists to diagnose. If the stack is | ||
| // genuinely gone, the log call below raises a fresh | ||
| // StackOverflowError that escapes publish anyway, which | ||
| // is the right outcome for that case. | ||
| // | ||
| // InternalError and UnknownError are contained for the | ||
| // same reason: both are documented as serious VM | ||
| // failures, and no subclass of either marks the fatal | ||
| // instance, so this code cannot tell a fatal one from a | ||
| // benign one and does not guess. Silence in the hierarchy | ||
| // is not evidence that these are harmless, and is not read | ||
| // as such. What the hierarchy does settle is the family | ||
| // they belong to: neither is an OutOfMemoryError, so no | ||
| // instance of either reaches the rethrow above, and | ||
| // containment is the whole of the decision. The image is | ||
| // scanned across every module in the boot layer for that | ||
| // answer, not java.base alone, because catch | ||
| // (OutOfMemoryError) matches subclasses from any module: | ||
| // UnknownError has no subclass there at all, and the sole | ||
| // InternalError subclass is java.util.zip.ZipError. That is | ||
| // named as a fact about the family, not as a hazard this | ||
| // hook will meet -- the JDK documents ZipError as no longer | ||
| // used and obsolete, superseded by ZipException, so a | ||
| // corrupt archive raises something else today. What can be | ||
| // told is that the bus itself is fine: the fault is inside | ||
| // one listener's frame, and the other listeners and the | ||
| // subscribers have no dependence on it. Re-throwing a | ||
| // VirtualMachineError merely because of its type was the | ||
| // bug; the narrow case above is the one that is actually | ||
| // unrecoverable. | ||
| log.error(name, "publish listener threw", t); | ||
| } | ||
| } | ||
| } | ||
|
Comment on lines
+251
to
+344
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win Listeners fire before topic type validation, so rejected publishes are still observed. The listener loop runs before the 🤖 Prompt for AI Agents |
||
|
|
||
| // Lazily create the topic from the value's runtime type. This matches | ||
| // Heron's behavior: publishers don't have to pre-register topics. | ||
| Class<?> valueType = value.getClass(); | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.