Skip to content

Commit eeef85d

Browse files
committed
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<Topic.Latest<T>> 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.
1 parent 54d881e commit eeef85d

5 files changed

Lines changed: 229 additions & 60 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 25 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -26,28 +26,41 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
2626
`publish` is reachable from the OpMode loop, the hardware thread, and the callback
2727
pool, so two publishers really can overlap.
2828

29-
The write path still takes the topic monitor, but only around the clock sample and
30-
the store. Two separate volatile fields would let two publishers interleave
31-
between the value write and the timestamp write; sampling the clock outside
32-
mutual exclusion would let a preempted publisher install an older pair after a
33-
newer one. Both regressions were reproduced and fixed; keeping the monitor on the
34-
write path preserves the ordering the previous `synchronized` body provided.
29+
The write path still takes the topic monitor, but only around the clock sample,
30+
one short-lived allocation and the store. Two separate volatile fields would let
31+
two publishers interleave between the value write and the timestamp write;
32+
sampling the clock outside mutual exclusion would let a preempted publisher
33+
install an older pair after a newer one. Both regressions were reproduced and
34+
fixed; keeping the monitor on the write path preserves the ordering the previous
35+
`synchronized` body provided.
36+
37+
Each publish therefore allocates one small short-lived pair on the write path.
38+
The read path is unchanged in allocation terms: `latestValueOr()` allocates
39+
nothing.
3540

3641
The topic's declared type is now normalized to its wrapper class **once**, at
3742
construction, for the per-publish type check. `type()` still reports the type the
3843
topic was created with — only the internal comparison field is boxed.
3944

4045
`latestPublishNanos()` still returns `0` before the first publish, unchanged.
4146

42-
**Read the timestamp before the value** when you use the two together as a
43-
staleness check. Each accessor now returns a self-consistent pair, but two separate
44-
calls can still straddle a publish, so the ordering rule remains — it is now
45-
documented on the public accessors and in the topics guide:
47+
**Prefer `Topic.latest()` when you need both the value and its age.** It returns the
48+
two from one snapshot read, so they provably come from the same publish:
4649

4750
```java
48-
long stamp = topic.latestPublishNanos(); // first
51+
Optional<Topic.Latest<T>> snap = topic.latest();
52+
T v = snap.map(Topic.Latest::value).orElse(defaultValue);
53+
long age = snap.map(Topic.Latest::ageNanos).orElse(Long.MAX_VALUE);
54+
```
55+
56+
`latestPublishNanos()` and `latestValueOr()` are unchanged and still correct. Composing
57+
them takes two calls, which can straddle a publish and leave the pair describing two
58+
different publishes; if you do compose them, **read the timestamp first**:
59+
60+
```java
61+
long stamp = topic.latestPublishNanos(); // first; 0 means nothing published yet
4962
T v = topic.latestValueOr(null); // then
50-
long age = System.nanoTime() - stamp;
63+
long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
5164
```
5265

5366
Reading the value first and the timestamp second can pair an older value with a

‎src/main/java/com/aaravlabs/synapse/Topic.java‎

Lines changed: 97 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -46,15 +46,50 @@ public final class Topic<T> {
4646
this.boxedType = box(type);
4747
}
4848

49-
/** One publish's value and the {@link System#nanoTime()} at which it was recorded. */
50-
private static final class Latest<T> {
51-
final T value;
52-
final long publishNanos;
49+
/**
50+
* One publish's value and the {@link System#nanoTime()} at which it was recorded.
51+
*
52+
* <p>Immutable. Returned by {@link Topic#latest()} so a caller can read the value and
53+
* its timestamp without risking a publish landing between two calls.
54+
*
55+
* @param <T> the message type carried by this topic
56+
*/
57+
public static final class Latest<T> {
58+
59+
private final T value;
60+
private final long publishNanos;
5361

5462
Latest(T value, long publishNanos) {
5563
this.value = value;
5664
this.publishNanos = publishNanos;
5765
}
66+
67+
/** @return the published value */
68+
public T value() {
69+
return value;
70+
}
71+
72+
/**
73+
* @return the {@link System#nanoTime()} at which {@link #value()} was recorded, on
74+
* the monotonic clock (arbitrary origin — only differences are meaningful)
75+
*/
76+
public long publishNanos() {
77+
return publishNanos;
78+
}
79+
80+
/**
81+
* @return nanos elapsed since {@link #publishNanos()}, on the monotonic clock.
82+
* Unlike {@code System.nanoTime() - publishNanos()}, this is always the age
83+
* of <i>this</i> value and needs no {@code 0} sentinel handling.
84+
*/
85+
public long ageNanos() {
86+
return System.nanoTime() - publishNanos;
87+
}
88+
89+
@Override
90+
public String toString() {
91+
return "Latest[" + value + " @" + publishNanos + "]";
92+
}
5893
}
5994

6095
/**
@@ -74,9 +109,9 @@ public Class<T> type() {
74109
/**
75110
* The most recently published value, or empty if nothing has been published yet.
76111
*
77-
* <p><b>Reading a value together with its age:</b> each accessor reads a
78-
* self-consistent pair, but two separate calls can still straddle a publish. Sample
79-
* the timestamp <i>first</i>, which can only over-report the value's age:
112+
* <p><b>Reading a value together with its age:</b> prefer {@link #latest()}, which
113+
* returns both from a single snapshot read. Composing them from two accessors works
114+
* only if you sample the timestamp <i>first</i>, which can only over-report the age:
80115
*
81116
* <pre>{@code
82117
* long stamp = topic.latestPublishNanos(); // first; 0 means nothing published yet
@@ -88,7 +123,7 @@ public Class<T> type() {
88123
* {@link #latestPublishNanos()} is {@code 0}, and subtracting it would yield the raw
89124
* {@link System#nanoTime()} reading — seconds to days — rather than an age.
90125
*
91-
* Reading the value first is the unsafe order: a publish landing between the two
126+
* <p>Reading the value first is the unsafe order: a publish landing between the two
92127
* calls pairs the older value with the newer timestamp, so an age check on that
93128
* pair passes even though the value is stale.
94129
*
@@ -110,7 +145,7 @@ public Optional<T> latestValue() {
110145
*/
111146
public long latestPublishNanos() {
112147
Latest<T> snap = latest;
113-
return snap == null ? 0L : snap.publishNanos;
148+
return snap == null ? 0L : snap.publishNanos();
114149
}
115150

116151
/**
@@ -122,8 +157,9 @@ public long latestPublishNanos() {
122157
* concurrent reader always sees a value and a timestamp from the same publish.
123158
* Splitting them across two volatile fields would let two concurrent publishers
124159
* interleave between the writes and leave the pair describing two different
125-
* publishes. The <b>write path takes the topic monitor</b>, around the clock sample
126-
* and the store — see below.
160+
* publishes. The <b>write path takes the topic monitor</b>, around the clock sample,
161+
* one short-lived allocation and the store — see below. Each publish therefore
162+
* allocates; {@link #latestValueOr(Object)} does not.
127163
*
128164
* @param value the value to record
129165
*/
@@ -133,9 +169,17 @@ void recordLatest(T value) {
133169
// nanoTime() and is then preempted installs an OLDER snapshot after a newer
134170
// one, and a reader can end up holding an older value beside a newer
135171
// timestamp -- an age computed from that stamp is under-reported, so a
136-
// staleness check can pass a value that is already stale. The read path stays
137-
// lock-free: only this write path takes the monitor, and it is held for two
138-
// field writes plus a clock read.
172+
// staleness check can pass a value that is already stale.
173+
//
174+
// The critical section is a clock read, an allocation and a volatile store --
175+
// NOT "two field writes and a clock read". Each publish allocates one 24-byte
176+
// Latest, which is not scalar-replaceable because it escapes into a volatile
177+
// field, so garbage scales with publish rate across every topic.
178+
//
179+
// Hoisting the allocation above the monitor was measured and is not a win: it
180+
// gains ~12 ns/publish uncontended but loses ~10-20 ns under four publishers,
181+
// since it lengthens the window in which a stalled publisher holds a snapshot
182+
// whose stamp is already stale. The allocation stays inside.
139183
synchronized (this) {
140184
this.latest = new Latest<>(value, System.nanoTime());
141185
}
@@ -158,14 +202,51 @@ boolean acceptsType(Class<?> other) {
158202
* Returns the most recent value, falling back to {@code defaultValue} if nothing has
159203
* been published yet. Convenience for {@code topic.latestValue().orElse(default)}.
160204
*
161-
* <p>Allocation-free — prefer this over {@link #latestValue()} on hot paths.
205+
* <p>Allocation-free — prefer this over {@link #latestValue()} on hot paths, and over
206+
* {@link #latest()} if you do not need the timestamp.
162207
*
163208
* @param defaultValue the value to return before the first publish
164209
* @return the latest value, or {@code defaultValue}
165210
*/
166211
public T latestValueOr(T defaultValue) {
167212
Latest<T> snap = latest;
168-
return snap == null ? defaultValue : snap.value;
213+
return snap == null ? defaultValue : snap.value();
214+
}
215+
216+
/**
217+
* The most recent value and the {@link System#nanoTime()} at which it was published,
218+
* as one consistent pair.
219+
*
220+
* <p><b>Prefer this over {@link #latestValue()} plus {@link #latestPublishNanos()}
221+
* when you are checking whether a value is stale.</b> The two accessors each read the
222+
* snapshot correctly, but a caller making two calls can straddle a publish between
223+
* them and end up holding a value from one publish and a timestamp from another:
224+
*
225+
* <pre>{@code
226+
* long stamp = topic.latestPublishNanos(); // publish A
227+
* ...publish B lands...
228+
* T v = topic.latestValueOr(null); // publish B's value, with A's stamp
229+
* long age = System.nanoTime() - stamp; // over-reports by a full publish interval
230+
* }</pre>
231+
*
232+
* <p>Reading the timestamp first makes that safe in one direction — the age can only
233+
* be over-reported, never under-reported — but over-reporting still costs you: a
234+
* fresh value gets rejected as stale. How often that happens depends on how long the
235+
* two reads take. Normally they are nanoseconds apart and a publish landing between
236+
* them is vanishingly rare. It stops being rare when the reader is preempted between
237+
* the two calls by a GC pause or the scheduler, since the window becomes
238+
* milliseconds. This accessor removes the window entirely.
239+
*
240+
* <p>Returns empty before the first publish.
241+
*
242+
* <p><b>Allocates one small {@link Optional} wrapper per call.</b> If you only want
243+
* the value and never check its age, use {@link #latestValueOr(Object)}, which is
244+
* allocation-free.
245+
*
246+
* @return an {@link Optional} holding the latest value and its publish timestamp
247+
*/
248+
public Optional<Latest<T>> latest() {
249+
return Optional.ofNullable(latest);
169250
}
170251

171252
/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */

‎src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java‎

Lines changed: 81 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,9 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception
7878
// not evidence that the code is correct. What this test does establish is that no
7979
// reader ever observes a value outside the published range, and that the final
8080
// state is some publisher's last value.
81+
// Timestamp and value from ONE snapshot read, so there is no window for a
82+
// publish to land in between. This is the property the two-call accessors cannot
83+
// offer, and the reason latest() exists.
8184
final int publishers = 4;
8285
final int perPublisher = 25_000;
8386
final int ids = publishers * perPublisher;
@@ -106,32 +109,38 @@ void concurrentPublishAndReadNeverLosesOrTearsTheLatestValue() throws Exception
106109
Thread reader = new Thread(() -> {
107110
try {
108111
while (readersRunning.get()) {
109-
// Timestamp FIRST, then value. On this order a publish landing
110-
// between the two reads can only make the reported age too large,
111-
// never too small, so a staleness check can never pass a stale value.
112-
long ts = topic.latestPublishNanos();
113-
Integer v = topic.latestValueOr(null);
114-
if (v == null) {
115-
if (ts != 0L) {
116-
failure.compareAndSet(null, new AssertionError(
117-
"timestamp recorded with no value: " + ts));
112+
// ONE snapshot read, so the value and the timestamp provably come
113+
// from the same publish. There is no window for a publish to land in,
114+
// so this cannot over-report age the way two calls can.
115+
Topic.Latest<Integer> snap = topic.latest().orElse(null);
116+
if (snap == null) {
117+
// Nothing published yet. There is deliberately no cross-check
118+
// against latestPublishNanos() here: a publish can land between
119+
// the two calls, so a nonzero stamp here would be that publish,
120+
// not an inconsistency.
121+
} else {
122+
Integer v = snap.value();
123+
if (v < 0 || v >= ids) {
124+
failure.compareAndSet(null,
125+
new AssertionError("torn/garbage latest value: " + v));
118126
return;
119127
}
120-
} else if (v < 0 || v >= ids) {
121-
failure.compareAndSet(null,
122-
new AssertionError("torn/garbage latest value: " + v));
123-
return;
124-
} else {
125-
// Only a timestamp AFTER the window ends is a tear: the reader
126-
// would be holding an older value beside a newer timestamp. A
127-
// timestamp before the window is the safe, expected straddle --
128-
// the reader sampled the timestamp, a publish landed, then it read
129-
// the newer value.
128+
// The pair is self-consistent by construction, so the only way
129+
// to fail is a stamp that predates the value it came with --
130+
// which the publish-window check bounds.
130131
long hi = windowHi.get(v);
131-
if (hi != Long.MIN_VALUE && ts > hi) {
132+
if (hi != Long.MIN_VALUE && snap.publishNanos() > hi) {
133+
failure.compareAndSet(null, new AssertionError(
134+
"value " + v + " paired with stamp " + snap.publishNanos()
135+
+ " later than its publish end " + hi));
136+
return;
137+
}
138+
// ageNanos must agree with the stamp it was taken from. Zero is
139+
// legitimate (same clock tick); only negative is impossible.
140+
long age = snap.ageNanos();
141+
if (age < 0 || age > System.nanoTime() - snap.publishNanos()) {
132142
failure.compareAndSet(null, new AssertionError(
133-
"value " + v + " is older than the timestamp " + ts
134-
+ " the reader sampled; its publish ended at " + hi));
143+
"implausible age " + age + " for stamp " + snap.publishNanos()));
135144
return;
136145
}
137146
}
@@ -240,6 +249,56 @@ void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() {
240249
assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic");
241250
}
242251

252+
@Test
253+
void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception {
254+
// latest() exists so a staleness check reads one snapshot instead of two. It must
255+
// be empty before the first publish, agree with both single-field accessors, and
256+
// never need a 0 sentinel.
257+
Topic<String> t = orchestrator.getOrCreateTopic("snap", String.class);
258+
assertFalse(t.latest().isPresent(), "latest() must be empty before the first publish");
259+
260+
orchestrator.publish("snap", "a");
261+
Topic.Latest<String> first = t.latest().orElseThrow();
262+
assertEquals("a", first.value());
263+
assertEquals(first.publishNanos(), t.latestPublishNanos(),
264+
"latest() and latestPublishNanos() must describe the same publish");
265+
assertEquals("a", t.latestValueOr(null));
266+
assertTrue(first.publishNanos() != 0L,
267+
"a published stamp must not be the pre-publish 0 sentinel");
268+
269+
// ageNanos is the age of THIS value, computed from its own stamp. Taking a fresh
270+
// clock reading afterwards must give an equal or LARGER delta, never a smaller
271+
// one -- a smaller delta would mean the age came from somewhere other than this
272+
// stamp, which is exactly the straddle latest() exists to prevent.
273+
// age >= 0, not > 0: a snapshot read in the same clock tick as its own publish
274+
// stamp legitimately reports age 0. Only a negative age is impossible here, and
275+
// that is all this check needs to catch.
276+
long age = first.ageNanos();
277+
assertTrue(age >= 0 && age < 1_000_000_000L, "ageNanos out of range: " + age);
278+
assertTrue(age <= System.nanoTime() - first.publishNanos(),
279+
"ageNanos " + age + " exceeds a delta measured later from the same stamp");
280+
281+
// A later publish replaces both halves together; there is no observable state in
282+
// which one half is from this publish and the other from the previous one.
283+
//
284+
// Wait for the clock to pass first.publishNanos() rather than sleeping a fixed
285+
// interval: System.nanoTime() may return the same tick for adjacent reads, so
286+
// Thread.sleep(2) does not guarantee the next publish gets a later stamp. Same
287+
// bounded-wait pattern as latestValueAndTimestampAdvanceTogether, so a broken
288+
// clock fails the assert instead of hanging.
289+
long deadline = System.nanoTime() + 50_000_000L;
290+
while (first.publishNanos() >= System.nanoTime() && System.nanoTime() <= deadline) {
291+
Thread.sleep(1);
292+
}
293+
orchestrator.publish("snap", "b");
294+
Topic.Latest<String> second = t.latest().orElseThrow();
295+
assertEquals("b", second.value());
296+
assertTrue(second.publishNanos() > first.publishNanos(),
297+
"a later publish must carry a later stamp");
298+
assertTrue(second.ageNanos() < first.ageNanos(),
299+
"the newer value must report a younger age");
300+
}
301+
243302
@Test
244303
void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() {
245304
// acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at

0 commit comments

Comments
 (0)