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
44 changes: 44 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,50 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- **`@SubscribedTo` dispatch no longer re-boxes the parameter type per message.**
The primitive-to-wrapper normalisation is computed once at bind time instead
of on every delivered message. Equivalent to the previous conditional check.
- **`Topic` latest-value reads no longer synchronize.** `latestValue()` and
`latestPublishNanos()` are plain volatile reads. The latest value and its
timestamp are published together as one immutable pair through a single volatile
field, so a reader always sees a value and a timestamp from the same publish.
`publish` is reachable from the OpMode loop, the hardware thread, and the callback
pool, so two publishers really can overlap.

The write path still takes the topic monitor, but only around the clock sample,
one short-lived allocation and the store. Two separate volatile fields would let
two publishers interleave between the value write and the timestamp write;
sampling the clock outside mutual exclusion would let a preempted publisher
install an older pair after a newer one. Both regressions were reproduced and
fixed; keeping the monitor on the write path preserves the ordering the previous
`synchronized` body provided.

Each publish therefore allocates one small short-lived pair on the write path.
The read path is unchanged in allocation terms: `latestValueOr()` allocates
nothing, while `latestValue()` wraps its result in an `Optional` as it did
before.

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

`latestPublishNanos()` still returns `0` before the first publish, unchanged. That
`0` is a heuristic rather than a proof — `System.nanoTime()` is permitted to return
`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:

```java
long stamp = topic.latestPublishNanos(); // first
Comment thread
kody-ai[bot] marked this conversation as resolved.
T v = topic.latestValueOr(null); // then
long age = stamp == 0L ? Long.MAX_VALUE : System.nanoTime() - stamp;
```

Reading the value first and the timestamp second can pair an older value with a
newer timestamp, so an age check on that pair passes even though the value is
stale. Reading the timestamp first can only over-report the age, never under-report
it.

## [0.4.0] - 2026-09-12

Expand Down
31 changes: 7 additions & 24 deletions src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -147,11 +147,9 @@ public <T> Topic<T> getOrCreateTopic(String topicName, Class<T> type) {
Objects.requireNonNull(topicName, "topicName");
Objects.requireNonNull(type, "type");

Class<?> normalized = boxed(type);

Topic<?> existing = topics.get(topicName);
if (existing != null) {
if (!boxed(existing.type()).isAssignableFrom(normalized)) {
if (!existing.acceptsType(type)) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' already exists with type "
+ existing.type().getName() + ", cannot re-create as "
Expand All @@ -163,7 +161,7 @@ public <T> Topic<T> getOrCreateTopic(String topicName, Class<T> type) {
Topic<T> created = new Topic<>(topicName, type);
Topic<?> prior = topics.putIfAbsent(topicName, created);
if (prior != null) {
if (!boxed(prior.type()).isAssignableFrom(normalized)) {
if (!prior.acceptsType(type)) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' already exists with type "
+ prior.type().getName() + ", cannot re-create as "
Expand All @@ -175,20 +173,6 @@ public <T> Topic<T> getOrCreateTopic(String topicName, Class<T> type) {
return created;
}

/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */
private static Class<?> boxed(Class<?> c) {
if (!c.isPrimitive()) return c;
if (c == int.class) return Integer.class;
if (c == long.class) return Long.class;
if (c == double.class) return Double.class;
if (c == float.class) return Float.class;
if (c == boolean.class) return Boolean.class;
if (c == byte.class) return Byte.class;
if (c == short.class) return Short.class;
if (c == char.class) return Character.class;
return c;
}

@Override
public Optional<Topic<?>> findTopic(String topicName) {
return Optional.ofNullable(topics.get(topicName));
Expand All @@ -198,7 +182,7 @@ public Optional<Topic<?>> findTopic(String topicName) {
@SuppressWarnings("unchecked")
public <T> Optional<Topic<T>> findTopic(String topicName, Class<T> type) {
Topic<?> t = topics.get(topicName);
if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty();
if (t == null || !t.acceptsType(type)) return Optional.empty();
return Optional.of((Topic<T>) t);
}

Expand All @@ -217,14 +201,13 @@ public <T> void publish(String topicName, T value) {

// Lazily create the topic from the value's runtime type. This matches
// Heron's behavior: publishers don't have to pre-register topics.
Class<?> valueType = value.getClass();
Topic<?> topic = topics.get(topicName);
if (topic == null) {
topic = getOrCreateTopic(topicName, valueType);
} else if (!boxed(topic.type()).isAssignableFrom(boxed(valueType))) {
topic = getOrCreateTopic(topicName, value.getClass());
} else if (!topic.acceptsValueClass(value.getClass())) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' is typed " + topic.type().getName()
+ " but publish got " + valueType.getName());
+ " but publish got " + value.getClass().getName());
}

((Topic<Object>) topic).recordLatest(value);
Expand Down Expand Up @@ -312,7 +295,7 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) {
@SuppressWarnings("unchecked")
public <T> Optional<T> getLatestValue(String topicName, Class<T> type) {
Topic<?> t = topics.get(topicName);
if (t == null || !boxed(t.type()).isAssignableFrom(boxed(type))) return Optional.empty();
if (t == null || !t.acceptsType(type)) return Optional.empty();
return (Optional<T>) t.latestValue();
}

Expand Down
147 changes: 133 additions & 14 deletions src/main/java/com/aaravlabs/synapse/Topic.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,43 @@ public final class Topic<T> {
private final String name;
private final Class<T> type;

// Guarded by `this` for write, volatile for the read in publish().
private volatile T latest;
private volatile long latestPublishNanos;
/**
* {@code type} normalized to its wrapper class. Cached at construction so the
* per-publish type check is a single field read instead of a primitive-boxing
* chain.
*/
private final Class<?> boxedType;

/**
* The most recently published value and its timestamp, as one immutable pair.
*
* <p>A single volatile reference rather than two volatile fields: with two, two
* concurrent publishers can interleave between the value write and the timestamp
* write and leave the pair describing two different publishes. Since
* {@link Orchestrator#publish(String, Object)} is reachable from the OpMode
* loop, the hardware thread, and the callback pool, that is not a theoretical
* interleaving.
*
* <p>Immutable so that a reader holding a reference sees a self-consistent pair
* with no further synchronization.
*/
private volatile Latest<T> latest;

Topic(String name, Class<T> type) {
this.name = name;
this.type = type;
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;

Latest(T value, long publishNanos) {
this.value = value;
this.publishNanos = publishNanos;
}
}

/**
Expand All @@ -44,42 +74,131 @@ 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:
*
* <pre>{@code
* long stamp = topic.latestPublishNanos(); // read the timestamp FIRST
* T v = topic.latestValueOr(null); // then the value
* 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
* calls pairs the older value with the newer timestamp, so an age check on that
* pair passes even though the value is stale.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
*
* @return an {@link Optional} holding the latest value
*/
public synchronized Optional<T> latestValue() {
return Optional.ofNullable(latest);
public Optional<T> latestValue() {
return Optional.ofNullable(latestValueOr(null));
}

/**
* Wall-clock nanos at which {@link #latestValue()} was last updated.
* Nanos at which {@link #latestValue()} was last updated, on the
* {@link System#nanoTime()} monotonic clock (arbitrary origin — only differences
* are meaningful).
*
* @return publish timestamp in nanoseconds ({@code System.nanoTime()} clock)
* <p>Sample this before {@link #latestValue()} when the two are used together as a
* staleness check — see {@link #latestValue()} for the ordering rule.
*
* @return publish timestamp in nanoseconds, or {@code 0} if nothing has been published
* yet. Note that {@code 0} is also a value {@link System#nanoTime()} is
* permitted to return, so this cannot be used on its own to prove that no
* publish has occurred — see {@link #latestValue()}.
*/
public synchronized long latestPublishNanos() {
return latestPublishNanos;
public long latestPublishNanos() {
Latest<T> snap = latest;
return snap == null ? 0L : snap.publishNanos;
}

/**
* Record a new latest value. Called by the orchestrator immediately before notifying
* subscribers.
*
* <p>The <b>read path is lock-free</b>: the value and its timestamp are published
* together as one immutable {@link Latest} through a single volatile reference, so a
* concurrent reader always sees a value and a timestamp from the same publish.
* Splitting them across two volatile fields would let two concurrent publishers
* interleave between the writes and leave the pair describing two different
* 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}.
*
* @param value the value to record
*/
synchronized void recordLatest(T value) {
this.latest = value;
this.latestPublishNanos = System.nanoTime();
void recordLatest(T value) {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
// The snapshot is sampled and installed under this topic's monitor so that
// install order matches timestamp order. Without it, a publisher that samples
// nanoTime() and is then preempted installs an OLDER snapshot after a newer
// one, and a reader can end up holding an older value beside a newer
// timestamp -- an age computed from that stamp is under-reported, so a
// staleness check can pass a value that is already stale.
//
// The critical section is a clock read, an allocation and a volatile store --
// NOT "two field writes and a clock read". Each publish allocates one Latest,
// which is not scalar-replaceable because it escapes into a volatile field, so
// garbage scales with publish rate across every topic.
//
// Hoisting the allocation above the monitor was measured and is not a win: it
// gains ~12 ns/publish uncontended but loses ~10-20 ns under four publishers,
// since it lengthens the window in which a stalled publisher holds a snapshot
// whose stamp is already stale. The allocation stays inside.
synchronized (this) {
Comment thread
IamCoder18 marked this conversation as resolved.
this.latest = new Latest<>(value, System.nanoTime());
Comment thread
kody-ai[bot] marked this conversation as resolved.
}
}
Comment thread
kody-ai[bot] marked this conversation as resolved.
Comment thread
kody-ai[bot] marked this conversation as resolved.

/**
* @return true if a value of runtime class {@code actual} may be published here.
* Delegates to {@link #acceptsType} so the two cannot drift apart.
*/
boolean acceptsValueClass(Class<?> actual) {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
return acceptsType(actual);
}

/** @return true if {@code other} is type-compatible with this topic. */
boolean acceptsType(Class<?> other) {
return boxedType.isAssignableFrom(box(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.
*
* @param defaultValue the value to return before the first publish
* @return the latest value, or {@code defaultValue}
*/
public T latestValueOr(T defaultValue) {
T v = latest;
return v != null ? v : defaultValue;
Latest<T> snap = latest;
return snap == null ? defaultValue : snap.value;
}

/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */
private static Class<?> box(Class<?> c) {
if (!c.isPrimitive()) return c;
if (c == int.class) return Integer.class;
if (c == long.class) return Long.class;
if (c == double.class) return Double.class;
if (c == float.class) return Float.class;
if (c == boolean.class) return Boolean.class;
if (c == byte.class) return Byte.class;
if (c == short.class) return Short.class;
if (c == char.class) return Character.class;
return c;
}

@Override
Expand Down
Loading
Loading