From 54d881e62c8ae9fd6a531f8c8a91c664a27b62b8 Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Tue, 29 Sep 2026 21:20:48 -0600 Subject: [PATCH 1/4] perf(topic): lock-free latest-value tracking `latestValue` and `latestPublishNanos` were each `synchronized` on the Topic instance. That monitor is the most-used read on the bus, so it serialised every reader against every writer. The value and its timestamp are now published together as one immutable `Latest` pair behind a single volatile field, which makes both accessors plain volatile reads. The write path keeps a monitor, but only around the `nanoTime()` sample and the store. Two volatile fields would not do: `publish` is reachable from the OpMode loop, the hardware thread and the callback pool, so two publishers can interleave between a value write and a timestamp write and leave the pair describing two different publishes -- a regression against the old monitor, which ordered those writes. Sampling the clock outside mutual exclusion is not sufficient either: a publisher that samples and is then preempted installs an OLDER snapshot after a newer one, so a reader can hold an older value beside a newer timestamp and an age computed from that stamp is under-reported. Holding the monitor across the sample and the store keeps install order equal to timestamp order, which is the property the old body had. The type work is hoisted out of the publish path: - `boxedType` is computed once in the constructor instead of re-normalising the declared type on every publish and every lookup. `Topic.type()` still reports the type the topic was created with; only the internal comparison field is boxed. - `acceptsType(Class)` is the old `boxed(t.type()).isAssignableFrom(boxed(other))` with the topic side hoisted, replacing that expression at five call sites. `acceptsValueClass(Class)` delegates to it, so the two cannot drift. `boxed()` moves from `OrchestratorImpl` to `Topic`, next to the field it reads. Two separate accessors can still straddle a publish, so the read order remains: sample the timestamp first. That can only over-report the value's age, never under-report it. This is documented on the public accessors, in the CHANGELOG and in the topics guide, and the guide's example had the unsafe order. No API or behaviour change; `latestPublishNanos()` still returns 0 before the first publish. Measured on the repo benchmark (`gradle :benchmarks:run --args='run --quick --scenarios S0 --styles synapse'`, 8 alternating runs of baseline and branch on one box, p50 ns/op, median): micro.topic.latestValue 44 ns -> 28 ns (-36%) micro.topic.recordLatest.1p1c 360 ns -> 192 ns (-47%) micro.topic.recordLatest.4p4c 340 ns -> 268 ns (-21%) micro.topic.recordLatest 89 ns -> 83 ns (-7%) micro.publish.subscribers0 140 ns -> 161 ns (+15%) The read-path win is the durable one. The contended publish numbers are much smaller than a fully lock-free write path would give, because the write path keeps its monitor; that is the cost of not regressing the timestamp ordering the old code had. `publish.subscribers0` and the `rawDirectCall` control both moved by more than this box's noise floor in the same runs, so read those two as inconclusive rather than as a regression. Treat every number as a ratio on the box that produced it. --- CHANGELOG.md | 35 +++ .../aaravlabs/synapse/OrchestratorImpl.java | 31 +- .../java/com/aaravlabs/synapse/Topic.java | 128 ++++++++- .../aaravlabs/synapse/TopicLockFreeTest.java | 266 ++++++++++++++++++ website/src/content/docs/concepts/topics.mdx | 12 +- 5 files changed, 433 insertions(+), 39 deletions(-) create mode 100644 src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 7cded3b..ef840b9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,41 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **`@SubscribedTo` dispatch no longer re-boxes the parameter type per message.** The primitive-to-wrapper normalisation is computed once at bind time instead of on every delivered message. Equivalent to the previous conditional check. +- **`Topic` latest-value reads no longer synchronize.** `latestValue()` and + `latestPublishNanos()` are plain volatile reads. The latest value and its + timestamp are published together as one immutable pair through a single volatile + field, so a reader always sees a value and a timestamp from the same publish. + `publish` is reachable from the OpMode loop, the hardware thread, and the callback + pool, so two publishers really can overlap. + + The write path still takes the topic monitor, but only around the clock sample and + the store. Two separate volatile fields would let two publishers interleave + between the value write and the timestamp write; sampling the clock outside + mutual exclusion would let a preempted publisher install an older pair after a + newer one. Both regressions were reproduced and fixed; keeping the monitor on the + write path preserves the ordering the previous `synchronized` body provided. + + The topic's declared type is now normalized to its wrapper class **once**, at + construction, for the per-publish type check. `type()` still reports the type the + topic was created with — only the internal comparison field is boxed. + + `latestPublishNanos()` still returns `0` before the first publish, unchanged. + + **Read the timestamp before the value** when you use the two together as a + staleness check. Each accessor now returns a self-consistent pair, but two separate + calls can still straddle a publish, so the ordering rule remains — it is now + documented on the public accessors and in the topics guide: + + ```java + long stamp = topic.latestPublishNanos(); // first + T v = topic.latestValueOr(null); // then + long age = System.nanoTime() - stamp; + ``` + + Reading the value first and the timestamp second can pair an older value with a + newer timestamp, so an age check on that pair passes even though the value is + stale. Reading the timestamp first can only over-report the age, never under-report + it. ## [0.4.0] - 2026-09-12 diff --git a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java index 6aa2f54..695e704 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -147,11 +147,9 @@ public Topic getOrCreateTopic(String topicName, Class type) { Objects.requireNonNull(topicName, "topicName"); Objects.requireNonNull(type, "type"); - Class normalized = boxed(type); - Topic existing = topics.get(topicName); if (existing != null) { - if (!boxed(existing.type()).isAssignableFrom(normalized)) { + if (!existing.acceptsType(type)) { throw new IllegalArgumentException( "Topic '" + topicName + "' already exists with type " + existing.type().getName() + ", cannot re-create as " @@ -163,7 +161,7 @@ public Topic getOrCreateTopic(String topicName, Class type) { Topic created = new Topic<>(topicName, type); Topic prior = topics.putIfAbsent(topicName, created); if (prior != null) { - if (!boxed(prior.type()).isAssignableFrom(normalized)) { + if (!prior.acceptsType(type)) { throw new IllegalArgumentException( "Topic '" + topicName + "' already exists with type " + prior.type().getName() + ", cannot re-create as " @@ -175,20 +173,6 @@ public Topic getOrCreateTopic(String topicName, Class type) { return created; } - /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ - private static Class boxed(Class c) { - if (!c.isPrimitive()) return c; - if (c == int.class) return Integer.class; - if (c == long.class) return Long.class; - if (c == double.class) return Double.class; - if (c == float.class) return Float.class; - if (c == boolean.class) return Boolean.class; - if (c == byte.class) return Byte.class; - if (c == short.class) return Short.class; - if (c == char.class) return Character.class; - return c; - } - @Override public Optional> findTopic(String topicName) { return Optional.ofNullable(topics.get(topicName)); @@ -198,7 +182,7 @@ public Optional> findTopic(String topicName) { @SuppressWarnings("unchecked") public Optional> findTopic(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty(); + if (t == null || !t.acceptsType(type)) return Optional.empty(); return Optional.of((Topic) t); } @@ -217,14 +201,13 @@ 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. - Class valueType = value.getClass(); Topic topic = topics.get(topicName); if (topic == null) { - topic = getOrCreateTopic(topicName, valueType); - } else if (!boxed(topic.type()).isAssignableFrom(boxed(valueType))) { + topic = getOrCreateTopic(topicName, value.getClass()); + } else if (!topic.acceptsValueClass(value.getClass())) { throw new IllegalArgumentException( "Topic '" + topicName + "' is typed " + topic.type().getName() - + " but publish got " + valueType.getName()); + + " but publish got " + value.getClass().getName()); } ((Topic) topic).recordLatest(value); @@ -312,7 +295,7 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { @SuppressWarnings("unchecked") public Optional getLatestValue(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty(); + if (t == null || !t.acceptsType(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 5b2302e..68e88fb 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -18,13 +18,43 @@ public final class Topic { private final String name; private final Class type; - // Guarded by `this` for write, volatile for the read in publish(). - private volatile T latest; - private volatile long latestPublishNanos; + /** + * {@code type} normalized to its wrapper class. Cached at construction so the + * per-publish type check is a single field read instead of a primitive-boxing + * chain. + */ + private final Class boxedType; + + /** + * The most recently published value and its timestamp, as one immutable pair. + * + *

A single volatile reference rather than two volatile fields: with two, two + * concurrent publishers can interleave between the value write and the timestamp + * write and leave the pair describing two different publishes. Since + * {@link Orchestrator#publish(String, Object)} is reachable from the OpMode + * loop, the hardware thread, and the callback pool, that is not a theoretical + * interleaving. + * + *

Immutable so that a reader holding a reference sees a self-consistent pair + * with no further synchronization. + */ + private volatile Latest latest; Topic(String name, Class type) { this.name = name; this.type = type; + this.boxedType = box(type); + } + + /** One publish's value and the {@link System#nanoTime()} at which it was recorded. */ + private static final class Latest { + final T value; + final long publishNanos; + + Latest(T value, long publishNanos) { + this.value = value; + this.publishNanos = publishNanos; + } } /** @@ -44,42 +74,112 @@ public Class type() { /** * The most recently published value, or empty if nothing has been published yet. * + *

Reading a value together with its age: each accessor reads a + * self-consistent pair, but two separate calls can still straddle a publish. Sample + * the timestamp first, which can only over-report the value's age: + * + *

{@code
+     * long stamp = topic.latestPublishNanos();  // first; 0 means nothing published yet
+     * T v = topic.latestValueOr(null);         // then
+     * long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
+     * }
+ * + *

The {@code 0} guard matters: before the first publish + * {@link #latestPublishNanos()} is {@code 0}, and subtracting it would yield the raw + * {@link System#nanoTime()} reading — seconds to days — rather than an age. + * + * Reading the value first is the unsafe order: a publish landing between the two + * calls pairs the older value with the newer timestamp, so an age check on that + * pair passes even though the value is stale. + * * @return an {@link Optional} holding the latest value */ - public synchronized Optional latestValue() { - return Optional.ofNullable(latest); + public Optional latestValue() { + return Optional.ofNullable(latestValueOr(null)); } /** - * Wall-clock nanos at which {@link #latestValue()} was last updated. + * Nanos at which {@link #latestValue()} was last updated, on the + * {@link System#nanoTime()} monotonic clock (arbitrary origin — only differences + * are meaningful). + * + *

Sample this before {@link #latestValue()} when the two are used together as a + * staleness check — see {@link #latestValue()} for the ordering rule. * - * @return publish timestamp in nanoseconds ({@code System.nanoTime()} clock) + * @return publish timestamp in nanoseconds, or 0 if nothing has been published */ - public synchronized long latestPublishNanos() { - return latestPublishNanos; + public long latestPublishNanos() { + Latest snap = latest; + return snap == null ? 0L : snap.publishNanos; } /** * Record a new latest value. Called by the orchestrator immediately before notifying * subscribers. * + *

The read path is lock-free: the value and its timestamp are published + * together as one immutable {@link Latest} through a single volatile reference, so a + * concurrent reader always sees a value and a timestamp from the same publish. + * Splitting them across two volatile fields would let two concurrent publishers + * interleave between the writes and leave the pair describing two different + * publishes. The write path takes the topic monitor, around the clock sample + * and the store — see below. + * * @param value the value to record */ - synchronized void recordLatest(T value) { - this.latest = value; - this.latestPublishNanos = System.nanoTime(); + void recordLatest(T value) { + // The snapshot is sampled and installed under this topic's monitor so that + // install order matches timestamp order. Without it, a publisher that samples + // nanoTime() and is then preempted installs an OLDER snapshot after a newer + // one, and a reader can end up holding an older value beside a newer + // timestamp -- an age computed from that stamp is under-reported, so a + // staleness check can pass a value that is already stale. The read path stays + // lock-free: only this write path takes the monitor, and it is held for two + // field writes plus a clock read. + synchronized (this) { + this.latest = new Latest<>(value, System.nanoTime()); + } + } + + /** + * @return true if a value of runtime class {@code actual} may be published here. + * Delegates to {@link #acceptsType} so the two cannot drift apart. + */ + boolean acceptsValueClass(Class actual) { + return acceptsType(actual); + } + + /** @return true if {@code other} is type-compatible with this topic. */ + boolean acceptsType(Class other) { + return boxedType.isAssignableFrom(box(other)); } /** * Returns the most recent value, falling back to {@code defaultValue} if nothing has * been published yet. Convenience for {@code topic.latestValue().orElse(default)}. * + *

Allocation-free — prefer this over {@link #latestValue()} on hot paths. + * * @param defaultValue the value to return before the first publish * @return the latest value, or {@code defaultValue} */ public T latestValueOr(T defaultValue) { - T v = latest; - return v != null ? v : defaultValue; + Latest snap = latest; + return snap == null ? defaultValue : snap.value; + } + + /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ + private static Class box(Class c) { + if (!c.isPrimitive()) return c; + if (c == int.class) return Integer.class; + if (c == long.class) return Long.class; + if (c == double.class) return Double.class; + if (c == float.class) return Float.class; + if (c == boolean.class) return Boolean.class; + if (c == byte.class) return Byte.class; + if (c == short.class) return Short.class; + if (c == char.class) return Character.class; + return c; } @Override diff --git a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java new file mode 100644 index 0000000..f54c26e --- /dev/null +++ b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java @@ -0,0 +1,266 @@ +package com.aaravlabs.synapse; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicLongArray; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Covers the lock-free {@link Topic} latest-value tracking: one volatile snapshot + * instead of a monitor, and a boxed type cached at construction so the per-publish type + * check is a single field read. + */ +class TopicLockFreeTest { + + private Orchestrator orchestrator; + private Topic topic; + + @BeforeEach void setUp() { + orchestrator = Orchestrator.create("topic-lock-free"); + topic = orchestrator.getOrCreateTopic("concurrent", Integer.class); + } + @AfterEach void tearDown() { orchestrator.close(); } + + @Test + void latestValueAndTimestampAdvanceTogether() throws Exception { + // Deliberately does not assume System.nanoTime() > 0: the JLS allows any + // origin, so only *relative* comparisons are asserted. + Topic t = orchestrator.getOrCreateTopic("t", String.class); + + orchestrator.publish("t", "a"); + assertEquals("a", t.latestValueOr(null)); + long afterFirst = t.latestPublishNanos(); + + // Do not assert that the first timestamp differs from a pre-publish reading: + // System.nanoTime() may return the same tick for both reads on a coarse-resolution + // platform, and the publish did record a timestamp. The value assertion above is + // what proves the publish landed; the timestamp comparison is made below, after + // the clock has demonstrably moved. + // + // Same value published again must still advance the timestamp — a cached + // "unchanged" shortcut would break staleness checks. Wait for the clock to pass + // the recorded stamp first (bounded, so a broken clock fails the assert instead + // of hanging). + long deadline = System.nanoTime() + 50_000_000L; + while (afterFirst >= System.nanoTime() && System.nanoTime() <= deadline) Thread.sleep(1); + orchestrator.publish("t", "a"); + assertTrue(t.latestPublishNanos() > afterFirst, + "republishing an identical value must still advance latestPublishNanos"); + } + + @Test + void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception { + // Topic reads are lock-free now (one volatile snapshot instead of a monitor on the + // read path). Hammer it from several publishers while a reader samples + // value+timestamp, and assert that a reader only ever observes a value that was + // actually published, that the pair is never torn, and that the final state is + // consistent. + // + // The pair is torn if a reader ends up holding a value OLDER than the timestamp it + // sampled, which would let an age check pass a value that is already stale. Each + // publisher publishes the globally unique ids in its own block and records the + // clock reading at the end of its own publish of that id; keyed by id those never + // move, so a timestamp later than the end of the value's own publish is + // unambiguous evidence of a split pair. + // + // A timestamp EARLIER than the value's own publish is not a defect: that is the + // reader sampling the timestamp, a publish landing, then reading the newer value. + // It is the documented order and it over-reports age, which is safe. + // + // This is a detector, not a gate, and its sensitivity was measured rather than + // assumed. Against a control that deliberately widens the interleaving window it + // fires every run. Against the two-volatile design this branch replaces it fires + // in roughly half the runs (~0.005% of samples), so a single green run here is + // not evidence that the code is correct. What this test does establish is that no + // reader ever observes a value outside the published range, and that the final + // state is some publisher's last value. + final int publishers = 4; + final int perPublisher = 25_000; + final int ids = publishers * perPublisher; + AtomicReference failure = new AtomicReference<>(); + AtomicReference readersRunning = new AtomicReference<>(Boolean.TRUE); + 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(), + // so a negative reading is legitimate and must stay distinguishable from + // "not recorded yet". + windowHi.set(i, Long.MIN_VALUE); + } + + // One reader sampling as tightly as it can: sampling density is what makes the + // tear detectable, so this loop deliberately does NOT yield. Thread.onSpinWait() + // is a CPU-relief hint that never gives up the timeslice, so the reader does hold + // a core while the publishers run, but only for the ~100ms they need, and the + // explicit isAlive() assertions below turn a genuinely starved publisher into a + // named failure rather than a misleading final-value assertion. + // + // Thread.yield() here would fix that starvation but weaken the detector, because + // yielding is part of what closes the interleaving window being observed. So the + // loop spins and the mitigation is a bounded spin plus the liveness assertions + // below: a publisher that genuinely cannot finish now fails by name instead of + // producing a misleading final-value failure. + Thread reader = new Thread(() -> { + try { + while (readersRunning.get()) { + // Timestamp FIRST, then value. On this order a publish landing + // between the two reads can only make the reported age too large, + // never too small, so a staleness check can never pass a stale value. + long ts = topic.latestPublishNanos(); + Integer v = topic.latestValueOr(null); + if (v == null) { + if (ts != 0L) { + failure.compareAndSet(null, new AssertionError( + "timestamp recorded with no value: " + ts)); + return; + } + } else if (v < 0 || v >= ids) { + failure.compareAndSet(null, + new AssertionError("torn/garbage latest value: " + v)); + return; + } else { + // Only a timestamp AFTER the window ends is a tear: the reader + // would be holding an older value beside a newer timestamp. A + // timestamp before the window is the safe, expected straddle -- + // the reader sampled the timestamp, a publish landed, then it read + // the newer value. + long hi = windowHi.get(v); + if (hi != Long.MIN_VALUE && ts > hi) { + failure.compareAndSet(null, new AssertionError( + "value " + v + " is older than the timestamp " + ts + + " the reader sampled; its publish ended at " + hi)); + return; + } + } + Thread.onSpinWait(); + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } + }, "topic-reader"); + reader.setDaemon(true); + reader.start(); + + Thread[] pubThreads = new Thread[publishers]; + for (int p = 0; p < publishers; p++) { + final int base = p * perPublisher; + pubThreads[p] = new Thread(() -> { + try { + 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. + windowHi.set(base + i, System.nanoTime()); + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } + }, "topic-publisher-" + p); + pubThreads[p].start(); + } + // Cleanup runs in finally: an assertion below can throw, and without this the + // spinning reader keeps a core busy for the rest of the JVM's life while the + // unjoined non-daemon publishers keep publishing -- slowing every later test and + // potentially hanging JVM exit. + try { + // A publisher that does not finish leaves a mid-block value as the final + // state, which would fail the assertion below for a reason that has nothing + // to do with the code under test. Check liveness explicitly so the failure + // names the thread. + for (Thread t : pubThreads) { + t.join(60_000); + assertFalse(t.isAlive(), "publisher " + t.getName() + " did not finish within 60s"); + } + + assertNull(failure.get(), "concurrent access failed: " + failure.get()); + // 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 + // numerically largest one. + int finalValue = orchestrator.getLatestValue("concurrent", Integer.class).orElse(-1); + assertTrue(finalValue >= 0 && (finalValue + 1) % perPublisher == 0, + "final latest value " + finalValue + + " is not the last value of any publisher"); + } finally { + // Signal the reader before joining it: it spins on this flag, so joining + // first would wait out the timeout on every run. Cleanup only -- deliberately + // no assertion here, because throwing from finally would mask whatever + // failure the try block is reporting. + readersRunning.set(Boolean.FALSE); + reader.join(5_000); + for (Thread t : pubThreads) { + t.join(60_000); + } + } + } + + @Test + void acceptsValueClassRejectsIncompatibleValues() { + // The per-publish type check must keep rejecting incompatible values, and a + // rejected value must not be recorded. Publishes reach acceptsValueClass with + // a runtime class only, so a topic typed Number has to accept several + // different concrete classes. + orchestrator.getOrCreateTopic("n", Number.class); + orchestrator.publish("n", 1); + orchestrator.publish("n", 2L); + orchestrator.publish("n", 3.5); + assertEquals(3.5, orchestrator.getLatestValue("n", Number.class).orElse(null)); + + assertThrows(IllegalArgumentException.class, () -> orchestrator.publish("n", "nope")); + // The rejected String must not have been recorded. + assertEquals(3.5, orchestrator.getLatestValue("n", Number.class).orElse(null)); + + // A topic typed Integer must still reject a String after accepting Integers. + orchestrator.getOrCreateTopic("i", Integer.class); + orchestrator.publish("i", 1); + orchestrator.publish("i", 2); + assertThrows(IllegalArgumentException.class, () -> orchestrator.publish("i", "nope")); + assertEquals(2, orchestrator.getLatestValue("i", Integer.class).orElse(null)); + } + + @Test + void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() { + // acceptsValueClass boxes its argument, so it agrees with acceptsType for the + // same pair even when handed a primitive class. + Topic t = orchestrator.getOrCreateTopic("d", Double.class); + 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()); + } + + Topic i = orchestrator.getOrCreateTopic("i", Integer.class); + assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic"); + } + + @Test + void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() { + // acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at + // four call sites. Primitive/wrapper equivalence has to survive. + orchestrator.getOrCreateTopic("t", Double.class); + orchestrator.publish("t", 0.5); + assertTrue(orchestrator.findTopic("t", double.class).isPresent(), + "Double topic must be findable as double"); + assertTrue(orchestrator.findTopic("t", Double.class).isPresent()); + assertEquals(0.5, orchestrator.getLatestValue("t", double.class).orElse(null)); + // Direction matters and is preserved: the check is + // `boxed(topicType).isAssignableFrom(boxed(requestedType))`, so you must ask + // with an equal-or-narrower type. Asking for the supertype Number finds + // nothing — same as before this change. + assertFalse(orchestrator.findTopic("t", Number.class).isPresent()); + assertFalse(orchestrator.findTopic("t", String.class).isPresent()); + assertFalse(orchestrator.findTopic("t", Object.class).isPresent()); + + // Re-creating with an equivalent type is allowed; an incompatible one throws. + assertNotNull(orchestrator.getOrCreateTopic("t", double.class)); + assertThrows(IllegalArgumentException.class, + () -> orchestrator.getOrCreateTopic("t", String.class)); + } +} diff --git a/website/src/content/docs/concepts/topics.mdx b/website/src/content/docs/concepts/topics.mdx index 384b130..0c9a389 100644 --- a/website/src/content/docs/concepts/topics.mdx +++ b/website/src/content/docs/concepts/topics.mdx @@ -51,8 +51,18 @@ Optional latest = orchestrator.getLatestValue("drive/power/fl", Double.c double fl = latest.orElse(0.0); Topic flTopic = orchestrator.findTopic("drive/power/fl", Double.class).orElseThrow(); + +// Read the timestamp first, then the value. Each accessor is internally consistent, +// but two separate calls can still straddle a publish, so this order is the only one +// that cannot pair an older value with a newer timestamp and report a fresh age for a +// value that is already stale. +// +// latestPublishNanos() is 0 until the topic's first publish, so treat 0 as "never +// published" rather than subtracting it -- that would hand you the raw nanoTime +// reading, which is seconds to days rather than an age. +long stamp = flTopic.latestPublishNanos(); +long nanosSinceUpdate = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; double fl2 = flTopic.latestValueOr(0.0); -long nanosSinceUpdate = System.nanoTime() - flTopic.latestPublishNanos(); ``` Fetching the latest value is the right tool inside periodic loops; subscribing is the right tool when the *event* matters, not the current state. See [Subscriptions](/docs/concepts/subscriptions) for the tradeoff. From eeef85d531244583e453230b04ebab0f90f54e6c Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Thu, 1 Oct 2026 19:24:24 -0600 Subject: [PATCH 2/4] feat(topic): add Topic.latest() for a single-read value + timestamp Checking whether a published value is stale today means two calls: long stamp = topic.latestPublishNanos(); T v = topic.latestValueOr(null); long age = System.nanoTime() - stamp; Both accessors are individually correct, but a publish can land between them, leaving the caller holding a value from one publish and a timestamp from another. Reading the timestamp first makes that safe in one direction -- the age can only be over-reported, never under-reported -- but over-reporting still discards a fresh value as stale, by up to a full publish interval. How often that happens depends on how long the two reads take. Normally they are nanoseconds apart and a publish landing between them is vanishingly rare. It stops being rare when the reader is preempted between the two calls by a GC pause or the scheduler, since the window becomes milliseconds. This library configures a small young gen, so young GCs are part of the deployment rather than a hypothetical. `latest()` returns the value and its timestamp from one snapshot read, so the pair provably comes from a single publish and the window does not exist: Optional> snap = topic.latest(); T v = snap.map(Topic.Latest::value).orElse(default); long age = snap.map(Topic.Latest::ageNanos).orElse(Long.MAX_VALUE); Additive only. No existing signature or behaviour changes, and callers that never check an age can keep using `latestValueOr` -- same single volatile read. `Topic.Latest` becomes public with `value()`, `publishNanos()` and `ageNanos()`. `ageNanos()` is the age of *this* value computed from its own stamp, which also removes the pre-publish `0` sentinel that the two-call form has to guard against by hand. `latestPublishNanos()` still returns 0 before the first publish, so the guard stays necessary for anyone composing it by hand. It is documented as such in the javadoc, the CHANGELOG and the topics guide. The concurrency test's reader now uses `latest()`, which is the stronger property: no straddle is representable, so the "timestamp first" ordering rule has nothing left to guard. That test previously carried a cross-check between `latest()` being empty and `latestPublishNanos()` being nonzero; it is gone, because a publish can legitimately land between those two reads and a nonzero stamp there says nothing about consistency. Suite: 61 tests, 0 failures. --- CHANGELOG.md | 37 ++++-- .../java/com/aaravlabs/synapse/Topic.java | 113 +++++++++++++++--- .../aaravlabs/synapse/TopicLockFreeTest.java | 103 ++++++++++++---- website/src/content/docs/api/index.mdx | 15 ++- website/src/content/docs/concepts/topics.mdx | 21 ++-- 5 files changed, 229 insertions(+), 60 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ef840b9..004b112 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,12 +26,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `publish` is reachable from the OpMode loop, the hardware thread, and the callback pool, so two publishers really can overlap. - The write path still takes the topic monitor, but only around the clock sample and - the store. Two separate volatile fields would let two publishers interleave - between the value write and the timestamp write; sampling the clock outside - mutual exclusion would let a preempted publisher install an older pair after a - newer one. Both regressions were reproduced and fixed; keeping the monitor on the - write path preserves the ordering the previous `synchronized` body provided. + The write path still takes the topic monitor, but only around the clock sample, + one short-lived allocation and the store. Two separate volatile fields would let + two publishers interleave between the value write and the timestamp write; + sampling the clock outside mutual exclusion would let a preempted publisher + install an older pair after a newer one. Both regressions were reproduced and + fixed; keeping the monitor on the write path preserves the ordering the previous + `synchronized` body provided. + + Each publish therefore allocates one small short-lived pair on the write path. + The read path is unchanged in allocation terms: `latestValueOr()` allocates + nothing. The topic's declared type is now normalized to its wrapper class **once**, at construction, for the per-publish type check. `type()` still reports the type the @@ -39,15 +44,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `latestPublishNanos()` still returns `0` before the first publish, unchanged. - **Read the timestamp before the value** when you use the two together as a - staleness check. Each accessor now returns a self-consistent pair, but two separate - calls can still straddle a publish, so the ordering rule remains — it is now - documented on the public accessors and in the topics guide: + **Prefer `Topic.latest()` when you need both the value and its age.** It returns the + two from one snapshot read, so they provably come from the same publish: ```java - long stamp = topic.latestPublishNanos(); // first + Optional> snap = topic.latest(); + T v = snap.map(Topic.Latest::value).orElse(defaultValue); + long age = snap.map(Topic.Latest::ageNanos).orElse(Long.MAX_VALUE); + ``` + + `latestPublishNanos()` and `latestValueOr()` are unchanged and still correct. Composing + them takes two calls, which can straddle a publish and leave the pair describing two + different publishes; if you do compose them, **read the timestamp first**: + + ```java + long stamp = topic.latestPublishNanos(); // first; 0 means nothing published yet T v = topic.latestValueOr(null); // then - long age = System.nanoTime() - stamp; + long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; ``` Reading the value first and the timestamp second can pair an older value with a diff --git a/src/main/java/com/aaravlabs/synapse/Topic.java b/src/main/java/com/aaravlabs/synapse/Topic.java index 68e88fb..c138032 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -46,15 +46,50 @@ public final class Topic { this.boxedType = box(type); } - /** One publish's value and the {@link System#nanoTime()} at which it was recorded. */ - private static final class Latest { - final T value; - final long publishNanos; + /** + * One publish's value and the {@link System#nanoTime()} at which it was recorded. + * + *

Immutable. Returned by {@link Topic#latest()} so a caller can read the value and + * its timestamp without risking a publish landing between two calls. + * + * @param the message type carried by this topic + */ + public static final class Latest { + + private final T value; + private final long publishNanos; Latest(T value, long publishNanos) { this.value = value; this.publishNanos = publishNanos; } + + /** @return the published value */ + public T value() { + return value; + } + + /** + * @return the {@link System#nanoTime()} at which {@link #value()} was recorded, on + * the monotonic clock (arbitrary origin — only differences are meaningful) + */ + public long publishNanos() { + return publishNanos; + } + + /** + * @return nanos elapsed since {@link #publishNanos()}, on the monotonic clock. + * Unlike {@code System.nanoTime() - publishNanos()}, this is always the age + * of this value and needs no {@code 0} sentinel handling. + */ + public long ageNanos() { + return System.nanoTime() - publishNanos; + } + + @Override + public String toString() { + return "Latest[" + value + " @" + publishNanos + "]"; + } } /** @@ -74,9 +109,9 @@ public Class type() { /** * The most recently published value, or empty if nothing has been published yet. * - *

Reading a value together with its age: each accessor reads a - * self-consistent pair, but two separate calls can still straddle a publish. Sample - * the timestamp first, which can only over-report the value's age: + *

Reading a value together with its age: prefer {@link #latest()}, which + * returns both from a single snapshot read. Composing them from two accessors works + * only if you sample the timestamp first, which can only over-report the age: * *

{@code
      * long stamp = topic.latestPublishNanos();  // first; 0 means nothing published yet
@@ -88,7 +123,7 @@ public Class type() {
      * {@link #latestPublishNanos()} is {@code 0}, and subtracting it would yield the raw
      * {@link System#nanoTime()} reading — seconds to days — rather than an age.
      *
-     * Reading the value first is the unsafe order: a publish landing between the two
+     * 

Reading the value first is the unsafe order: a publish landing between the two * calls pairs the older value with the newer timestamp, so an age check on that * pair passes even though the value is stale. * @@ -110,7 +145,7 @@ public Optional latestValue() { */ public long latestPublishNanos() { Latest snap = latest; - return snap == null ? 0L : snap.publishNanos; + return snap == null ? 0L : snap.publishNanos(); } /** @@ -122,8 +157,9 @@ public long latestPublishNanos() { * concurrent reader always sees a value and a timestamp from the same publish. * Splitting them across two volatile fields would let two concurrent publishers * interleave between the writes and leave the pair describing two different - * publishes. The write path takes the topic monitor, around the clock sample - * and the store — see below. + * publishes. The write path takes the topic monitor, around the clock sample, + * one short-lived allocation and the store — see below. Each publish therefore + * allocates; {@link #latestValueOr(Object)} does not. * * @param value the value to record */ @@ -133,9 +169,17 @@ void recordLatest(T value) { // nanoTime() and is then preempted installs an OLDER snapshot after a newer // one, and a reader can end up holding an older value beside a newer // timestamp -- an age computed from that stamp is under-reported, so a - // staleness check can pass a value that is already stale. The read path stays - // lock-free: only this write path takes the monitor, and it is held for two - // field writes plus a clock read. + // staleness check can pass a value that is already stale. + // + // The critical section is a clock read, an allocation and a volatile store -- + // NOT "two field writes and a clock read". Each publish allocates one 24-byte + // Latest, which is not scalar-replaceable because it escapes into a volatile + // field, so garbage scales with publish rate across every topic. + // + // Hoisting the allocation above the monitor was measured and is not a win: it + // gains ~12 ns/publish uncontended but loses ~10-20 ns under four publishers, + // since it lengthens the window in which a stalled publisher holds a snapshot + // whose stamp is already stale. The allocation stays inside. synchronized (this) { this.latest = new Latest<>(value, System.nanoTime()); } @@ -158,14 +202,51 @@ boolean acceptsType(Class other) { * Returns the most recent value, falling back to {@code defaultValue} if nothing has * been published yet. Convenience for {@code topic.latestValue().orElse(default)}. * - *

Allocation-free — prefer this over {@link #latestValue()} on hot paths. + *

Allocation-free — prefer this over {@link #latestValue()} on hot paths, and over + * {@link #latest()} if you do not need the timestamp. * * @param defaultValue the value to return before the first publish * @return the latest value, or {@code defaultValue} */ public T latestValueOr(T defaultValue) { Latest snap = latest; - return snap == null ? defaultValue : snap.value; + return snap == null ? defaultValue : snap.value(); + } + + /** + * The most recent value and the {@link System#nanoTime()} at which it was published, + * as one consistent pair. + * + *

Prefer this over {@link #latestValue()} plus {@link #latestPublishNanos()} + * when you are checking whether a value is stale. The two accessors each read the + * snapshot correctly, but a caller making two calls can straddle a publish between + * them and end up holding a value from one publish and a timestamp from another: + * + *

{@code
+     * long stamp = topic.latestPublishNanos();   // publish A
+     *                                                  ...publish B lands...
+     * T v = topic.latestValueOr(null);             // publish B's value, with A's stamp
+     * long age = System.nanoTime() - stamp;        // over-reports by a full publish interval
+     * }
+ * + *

Reading the timestamp first makes that safe in one direction — the age can only + * be over-reported, never under-reported — but over-reporting still costs you: a + * fresh value gets rejected as stale. How often that happens depends on how long the + * two reads take. Normally they are nanoseconds apart and a publish landing between + * them is vanishingly rare. It stops being rare when the reader is preempted between + * the two calls by a GC pause or the scheduler, since the window becomes + * milliseconds. This accessor removes the window entirely. + * + *

Returns empty before the first publish. + * + *

Allocates one small {@link Optional} wrapper per call. If you only want + * the value and never check its age, use {@link #latestValueOr(Object)}, which is + * allocation-free. + * + * @return an {@link Optional} holding the latest value and its publish timestamp + */ + public Optional> latest() { + return Optional.ofNullable(latest); } /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ diff --git a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java index f54c26e..9b8ba71 100644 --- a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java +++ b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java @@ -78,6 +78,9 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception // not evidence that the code is correct. What this test does establish is that no // reader ever observes a value outside the published range, and that the final // state is some publisher's last value. + // Timestamp and value from ONE snapshot read, so there is no window for a + // publish to land in between. This is the property the two-call accessors cannot + // offer, and the reason latest() exists. final int publishers = 4; final int perPublisher = 25_000; final int ids = publishers * perPublisher; @@ -106,32 +109,38 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception Thread reader = new Thread(() -> { try { while (readersRunning.get()) { - // Timestamp FIRST, then value. On this order a publish landing - // between the two reads can only make the reported age too large, - // never too small, so a staleness check can never pass a stale value. - long ts = topic.latestPublishNanos(); - Integer v = topic.latestValueOr(null); - if (v == null) { - if (ts != 0L) { - failure.compareAndSet(null, new AssertionError( - "timestamp recorded with no value: " + ts)); + // ONE snapshot read, so the value and the timestamp provably come + // from the same publish. There is no window for a publish to land in, + // so this cannot over-report age the way two calls can. + Topic.Latest snap = topic.latest().orElse(null); + if (snap == null) { + // Nothing published yet. There is deliberately no cross-check + // against latestPublishNanos() here: a publish can land between + // the two calls, so a nonzero stamp here would be that publish, + // not an inconsistency. + } else { + Integer v = snap.value(); + if (v < 0 || v >= ids) { + failure.compareAndSet(null, + new AssertionError("torn/garbage latest value: " + v)); return; } - } else if (v < 0 || v >= ids) { - failure.compareAndSet(null, - new AssertionError("torn/garbage latest value: " + v)); - return; - } else { - // Only a timestamp AFTER the window ends is a tear: the reader - // would be holding an older value beside a newer timestamp. A - // timestamp before the window is the safe, expected straddle -- - // the reader sampled the timestamp, a publish landed, then it read - // the newer value. + // 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. long hi = windowHi.get(v); - if (hi != Long.MIN_VALUE && ts > hi) { + if (hi != Long.MIN_VALUE && snap.publishNanos() > hi) { + failure.compareAndSet(null, new AssertionError( + "value " + v + " paired with stamp " + snap.publishNanos() + + " later than its publish end " + hi)); + return; + } + // ageNanos must agree with the stamp it was taken from. Zero is + // legitimate (same clock tick); only negative is impossible. + long age = snap.ageNanos(); + if (age < 0 || age > System.nanoTime() - snap.publishNanos()) { failure.compareAndSet(null, new AssertionError( - "value " + v + " is older than the timestamp " + ts - + " the reader sampled; its publish ended at " + hi)); + "implausible age " + age + " for stamp " + snap.publishNanos())); return; } } @@ -240,6 +249,56 @@ void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() { assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic"); } + @Test + void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception { + // latest() exists so a staleness check reads one snapshot instead of two. It must + // be empty before the first publish, agree with both single-field accessors, and + // never need a 0 sentinel. + Topic t = orchestrator.getOrCreateTopic("snap", String.class); + assertFalse(t.latest().isPresent(), "latest() must be empty before the first publish"); + + orchestrator.publish("snap", "a"); + 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)); + assertTrue(first.publishNanos() != 0L, + "a published stamp must not be the pre-publish 0 sentinel"); + + // ageNanos is the age of THIS value, computed from its own stamp. Taking a fresh + // clock reading afterwards must give an equal or LARGER delta, never a smaller + // one -- a smaller delta would mean the age came from somewhere other than this + // stamp, which is exactly the straddle latest() exists to prevent. + // age >= 0, not > 0: a snapshot read in the same clock tick as its own publish + // stamp legitimately reports age 0. Only a negative age is impossible here, and + // that is all this check needs to catch. + long age = first.ageNanos(); + assertTrue(age >= 0 && age < 1_000_000_000L, "ageNanos out of range: " + age); + assertTrue(age <= System.nanoTime() - first.publishNanos(), + "ageNanos " + age + " exceeds a delta measured later from the same stamp"); + + // A later publish replaces both halves together; there is no observable state in + // which one half is from this publish and the other from the previous one. + // + // Wait for the clock to pass first.publishNanos() rather than sleeping a fixed + // interval: System.nanoTime() may return the same tick for adjacent reads, so + // Thread.sleep(2) does not guarantee the next publish gets a later stamp. Same + // bounded-wait pattern as latestValueAndTimestampAdvanceTogether, so a broken + // clock fails the assert instead of hanging. + long deadline = System.nanoTime() + 50_000_000L; + while (first.publishNanos() >= System.nanoTime() && System.nanoTime() <= deadline) { + Thread.sleep(1); + } + orchestrator.publish("snap", "b"); + Topic.Latest second = t.latest().orElseThrow(); + assertEquals("b", second.value()); + assertTrue(second.publishNanos() > first.publishNanos(), + "a later publish must carry a later stamp"); + assertTrue(second.ageNanos() < first.ageNanos(), + "the newer value must report a younger age"); + } + @Test void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() { // acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at diff --git a/website/src/content/docs/api/index.mdx b/website/src/content/docs/api/index.mdx index c962550..c6f3ffb 100644 --- a/website/src/content/docs/api/index.mdx +++ b/website/src/content/docs/api/index.mdx @@ -48,8 +48,19 @@ Base class for robot logic. `protected final Orchestrator orchestrator` is avail | --- | --- | | `name()` / `type()` | identity | | `latestValue()` | `Optional` | -| `latestValueOr(default)` | `T` | -| `latestPublishNanos()` | `long`, `System.nanoTime()` clock | +| `latestValueOr(default)` | `T`; allocation-free, prefer on hot paths | +| `latestPublishNanos()` | `long`, `System.nanoTime()` clock, `0` before the first publish | +| `latest()` | `Optional>`; value and timestamp from one snapshot read | + +### `Topic.Latest` + +The pair returned by `Topic.latest()`. Immutable. + +| Method | Returns | +| --- | --- | +| `value()` | `T` | +| `publishNanos()` | `long`, `System.nanoTime()` clock | +| `ageNanos()` | `long`, age of this value; no `0` sentinel handling | ### Subscription diff --git a/website/src/content/docs/concepts/topics.mdx b/website/src/content/docs/concepts/topics.mdx index 0c9a389..0397db8 100644 --- a/website/src/content/docs/concepts/topics.mdx +++ b/website/src/content/docs/concepts/topics.mdx @@ -52,14 +52,19 @@ double fl = latest.orElse(0.0); Topic flTopic = orchestrator.findTopic("drive/power/fl", Double.class).orElseThrow(); -// Read the timestamp first, then the value. Each accessor is internally consistent, -// but two separate calls can still straddle a publish, so this order is the only one -// that cannot pair an older value with a newer timestamp and report a fresh age for a -// value that is already stale. -// -// latestPublishNanos() is 0 until the topic's first publish, so treat 0 as "never -// published" rather than subtracting it -- that would hand you the raw nanoTime -// reading, which is seconds to days rather than an age. +// One snapshot read gives the value and its timestamp together, so the two can never +// come from different publishes. ageNanos() is the age of this exact value, and needs no +// special handling before the first publish -- latest() is simply empty then. +Optional> flSnapshot = flTopic.latest(); +double fl2 = flSnapshot.map(Topic.Latest::value).orElse(0.0); +long nanosSinceUpdate = flSnapshot.map(Topic.Latest::ageNanos).orElse(Long.MAX_VALUE); +``` + +If you only want the value and never check its age, keep using `latestValueOr(0.0)` — it is the same single volatile read, reads better at a glance, and allocates nothing. `latest()` allocates one small `Optional` wrapper per call. + +Composing `latestPublishNanos()` and `latestValueOr()` still works. It takes two calls, so a publish can land between them; if you do it, read the timestamp first, and note that it is `0` until the topic's first publish: + +```java long stamp = flTopic.latestPublishNanos(); long nanosSinceUpdate = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; double fl2 = flTopic.latestValueOr(0.0); From b4dcf33629b0d31afb6e36e12c913f115caa957a Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Tue, 29 Sep 2026 21:20:48 -0600 Subject: [PATCH 3/4] perf(topic): lock-free latest-value tracking `latestValue` and `latestPublishNanos` were each `synchronized` on the Topic instance. That monitor is the most-used read on the bus, so it serialised every reader against every writer. The value and its timestamp are now published together as one immutable `Latest` pair behind a single volatile field, which makes both accessors plain volatile reads. The write path keeps a monitor, but only around the `nanoTime()` sample and the store. Two volatile fields would not do: `publish` is reachable from the OpMode loop, the hardware thread and the callback pool, so two publishers can interleave between a value write and a timestamp write and leave the pair describing two different publishes -- a regression against the old monitor, which ordered those writes. Sampling the clock outside mutual exclusion is not sufficient either: a publisher that samples and is then preempted installs an OLDER snapshot after a newer one, so a reader can hold an older value beside a newer timestamp and an age computed from that stamp is under-reported. Holding the monitor across the sample and the store keeps install order equal to timestamp order, which is the property the old body had. The type work is hoisted out of the publish path: - `boxedType` is computed once in the constructor instead of re-normalising the declared type on every publish and every lookup. `Topic.type()` still reports the type the topic was created with; only the internal comparison field is boxed. - `acceptsType(Class)` is the old `boxed(t.type()).isAssignableFrom(boxed(other))` with the topic side hoisted, replacing that expression at five call sites. `acceptsValueClass(Class)` delegates to it, so the two cannot drift. `boxed()` moves from `OrchestratorImpl` to `Topic`, next to the field it reads. Two separate accessors can still straddle a publish, so the read order remains: sample the timestamp first. That can only over-report the value's age, never under-report it. This is documented on the public accessors, in the CHANGELOG and in the topics guide, and the guide's example had the unsafe order. No API or behaviour change; `latestPublishNanos()` still returns 0 before the first publish. Measured on the repo benchmark (`gradle :benchmarks:run --args='run --quick --scenarios S0 --styles synapse'`, 8 alternating runs of baseline and branch on one box, p50 ns/op, median): micro.topic.latestValue 44 ns -> 28 ns (-36%) micro.topic.recordLatest.1p1c 360 ns -> 192 ns (-47%) micro.topic.recordLatest.4p4c 340 ns -> 268 ns (-21%) micro.topic.recordLatest 89 ns -> 83 ns (-7%) micro.publish.subscribers0 140 ns -> 161 ns (+15%) The read-path win is the durable one. The contended publish numbers are much smaller than a fully lock-free write path would give, because the write path keeps its monitor; that is the cost of not regressing the timestamp ordering the old code had. `publish.subscribers0` and the `rawDirectCall` control both moved by more than this box's noise floor in the same runs, so read those two as inconclusive rather than as a regression. Treat every number as a ratio on the box that produced it. --- CHANGELOG.md | 43 +++ .../aaravlabs/synapse/OrchestratorImpl.java | 31 +- .../java/com/aaravlabs/synapse/Topic.java | 146 +++++++++- .../aaravlabs/synapse/TopicLockFreeTest.java | 266 ++++++++++++++++++ website/src/content/docs/concepts/topics.mdx | 14 +- 5 files changed, 461 insertions(+), 39 deletions(-) create mode 100644 src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 7cded3b..15d0e65 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,49 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **`@SubscribedTo` dispatch no longer re-boxes the parameter type per message.** The primitive-to-wrapper normalisation is computed once at bind time instead of on every delivered message. Equivalent to the previous conditional check. +- **`Topic` latest-value reads no longer synchronize.** `latestValue()` and + `latestPublishNanos()` are plain volatile reads. The latest value and its + timestamp are published together as one immutable pair through a single volatile + field, so a reader always sees a value and a timestamp from the same publish. + `publish` is reachable from the OpMode loop, the hardware thread, and the callback + pool, so two publishers really can overlap. + + The write path still takes the topic monitor, but only around the clock sample, + one short-lived allocation and the store. Two separate volatile fields would let + two publishers interleave between the value write and the timestamp write; + sampling the clock outside mutual exclusion would let a preempted publisher + install an older pair after a newer one. Both regressions were reproduced and + fixed; keeping the monitor on the write path preserves the ordering the previous + `synchronized` body provided. + + Each publish therefore allocates one small short-lived pair on the write path. + The read path is unchanged in allocation terms: `latestValueOr()` allocates + nothing. + + The topic's declared type is now normalized to its wrapper class **once**, at + construction, for the per-publish type check. `type()` still reports the type the + topic was created with — only the internal comparison field is boxed. + + `latestPublishNanos()` still returns `0` before the first publish, unchanged. That + `0` is a heuristic rather than a proof — `System.nanoTime()` is permitted to return + `0`, so a genuine publish can carry it too. The resulting over-reported age rejects + a fresh value rather than admitting a stale one, so the failure direction is safe. + + **Read the timestamp before the value** when you use the two together as a + staleness check. Each accessor now returns a self-consistent pair, but two separate + calls can still straddle a publish, so the ordering rule remains — it is now + documented on the public accessors and in the topics guide: + + ```java + long stamp = topic.latestPublishNanos(); // first + T v = topic.latestValueOr(null); // then + long age = System.nanoTime() - stamp; + ``` + + Reading the value first and the timestamp second can pair an older value with a + newer timestamp, so an age check on that pair passes even though the value is + stale. Reading the timestamp first can only over-report the age, never under-report + it. ## [0.4.0] - 2026-09-12 diff --git a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java index 6aa2f54..695e704 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -147,11 +147,9 @@ public Topic getOrCreateTopic(String topicName, Class type) { Objects.requireNonNull(topicName, "topicName"); Objects.requireNonNull(type, "type"); - Class normalized = boxed(type); - Topic existing = topics.get(topicName); if (existing != null) { - if (!boxed(existing.type()).isAssignableFrom(normalized)) { + if (!existing.acceptsType(type)) { throw new IllegalArgumentException( "Topic '" + topicName + "' already exists with type " + existing.type().getName() + ", cannot re-create as " @@ -163,7 +161,7 @@ public Topic getOrCreateTopic(String topicName, Class type) { Topic created = new Topic<>(topicName, type); Topic prior = topics.putIfAbsent(topicName, created); if (prior != null) { - if (!boxed(prior.type()).isAssignableFrom(normalized)) { + if (!prior.acceptsType(type)) { throw new IllegalArgumentException( "Topic '" + topicName + "' already exists with type " + prior.type().getName() + ", cannot re-create as " @@ -175,20 +173,6 @@ public Topic getOrCreateTopic(String topicName, Class type) { return created; } - /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ - private static Class boxed(Class c) { - if (!c.isPrimitive()) return c; - if (c == int.class) return Integer.class; - if (c == long.class) return Long.class; - if (c == double.class) return Double.class; - if (c == float.class) return Float.class; - if (c == boolean.class) return Boolean.class; - if (c == byte.class) return Byte.class; - if (c == short.class) return Short.class; - if (c == char.class) return Character.class; - return c; - } - @Override public Optional> findTopic(String topicName) { return Optional.ofNullable(topics.get(topicName)); @@ -198,7 +182,7 @@ public Optional> findTopic(String topicName) { @SuppressWarnings("unchecked") public Optional> findTopic(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty(); + if (t == null || !t.acceptsType(type)) return Optional.empty(); return Optional.of((Topic) t); } @@ -217,14 +201,13 @@ 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. - Class valueType = value.getClass(); Topic topic = topics.get(topicName); if (topic == null) { - topic = getOrCreateTopic(topicName, valueType); - } else if (!boxed(topic.type()).isAssignableFrom(boxed(valueType))) { + topic = getOrCreateTopic(topicName, value.getClass()); + } else if (!topic.acceptsValueClass(value.getClass())) { throw new IllegalArgumentException( "Topic '" + topicName + "' is typed " + topic.type().getName() - + " but publish got " + valueType.getName()); + + " but publish got " + value.getClass().getName()); } ((Topic) topic).recordLatest(value); @@ -312,7 +295,7 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) { @SuppressWarnings("unchecked") public Optional getLatestValue(String topicName, Class type) { Topic t = topics.get(topicName); - if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty(); + if (t == null || !t.acceptsType(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 5b2302e..ecda38b 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -18,13 +18,43 @@ public final class Topic { private final String name; private final Class type; - // Guarded by `this` for write, volatile for the read in publish(). - private volatile T latest; - private volatile long latestPublishNanos; + /** + * {@code type} normalized to its wrapper class. Cached at construction so the + * per-publish type check is a single field read instead of a primitive-boxing + * chain. + */ + private final Class boxedType; + + /** + * The most recently published value and its timestamp, as one immutable pair. + * + *

A single volatile reference rather than two volatile fields: with two, two + * concurrent publishers can interleave between the value write and the timestamp + * write and leave the pair describing two different publishes. Since + * {@link Orchestrator#publish(String, Object)} is reachable from the OpMode + * loop, the hardware thread, and the callback pool, that is not a theoretical + * interleaving. + * + *

Immutable so that a reader holding a reference sees a self-consistent pair + * with no further synchronization. + */ + private volatile Latest latest; Topic(String name, Class type) { this.name = name; this.type = type; + this.boxedType = box(type); + } + + /** One publish's value and the {@link System#nanoTime()} at which it was recorded. */ + private static final class Latest { + final T value; + final long publishNanos; + + Latest(T value, long publishNanos) { + this.value = value; + this.publishNanos = publishNanos; + } } /** @@ -44,42 +74,130 @@ public Class type() { /** * The most recently published value, or empty if nothing has been published yet. * + *

Reading a value together with its age: each accessor reads a + * self-consistent pair, but two separate calls can still straddle a publish. Sample + * the timestamp first, which can only over-report the value's age: + * + *

{@code
+     * long stamp = topic.latestPublishNanos();  // read the timestamp FIRST
+     * T v = topic.latestValueOr(null);         // then the value
+     * long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
+     * }
+ * + *

Reading the timestamp first is what makes the pair safe: a publish landing + * between the two calls can only over-report the age, never under-report it. + * + *

The {@code 0} guard handles "nothing published yet", where the stamp is still + * its initial {@code 0} and subtracting it would yield the raw + * {@link System#nanoTime()} reading. It is a heuristic, not a proof: the JLS permits + * {@link System#nanoTime()} to return {@code 0}, so a genuine publish can carry a + * {@code 0} stamp too. That case only over-reports the age, which rejects a fresh + * value rather than admitting a stale one, so the failure direction is safe. + * + * Reading the value first is the unsafe order: a publish landing between the two + * calls pairs the older value with the newer timestamp, so an age check on that + * pair passes even though the value is stale. + * * @return an {@link Optional} holding the latest value */ - public synchronized Optional latestValue() { - return Optional.ofNullable(latest); + public Optional latestValue() { + return Optional.ofNullable(latestValueOr(null)); } /** - * Wall-clock nanos at which {@link #latestValue()} was last updated. + * Nanos at which {@link #latestValue()} was last updated, on the + * {@link System#nanoTime()} monotonic clock (arbitrary origin — only differences + * are meaningful). * - * @return publish timestamp in nanoseconds ({@code System.nanoTime()} clock) + *

Sample this before {@link #latestValue()} when the two are used together as a + * staleness check — see {@link #latestValue()} for the ordering rule. + * + * @return publish timestamp in nanoseconds, or {@code 0} if nothing has been published + * yet. Note that {@code 0} is also a value {@link System#nanoTime()} is + * permitted to return, so this cannot be used on its own to prove that no + * publish has occurred — see {@link #latestValue()}. */ - public synchronized long latestPublishNanos() { - return latestPublishNanos; + public long latestPublishNanos() { + Latest snap = latest; + return snap == null ? 0L : snap.publishNanos; } /** * Record a new latest value. Called by the orchestrator immediately before notifying * subscribers. * + *

The read path is lock-free: the value and its timestamp are published + * together as one immutable {@link Latest} through a single volatile reference, so a + * concurrent reader always sees a value and a timestamp from the same publish. + * Splitting them across two volatile fields would let two concurrent publishers + * interleave between the writes and leave the pair describing two different + * publishes. The write path takes the topic monitor, around the clock sample, + * one short-lived allocation and the store — see below. Each publish therefore + * allocates; the read accessors do not. + * * @param value the value to record */ - synchronized void recordLatest(T value) { - this.latest = value; - this.latestPublishNanos = System.nanoTime(); + void recordLatest(T value) { + // The snapshot is sampled and installed under this topic's monitor so that + // install order matches timestamp order. Without it, a publisher that samples + // nanoTime() and is then preempted installs an OLDER snapshot after a newer + // one, and a reader can end up holding an older value beside a newer + // timestamp -- an age computed from that stamp is under-reported, so a + // staleness check can pass a value that is already stale. + // + // The critical section is a clock read, an allocation and a volatile store -- + // NOT "two field writes and a clock read". Each publish allocates one 24-byte + // Latest, which is not scalar-replaceable because it escapes into a volatile + // field, so garbage scales with publish rate across every topic. + // + // Hoisting the allocation above the monitor was measured and is not a win: it + // gains ~12 ns/publish uncontended but loses ~10-20 ns under four publishers, + // since it lengthens the window in which a stalled publisher holds a snapshot + // whose stamp is already stale. The allocation stays inside. + synchronized (this) { + this.latest = new Latest<>(value, System.nanoTime()); + } + } + + /** + * @return true if a value of runtime class {@code actual} may be published here. + * Delegates to {@link #acceptsType} so the two cannot drift apart. + */ + boolean acceptsValueClass(Class actual) { + return acceptsType(actual); + } + + /** @return true if {@code other} is type-compatible with this topic. */ + boolean acceptsType(Class other) { + return boxedType.isAssignableFrom(box(other)); } /** * Returns the most recent value, falling back to {@code defaultValue} if nothing has * been published yet. Convenience for {@code topic.latestValue().orElse(default)}. * + *

Allocation-free — prefer this over {@link #latestValue()} on hot paths. + * * @param defaultValue the value to return before the first publish * @return the latest value, or {@code defaultValue} */ public T latestValueOr(T defaultValue) { - T v = latest; - return v != null ? v : defaultValue; + Latest snap = latest; + return snap == null ? defaultValue : snap.value; + } + + /** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */ + private static Class box(Class c) { + if (!c.isPrimitive()) return c; + if (c == int.class) return Integer.class; + if (c == long.class) return Long.class; + if (c == double.class) return Double.class; + if (c == float.class) return Float.class; + if (c == boolean.class) return Boolean.class; + if (c == byte.class) return Byte.class; + if (c == short.class) return Short.class; + if (c == char.class) return Character.class; + return c; } @Override diff --git a/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java new file mode 100644 index 0000000..f54c26e --- /dev/null +++ b/src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java @@ -0,0 +1,266 @@ +package com.aaravlabs.synapse; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicLongArray; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Covers the lock-free {@link Topic} latest-value tracking: one volatile snapshot + * instead of a monitor, and a boxed type cached at construction so the per-publish type + * check is a single field read. + */ +class TopicLockFreeTest { + + private Orchestrator orchestrator; + private Topic topic; + + @BeforeEach void setUp() { + orchestrator = Orchestrator.create("topic-lock-free"); + topic = orchestrator.getOrCreateTopic("concurrent", Integer.class); + } + @AfterEach void tearDown() { orchestrator.close(); } + + @Test + void latestValueAndTimestampAdvanceTogether() throws Exception { + // Deliberately does not assume System.nanoTime() > 0: the JLS allows any + // origin, so only *relative* comparisons are asserted. + Topic t = orchestrator.getOrCreateTopic("t", String.class); + + orchestrator.publish("t", "a"); + assertEquals("a", t.latestValueOr(null)); + long afterFirst = t.latestPublishNanos(); + + // Do not assert that the first timestamp differs from a pre-publish reading: + // System.nanoTime() may return the same tick for both reads on a coarse-resolution + // platform, and the publish did record a timestamp. The value assertion above is + // what proves the publish landed; the timestamp comparison is made below, after + // the clock has demonstrably moved. + // + // Same value published again must still advance the timestamp — a cached + // "unchanged" shortcut would break staleness checks. Wait for the clock to pass + // the recorded stamp first (bounded, so a broken clock fails the assert instead + // of hanging). + long deadline = System.nanoTime() + 50_000_000L; + while (afterFirst >= System.nanoTime() && System.nanoTime() <= deadline) Thread.sleep(1); + orchestrator.publish("t", "a"); + assertTrue(t.latestPublishNanos() > afterFirst, + "republishing an identical value must still advance latestPublishNanos"); + } + + @Test + void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception { + // Topic reads are lock-free now (one volatile snapshot instead of a monitor on the + // read path). Hammer it from several publishers while a reader samples + // value+timestamp, and assert that a reader only ever observes a value that was + // actually published, that the pair is never torn, and that the final state is + // consistent. + // + // The pair is torn if a reader ends up holding a value OLDER than the timestamp it + // sampled, which would let an age check pass a value that is already stale. Each + // publisher publishes the globally unique ids in its own block and records the + // clock reading at the end of its own publish of that id; keyed by id those never + // move, so a timestamp later than the end of the value's own publish is + // unambiguous evidence of a split pair. + // + // A timestamp EARLIER than the value's own publish is not a defect: that is the + // reader sampling the timestamp, a publish landing, then reading the newer value. + // It is the documented order and it over-reports age, which is safe. + // + // This is a detector, not a gate, and its sensitivity was measured rather than + // assumed. Against a control that deliberately widens the interleaving window it + // fires every run. Against the two-volatile design this branch replaces it fires + // in roughly half the runs (~0.005% of samples), so a single green run here is + // not evidence that the code is correct. What this test does establish is that no + // reader ever observes a value outside the published range, and that the final + // state is some publisher's last value. + final int publishers = 4; + final int perPublisher = 25_000; + final int ids = publishers * perPublisher; + AtomicReference failure = new AtomicReference<>(); + AtomicReference readersRunning = new AtomicReference<>(Boolean.TRUE); + 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(), + // so a negative reading is legitimate and must stay distinguishable from + // "not recorded yet". + windowHi.set(i, Long.MIN_VALUE); + } + + // One reader sampling as tightly as it can: sampling density is what makes the + // tear detectable, so this loop deliberately does NOT yield. Thread.onSpinWait() + // is a CPU-relief hint that never gives up the timeslice, so the reader does hold + // a core while the publishers run, but only for the ~100ms they need, and the + // explicit isAlive() assertions below turn a genuinely starved publisher into a + // named failure rather than a misleading final-value assertion. + // + // Thread.yield() here would fix that starvation but weaken the detector, because + // yielding is part of what closes the interleaving window being observed. So the + // loop spins and the mitigation is a bounded spin plus the liveness assertions + // below: a publisher that genuinely cannot finish now fails by name instead of + // producing a misleading final-value failure. + Thread reader = new Thread(() -> { + try { + while (readersRunning.get()) { + // Timestamp FIRST, then value. On this order a publish landing + // between the two reads can only make the reported age too large, + // never too small, so a staleness check can never pass a stale value. + long ts = topic.latestPublishNanos(); + Integer v = topic.latestValueOr(null); + if (v == null) { + if (ts != 0L) { + failure.compareAndSet(null, new AssertionError( + "timestamp recorded with no value: " + ts)); + return; + } + } else if (v < 0 || v >= ids) { + failure.compareAndSet(null, + new AssertionError("torn/garbage latest value: " + v)); + return; + } else { + // Only a timestamp AFTER the window ends is a tear: the reader + // would be holding an older value beside a newer timestamp. A + // timestamp before the window is the safe, expected straddle -- + // the reader sampled the timestamp, a publish landed, then it read + // the newer value. + long hi = windowHi.get(v); + if (hi != Long.MIN_VALUE && ts > hi) { + failure.compareAndSet(null, new AssertionError( + "value " + v + " is older than the timestamp " + ts + + " the reader sampled; its publish ended at " + hi)); + return; + } + } + Thread.onSpinWait(); + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } + }, "topic-reader"); + reader.setDaemon(true); + reader.start(); + + Thread[] pubThreads = new Thread[publishers]; + for (int p = 0; p < publishers; p++) { + final int base = p * perPublisher; + pubThreads[p] = new Thread(() -> { + try { + 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. + windowHi.set(base + i, System.nanoTime()); + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } + }, "topic-publisher-" + p); + pubThreads[p].start(); + } + // Cleanup runs in finally: an assertion below can throw, and without this the + // spinning reader keeps a core busy for the rest of the JVM's life while the + // unjoined non-daemon publishers keep publishing -- slowing every later test and + // potentially hanging JVM exit. + try { + // A publisher that does not finish leaves a mid-block value as the final + // state, which would fail the assertion below for a reason that has nothing + // to do with the code under test. Check liveness explicitly so the failure + // names the thread. + for (Thread t : pubThreads) { + t.join(60_000); + assertFalse(t.isAlive(), "publisher " + t.getName() + " did not finish within 60s"); + } + + assertNull(failure.get(), "concurrent access failed: " + failure.get()); + // 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 + // numerically largest one. + int finalValue = orchestrator.getLatestValue("concurrent", Integer.class).orElse(-1); + assertTrue(finalValue >= 0 && (finalValue + 1) % perPublisher == 0, + "final latest value " + finalValue + + " is not the last value of any publisher"); + } finally { + // Signal the reader before joining it: it spins on this flag, so joining + // first would wait out the timeout on every run. Cleanup only -- deliberately + // no assertion here, because throwing from finally would mask whatever + // failure the try block is reporting. + readersRunning.set(Boolean.FALSE); + reader.join(5_000); + for (Thread t : pubThreads) { + t.join(60_000); + } + } + } + + @Test + void acceptsValueClassRejectsIncompatibleValues() { + // The per-publish type check must keep rejecting incompatible values, and a + // rejected value must not be recorded. Publishes reach acceptsValueClass with + // a runtime class only, so a topic typed Number has to accept several + // different concrete classes. + orchestrator.getOrCreateTopic("n", Number.class); + orchestrator.publish("n", 1); + orchestrator.publish("n", 2L); + orchestrator.publish("n", 3.5); + assertEquals(3.5, orchestrator.getLatestValue("n", Number.class).orElse(null)); + + assertThrows(IllegalArgumentException.class, () -> orchestrator.publish("n", "nope")); + // The rejected String must not have been recorded. + assertEquals(3.5, orchestrator.getLatestValue("n", Number.class).orElse(null)); + + // A topic typed Integer must still reject a String after accepting Integers. + orchestrator.getOrCreateTopic("i", Integer.class); + orchestrator.publish("i", 1); + orchestrator.publish("i", 2); + assertThrows(IllegalArgumentException.class, () -> orchestrator.publish("i", "nope")); + assertEquals(2, orchestrator.getLatestValue("i", Integer.class).orElse(null)); + } + + @Test + void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() { + // acceptsValueClass boxes its argument, so it agrees with acceptsType for the + // same pair even when handed a primitive class. + Topic t = orchestrator.getOrCreateTopic("d", Double.class); + 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()); + } + + Topic i = orchestrator.getOrCreateTopic("i", Integer.class); + assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic"); + } + + @Test + void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() { + // acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at + // four call sites. Primitive/wrapper equivalence has to survive. + orchestrator.getOrCreateTopic("t", Double.class); + orchestrator.publish("t", 0.5); + assertTrue(orchestrator.findTopic("t", double.class).isPresent(), + "Double topic must be findable as double"); + assertTrue(orchestrator.findTopic("t", Double.class).isPresent()); + assertEquals(0.5, orchestrator.getLatestValue("t", double.class).orElse(null)); + // Direction matters and is preserved: the check is + // `boxed(topicType).isAssignableFrom(boxed(requestedType))`, so you must ask + // with an equal-or-narrower type. Asking for the supertype Number finds + // nothing — same as before this change. + assertFalse(orchestrator.findTopic("t", Number.class).isPresent()); + assertFalse(orchestrator.findTopic("t", String.class).isPresent()); + assertFalse(orchestrator.findTopic("t", Object.class).isPresent()); + + // Re-creating with an equivalent type is allowed; an incompatible one throws. + assertNotNull(orchestrator.getOrCreateTopic("t", double.class)); + assertThrows(IllegalArgumentException.class, + () -> orchestrator.getOrCreateTopic("t", String.class)); + } +} diff --git a/website/src/content/docs/concepts/topics.mdx b/website/src/content/docs/concepts/topics.mdx index 384b130..8a85f94 100644 --- a/website/src/content/docs/concepts/topics.mdx +++ b/website/src/content/docs/concepts/topics.mdx @@ -51,8 +51,20 @@ Optional latest = orchestrator.getLatestValue("drive/power/fl", Double.c double fl = latest.orElse(0.0); Topic flTopic = orchestrator.findTopic("drive/power/fl", Double.class).orElseThrow(); + +// Read the timestamp first, then the value. Each accessor is internally consistent, +// but two separate calls can still straddle a publish, so this order is the only one +// that cannot pair an older value with a newer timestamp and report a fresh age for a +// value that is already stale. +// +// latestPublishNanos() is 0 until the topic's first publish, so guard before +// subtracting it -- otherwise you get the raw nanoTime reading, which is seconds to +// days rather than an age. It is a heuristic, not a proof: nanoTime() is permitted to +// return 0, so a real publish can carry a 0 stamp too. That case over-reports the age +// and rejects a fresh value, which is the safe direction. +long stamp = flTopic.latestPublishNanos(); +long nanosSinceUpdate = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; double fl2 = flTopic.latestValueOr(0.0); -long nanosSinceUpdate = System.nanoTime() - flTopic.latestPublishNanos(); ``` Fetching the latest value is the right tool inside periodic loops; subscribing is the right tool when the *event* matters, not the current state. See [Subscriptions](/docs/concepts/subscriptions) for the tradeoff. From ccfabe2d33235fd70eb687315895c0c11c77ba1e Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Thu, 1 Oct 2026 20:51:57 -0600 Subject: [PATCH 4/4] docs(topic): mark the pre-publish 0 stamp as a heuristic, not proof System.nanoTime() is permitted to return 0, so a 0 stamp is not proof that nothing has been published -- a genuine publish can carry one too. Stated in the latestPublishNanos() javadoc, the latestValue() example, the topics guide and the CHANGELOG, along with which way the resulting over-reported age fails: it rejects a fresh value rather than admitting a stale one. latest() on this branch removes the question entirely, since one snapshot read needs no sentinel. --- CHANGELOG.md | 7 ++-- .../java/com/aaravlabs/synapse/Topic.java | 35 ++++++++++++++----- website/src/content/docs/concepts/topics.mdx | 2 +- 3 files changed, 32 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 004b112..fe598ce 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,7 +42,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 construction, for the per-publish type check. `type()` still reports the type the topic was created with — only the internal comparison field is boxed. - `latestPublishNanos()` still returns `0` before the first publish, unchanged. + `latestPublishNanos()` still returns `0` before the first publish, unchanged. That + `0` is a heuristic rather than a proof — `System.nanoTime()` is permitted to return + `0`, so a genuine publish can carry it too. The resulting over-reported age rejects + a fresh value rather than admitting a stale one, so the failure direction is safe. **Prefer `Topic.latest()` when you need both the value and its age.** It returns the two from one snapshot read, so they provably come from the same publish: @@ -58,7 +61,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 different publishes; if you do compose them, **read the timestamp first**: ```java - long stamp = topic.latestPublishNanos(); // first; 0 means nothing published yet + long stamp = topic.latestPublishNanos(); // read the timestamp FIRST T v = topic.latestValueOr(null); // then long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; ``` diff --git a/src/main/java/com/aaravlabs/synapse/Topic.java b/src/main/java/com/aaravlabs/synapse/Topic.java index c138032..e3b6487 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -64,12 +64,18 @@ public static final class Latest { this.publishNanos = publishNanos; } - /** @return the published value */ + /** + * The published value. + * + * @return the value published in this snapshot + */ public T value() { return value; } /** + * When this value was recorded. + * * @return the {@link System#nanoTime()} at which {@link #value()} was recorded, on * the monotonic clock (arbitrary origin — only differences are meaningful) */ @@ -78,9 +84,12 @@ public long publishNanos() { } /** - * @return nanos elapsed since {@link #publishNanos()}, on the monotonic clock. - * Unlike {@code System.nanoTime() - publishNanos()}, this is always the age - * of this value and needs no {@code 0} sentinel handling. + * How long ago this value was recorded, measured from its own stamp. + * + *

Unlike {@code System.nanoTime() - publishNanos()}, this is always the age of + * this value and needs no {@code 0} sentinel handling. + * + * @return nanos elapsed since {@link #publishNanos()}, on the monotonic clock */ public long ageNanos() { return System.nanoTime() - publishNanos; @@ -114,14 +123,19 @@ public Class type() { * only if you sample the timestamp first, which can only over-report the age: * *

{@code
-     * long stamp = topic.latestPublishNanos();  // first; 0 means nothing published yet
+     * long stamp = topic.latestPublishNanos();  // read the timestamp FIRST
      * T v = topic.latestValueOr(null);         // then
      * long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
      * }
* - *

The {@code 0} guard matters: before the first publish - * {@link #latestPublishNanos()} is {@code 0}, and subtracting it would yield the raw - * {@link System#nanoTime()} reading — seconds to days — rather than an age. + *

The {@code 0} guard handles "nothing published yet", where the stamp is still + * its initial {@code 0} and subtracting it would yield the raw + * {@link System#nanoTime()} reading. It is a heuristic, not a proof: the JLS permits + * {@link System#nanoTime()} to return {@code 0}, so a genuine publish can carry a + * {@code 0} stamp too. That case only over-reports the age, which rejects a fresh + * value rather than admitting a stale one, so the failure direction is safe. + * + *

{@link #latest()} removes the question entirely: one snapshot read, no sentinel. * *

Reading the value first is the unsafe order: a publish landing between the two * calls pairs the older value with the newer timestamp, so an age check on that @@ -141,7 +155,10 @@ public Optional latestValue() { *

Sample this before {@link #latestValue()} when the two are used together as a * staleness check — see {@link #latestValue()} for the ordering rule. * - * @return publish timestamp in nanoseconds, or 0 if nothing has been published + * @return publish timestamp in nanoseconds, or {@code 0} if nothing has been published + * yet. Note that {@code 0} is also a value {@link System#nanoTime()} is + * permitted to return, so this cannot be used on its own to prove that no + * publish has occurred — see {@link #latestValue()}. */ public long latestPublishNanos() { Latest snap = latest; diff --git a/website/src/content/docs/concepts/topics.mdx b/website/src/content/docs/concepts/topics.mdx index 0397db8..e1d5b08 100644 --- a/website/src/content/docs/concepts/topics.mdx +++ b/website/src/content/docs/concepts/topics.mdx @@ -62,7 +62,7 @@ long nanosSinceUpdate = flSnapshot.map(Topic.Latest::ageNanos).orElse(Long.MAX_V If you only want the value and never check its age, keep using `latestValueOr(0.0)` — it is the same single volatile read, reads better at a glance, and allocates nothing. `latest()` allocates one small `Optional` wrapper per call. -Composing `latestPublishNanos()` and `latestValueOr()` still works. It takes two calls, so a publish can land between them; if you do it, read the timestamp first, and note that it is `0` until the topic's first publish: +Composing `latestPublishNanos()` and `latestValueOr()` still works. It takes two calls, so a publish can land between them; if you do it, read the timestamp first, and note that the stamp is `0` until the topic's first publish, so guard before subtracting it — otherwise you get the raw `nanoTime()` reading, which is seconds to days rather than an age. That guard is a heuristic, not a proof: `nanoTime()` is permitted to return `0`, so a real publish can carry a `0` stamp too, which over-reports the age and rejects a fresh value — the safe direction. ```java long stamp = flTopic.latestPublishNanos();