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..d6fb568 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -131,6 +131,26 @@ 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. `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. + ## [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..b2b5b9f 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -145,35 +145,72 @@ 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()); - } - 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()); - } - return (Topic) prior; + return (Topic) validateExisting(topicName, type, prior, castSafe); } log.info(name, "created topic " + created); return created; } + /** + * 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 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()); + } + 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()); + } + return existing; + } + @Override public Optional> findTopic(String topicName) { return Optional.ofNullable(topics.get(topicName)); @@ -183,7 +220,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); } @@ -329,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() @@ -365,7 +411,16 @@ 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() 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 + // 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 +429,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 +443,22 @@ 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. + // + // 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); + } } /** @@ -407,14 +481,34 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { } }); }; - // Replace in the list. Subscription object's handler() returns the - // original so removeSubscription still works correctly. - 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); + // 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. + // + // 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 ---------------------------------------------------- @@ -423,7 +517,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(); } @@ -687,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 8c38be3..9a03942 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -217,6 +217,29 @@ 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. + * + *

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); + } + /** * 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 +311,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..828f731 --- /dev/null +++ b/src/test/java/com/aaravlabs/synapse/SubscriptionLifetimeTest.java @@ -0,0 +1,209 @@ +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.Set; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.IntSupplier; + +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() 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); + List got = Collections.synchronizedList(new ArrayList<>()); + orchestrator.subscribe("t", int.class, got::add); + + orchestrator.publish("t", 1); + orchestrator.publish("t", 2); + await(() -> got.size(), "both published values must be delivered"); + + // 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"); + } + + @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 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 + // 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(); } + }; + // 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"); + // 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"); + int settled = hits.get(); + orchestrator.publish("hw", "after-unregister"); + 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() { + 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, "the hardware-thread callback must run before removal"); + await(control::get, "the control callback must run before removal"); + + sub.unsubscribe(); + int settled = hits.get(); + orchestrator.publish("hw2", "after-unsub"); + 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"); + } + + /** + * 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.getAsInt() > 0) return; + try { + Thread.sleep(5); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + 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 8b4b1bf..94f9e0d 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,21 @@ 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 from one clock reading + // would reduce to a restatement of the strictly-later-stamp assertion above and + // 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