diff --git a/CHANGELOG.md b/CHANGELOG.md index fa5388a..faf6e99 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -48,13 +48,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `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: + **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(); // 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 d3818a2..da742c9 100644 --- a/src/main/java/com/aaravlabs/synapse/Topic.java +++ b/src/main/java/com/aaravlabs/synapse/Topic.java @@ -46,15 +46,59 @@ 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; } + + /** + * 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) + */ + public long publishNanos() { + return publishNanos; + } + + /** + * 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; + } + + @Override + public String toString() { + return "Latest[" + value + " @" + publishNanos + "]"; + } } /** @@ -74,19 +118,16 @@ 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();  // read the timestamp FIRST
-     * T v = topic.latestValueOr(null);         // then the value
+     * T v = topic.latestValueOr(null);         // then
      * 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 @@ -94,7 +135,9 @@ public Class type() { * {@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 + *

{@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 * pair passes even though the value is stale. * @@ -119,7 +162,7 @@ public Optional latestValue() { */ public long latestPublishNanos() { Latest snap = latest; - return snap == null ? 0L : snap.publishNanos; + return snap == null ? 0L : snap.publishNanos(); } /** @@ -134,7 +177,8 @@ public long latestPublishNanos() { * 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}. + * while {@link #latestValue()} and {@link #latest()} each wrap their result in an + * {@link Optional}. * * @param value the value to record */ @@ -177,14 +221,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 8a85f94..1c5fc14 100644 --- a/website/src/content/docs/concepts/topics.mdx +++ b/website/src/content/docs/concepts/topics.mdx @@ -52,16 +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 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. +// 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 +// sentinel 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. 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(); long nanosSinceUpdate = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp; double fl2 = flTopic.latestValueOr(0.0);