Skip to content
Open
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
10 changes: 10 additions & 0 deletions .github/workflows/docker.yml
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,16 @@ jobs:
username: ${{ github.actor }}
password: ${{ secrets.GITHUB_TOKEN }}

# IMAGE is set through GITHUB_ENV by each build job's own prepare step, so it does
# not cross the job boundary. Without recomputing it here the manifest step sees an
# empty IMAGE and buildx rejects "@sha256:..." as an invalid reference, which is why
# the published multi-arch tag silently stopped updating.
- name: Prepare image name
env:
SOURCE: ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}
run: |
echo "IMAGE=${SOURCE,,}" >> "$GITHUB_ENV"

- name: Extract metadata
id: meta
uses: docker/metadata-action@902fa8ec7d6ecbf8d84d538b9b233a880e428804 # v5.7.0
Expand Down
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,26 @@ listeners registered, `publish` costs a single volatile read.
stale. Reading the timestamp first can only over-report the age, never under-report
it.

### Fixed

- **`unregisterNode` and `unsubscribe` now stop `@OnHardwareThread` subscribers.**
Routing a subscription to the hardware thread replaced the entry in the subscriber
list, but the subscription still reported the unwrapped handler, so removal searched
for something the list no longer held and silently did nothing. An unregistered node
kept receiving publishes — it could still drive hardware after being removed.
- **`subscribe` works with primitive classes.** `int.class.isInstance(x)` is false for
every argument, so casting with the raw type made every delivery throw
`ClassCastException`, which was swallowed into a log line: the subscriber registered,
never fired, and nothing said why. Primitive and wrapper types are now treated
interchangeably here, as they already were on the publish path.
- **Typed lookups no longer hand back a topic that cannot be cast to the requested
type.** An `Object` topic accepts a request for `String`, and the unchecked cast then
failed later at the caller's own line — after `isPresent()` had reported a value. The
request is now refused. `getOrCreateTopic` throws with a message naming the actual
type; `findTopic` and `getLatestValue` return empty. Narrowing to a type the
topic really holds still works, and `subscribeRaw` (used by the annotation binder)
still accepts an `Object` topic.

## [0.4.0] - 2026-09-12

### Changed
Expand Down
165 changes: 136 additions & 29 deletions src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -145,35 +145,72 @@ private ThreadFactory makeThreadFactory(String prefix) {
@Override
@SuppressWarnings("unchecked")
public <T> Topic<T> getOrCreateTopic(String topicName, Class<T> type) {
return (Topic<T>) getOrCreateTopic(topicName, type, true);
}

/**
* Topic lookup that skips the cast-safety check for callers whose handler receives
* {@code Object}. Used by the annotation binder and by {@link #publish}, where an
* Object-typed topic is a supported configuration rather than a mistake, and where
* the caller never receives the topic as a {@code Topic<T>} it could cast wrongly.
*/
private Topic<?> getOrCreateTopicUnchecked(String topicName, Class<?> type) {
return getOrCreateTopic(topicName, type, false);
}

/**
* The single topic-creation path. {@code castSafe} selects whether an existing topic
* that accepts {@code type} may still be handed back.
*
* <p>One method rather than two: the {@code exists} / {@code putIfAbsent} race
* branches, the error strings, and the creation log have to stay identical, and a
* duplicated body is where the two copies silently drift.
*/
@SuppressWarnings("unchecked")
private <T> Topic<T> getOrCreateTopic(String topicName, Class<T> type, boolean castSafe) {
Objects.requireNonNull(topicName, "topicName");
Objects.requireNonNull(type, "type");

Topic<?> existing = topics.get(topicName);
if (existing != null) {
if (!existing.acceptsType(type)) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' already exists with type "
+ existing.type().getName() + ", cannot re-create as "
+ type.getName());
}
return (Topic<T>) existing;
return (Topic<T>) validateExisting(topicName, type, existing, castSafe);
}

Topic<T> created = new Topic<>(topicName, type);
Topic<?> prior = topics.putIfAbsent(topicName, created);
if (prior != null) {
if (!prior.acceptsType(type)) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' already exists with type "
+ prior.type().getName() + ", cannot re-create as "
+ type.getName());
}
return (Topic<T>) prior;
return (Topic<T>) validateExisting(topicName, type, prior, castSafe);
}
log.info(name, "created topic " + created);
return created;
}

/**
* Check an already-registered topic against a requested type. Shared by the
* {@code topics.get} and {@code putIfAbsent} branches of {@link #getOrCreateTopic},
* which must agree on the error text.
*/
private Topic<?> validateExisting(String topicName, Class<?> type,
Topic<?> existing, boolean castSafe) {
if (!existing.acceptsType(type)) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' already exists with type "
+ existing.type().getName() + ", cannot re-create as "
+ type.getName());
}
if (castSafe && !existing.safelyReturnsAs(type)) {
// compatible, but this topic may hold values the caller cannot cast to T.
// Returning it anyway defers a ClassCastException to an unrelated line.
throw new IllegalArgumentException(
"Topic '" + topicName + "' holds "
+ existing.type().getSimpleName() + ", which cannot be "
+ "returned as " + type.getName()
+ "; read it via findTopic(name) and treat the value as "
+ existing.type().getName());
}
return existing;
}

@Override
public Optional<Topic<?>> findTopic(String topicName) {
return Optional.ofNullable(topics.get(topicName));
Expand All @@ -183,7 +220,9 @@ 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 || !t.acceptsType(type)) return Optional.empty();
if (t == null || !t.acceptsType(type) || !t.safelyReturnsAs(type)) {
return Optional.empty();
}
return Optional.of((Topic<T>) t);
}

Expand Down Expand Up @@ -329,9 +368,16 @@ 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.
// Unchecked on purpose: publish immediately downcasts to Topic<Object> and never
// exposes the topic to a caller as a Topic<T>, so the cast-safety rule that
// protects typed lookups would only make this path throw. Concretely, a binder
// registering a zero-argument @SubscribedTo handler installs an Object-typed
// topic; if that lands between topics.get and the create below, a cast-safe
// getOrCreateTopic would refuse String.isAssignableFrom(Object) == false and
// silently drop a publish that used to succeed.
Topic<?> topic = topics.get(topicName);
if (topic == null) {
topic = getOrCreateTopic(topicName, value.getClass());
topic = getOrCreateTopicUnchecked(topicName, value.getClass());
} else if (!topic.acceptsValueClass(value.getClass())) {
throw new IllegalArgumentException(
"Topic '" + topicName + "' is typed " + topic.type().getName()
Expand Down Expand Up @@ -365,7 +411,16 @@ public <T> Subscription subscribe(String topicName, Class<T> type,
Consumer<? super T> handler) {
Objects.requireNonNull(handler, "handler");
Topic<T> topic = getOrCreateTopic(topicName, type);
MessageHandler wrapped = msg -> handler.accept(type.cast(msg));
// Class.cast() throws for any non-null argument when the class is primitive -- it
// performs an isInstance check, and isInstance is false for every wrapper value --
// so casting with the raw type made every delivery throw ClassCastException, which
// dispatchCallback swallowed into a log line -- the subscriber was registered,
// never fired, and nothing else indicated why. Box the type for the runtime check;
// Consumer<? super T> erases to accept(Object), so no cast on T is needed and
// none would survive erasure anyway.
final Class<?> boxedType = Topic.box(type);
@SuppressWarnings("unchecked")
MessageHandler wrapped = msg -> ((Consumer<Object>) handler).accept(boxedType.cast(msg));
SubscriberList list = subscribers.computeIfAbsent(topic, k -> new SubscriberList());
list.add(wrapped);
return new Subscription(topic, wrapped, this);
Expand All @@ -374,7 +429,11 @@ public <T> Subscription subscribe(String topicName, Class<T> type,
/** Raw subscribe used by the annotation binder where the parameter type is reflective. */
public Subscription subscribeRaw(String topicName, Class<?> type,
java.util.function.Consumer<Object> handler) {
Topic<?> topic = getOrCreateTopic(topicName, type);
// Deliberately no safelyReturnsAs check. This path exists for the annotation
// binder, which uses Object as the topic type when handlers of differing
// parameter types share a topic and does its own per-message isInstance filter.
// The handler receives Object, so a wider topic is correct rather than unsafe.
Topic<?> topic = getOrCreateTopicUnchecked(topicName, type);
MessageHandler wrapped = handler::accept;
SubscriberList list = subscribers.computeIfAbsent(topic, k -> new SubscriberList());
list.add(wrapped);
Expand All @@ -384,7 +443,22 @@ public Subscription subscribeRaw(String topicName, Class<?> type,
void removeSubscription(Subscription sub) {
Topic<?> t = sub.topic();
SubscriberList list = subscribers.get(t);
if (list != null) list.remove(sub.handler());
// markSubscriptionAsHardwareThreaded may have replaced the list entry with a
// wrapper, so look for whichever handler is actually registered. Searching only
// sub.handler() missed the wrapper and left an unregistered @OnHardwareThread
// subscriber running -- it kept receiving publishes, so an unregistered node
// could still drive the hardware.
//
// The lookup and the removal share hardwareRerouteLock with the mark, because
// the wrapper is recorded in two steps (list.replace, then the map put) and a
// removal landing between them would read an empty map, fall back to the
// original handler the replace had already removed, and no-op -- re-registering
// the exact leak this lookup exists to close.
synchronized (hardwareRerouteLock) {
MessageHandler registered = hardwareRerouted.getOrDefault(sub, sub.handler());
if (list != null) list.remove(registered);
hardwareRerouted.remove(sub);
}
}

/**
Expand All @@ -407,14 +481,34 @@ public void markSubscriptionAsHardwareThreaded(Subscription sub) {
}
});
};
// Replace in the list. Subscription object's handler() returns the
// original so removeSubscription still works correctly.
list.replace(original, wrapped);
// Stash the wrapped handler so subscribers still see the original via
// sub.handler() (which is unchanged), but the dispatch uses wrapped.
hardwareRerouted.put(sub, wrapped);
// Replace in the list. sub.handler() keeps returning the original so callers see
// a stable identity, and hardwareRerouted records which handler actually sits in
// the list so removeSubscription can find and remove exactly that one.
//
// Record only when the replacement actually happened. The method is public and
// the Subscription stays reachable afterwards, so a second mark is reachable:
// replace() looks for sub.handler(), which the first mark already took out of the
// list, so it no-ops. Recording the second wrapper anyway would point
// removeSubscription at a handler the list never held, leaving the first wrapper
// registered forever.
//
// Same lock as removeSubscription, so a concurrent removal cannot observe the
// list with the wrapper installed but the map still holding the original.
synchronized (hardwareRerouteLock) {
if (list.replace(original, wrapped)) {
hardwareRerouted.put(sub, wrapped);
}
}
}

/**
* Guards the two-step hardware reroute bookkeeping ({@link SubscriberList#replace}
* plus the {@link #hardwareRerouted} entry) against {@link #removeSubscription}.
* Rare by construction -- registration and removal both happen off the publish path
* -- and held only for the duration of that bookkeeping, never across dispatch.
*/
private final Object hardwareRerouteLock = new Object();

private final java.util.Map<Subscription, MessageHandler> hardwareRerouted = new ConcurrentHashMap<>();

// ---- fetch latest ----------------------------------------------------
Expand All @@ -423,7 +517,13 @@ 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 || !t.acceptsType(type)) return Optional.empty();
// safelyReturnsAs, not just acceptsType: an Object topic accepts a request for
// Integer, but handing back Optional<Integer> when the stored value is a String
// moves a ClassCastException to the caller's own site, after isPresent() has
// already told them the value is there.
if (t == null || !t.acceptsType(type) || !t.safelyReturnsAs(type)) {
return Optional.empty();
}
return (Optional<T>) t.latestValue();
}

Expand Down Expand Up @@ -687,17 +787,24 @@ void remove(MessageHandler h) {
snapshot = next;
}
}
/** Atomically replace {@code old} with {@code next}. No-op if not found. */
void replace(MessageHandler old, MessageHandler next) {
/**
* Atomically replace {@code old} with {@code next}.
*
* @return true if {@code old} was present and is now {@code next}; false when it
* was absent, in which case nothing changed. Callers use the result to
* avoid recording a wrapper the list never accepted.
*/
boolean replace(MessageHandler old, MessageHandler next) {
synchronized (lock) {
MessageHandler[] cur = snapshot;
int idx = -1;
for (int i = 0; i < cur.length; i++) if (cur[i] == old) { idx = i; break; }
if (idx < 0) return;
if (idx < 0) return false;
MessageHandler[] out = new MessageHandler[cur.length];
System.arraycopy(cur, 0, out, 0, cur.length);
out[idx] = next;
snapshot = out;
return true;
}
}
MessageHandler[] snapshot() { return snapshot; }
Expand Down
25 changes: 24 additions & 1 deletion src/main/java/com/aaravlabs/synapse/Topic.java
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,29 @@ boolean acceptsType(Class<?> other) {
return boxedType.isAssignableFrom(box(other));
}

/**
* Whether a caller asking for {@code requested} can be handed this topic as a
* {@code Topic<requested>} without a ClassCastException later.
*
* <p>Checks the <b>opposite</b> assignability direction from {@link #acceptsType}, so
* the two are complementary rather than one being stronger: {@code acceptsType} asks
* whether the topic can carry {@code requested} values
* ({@code topicType.isAssignableFrom(requested)}), while this asks whether this
* topic's values can be <i>handed to</i> a caller holding a {@code requested}
* ({@code requested.isAssignableFrom(topicType)}). Neither implies the other, which is
* why callers that hand back a {@code Topic<requested>} need both.
*
* <p>The gap this closes: {@code Object} topic, request for {@code String} — {@code
* acceptsType} passes it, so an {@code Object} topic accepts a request for {@code
* String}. Handing back {@code Topic<String>} for a topic that may hold any {@code
* Object} is not safe — the cast is erased, so the failure surfaces later at the
* caller's own site, possibly on the hardware thread. Callers that receive the topic as
* {@code Topic<Object>} are unaffected.
*/
boolean safelyReturnsAs(Class<?> requested) {
return box(requested).isAssignableFrom(boxedType);
}

/**
* Returns the most recent value, falling back to {@code defaultValue} if nothing has
* been published yet. Convenience for {@code topic.latestValue().orElse(default)}.
Expand Down Expand Up @@ -288,7 +311,7 @@ public Optional<Latest<T>> latest() {
}

/** Treat primitive {@code double.class} and wrapper {@code Double.class} as the same type. */
private static Class<?> box(Class<?> c) {
static Class<?> box(Class<?> c) {
if (!c.isPrimitive()) return c;
if (c == int.class) return Integer.class;
if (c == long.class) return Long.class;
Expand Down
Loading
Loading