Skip to content
Closed
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
18 changes: 13 additions & 5 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<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);
```

`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
113 changes: 97 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,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)}.
*
* <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>Returns empty before the first publish.
*
* <p><b>Allocates one small {@link Optional} wrapper per call.</b> 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<T>> latest() {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
return Optional.ofNullable(latest);
Comment thread
sourcery-ai[bot] marked this conversation as resolved.
}

/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */
Expand Down
103 changes: 81 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,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<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. 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()));
Comment thread
sourcery-ai[bot] marked this conversation as resolved.
return;
}
}
Expand Down Expand Up @@ -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<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));
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<String> 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
Expand Down
15 changes: 13 additions & 2 deletions website/src/content/docs/api/index.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,19 @@ Base class for robot logic. `protected final Orchestrator orchestrator` is avail
| --- | --- |
| `name()` / `type()` | identity |
| `latestValue()` | `Optional<T>` |
| `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<Topic.Latest<T>>`; value and timestamp from one snapshot read |

### `Topic.Latest<T>`

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

Expand Down
Loading