From 3111a9d8a6f89d6234237c42dc336f8406678486 Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Fri, 2 Oct 2026 12:56:12 -0600 Subject: [PATCH 1/2] fix(bus): stop hardware-thread subscribers, primitive subscribe, widened casts Three defects found by auditing the publish path, each reproduced against the real classes before being fixed. 1. unregisterNode / unsubscribe did not stop an @OnHardwareThread subscriber. markSubscriptionAsHardwareThreaded replaces the list entry with a wrapper, but Subscription.handler() still returns the unwrapped handler, so removeSubscription searched for something the list no longer held and returned silently. An unregistered node kept receiving publishes, so it could still drive hardware. hardwareRerouted already recorded the registered wrapper; it was written but never read. 2. subscribe(String, int.class, handler) registered and then never fired. Class.cast() returns false for every argument when the class is primitive, so every delivery threw ClassCastException and dispatchCallback swallowed it into a log line. The type is boxed before the runtime check, matching the primitive/wrapper equivalence the publish path already applies. 3. getOrCreateTopic / findTopic / getLatestValue handed back a value the caller could not cast to the requested type. Object.isAssignableFrom( String) is true, so an Object topic satisfied a request for String, and the erased cast failed later at the caller's own line -- after isPresent() had already reported a value present. New Topic.safelyReturnsAs refuses the request with a message naming the actual type. Narrowing to a type the topic really holds still works, and subscribeRaw keeps accepting an Object topic because its handler receives Object, which is the binder's supported "several handlers, one topic" configuration. Tests: new SubscriptionLifetimeTest covers all three; each fails when its fix is reverted. Also in TopicLockFreeTest: - The concurrency test counted no samples, so it passed even if latest() returned empty always and every check inside the reader was skipped. It now counts and asserts a non-zero observation count. - Two assertions compared a value with itself: latest().publishNanos() against latestPublishNanos() (both read the same volatile field), and acceptsValueClass against acceptsType (which delegates to it). Replaced with literal expectations and with an externally-clocked publish window, which is the only form that cannot be satisfied by a skewed timestamp. - latestPublishNanos() now has its own window-bracketed assertion, so replacing it with a live clock reading fails the test. Docs: a comment records why the install-ordering guarantee is not covered by a test. The sample-to-store window is nanoseconds wide; a spinning detector catches it 3 runs in 6 but starves publishers and hung the suite, and a yielding detector never catches it. A test that cannot fail reliably is worse than none. CI: the Docker workflow's merge job never prepared IMAGE, so buildx got "@sha256:..." with no repository and the multi-arch manifest step failed. The published image has not updated since at least ce0f407; every build job recomputes IMAGE locally via GITHUB_ENV and that does not cross the job boundary. Suite: 90 tests, 0 failures across 10 consecutive runs and on a 2-core-pinned runner. javadoc: 0 errors, no new warnings. --- .github/workflows/docker.yml | 10 ++ CHANGELOG.md | 19 +++ .../aaravlabs/synapse/OrchestratorImpl.java | 95 ++++++++++-- .../java/com/aaravlabs/synapse/Topic.java | 17 ++- .../synapse/SubscriptionLifetimeTest.java | 140 ++++++++++++++++++ .../aaravlabs/synapse/TopicLockFreeTest.java | 126 +++++++++++----- 6 files changed, 359 insertions(+), 48 deletions(-) create mode 100644 src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java diff --git a/.github/workflows/docker.yml b/.github/workflows/docker.yml index 04e3a06..a41939a 100644 --- a/.github/workflows/docker.yml +++ b/.github/workflows/docker.yml @@ -201,6 +201,16 @@ jobs: username: ${{ github.actor }} password: ${{ secrets.GITHUB_TOKEN }} + # IMAGE is set through GITHUB_ENV by each build job's own prepare step, so it does + # not cross the job boundary. Without recomputing it here the manifest step sees an + # empty IMAGE and buildx rejects "@sha256:..." as an invalid reference, which is why + # the published multi-arch tag silently stopped updating. + - name: Prepare image name + env: + SOURCE: ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }} + run: | + echo "IMAGE=${SOURCE,,}" >> "$GITHUB_ENV" + - name: Extract metadata id: meta uses: docker/metadata-action@902fa8ec7d6ecbf8d84d538b9b233a880e428804 # v5.7.0 diff --git a/CHANGELOG.md b/CHANGELOG.md index aa2cd1d..d91f3b8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -131,6 +131,25 @@ listeners registered, `publish` costs a single volatile read. stale. Reading the timestamp first can only over-report the age, never under-report it. +### Fixed + +- **`unregisterNode` and `unsubscribe` now stop `@OnHardwareThread` subscribers.** + Routing a subscription to the hardware thread replaced the entry in the subscriber + list, but the subscription still reported the unwrapped handler, so removal searched + for something the list no longer held and silently did nothing. An unregistered node + kept receiving publishes — it could still drive hardware after being removed. +- **`subscribe` works with primitive classes.** `int.class.isInstance(x)` is false for + every argument, so casting with the raw type made every delivery throw + `ClassCastException`, which was swallowed into a log line: the subscriber registered, + never fired, and nothing said why. Primitive and wrapper types are now treated + interchangeably here, as they already were on the publish path. +- **Typed lookups no longer hand back a topic that cannot be cast to the requested + type.** An `Object` topic accepts a request for `String`, and the unchecked cast then + failed later at the caller's own line — after `isPresent()` had reported a value. The + request is now refused with a message naming the actual type. Narrowing to a type the + topic really holds still works, and `subscribeRaw` (used by the annotation binder) + still accepts an `Object` topic. + ## [0.4.0] - 2026-09-12 ### Changed diff --git a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java index 921af24..6a60bbc 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -156,6 +156,16 @@ public Topic getOrCreateTopic(String topicName, Class type) { + existing.type().getName() + ", cannot re-create as " + type.getName()); } + if (!existing.safelyReturnsAs(type)) { + // compatible, but this topic may hold values the caller cannot cast to T. + // Returning it anyway defers a ClassCastException to an unrelated line. + throw new IllegalArgumentException( + "Topic '" + topicName + "' holds " + + existing.type().getSimpleName() + ", which cannot be " + + "returned as " + type.getName() + + "; read it via findTopic(name) and treat the value as " + + existing.type().getName()); + } return (Topic) existing; } @@ -168,12 +178,53 @@ public Topic getOrCreateTopic(String topicName, Class type) { + prior.type().getName() + ", cannot re-create as " + type.getName()); } + if (!prior.safelyReturnsAs(type)) { + throw new IllegalArgumentException( + "Topic '" + topicName + "' holds " + + prior.type().getSimpleName() + ", which cannot be " + + "returned as " + type.getName() + + "; read it via findTopic(name) and treat the value as " + + prior.type().getName()); + } return (Topic) prior; } log.info(name, "created topic " + created); return created; } + /** + * Topic lookup that skips the cast-safety check for callers whose handler receives + * {@code Object}. Used by the annotation binder, where an Object-typed topic is a + * supported configuration rather than a mistake. + */ + private Topic getOrCreateTopicUnchecked(String topicName, Class type) { + Objects.requireNonNull(topicName, "topicName"); + Objects.requireNonNull(type, "type"); + Topic existing = topics.get(topicName); + if (existing != null) { + if (!existing.acceptsType(type)) { + throw new IllegalArgumentException( + "Topic '" + topicName + "' already exists with type " + + existing.type().getName() + ", cannot re-create as " + + type.getName()); + } + return existing; + } + Topic created = new Topic<>(topicName, type); + Topic prior = topics.putIfAbsent(topicName, created); + if (prior != null) { + if (!prior.acceptsType(type)) { + throw new IllegalArgumentException( + "Topic '" + topicName + "' already exists with type " + + prior.type().getName() + ", cannot re-create as " + + type.getName()); + } + return prior; + } + log.info(name, "created topic " + created); + return created; + } + @Override public Optional> findTopic(String topicName) { return Optional.ofNullable(topics.get(topicName)); @@ -183,7 +234,9 @@ public Optional> findTopic(String topicName) { @SuppressWarnings("unchecked") public Optional> findTopic(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !t.acceptsType(type)) return Optional.empty(); + if (t == null || !t.acceptsType(type) || !t.safelyReturnsAs(type)) { + return Optional.empty(); + } return Optional.of((Topic) t); } @@ -365,7 +418,15 @@ public Subscription subscribe(String topicName, Class type, Consumer handler) { Objects.requireNonNull(handler, "handler"); Topic topic = getOrCreateTopic(topicName, type); - MessageHandler wrapped = msg -> handler.accept(type.cast(msg)); + // Class.cast() returns false for every argument when the class is primitive, so + // casting with the raw type made every delivery throw ClassCastException, which + // dispatchCallback swallowed into a log line -- the subscriber was registered, + // never fired, and nothing else indicated why. Box the type for the runtime check; + // Consumer erases to accept(Object), so no cast on T is needed and + // none would survive erasure anyway. + final Class boxedType = Topic.box(type); + @SuppressWarnings("unchecked") + MessageHandler wrapped = msg -> ((Consumer) handler).accept(boxedType.cast(msg)); SubscriberList list = subscribers.computeIfAbsent(topic, k -> new SubscriberList()); list.add(wrapped); return new Subscription(topic, wrapped, this); @@ -374,7 +435,11 @@ public Subscription subscribe(String topicName, Class type, /** Raw subscribe used by the annotation binder where the parameter type is reflective. */ public Subscription subscribeRaw(String topicName, Class type, java.util.function.Consumer handler) { - Topic topic = getOrCreateTopic(topicName, type); + // Deliberately no safelyReturnsAs check. This path exists for the annotation + // binder, which uses Object as the topic type when handlers of differing + // parameter types share a topic and does its own per-message isInstance filter. + // The handler receives Object, so a wider topic is correct rather than unsafe. + Topic topic = getOrCreateTopicUnchecked(topicName, type); MessageHandler wrapped = handler::accept; SubscriberList list = subscribers.computeIfAbsent(topic, k -> new SubscriberList()); list.add(wrapped); @@ -384,7 +449,14 @@ public Subscription subscribeRaw(String topicName, Class type, void removeSubscription(Subscription sub) { Topic t = sub.topic(); SubscriberList list = subscribers.get(t); - if (list != null) list.remove(sub.handler()); + // markSubscriptionAsHardwareThreaded may have replaced the list entry with a + // wrapper, so look for whichever handler is actually registered. Searching only + // sub.handler() missed the wrapper and left an unregistered @OnHardwareThread + // subscriber running -- it kept receiving publishes, so an unregistered node + // could still drive the hardware. + MessageHandler registered = hardwareRerouted.getOrDefault(sub, sub.handler()); + if (list != null) list.remove(registered); + hardwareRerouted.remove(sub); } /** @@ -407,11 +479,10 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { } }); }; - // Replace in the list. Subscription object's handler() returns the - // original so removeSubscription still works correctly. + // Replace in the list. sub.handler() keeps returning the original so callers see + // a stable identity, and hardwareRerouted records which handler actually sits in + // the list so removeSubscription can find and remove exactly that one. list.replace(original, wrapped); - // Stash the wrapped handler so subscribers still see the original via - // sub.handler() (which is unchanged), but the dispatch uses wrapped. hardwareRerouted.put(sub, wrapped); } @@ -423,7 +494,13 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { @SuppressWarnings("unchecked") public Optional getLatestValue(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !t.acceptsType(type)) return Optional.empty(); + // safelyReturnsAs, not just acceptsType: an Object topic accepts a request for + // Integer, but handing back Optional when the stored value is a String + // moves a ClassCastException to the caller's own site, after isPresent() has + // already told them the value is there. + if (t == null || !t.acceptsType(type) || !t.safelyReturnsAs(type)) { + return Optional.empty(); + } return (Optional) t.latestValue(); } diff --git a/src/main/java/com/aaravlabs/synapse/Topic.java b/src/main/java/com/aaravlabs/synapse/Topic.java index 8c38be3..4033e55 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -217,6 +217,21 @@ boolean acceptsType(Class other) { return boxedType.isAssignableFrom(box(other)); } + /** + * Whether a caller asking for {@code requested} can be handed this topic as a + * {@code Topic} without a ClassCastException later. + * + *

Strictly stronger than {@link #acceptsType}: that only asks whether the + * requested type is compatible with this topic, so an {@code Object} topic accepts a + * request for {@code String}. Handing back {@code Topic} for a topic that may + * hold any {@code Object} is not safe — the cast is erased, so the failure surfaces + * later at the caller's own site, possibly on the hardware thread. Callers that + * receive the topic as {@code Topic} are unaffected. + */ + boolean safelyReturnsAs(Class requested) { + return box(requested).isAssignableFrom(boxedType); + } + /** * Returns the most recent value, falling back to {@code defaultValue} if nothing has * been published yet. Convenience for {@code topic.latestValue().orElse(default)}. @@ -288,7 +303,7 @@ public Optional> latest() { } /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ - private static Class box(Class c) { + static Class box(Class c) { if (!c.isPrimitive()) return c; if (c == int.class) return Integer.class; if (c == long.class) return Long.class; diff --git a/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java b/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java new file mode 100644 index 0000000..7a93c95 --- /dev/null +++ b/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java @@ -0,0 +1,140 @@ +package com.aaravlabs.synapse; + +import com.aaravlabs.synapse.annotation.OnHardwareThread; +import com.aaravlabs.synapse.annotation.SubscribedTo; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Guards three ways a subscription or typed lookup could outlive or misrepresent the + * value it was given. All three were defects found by audit and reproduced here first. + */ +class SubscriptionLifetimeTest { + + private Orchestrator orchestrator; + + @BeforeEach void setUp() { orchestrator = Orchestrator.create("lifetime"); } + @AfterEach void tearDown() { orchestrator.close(); } + + @Test + void subscribeWithAPrimitiveClassStillDelivers() { + // The library treats primitive and wrapper types as interchangeable everywhere -- + // getOrCreateTopic accepts int.class for an Integer topic, and acceptsType boxes + // before comparing. subscribe used to cast with the raw type instead, and + // Class.cast() returns false for every argument when the class is primitive, so + // every delivery threw ClassCastException. dispatchCallback swallowed it into a log + // line: the subscriber was registered, never fired, and nothing said why. + orchestrator.getOrCreateTopic("t", int.class); + List got = Collections.synchronizedList(new ArrayList<>()); + orchestrator.subscribe("t", int.class, got::add); + + orchestrator.publish("t", 1); + orchestrator.publish("t", 2); + await(() -> got.size() >= 2); + + assertEquals(List.of(1, 2), got, + "subscribe with a primitive class must deliver, not silently drop"); + } + + @Test + void unwrappingAnObjectTopicAsANarrowerTypeIsRefused() { + // Object.isAssignableFrom(String) is true, so the types are "compatible" and the + // old check let it through. The caller received a Topic for a topic that + // may hold any Object, and the erased cast failed later at their own call site -- + // after isPresent() had already reported a value. Refused up front instead. + orchestrator.getOrCreateTopic("n", Object.class); + + IllegalArgumentException e = assertThrows(IllegalArgumentException.class, + () -> orchestrator.getOrCreateTopic("n", String.class), + "an Object topic must not be handed back as Topic"); + assertTrue(e.getMessage().contains("Object"), + "the message should name the actual type, was: " + e.getMessage()); + + assertFalse(orchestrator.findTopic("n", String.class).isPresent(), + "findTopic must not hand back a Topic for an Object topic"); + assertTrue(orchestrator.findTopic("n", Object.class).isPresent(), + "asking for the actual type must still work"); + + orchestrator.publish("n", 1.5); + assertFalse(orchestrator.getLatestValue("n", String.class).isPresent(), + "getLatestValue must not report a Double as a String"); + assertEquals(1.5, orchestrator.getLatestValue("n", Object.class).orElse(null), + "asking for the actual type must still return the value"); + + // Narrowing to a type the topic can actually hold is still allowed. + assertNotNull(orchestrator.getOrCreateTopic("n2", String.class)); + orchestrator.publish("n2", "ok"); + assertEquals("ok", orchestrator.getLatestValue("n2", String.class).orElse(null)); + } + + @Test + void unregisterNodeStopsAnOnHardwareThreadSubscriber() throws Exception { + // The hardware-thread wrapper replaces the list entry, but the Subscription still + // reported the unwrapped handler, so removeSubscription searched for something the + // list no longer held and returned without removing anything. An unregistered node + // kept receiving publishes, so it could still drive the hardware -- and the pool is + // never shut down by unregisterNode. + AtomicInteger hits = new AtomicInteger(); + Node n = new Node(orchestrator) { + @SubscribedTo(topic = "hw") + @OnHardwareThread + public void on(String s) { hits.incrementAndGet(); } + }; + orchestrator.registerNode("n", n); + orchestrator.publish("hw", "before"); + await(() -> hits.get() >= 1); + + orchestrator.unregisterNode("n"); + // Give any already-queued work time to drain before counting. + Thread.sleep(200); + int settled = hits.get(); + orchestrator.publish("hw", "after-unregister"); + Thread.sleep(400); + + assertEquals(settled, hits.get(), + "an unregistered node must stop receiving, including on the hardware thread"); + } + + @Test + void unsubscribeStopsAnOnHardwareThreadSubscriber() throws Exception { + AtomicInteger hits = new AtomicInteger(); + Topic t = orchestrator.getOrCreateTopic("hw2", String.class); + Subscription sub = orchestrator.subscribe("hw2", String.class, s -> hits.incrementAndGet()); + // Route it through the same path the binder uses for @OnHardwareThread. + ((OrchestratorImpl) orchestrator).markSubscriptionAsHardwareThreaded(sub); + + orchestrator.publish("hw2", "before"); + await(() -> hits.get() >= 1); + + sub.unsubscribe(); + Thread.sleep(200); + int settled = hits.get(); + orchestrator.publish("hw2", "after-unsub"); + Thread.sleep(400); + + assertEquals(settled, hits.get(), + "unsubscribe must stop a subscriber that was rerouted to the hardware thread"); + assertTrue(t.latestValue().isPresent(), "publish still records the latest value"); + } + + private static void await(java.util.function.BooleanSupplier cond) { + long deadline = System.nanoTime() + 3_000_000_000L; + while (System.nanoTime() < deadline) { + if (cond.getAsBoolean()) return; + try { + Thread.sleep(5); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + } + } +} diff --git a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java index 8b4b1bf..6e362bf 100644 --- a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java +++ b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java @@ -4,6 +4,8 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicLongArray; import java.util.concurrent.atomic.AtomicReference; @@ -86,6 +88,11 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception final int ids = publishers * perPublisher; AtomicReference failure = new AtomicReference<>(); AtomicReference readersRunning = new AtomicReference<>(Boolean.TRUE); + // Counted, and asserted non-zero below. Without this the whole reader loop can be + // skipped forever: make latest() return empty and snap is always null, so no + // assertion inside the loop ever runs and the test still passes on its final-state + // checks. A detector that cannot be observed to have fired is not evidence. + AtomicLong samplesWithValue = new AtomicLong(); AtomicLongArray windowHi = new AtomicLongArray(ids); for (int i = 0; i < ids; i++) { // Long.MIN_VALUE, not -1: the JLS allows any origin for System.nanoTime(), @@ -119,15 +126,15 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception // the two calls, so a nonzero stamp here would be that publish, // not an inconsistency. } else { + samplesWithValue.incrementAndGet(); Integer v = snap.value(); if (v < 0 || v >= ids) { failure.compareAndSet(null, new AssertionError("torn/garbage latest value: " + v)); return; } - // The pair is self-consistent by construction, so the only way - // to fail is a stamp that predates the value it came with -- - // which the publish-window check bounds. + // Only a stamp AFTER the publish ended indicates a torn pair: the + // reader would be holding an older value beside a newer timestamp. long hi = windowHi.get(v); if (hi != Long.MIN_VALUE && snap.publishNanos() > hi) { failure.compareAndSet(null, new AssertionError( @@ -167,8 +174,8 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception for (int i = 0; i < perPublisher; i++) { orchestrator.publish("concurrent", base + i); // Recorded after the publish returns. Recording it before would - // make the detector blind: a reader would essentially never see a - // value whose window had closed. + // make the tear detector blind: a reader would essentially never + // see a value whose window had closed. windowHi.set(base + i, System.nanoTime()); } } catch (Throwable t) { @@ -192,6 +199,10 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception } assertNull(failure.get(), "concurrent access failed: " + failure.get()); + // The reader must actually have observed the topic. If it saw nothing, every + // check above was vacuously skipped and this green run means nothing. + assertTrue(samplesWithValue.get() > 0, + "the reader never observed a value, so the tear checks never ran"); // Concurrent publishers each write a contiguous ascending block, so which // publisher's final write lands last is non-deterministic. The invariant is // that the surviving value is some publisher's last value, not that it is the @@ -213,6 +224,29 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception } } + /* + * NOT a test: a record of why there is no test for the install-ordering guarantee. + * + * recordLatest holds the topic monitor across the clock sample and the store, so + * install order equals timestamp order. Without it, a publisher that samples the clock + * and is then preempted installs an OLDER snapshot after a newer one, and the visible + * stamp moves backwards. That regression was reproduced against the real class by a + * standalone probe, but it is not covered here, and this is the reasoning: + * + * - The window between sampling and storing is nanoseconds wide in the steady state. + * Sampling `now` immediately before the store and asserting install order directly + * would test the test, not the code. + * - Detecting it by observation needs a reader tight enough to land inside that + * window. A spinning reader detects it 3 runs in 6, but starves the publishers on a + * 1-2 core runner and made the suite hang. A yielding reader never detects it at all, + * because it samples orders of magnitude too slowly. + * - So every arrangement either flakes, hangs, or cannot fire. A test that cannot fail + * reliably is worse than no test: it reports coverage it does not have. + * + * The guarantee rests on the monitor in recordLatest plus the reasoning above, not on + * an assertion. If that ever changes, this is the gap. + */ + @Test void acceptsValueClassRejectsIncompatibleValues() { // The per-publish type check must keep rejecting incompatible values, and a @@ -245,11 +279,13 @@ void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() { assertTrue(t.acceptsValueClass(Double.class)); assertTrue(t.acceptsValueClass(double.class), "primitive must match wrapper topic"); assertTrue(t.acceptsType(double.class)); - // acceptsValueClass delegates to acceptsType, so the pair cannot disagree. - for (Class c : new Class[]{Double.class, double.class, Number.class, Object.class}) { - assertEquals(t.acceptsType(c), t.acceptsValueClass(c), - "acceptsValueClass must agree with acceptsType for " + c.getSimpleName()); - } + // Literal expectations, not a comparison of acceptsValueClass against acceptsType: + // acceptsValueClass delegates to acceptsType, so asserting they agree restates the + // delegation rather than testing a truth value. These pin the actual semantics. + assertFalse(t.acceptsValueClass(Number.class), + "a supertype of the topic type must not be accepted as a publish type"); + assertFalse(t.acceptsValueClass(Object.class), + "Object must not be accepted for a Double topic"); Topic i = orchestrator.getOrCreateTopic("i", Integer.class); assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic"); @@ -263,33 +299,48 @@ void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception { Topic t = orchestrator.getOrCreateTopic("snap", String.class); assertFalse(t.latest().isPresent(), "latest() must be empty before the first publish"); + // Clock readings taken OUTSIDE the library, around the publish. These are the only + // independent reference for the stamp: comparing latest().publishNanos() with + // latestPublishNanos() reads the same volatile field twice and holds for any value, + // and an age bracket is offset-invariant -- a Latest stamped with publishNanos + K + // shifts the age and both bounds by K and still passes. A wrong-but-consistent + // timestamp is exactly the staleness failure latest() exists to prevent, so the + // stamp has to be pinned against a clock the library did not touch. + long publishBefore = System.nanoTime(); orchestrator.publish("snap", "a"); + long publishAfter = System.nanoTime(); + Topic.Latest first = t.latest().orElseThrow(); assertEquals("a", first.value()); - assertEquals(first.publishNanos(), t.latestPublishNanos(), - "latest() and latestPublishNanos() must describe the same publish"); assertEquals("a", t.latestValueOr(null)); + // Compare the SNAPSHOT identity, not the stamp. latest() and latestPublishNanos() + // both read the same volatile field, so comparing their stamps holds for any value + // an implementation could produce and would catch nothing. latestPublishNanos() is + // pinned to an externally-clocked window below instead. - // No assertion that publishNanos() is nonzero: System.nanoTime() has an arbitrary - // origin, so a published snapshot may legitimately carry 0. The present Optional - // already proves a publish happened, which is the property worth checking here. + // The stamp must fall inside the window this publish occupied, which no + // self-consistent offset can satisfy. Pinned for both accessors, since neither + // reading is derived from the other. + assertTrue(first.publishNanos() >= publishBefore + && first.publishNanos() <= publishAfter, + "publish stamp " + first.publishNanos() + " lies outside the publish window [" + + publishBefore + "," + publishAfter + "]"); + long recorded = t.latestPublishNanos(); + assertTrue(recorded >= publishBefore && recorded <= publishAfter, + "latestPublishNanos() " + recorded + " must be the stamp of this publish," + + " not a live clock reading, which would be outside [" + + publishBefore + "," + publishAfter + "]"); - // ageNanos must be the age of THIS value, from ITS OWN stamp. - // - // Comparing ageNanos() against `now - publishNanos()` cannot establish that: - // both sides recompute the same subtraction, so the relation holds for any pair - // and a Latest carrying a foreign stamp still passes. What actually pins the - // stamp to this value is the cross-check above -- latest().publishNanos() equals - // latestPublishNanos() equals the stamp of the "a" publish -- combined with a - // non-negative age and a stamp that is not in the future. - // - // Both bounds are relative to the clock, not absolute wall-clock deadlines: a - // snapshot read in the same tick as its own publish legitimately reports age 0, - // and a GC pause between the publish and this read can legitimately make the - // age arbitrarily large. - // The clock reads must bracket the ageNanos() call itself; a reading taken - // outside it cannot bound the result in either direction, because a pause - // between the reading and the call would break the relation. + // No assertion that the stamp is nonzero: System.nanoTime() has an arbitrary + // origin, so a published snapshot may legitimately carry 0. + + // ageNanos must be the age of THIS value, from ITS OWN stamp. The clock reads + // bracket the call itself -- a reading sampled outside it cannot bound the result + // in either direction, because a pause between the reading and the call would + // break the relation. Both bounds are relative to the clock rather than absolute + // wall-clock deadlines: a snapshot read in the same tick as its own publish + // legitimately reports age 0, and a GC pause between the publish and this read can + // legitimately make the age arbitrarily large. long ageBefore = System.nanoTime(); long age = first.ageNanos(); long ageAfter = System.nanoTime(); @@ -316,13 +367,12 @@ void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception { assertEquals("b", second.value()); assertTrue(second.publishNanos() > first.publishNanos(), "a later publish must carry a later stamp"); - // Derive both ages from ONE clock reading rather than calling ageNanos() twice. - // Sampled at different moments the comparison depends on how long each call took, - // so a pause between the publish and the second read could make the newer value - // look older. From a single now, a strictly later stamp is a strictly smaller age. - long laterNow = System.nanoTime(); - assertTrue(laterNow - second.publishNanos() < laterNow - first.publishNanos(), - "the newer value must report a younger age"); + // Actually call ageNanos() on both snapshots. Deriving both ages from one clock reading + // would reduce to a restatement of the strictly-later-stamp assertion above and + // would never exercise ageNanos at all. + assertTrue(second.ageNanos() <= first.ageNanos(), + "the newer value must not report a younger age than the one it replaced: " + + second.ageNanos() + " vs " + first.ageNanos()); } @Test From 302099df88b5833a24a2e26b3d1d28d79edfc4e1 Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Fri, 2 Oct 2026 13:55:01 -0600 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20address=20PR=2032=20review=20?= =?UTF-8?q?=E2=80=94=20reroute/removal=20race,=20publish=20cast-safety,=20?= =?UTF-8?q?test=20calibration?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Main source: - Serialize markSubscriptionAsHardwareThreaded and removeSubscription on a shared hardwareRerouteLock. The wrapper is recorded in two steps (list.replace, then the map put) and a removal landing between them read an empty map, fell back to the original handler the replace had already removed, and no-opped — re-opening the exact hardware-thread leak this PR closes. - SubscriberList.replace now reports whether it replaced anything, and the wrapper is recorded only on success. A second mark on an already-rerouted Subscription no-ops in replace(); recording its wrapper would point removeSubscription at a handler the list never held. - Route publish's lazy auto-create through getOrCreateTopicUnchecked. publish downcasts to Topic immediately and never exposes the topic as a Topic, so the cast-safety rule would only make this path throw: a binder installing an Object-typed topic for a zero-argument @SubscribedTo handler between topics.get and the create made String.isAssignableFrom(Object) == false, dropping a publish that had always succeeded. - Collapse getOrCreateTopicUnchecked into getOrCreateTopic(name, type, castSafe) with a shared validateExisting. The duplicated body carried identical error strings, race branches, and the creation log, and was free to drift. - Correct the safelyReturnsAs javadoc: it checks the opposite assignability direction from acceptsType, so the two are complementary rather than one being stronger. - Correct the Class.cast comment: it throws for any non-null argument when the class is primitive (it is an isInstance check, and isInstance is false for wrappers), rather than "returning false". Tests: - await() fails on expiry instead of returning silently, and takes a named condition. A silent timeout let the removal tests compare zero hits against zero hits and pass with the callback never having run. - Add a control subscriber to both removal tests, never unregistered. Its arrival proves the post-removal publish was dispatched and the hardware thread was live, so "nothing arrived" means the subject was removed rather than the bus having gone quiet. Replaces the Thread.sleep(200)/sleep(400) scheduling assumptions, and asserts the positive control that a pre-removal delivery actually ran. - Assert set membership rather than delivery order in the primitive-subscribe test; publish dispatches to a four-worker pool, so ordering was never guaranteed. - Bound the ageNanos comparison by the measured span between the two clock samples, so independent sampling cannot produce a spurious ordering failure, and fix the inverted failure message. - Add coverage for publishing to a binder-created Object-typed topic, and state in the test what it does and does not pin (the defect itself needed a race to reach). 91 tests, 0 failures across 6 full-suite runs plus a 2-core pinned run. javadoc: 0 errors, no new warnings. --- CHANGELOG.md | 3 +- .../aaravlabs/synapse/OrchestratorImpl.java | 170 ++++++++++-------- .../java/com/aaravlabs/synapse/Topic.java | 20 ++- .../synapse/SubscriptionLifetimeTest.java | 99 ++++++++-- .../aaravlabs/synapse/TopicLockFreeTest.java | 19 +- 5 files changed, 214 insertions(+), 97 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d91f3b8..d6fb568 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -146,7 +146,8 @@ listeners registered, `publish` costs a single volatile read. - **Typed lookups no longer hand back a topic that cannot be cast to the requested type.** An `Object` topic accepts a request for `String`, and the unchecked cast then failed later at the caller's own line — after `isPresent()` had reported a value. The - request is now refused with a message naming the actual type. Narrowing to a type the + request is now refused. `getOrCreateTopic` throws with a message naming the actual + type; `findTopic` and `getLatestValue` return empty. Narrowing to a type the topic really holds still works, and `subscribeRaw` (used by the annotation binder) still accepts an `Object` topic. diff --git a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java index 6a60bbc..b2b5b9f 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -145,84 +145,70 @@ private ThreadFactory makeThreadFactory(String prefix) { @Override @SuppressWarnings("unchecked") public Topic getOrCreateTopic(String topicName, Class type) { + return (Topic) getOrCreateTopic(topicName, type, true); + } + + /** + * Topic lookup that skips the cast-safety check for callers whose handler receives + * {@code Object}. Used by the annotation binder and by {@link #publish}, where an + * Object-typed topic is a supported configuration rather than a mistake, and where + * the caller never receives the topic as a {@code Topic} it could cast wrongly. + */ + private Topic getOrCreateTopicUnchecked(String topicName, Class type) { + return getOrCreateTopic(topicName, type, false); + } + + /** + * The single topic-creation path. {@code castSafe} selects whether an existing topic + * that accepts {@code type} may still be handed back. + * + *

One method rather than two: the {@code exists} / {@code putIfAbsent} race + * branches, the error strings, and the creation log have to stay identical, and a + * duplicated body is where the two copies silently drift. + */ + @SuppressWarnings("unchecked") + private Topic getOrCreateTopic(String topicName, Class type, boolean castSafe) { Objects.requireNonNull(topicName, "topicName"); Objects.requireNonNull(type, "type"); Topic existing = topics.get(topicName); if (existing != null) { - if (!existing.acceptsType(type)) { - throw new IllegalArgumentException( - "Topic '" + topicName + "' already exists with type " - + existing.type().getName() + ", cannot re-create as " - + type.getName()); - } - if (!existing.safelyReturnsAs(type)) { - // compatible, but this topic may hold values the caller cannot cast to T. - // Returning it anyway defers a ClassCastException to an unrelated line. - throw new IllegalArgumentException( - "Topic '" + topicName + "' holds " - + existing.type().getSimpleName() + ", which cannot be " - + "returned as " + type.getName() - + "; read it via findTopic(name) and treat the value as " - + existing.type().getName()); - } - return (Topic) existing; + return (Topic) validateExisting(topicName, type, existing, castSafe); } Topic created = new Topic<>(topicName, type); Topic prior = topics.putIfAbsent(topicName, created); if (prior != null) { - if (!prior.acceptsType(type)) { - throw new IllegalArgumentException( - "Topic '" + topicName + "' already exists with type " - + prior.type().getName() + ", cannot re-create as " - + type.getName()); - } - if (!prior.safelyReturnsAs(type)) { - throw new IllegalArgumentException( - "Topic '" + topicName + "' holds " - + prior.type().getSimpleName() + ", which cannot be " - + "returned as " + type.getName() - + "; read it via findTopic(name) and treat the value as " - + prior.type().getName()); - } - return (Topic) prior; + return (Topic) validateExisting(topicName, type, prior, castSafe); } log.info(name, "created topic " + created); return created; } /** - * Topic lookup that skips the cast-safety check for callers whose handler receives - * {@code Object}. Used by the annotation binder, where an Object-typed topic is a - * supported configuration rather than a mistake. + * Check an already-registered topic against a requested type. Shared by the + * {@code topics.get} and {@code putIfAbsent} branches of {@link #getOrCreateTopic}, + * which must agree on the error text. */ - private Topic getOrCreateTopicUnchecked(String topicName, Class type) { - Objects.requireNonNull(topicName, "topicName"); - Objects.requireNonNull(type, "type"); - Topic existing = topics.get(topicName); - if (existing != null) { - if (!existing.acceptsType(type)) { - throw new IllegalArgumentException( - "Topic '" + topicName + "' already exists with type " - + existing.type().getName() + ", cannot re-create as " - + type.getName()); - } - return existing; + private Topic validateExisting(String topicName, Class type, + Topic existing, boolean castSafe) { + if (!existing.acceptsType(type)) { + throw new IllegalArgumentException( + "Topic '" + topicName + "' already exists with type " + + existing.type().getName() + ", cannot re-create as " + + type.getName()); } - Topic created = new Topic<>(topicName, type); - Topic prior = topics.putIfAbsent(topicName, created); - if (prior != null) { - if (!prior.acceptsType(type)) { - throw new IllegalArgumentException( - "Topic '" + topicName + "' already exists with type " - + prior.type().getName() + ", cannot re-create as " - + type.getName()); - } - return prior; + if (castSafe && !existing.safelyReturnsAs(type)) { + // compatible, but this topic may hold values the caller cannot cast to T. + // Returning it anyway defers a ClassCastException to an unrelated line. + throw new IllegalArgumentException( + "Topic '" + topicName + "' holds " + + existing.type().getSimpleName() + ", which cannot be " + + "returned as " + type.getName() + + "; read it via findTopic(name) and treat the value as " + + existing.type().getName()); } - log.info(name, "created topic " + created); - return created; + return existing; } @Override @@ -382,9 +368,16 @@ public void publish(String topicName, T value) { // Lazily create the topic from the value's runtime type. This matches // Heron's behavior: publishers don't have to pre-register topics. + // Unchecked on purpose: publish immediately downcasts to Topic and never + // exposes the topic to a caller as a Topic, so the cast-safety rule that + // protects typed lookups would only make this path throw. Concretely, a binder + // registering a zero-argument @SubscribedTo handler installs an Object-typed + // topic; if that lands between topics.get and the create below, a cast-safe + // getOrCreateTopic would refuse String.isAssignableFrom(Object) == false and + // silently drop a publish that used to succeed. Topic topic = topics.get(topicName); if (topic == null) { - topic = getOrCreateTopic(topicName, value.getClass()); + topic = getOrCreateTopicUnchecked(topicName, value.getClass()); } else if (!topic.acceptsValueClass(value.getClass())) { throw new IllegalArgumentException( "Topic '" + topicName + "' is typed " + topic.type().getName() @@ -418,8 +411,9 @@ public Subscription subscribe(String topicName, Class type, Consumer handler) { Objects.requireNonNull(handler, "handler"); Topic topic = getOrCreateTopic(topicName, type); - // Class.cast() returns false for every argument when the class is primitive, so - // casting with the raw type made every delivery throw ClassCastException, which + // Class.cast() throws for any non-null argument when the class is primitive -- it + // performs an isInstance check, and isInstance is false for every wrapper value -- + // so casting with the raw type made every delivery throw ClassCastException, which // dispatchCallback swallowed into a log line -- the subscriber was registered, // never fired, and nothing else indicated why. Box the type for the runtime check; // Consumer erases to accept(Object), so no cast on T is needed and @@ -454,9 +448,17 @@ void removeSubscription(Subscription sub) { // sub.handler() missed the wrapper and left an unregistered @OnHardwareThread // subscriber running -- it kept receiving publishes, so an unregistered node // could still drive the hardware. - MessageHandler registered = hardwareRerouted.getOrDefault(sub, sub.handler()); - if (list != null) list.remove(registered); - hardwareRerouted.remove(sub); + // + // The lookup and the removal share hardwareRerouteLock with the mark, because + // the wrapper is recorded in two steps (list.replace, then the map put) and a + // removal landing between them would read an empty map, fall back to the + // original handler the replace had already removed, and no-op -- re-registering + // the exact leak this lookup exists to close. + synchronized (hardwareRerouteLock) { + MessageHandler registered = hardwareRerouted.getOrDefault(sub, sub.handler()); + if (list != null) list.remove(registered); + hardwareRerouted.remove(sub); + } } /** @@ -482,10 +484,31 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { // Replace in the list. sub.handler() keeps returning the original so callers see // a stable identity, and hardwareRerouted records which handler actually sits in // the list so removeSubscription can find and remove exactly that one. - list.replace(original, wrapped); - hardwareRerouted.put(sub, wrapped); + // + // Record only when the replacement actually happened. The method is public and + // the Subscription stays reachable afterwards, so a second mark is reachable: + // replace() looks for sub.handler(), which the first mark already took out of the + // list, so it no-ops. Recording the second wrapper anyway would point + // removeSubscription at a handler the list never held, leaving the first wrapper + // registered forever. + // + // Same lock as removeSubscription, so a concurrent removal cannot observe the + // list with the wrapper installed but the map still holding the original. + synchronized (hardwareRerouteLock) { + if (list.replace(original, wrapped)) { + hardwareRerouted.put(sub, wrapped); + } + } } + /** + * Guards the two-step hardware reroute bookkeeping ({@link SubscriberList#replace} + * plus the {@link #hardwareRerouted} entry) against {@link #removeSubscription}. + * Rare by construction -- registration and removal both happen off the publish path + * -- and held only for the duration of that bookkeeping, never across dispatch. + */ + private final Object hardwareRerouteLock = new Object(); + private final java.util.Map hardwareRerouted = new ConcurrentHashMap<>(); // ---- fetch latest ---------------------------------------------------- @@ -764,17 +787,24 @@ void remove(MessageHandler h) { snapshot = next; } } - /** Atomically replace {@code old} with {@code next}. No-op if not found. */ - void replace(MessageHandler old, MessageHandler next) { + /** + * Atomically replace {@code old} with {@code next}. + * + * @return true if {@code old} was present and is now {@code next}; false when it + * was absent, in which case nothing changed. Callers use the result to + * avoid recording a wrapper the list never accepted. + */ + boolean replace(MessageHandler old, MessageHandler next) { synchronized (lock) { MessageHandler[] cur = snapshot; int idx = -1; for (int i = 0; i < cur.length; i++) if (cur[i] == old) { idx = i; break; } - if (idx < 0) return; + if (idx < 0) return false; MessageHandler[] out = new MessageHandler[cur.length]; System.arraycopy(cur, 0, out, 0, cur.length); out[idx] = next; snapshot = out; + return true; } } MessageHandler[] snapshot() { return snapshot; } diff --git a/src/main/java/com/aaravlabs/synapse/Topic.java b/src/main/java/com/aaravlabs/synapse/Topic.java index 4033e55..9a03942 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -221,12 +221,20 @@ boolean acceptsType(Class other) { * Whether a caller asking for {@code requested} can be handed this topic as a * {@code Topic} without a ClassCastException later. * - *

Strictly stronger than {@link #acceptsType}: that only asks whether the - * requested type is compatible with this topic, so an {@code Object} topic accepts a - * request for {@code String}. Handing back {@code Topic} for a topic that may - * hold any {@code Object} is not safe — the cast is erased, so the failure surfaces - * later at the caller's own site, possibly on the hardware thread. Callers that - * receive the topic as {@code Topic} are unaffected. + *

Checks the opposite assignability direction from {@link #acceptsType}, so + * the two are complementary rather than one being stronger: {@code acceptsType} asks + * whether the topic can carry {@code requested} values + * ({@code topicType.isAssignableFrom(requested)}), while this asks whether this + * topic's values can be handed to a caller holding a {@code requested} + * ({@code requested.isAssignableFrom(topicType)}). Neither implies the other, which is + * why callers that hand back a {@code Topic} need both. + * + *

The gap this closes: {@code Object} topic, request for {@code String} — {@code + * acceptsType} passes it, so an {@code Object} topic accepts a request for {@code + * String}. Handing back {@code Topic} for a topic that may hold any {@code + * Object} is not safe — the cast is erased, so the failure surfaces later at the + * caller's own site, possibly on the hardware thread. Callers that receive the topic as + * {@code Topic} are unaffected. */ boolean safelyReturnsAs(Class requested) { return box(requested).isAssignableFrom(boxedType); diff --git a/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java b/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java index 7a93c95..828f731 100644 --- a/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java +++ b/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java @@ -9,7 +9,9 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.IntSupplier; import static org.junit.jupiter.api.Assertions.*; @@ -29,7 +31,7 @@ void subscribeWithAPrimitiveClassStillDelivers() { // The library treats primitive and wrapper types as interchangeable everywhere -- // getOrCreateTopic accepts int.class for an Integer topic, and acceptsType boxes // before comparing. subscribe used to cast with the raw type instead, and - // Class.cast() returns false for every argument when the class is primitive, so + // Class.cast() throws for any non-null argument when the class is primitive, so // every delivery threw ClassCastException. dispatchCallback swallowed it into a log // line: the subscriber was registered, never fired, and nothing said why. orchestrator.getOrCreateTopic("t", int.class); @@ -38,9 +40,12 @@ void subscribeWithAPrimitiveClassStillDelivers() { orchestrator.publish("t", 1); orchestrator.publish("t", 2); - await(() -> got.size() >= 2); + await(() -> got.size(), "both published values must be delivered"); - assertEquals(List.of(1, 2), got, + // Membership, not order: publish dispatches each callback to a four-worker pool, + // so the second task can run before the first. TopicTest documents the same. + assertEquals(2, got.size(), "both published values must be delivered exactly once"); + assertEquals(Set.of(1, 2), Set.copyOf(got), "subscribe with a primitive class must deliver, not silently drop"); } @@ -76,7 +81,41 @@ void unwrappingAnObjectTopicAsANarrowerTypeIsRefused() { } @Test - void unregisterNodeStopsAnOnHardwareThreadSubscriber() throws Exception { + void publishToABinderCreatedObjectTopicStillSucceeds() { + // The cast-safety rule protects typed lookups, but publish is not one: it + // downcasts to Topic immediately and never exposes the topic to a caller. + // A zero-argument @SubscribedTo handler makes the binder register an Object-typed + // topic, so enforcing the rule on publish's lazy auto-create would reject String + // against that topic and drop a publish that has always worked. + // + // Scope of this test: it pins the observable behaviour — publish to an + // Object-typed topic delivers and records. The specific defect was on the lazy + // auto-create branch, reachable only when a binder installs the Object topic + // between publish's topics.get and its create. This test reaches the topic via + // the map-hit branch instead, so it does not by itself detect a revert of that + // one line; it guards the behaviour the revert would break, and the line itself is + // covered by review rather than by a test that would have to win a race. + AtomicInteger hits = new AtomicInteger(); + Node n = new Node(orchestrator) { + @SubscribedTo(topic = "any") + @OnHardwareThread + public void on() { hits.incrementAndGet(); } + }; + orchestrator.registerNode("n", n); + assertEquals(Object.class, orchestrator.findTopic("any").orElseThrow().type(), + "a zero-argument handler must give the topic the Object type"); + + // No assertDoesNotThrow: the publish must simply deliver, and the await below is + // the assertion. A thrown IllegalArgumentException fails the test on its own. + orchestrator.publish("any", "a string on an Object topic"); + await(hits::get, "the Object-typed topic must still deliver to its subscriber"); + assertEquals("a string on an Object topic", + orchestrator.getLatestValue("any", Object.class).orElse(null), + "publish must still record the latest value"); + } + + @Test + void unregisterNodeStopsAnOnHardwareThreadSubscriber() { // The hardware-thread wrapper replaces the list entry, but the Subscription still // reported the unwrapped handler, so removeSubscription searched for something the // list no longer held and returned without removing anything. An unregistered node @@ -88,53 +127,83 @@ void unregisterNodeStopsAnOnHardwareThreadSubscriber() throws Exception { @OnHardwareThread public void on(String s) { hits.incrementAndGet(); } }; + // Control subscriber, never unregistered. It is what makes the negative assertion + // below falsifiable: its arrival proves the post-unregister publish was dispatched + // and the hardware thread was live, so a silent subject means the subject was + // removed rather than the bus having gone quiet. Without it, "nothing arrived" is + // also what a starved pool looks like, and the test would pass either way. + AtomicInteger control = new AtomicInteger(); + Node ctl = new Node(orchestrator) { + @SubscribedTo(topic = "hw") + @OnHardwareThread + public void on(String s) { control.incrementAndGet(); } + }; + orchestrator.registerNode("n", n); + orchestrator.registerNode("ctl", ctl); orchestrator.publish("hw", "before"); - await(() -> hits.get() >= 1); + // Positive control, asserted by await: the callback really ran before removal. + await(hits::get, "the hardware-thread callback must run before removal"); + await(control::get, "the control callback must run before removal"); + assertEquals(1, hits.get(), "exactly the one pre-unregister publish should have been seen"); orchestrator.unregisterNode("n"); - // Give any already-queued work time to drain before counting. - Thread.sleep(200); int settled = hits.get(); orchestrator.publish("hw", "after-unregister"); - Thread.sleep(400); + await(control::get, + "the control must observe the post-unregister publish, proving dispatch ran"); assertEquals(settled, hits.get(), "an unregistered node must stop receiving, including on the hardware thread"); } @Test - void unsubscribeStopsAnOnHardwareThreadSubscriber() throws Exception { + void unsubscribeStopsAnOnHardwareThreadSubscriber() { AtomicInteger hits = new AtomicInteger(); Topic t = orchestrator.getOrCreateTopic("hw2", String.class); Subscription sub = orchestrator.subscribe("hw2", String.class, s -> hits.incrementAndGet()); + // Control subscriber that is never unsubscribed, for the reason given above. + AtomicInteger control = new AtomicInteger(); + orchestrator.subscribe("hw2", String.class, s -> control.incrementAndGet()); // Route it through the same path the binder uses for @OnHardwareThread. ((OrchestratorImpl) orchestrator).markSubscriptionAsHardwareThreaded(sub); orchestrator.publish("hw2", "before"); - await(() -> hits.get() >= 1); + await(hits::get, "the hardware-thread callback must run before removal"); + await(control::get, "the control callback must run before removal"); sub.unsubscribe(); - Thread.sleep(200); int settled = hits.get(); orchestrator.publish("hw2", "after-unsub"); - Thread.sleep(400); + await(control::get, + "the control must observe the post-unsubscribe publish, proving dispatch ran"); assertEquals(settled, hits.get(), "unsubscribe must stop a subscriber that was rerouted to the hardware thread"); assertTrue(t.latestValue().isPresent(), "publish still records the latest value"); } - private static void await(java.util.function.BooleanSupplier cond) { + /** + * Poll {@code cond} until it is true, failing the test if it never becomes true. + * + *

Failing on expiry is the point: a silent timeout lets a test that only counts + * deliveries compare zero against zero and pass without the callback ever running. + * + * @param cond the condition to await; it must count real observations so the caller + * can assert on them afterwards + * @param message what the awaited condition was, for the failure output + */ + private static void await(java.util.function.IntSupplier cond, String message) { long deadline = System.nanoTime() + 3_000_000_000L; while (System.nanoTime() < deadline) { - if (cond.getAsBoolean()) return; + if (cond.getAsInt() > 0) return; try { Thread.sleep(5); } catch (InterruptedException e) { Thread.currentThread().interrupt(); - return; + fail("interrupted while waiting for: " + message, e); } } + fail("condition never became true within 3s: " + message); } } diff --git a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java index 6e362bf..94f9e0d 100644 --- a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java +++ b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java @@ -367,12 +367,21 @@ void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception { assertEquals("b", second.value()); assertTrue(second.publishNanos() > first.publishNanos(), "a later publish must carry a later stamp"); - // Actually call ageNanos() on both snapshots. Deriving both ages from one clock reading + // Actually call ageNanos() on both snapshots -- deriving both from one clock reading // would reduce to a restatement of the strictly-later-stamp assertion above and - // would never exercise ageNanos at all. - assertTrue(second.ageNanos() <= first.ageNanos(), - "the newer value must not report a younger age than the one it replaced: " - + second.ageNanos() + " vs " + first.ageNanos()); + // would never exercise ageNanos at all. The two calls sample the clock separately, + // so the elapsed span between them is measured and allowed for: that span is the + // sampling skew, and only skew of that size could let the newer value report the + // older age. Clock readings either side keep the allowance honest rather than + // granting an unbounded tolerance. + long skewBefore = System.nanoTime(); + long secondAge = second.ageNanos(); + long firstAge = first.ageNanos(); + long skewAfter = System.nanoTime(); + long maxSkew = skewAfter - skewBefore; + assertTrue(secondAge <= firstAge + maxSkew, + "the newer value must not report an older age than the one it replaced: " + + secondAge + " vs " + firstAge + " (clock skew allowance " + maxSkew + ")"); } @Test