Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 24 additions & 5 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -96,13 +96,32 @@ listeners registered, `publish` costs a single volatile read.
`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<Topic.Latest<T>> snap = topic.latest();
T v;
long age;
if (snap.isPresent()) {
Topic.Latest<T> s = snap.get();
v = s.value();
age = s.ageNanos();
} else {
v = defaultValue;
age = Long.MAX_VALUE;
}
```

Unwrap once rather than calling `map()` per field: each `map()` allocates its own
`Optional` and boxes the `long`, which is avoidable garbage in a periodic loop.

`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;
```
Expand Down
132 changes: 116 additions & 16 deletions src/main/java/com/aaravlabs/synapse/Topic.java
Original file line number Diff line number Diff line change
Expand Up @@ -46,15 +46,59 @@ public final class Topic<T> {
this.boxedType = box(type);
}

/** One publish's value and the {@link System#nanoTime()} at which it was recorded. */
private static final class Latest<T> {
final T value;
final long publishNanos;
/**
* One publish's value and the {@link System#nanoTime()} at which it was recorded.
*
* <p>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 <T> the message type carried by this topic
*/
public static final class Latest<T> {

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.
*
* <p>Unlike {@code System.nanoTime() - publishNanos()}, this is always the age of
* <i>this</i> 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 + "]";
}
}

/**
Expand All @@ -74,27 +118,26 @@ public Class<T> type() {
/**
* The most recently published value, or empty if nothing has been published yet.
*
* <p><b>Reading a value together with its age:</b> each accessor reads a
* self-consistent pair, but two separate calls can still straddle a publish. Sample
* the timestamp <i>first</i>, which can only over-report the value's age:
* <p><b>Reading a value together with its age:</b> prefer {@link #latest()}, which
* returns both from a single snapshot read. Composing them from two accessors works
* only if you sample the timestamp <i>first</i>, which can only over-report the age:
*
* <pre>{@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;
* }</pre>
*
* <p>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.
*
* <p>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
* <p>{@link #latest()} removes the question entirely: one snapshot read, no sentinel.
*
* <p>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.
*
Expand All @@ -119,7 +162,7 @@ public Optional<T> latestValue() {
*/
public long latestPublishNanos() {
Latest<T> snap = latest;
return snap == null ? 0L : snap.publishNanos;
return snap == null ? 0L : snap.publishNanos();
}

/**
Expand All @@ -134,7 +177,8 @@ public long latestPublishNanos() {
* publishes. The <b>write path takes the topic monitor</b>, 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
*/
Expand Down Expand Up @@ -177,14 +221,70 @@ 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)}.
*
* <p>Allocation-free — prefer this over {@link #latestValue()} on hot paths.
* <p>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<T> 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.
*
* <p><b>Prefer this over {@link #latestValue()} plus {@link #latestPublishNanos()}
* when you are checking whether a value is stale.</b> 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:
*
* <pre>{@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
* }</pre>
*
* <p>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.
*
* <p>Unwrap once rather than calling {@code map()} per field: each {@code map()} call
* allocates its own {@link Optional} and boxes the {@code long} age, which is
* avoidable garbage in a periodic loop.
*
* <pre>{@code
* Optional<Topic.Latest<T>> snap = topic.latest();
* T v;
* long age;
* if (snap.isPresent()) {
* Topic.Latest<T> s = snap.get();
* v = s.value();
* age = s.ageNanos();
* } else {
* v = defaultValue;
* age = Long.MAX_VALUE;
* }
* }</pre>
*
* <p>Returns empty before the first publish.
*
* <p><b>Non-empty calls allocate one small {@link Optional} wrapper.</b> Calls before
* the first publish return the shared {@link Optional#empty()} instance and allocate
* nothing. If you only want the value and never check its age, use
* {@link #latestValueOr(Object)}, which never allocates.
*
* @return an {@link Optional} holding the latest value and its publish timestamp
*/
public Optional<Latest<T>> latest() {
return Optional.ofNullable(latest);
}

/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */
Expand Down
129 changes: 107 additions & 22 deletions src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -106,32 +109,44 @@ 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<Integer> 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, so bracket
// the call between two clock reads and require the age to fall
// between both deltas. A one-sided check lets an implementation
// returning a constant 0 pass. Zero is legitimate (same tick).
long ageBefore = System.nanoTime();
long age = snap.ageNanos();
long ageAfter = System.nanoTime();
if (age < 0
|| age < ageBefore - snap.publishNanos()
|| age > ageAfter - 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;
}
}
Expand Down Expand Up @@ -240,6 +255,76 @@ 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<String> t = orchestrator.getOrCreateTopic("snap", String.class);
assertFalse(t.latest().isPresent(), "latest() must be empty before the first publish");

orchestrator.publish("snap", "a");
Topic.Latest<String> 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));

// No assertion that publishNanos() is nonzero: System.nanoTime() has an arbitrary
// origin, so a published snapshot may legitimately carry 0. The present Optional
// already proves a publish happened, which is the property worth checking here.

// ageNanos must be the age of THIS value, from ITS OWN stamp.
//
// Comparing ageNanos() against `now - publishNanos()` cannot establish that:
// both sides recompute the same subtraction, so the relation holds for any pair
// and a Latest carrying a foreign stamp still passes. What actually pins the
// stamp to this value is the cross-check above -- latest().publishNanos() equals
// latestPublishNanos() equals the stamp of the "a" publish -- combined with a
// non-negative age and a stamp that is not in the future.
Comment on lines +279 to +284

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

kody code-review Bug high

Tautological stamp validation in src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java leaves Latest.publishNanos() unverified for the "a" publish: latest() (Topic.java:287, Optional.ofNullable(latest)) and latestPublishNanos() (Topic.java:163-166, latest.publishNanos()) read the same volatile field, so the assertEquals at lines 269-270 passes for any stamp and the replacement age assertion cannot substitute. The bracket at lines 297-300 remains offset-invariant because a Latest built with stamp + K shifts age and both bounds by exactly K, letting a wrong-but-consistent timestamp—the staleness failure latest() exists to prevent—pass the whole method; bracket orchestrator.publish("snap", "a") with same-thread System.nanoTime() reads and assert that first.publishNanos() falls between them, the only check that rejects a constant offset.

// The cross-check above only proves the two accessors read the same field; it is
// tautological for the stamp's identity, so the stamp is pinned here instead by
// bracketing the publish itself between two clock reads on this thread.
long publishBefore = System.nanoTime();
orchestrator.publish("snap", "a");
Topic.Latest<String> first = t.latest().orElseThrow();
long publishAfter = System.nanoTime();
assertTrue(first.publishNanos() >= publishBefore
                && first.publishNanos() <= publishAfter,
        "stamp " + first.publishNanos() + " is not the clock reading taken during the publish"
                + " between " + publishBefore + " and " + publishAfter);

// ageNanos must be the age of THIS value, from ITS OWN stamp: a constant 0 or a
// constant offset would fall outside this bracket, while a GC pause between the
// reads only moves T inside it.
Prompt for LLM

File src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java:

Line 279 to 284:

Tautological stamp validation in `src/test/java/com/aaravlabs/synapse/TopicLockFreeTest.java` leaves `Latest.publishNanos()` unverified for the `"a"` publish: `latest()` (`Topic.java:287`, `Optional.ofNullable(latest)`) and `latestPublishNanos()` (`Topic.java:163-166`, `latest.publishNanos()`) read the same volatile field, so the `assertEquals` at lines 269-270 passes for any stamp and the replacement age assertion cannot substitute. The bracket at lines 297-300 remains offset-invariant because a `Latest` built with `stamp + K` shifts `age` and both bounds by exactly K, letting a wrong-but-consistent timestamp—the staleness failure `latest()` exists to prevent—pass the whole method; bracket `orchestrator.publish("snap", "a")` with same-thread `System.nanoTime()` reads and assert that `first.publishNanos()` falls between them, the only check that rejects a constant offset.

Suggested Code:

        // The cross-check above only proves the two accessors read the same field; it is
        // tautological for the stamp's identity, so the stamp is pinned here instead by
        // bracketing the publish itself between two clock reads on this thread.
        long publishBefore = System.nanoTime();
        orchestrator.publish("snap", "a");
        Topic.Latest<String> first = t.latest().orElseThrow();
        long publishAfter = System.nanoTime();
        assertTrue(first.publishNanos() >= publishBefore
                        && first.publishNanos() <= publishAfter,
                "stamp " + first.publishNanos() + " is not the clock reading taken during the publish"
                        + " between " + publishBefore + " and " + publishAfter);

        // ageNanos must be the age of THIS value, from ITS OWN stamp: a constant 0 or a
        // constant offset would fall outside this bracket, while a GC pause between the
        // reads only moves T inside it.

Talk to Kody by mentioning @kody

Was this suggestion helpful? React with 👍 or 👎 to help Kody learn from this interaction.

​

​

//
// Both bounds are relative to the clock, not absolute wall-clock deadlines: a
// snapshot read in the same tick as its own publish legitimately reports age 0,
// and a GC pause between the publish and this read can legitimately make the
// age arbitrarily large.
// The clock reads must bracket the ageNanos() call itself; a reading taken
// outside it cannot bound the result in either direction, because a pause
// between the reading and the call would break the relation.
long ageBefore = System.nanoTime();
long age = first.ageNanos();
long ageAfter = System.nanoTime();
assertTrue(age >= 0, "ageNanos must not be negative: " + age);
assertTrue(age >= ageBefore - first.publishNanos()
&& age <= ageAfter - first.publishNanos(),
"ageNanos " + age + " is not the age of stamp " + first.publishNanos()
+ " measured between " + ageBefore + " and " + ageAfter);

// 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<String> second = t.latest().orElseThrow();
assertEquals("b", second.value());
assertTrue(second.publishNanos() > first.publishNanos(),
"a later publish must carry a later stamp");
// Derive both ages from ONE clock reading rather than calling ageNanos() twice.
// Sampled at different moments the comparison depends on how long each call took,
// so a pause between the publish and the second read could make the newer value
// look older. From a single now, a strictly later stamp is a strictly smaller age.
long laterNow = System.nanoTime();
assertTrue(laterNow - second.publishNanos() < laterNow - first.publishNanos(),
"the newer value must report a younger age");
}

@Test
void acceptsTypeIsEquivalentToTheOldBoxedAssignableFromCheck() {
// acceptsType() replaced `boxed(t.type()).isAssignableFrom(boxed(type))` at
Expand Down
Loading
Loading