From 34f8ae8e6748d4d1f7f9609392e9b637a4fef9ae Mon Sep 17 00:00:00 2001 From: Aarav Sharma Date: Tue, 29 Sep 2026 21:20:48 -0600 Subject: [PATCH] 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 | 44 +++ .../aaravlabs/synapse/OrchestratorImpl.java | 31 +- .../java/com/aaravlabs/synapse/Topic.java | 147 +++++++++- .../aaravlabs/synapse/TopicLockFreeTest.java | 266 ++++++++++++++++++ website/src/content/docs/concepts/topics.mdx | 14 +- 5 files changed, 463 insertions(+), 39 deletions(-) create mode 100644 src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 7cded3b..fa5388a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,50 @@ 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, while `latestValue()` wraps its result in an `Optional` as it did + before. + + 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 = stamp == 0L ? Long.MAX_VALUE : 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..d3818a2 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,131 @@ 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; on the read side {@link #latestValueOr(Object)} allocates nothing, + * while {@link #latestValue()} wraps its result in an {@link Optional}. + * * @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 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.