Skip to content

Commit 9624924

Browse files
authored
feat(topic): add Topic.latest() for a single-read value + timestamp (#31)
* feat(topic): add Topic.latest() for a single-read value + timestamp Checking whether a published value is stale needs 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(defaultValue); long age = snap.map(Topic.Latest::ageNanos).orElse(Long.MAX_VALUE); Additive only. No existing signature or behaviour changes. Callers that never check an age keep using `latestValueOr` -- the same single volatile read, and allocation-free. `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 by hand. `latestPublishNanos()` still returns 0 before the first publish, so that guard stays necessary for anyone composing by hand, and it is documented as a heuristic rather than a proof -- System.nanoTime() is permitted to return 0 too. It is stated 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. Its cross-check between `latest()` being empty and `latestPublishNanos()` being nonzero is gone, because a publish can legitimately land between those two reads. Suite: 86 tests, 0 failures. javadoc: no errors, no new warnings. * test(topic): drop absolute time bounds; document empty-read allocation - Remove the one-second age ceiling. A GC pause between the publish and the ageNanos() read legitimately makes the age arbitrarily large, so the bound imposed a scheduling deadline on the test. Both bounds are now relative to the clock. - Remove the nonzero publishNanos() assertion. System.nanoTime() has an arbitrary origin, so a published snapshot may carry 0; the present Optional already proves a publish happened. - latest() allocates an Optional wrapper only when it returns a value. Empty reads return the shared Optional.empty() and allocate nothing, which is now what the javadoc and the topics guide say. Suite 86/86, javadoc clean. * test(topic): bracket ageNanos() between two clock reads Three review findings on the snapshot tests: - ageNanos() was only bounded from above, so an implementation returning a constant 0 satisfied the check while reporting nothing about the real age. Bracket the call between two clock reads in both the unit test and the concurrency reader, and require the age to fall between both deltas from the snapshot's own stamp. - second.ageNanos() < first.ageNanos() sampled at two different moments, so a pause between the publish and the second read could make the newer value look older. Derive both ages from a single clock reading: from one now, a strictly later stamp is a strictly smaller age. The one-second wall-clock ceiling and the nonzero-stamp assertion reported earlier were already removed in the previous commit; these reviews landed on the older SHA. Suite 86/86, tear test green on a 2-core-pinned runner, javadoc clean. * docs(topic): unwrap the snapshot once in every example The value-and-age example called Optional.map() once per field, so each non-empty read allocated two more Optional wrappers and boxed the long age -- in a guide aimed at periodic loops, teaching the allocating form three times over (the latest() wrapper plus two map() calls). Every example now unwraps with isPresent()/get() and calls value() and ageNanos() directly: the topics guide, the javadoc and the CHANGELOG. * test(topic): bracket ageNanos() between clock reads taken around the call The previous bounds sampled a clock reading outside the ageNanos() call, so they could not bound the result in either direction: a pause between the sampling and the call breaks the lower bound, and the age correctly exceeds a reading taken before it. Both the unit test and the concurrency reader now take one clock reading immediately before the call and one immediately after, and require the age to fall between the two deltas from the snapshot's own stamp. That is the only bracketing that constrains the value, and it holds regardless of how long the call itself takes. Suite 86/86 across 8 consecutive runs, tear test green pinned to two cores, javadoc clean.
1 parent 2a8f572 commit 9624924

5 files changed

Lines changed: 281 additions & 55 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -96,13 +96,32 @@ listeners registered, `publish` costs a single volatile read.
9696
`0`, so a genuine publish can carry it too. The resulting over-reported age rejects
9797
a fresh value rather than admitting a stale one, so the failure direction is safe.
9898

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

104102
```java
105-
long stamp = topic.latestPublishNanos(); // first
103+
Optional<Topic.Latest<T>> snap = topic.latest();
104+
T v;
105+
long age;
106+
if (snap.isPresent()) {
107+
Topic.Latest<T> s = snap.get();
108+
v = s.value();
109+
age = s.ageNanos();
110+
} else {
111+
v = defaultValue;
112+
age = Long.MAX_VALUE;
113+
}
114+
```
115+
116+
Unwrap once rather than calling `map()` per field: each `map()` allocates its own
117+
`Optional` and boxes the `long`, which is avoidable garbage in a periodic loop.
118+
119+
`latestPublishNanos()` and `latestValueOr()` are unchanged and still correct. Composing
120+
them takes two calls, which can straddle a publish and leave the pair describing two
121+
different publishes; if you do compose them, **read the timestamp first**:
122+
123+
```java
124+
long stamp = topic.latestPublishNanos(); // read the timestamp FIRST
106125
T v = topic.latestValueOr(null); // then
107126
long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
108127
```

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

Lines changed: 116 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -46,15 +46,59 @@ 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+
/**
68+
* The published value.
69+
*
70+
* @return the value published in this snapshot
71+
*/
72+
public T value() {
73+
return value;
74+
}
75+
76+
/**
77+
* When this value was recorded.
78+
*
79+
* @return the {@link System#nanoTime()} at which {@link #value()} was recorded, on
80+
* the monotonic clock (arbitrary origin — only differences are meaningful)
81+
*/
82+
public long publishNanos() {
83+
return publishNanos;
84+
}
85+
86+
/**
87+
* How long ago this value was recorded, measured from its own stamp.
88+
*
89+
* <p>Unlike {@code System.nanoTime() - publishNanos()}, this is always the age of
90+
* <i>this</i> value and needs no {@code 0} sentinel handling.
91+
*
92+
* @return nanos elapsed since {@link #publishNanos()}, on the monotonic clock
93+
*/
94+
public long ageNanos() {
95+
return System.nanoTime() - publishNanos;
96+
}
97+
98+
@Override
99+
public String toString() {
100+
return "Latest[" + value + " @" + publishNanos + "]";
101+
}
58102
}
59103

60104
/**
@@ -74,27 +118,26 @@ public Class<T> type() {
74118
/**
75119
* The most recently published value, or empty if nothing has been published yet.
76120
*
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:
121+
* <p><b>Reading a value together with its age:</b> prefer {@link #latest()}, which
122+
* returns both from a single snapshot read. Composing them from two accessors works
123+
* only if you sample the timestamp <i>first</i>, which can only over-report the age:
80124
*
81125
* <pre>{@code
82126
* long stamp = topic.latestPublishNanos(); // read the timestamp FIRST
83-
* T v = topic.latestValueOr(null); // then the value
127+
* T v = topic.latestValueOr(null); // then
84128
* long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
85129
* }</pre>
86130
*
87-
* <p>Reading the timestamp first is what makes the pair safe: a publish landing
88-
* between the two calls can only over-report the age, never under-report it.
89-
*
90131
* <p>The {@code 0} guard handles "nothing published yet", where the stamp is still
91132
* its initial {@code 0} and subtracting it would yield the raw
92133
* {@link System#nanoTime()} reading. It is a heuristic, not a proof: the JLS permits
93134
* {@link System#nanoTime()} to return {@code 0}, so a genuine publish can carry a
94135
* {@code 0} stamp too. That case only over-reports the age, which rejects a fresh
95136
* value rather than admitting a stale one, so the failure direction is safe.
96137
*
97-
* Reading the value first is the unsafe order: a publish landing between the two
138+
* <p>{@link #latest()} removes the question entirely: one snapshot read, no sentinel.
139+
*
140+
* <p>Reading the value first is the unsafe order: a publish landing between the two
98141
* calls pairs the older value with the newer timestamp, so an age check on that
99142
* pair passes even though the value is stale.
100143
*
@@ -119,7 +162,7 @@ public Optional<T> latestValue() {
119162
*/
120163
public long latestPublishNanos() {
121164
Latest<T> snap = latest;
122-
return snap == null ? 0L : snap.publishNanos;
165+
return snap == null ? 0L : snap.publishNanos();
123166
}
124167

125168
/**
@@ -134,7 +177,8 @@ public long latestPublishNanos() {
134177
* publishes. The <b>write path takes the topic monitor</b>, around the clock sample,
135178
* one short-lived allocation and the store — see below. Each publish therefore
136179
* allocates; on the read side {@link #latestValueOr(Object)} allocates nothing,
137-
* while {@link #latestValue()} wraps its result in an {@link Optional}.
180+
* while {@link #latestValue()} and {@link #latest()} each wrap their result in an
181+
* {@link Optional}.
138182
*
139183
* @param value the value to record
140184
*/
@@ -177,14 +221,70 @@ boolean acceptsType(Class<?> other) {
177221
* Returns the most recent value, falling back to {@code defaultValue} if nothing has
178222
* been published yet. Convenience for {@code topic.latestValue().orElse(default)}.
179223
*
180-
* <p>Allocation-free — prefer this over {@link #latestValue()} on hot paths.
224+
* <p>Allocation-free — prefer this over {@link #latestValue()} on hot paths, and over
225+
* {@link #latest()} if you do not need the timestamp.
181226
*
182227
* @param defaultValue the value to return before the first publish
183228
* @return the latest value, or {@code defaultValue}
184229
*/
185230
public T latestValueOr(T defaultValue) {
186231
Latest<T> snap = latest;
187-
return snap == null ? defaultValue : snap.value;
232+
return snap == null ? defaultValue : snap.value();
233+
}
234+
235+
/**
236+
* The most recent value and the {@link System#nanoTime()} at which it was published,
237+
* as one consistent pair.
238+
*
239+
* <p><b>Prefer this over {@link #latestValue()} plus {@link #latestPublishNanos()}
240+
* when you are checking whether a value is stale.</b> The two accessors each read the
241+
* snapshot correctly, but a caller making two calls can straddle a publish between
242+
* them and end up holding a value from one publish and a timestamp from another:
243+
*
244+
* <pre>{@code
245+
* long stamp = topic.latestPublishNanos(); // publish A
246+
* ...publish B lands...
247+
* T v = topic.latestValueOr(null); // publish B's value, with A's stamp
248+
* long age = System.nanoTime() - stamp; // over-reports by a full publish interval
249+
* }</pre>
250+
*
251+
* <p>Reading the timestamp first makes that safe in one direction — the age can only
252+
* be over-reported, never under-reported — but over-reporting still costs you: a
253+
* fresh value gets rejected as stale. How often that happens depends on how long the
254+
* two reads take. Normally they are nanoseconds apart and a publish landing between
255+
* them is vanishingly rare. It stops being rare when the reader is preempted between
256+
* the two calls by a GC pause or the scheduler, since the window becomes
257+
* milliseconds. This accessor removes the window entirely.
258+
*
259+
* <p>Unwrap once rather than calling {@code map()} per field: each {@code map()} call
260+
* allocates its own {@link Optional} and boxes the {@code long} age, which is
261+
* avoidable garbage in a periodic loop.
262+
*
263+
* <pre>{@code
264+
* Optional<Topic.Latest<T>> snap = topic.latest();
265+
* T v;
266+
* long age;
267+
* if (snap.isPresent()) {
268+
* Topic.Latest<T> s = snap.get();
269+
* v = s.value();
270+
* age = s.ageNanos();
271+
* } else {
272+
* v = defaultValue;
273+
* age = Long.MAX_VALUE;
274+
* }
275+
* }</pre>
276+
*
277+
* <p>Returns empty before the first publish.
278+
*
279+
* <p><b>Non-empty calls allocate one small {@link Optional} wrapper.</b> Calls before
280+
* the first publish return the shared {@link Optional#empty()} instance and allocate
281+
* nothing. If you only want the value and never check its age, use
282+
* {@link #latestValueOr(Object)}, which never allocates.
283+
*
284+
* @return an {@link Optional} holding the latest value and its publish timestamp
285+
*/
286+
public Optional<Latest<T>> latest() {
287+
return Optional.ofNullable(latest);
188288
}
189289

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

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

Lines changed: 107 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,44 @@ 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, so bracket
139+
// the call between two clock reads and require the age to fall
140+
// between both deltas. A one-sided check lets an implementation
141+
// returning a constant 0 pass. Zero is legitimate (same tick).
142+
long ageBefore = System.nanoTime();
143+
long age = snap.ageNanos();
144+
long ageAfter = System.nanoTime();
145+
if (age < 0
146+
|| age < ageBefore - snap.publishNanos()
147+
|| age > ageAfter - snap.publishNanos()) {
132148
failure.compareAndSet(null, new AssertionError(
133-
"value " + v + " is older than the timestamp " + ts
134-
+ " the reader sampled; its publish ended at " + hi));
149+
"implausible age " + age + " for stamp " + snap.publishNanos()));
135150
return;
136151
}
137152
}
@@ -240,6 +255,76 @@ void acceptsValueClassTreatsPrimitivesAndWrappersAsEquivalent() {
240255
assertFalse(i.acceptsValueClass(double.class), "Double must not match Integer topic");
241256
}
242257

258+
@Test
259+
void latestReturnsValueAndTimestampFromOneConsistentPublish() throws Exception {
260+
// latest() exists so a staleness check reads one snapshot instead of two. It must
261+
// be empty before the first publish, agree with both single-field accessors, and
262+
// never need a 0 sentinel.
263+
Topic<String> t = orchestrator.getOrCreateTopic("snap", String.class);
264+
assertFalse(t.latest().isPresent(), "latest() must be empty before the first publish");
265+
266+
orchestrator.publish("snap", "a");
267+
Topic.Latest<String> first = t.latest().orElseThrow();
268+
assertEquals("a", first.value());
269+
assertEquals(first.publishNanos(), t.latestPublishNanos(),
270+
"latest() and latestPublishNanos() must describe the same publish");
271+
assertEquals("a", t.latestValueOr(null));
272+
273+
// No assertion that publishNanos() is nonzero: System.nanoTime() has an arbitrary
274+
// origin, so a published snapshot may legitimately carry 0. The present Optional
275+
// already proves a publish happened, which is the property worth checking here.
276+
277+
// ageNanos must be the age of THIS value, from ITS OWN stamp.
278+
//
279+
// Comparing ageNanos() against `now - publishNanos()` cannot establish that:
280+
// both sides recompute the same subtraction, so the relation holds for any pair
281+
// and a Latest carrying a foreign stamp still passes. What actually pins the
282+
// stamp to this value is the cross-check above -- latest().publishNanos() equals
283+
// latestPublishNanos() equals the stamp of the "a" publish -- combined with a
284+
// non-negative age and a stamp that is not in the future.
285+
//
286+
// Both bounds are relative to the clock, not absolute wall-clock deadlines: a
287+
// snapshot read in the same tick as its own publish legitimately reports age 0,
288+
// and a GC pause between the publish and this read can legitimately make the
289+
// age arbitrarily large.
290+
// The clock reads must bracket the ageNanos() call itself; a reading taken
291+
// outside it cannot bound the result in either direction, because a pause
292+
// between the reading and the call would break the relation.
293+
long ageBefore = System.nanoTime();
294+
long age = first.ageNanos();
295+
long ageAfter = System.nanoTime();
296+
assertTrue(age >= 0, "ageNanos must not be negative: " + age);
297+
assertTrue(age >= ageBefore - first.publishNanos()
298+
&& age <= ageAfter - first.publishNanos(),
299+
"ageNanos " + age + " is not the age of stamp " + first.publishNanos()
300+
+ " measured between " + ageBefore + " and " + ageAfter);
301+
302+
// A later publish replaces both halves together; there is no observable state in
303+
// which one half is from this publish and the other from the previous one.
304+
//
305+
// Wait for the clock to pass first.publishNanos() rather than sleeping a fixed
306+
// interval: System.nanoTime() may return the same tick for adjacent reads, so
307+
// Thread.sleep(2) does not guarantee the next publish gets a later stamp. Same
308+
// bounded-wait pattern as latestValueAndTimestampAdvanceTogether, so a broken
309+
// clock fails the assert instead of hanging.
310+
long deadline = System.nanoTime() + 50_000_000L;
311+
while (first.publishNanos() >= System.nanoTime() && System.nanoTime() <= deadline) {
312+
Thread.sleep(1);
313+
}
314+
orchestrator.publish("snap", "b");
315+
Topic.Latest<String> second = t.latest().orElseThrow();
316+
assertEquals("b", second.value());
317+
assertTrue(second.publishNanos() > first.publishNanos(),
318+
"a later publish must carry a later stamp");
319+
// Derive both ages from ONE clock reading rather than calling ageNanos() twice.
320+
// Sampled at different moments the comparison depends on how long each call took,
321+
// so a pause between the publish and the second read could make the newer value
322+
// look older. From a single now, a strictly later stamp is a strictly smaller age.
323+
long laterNow = System.nanoTime();
324+
assertTrue(laterNow - second.publishNanos() < laterNow - first.publishNanos(),
325+
"the newer value must report a younger age");
326+
}
327+
243328
@Test
244329
void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() {
245330
// acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at

0 commit comments

Comments
 (0)