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
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Changed

- **`Orchestrator.hardware()` now returns a shared instance.** The
`HardwareActions` facade is constructed once per orchestrator instead of once
per call, and the `HardwareView` handed to `bulkRead` callbacks is likewise
shared across registrations rather than allocated per registration. Both are
immutable single-reference views over the orchestrator, so this is not a
behavioral change beyond object identity — it is now safe to capture the
facade once and reuse it on hot paths.
- **`@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.

## [0.4.0] - 2026-09-12

### Changed
Expand Down
18 changes: 15 additions & 3 deletions src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,11 @@ public final class OrchestratorImpl implements Orchestrator {
*/
private final java.util.concurrent.ScheduledExecutorService hardwareThread;

// Immutable one-field views over this orchestrator, so they are built once
// rather than allocated per hardware() call or per bulk-read registration.
private final com.aaravlabs.synapse.ftc.HardwareActions hardwareActions;
private final com.aaravlabs.synapse.ftc.HardwareView hardwareView;

private volatile boolean closed = false;

private OrchestratorImpl(String name, LogSink log) {
Expand Down Expand Up @@ -97,6 +102,9 @@ private OrchestratorImpl(String name, LogSink log) {
// (for hardware loops) go to the same executor / thread.
this.hardwareThread = java.util.concurrent.Executors
.newSingleThreadScheduledExecutor(hwTf);

this.hardwareActions = new com.aaravlabs.synapse.ftc.HardwareActions(this);
this.hardwareView = new com.aaravlabs.synapse.ftc.HardwareView(this);
}

private volatile Thread hardwareThreadThread;
Expand Down Expand Up @@ -434,8 +442,7 @@ public com.aaravlabs.synapse.ftc.HardwareActions.BulkReadHandle scheduleHardware
int hz, com.aaravlabs.synapse.ftc.BulkReader reader) {
if (hz <= 0) throw new IllegalArgumentException("hz must be > 0");
long delayMs = Math.max(1, 1000L / hz);
com.aaravlabs.synapse.ftc.HardwareView view =
new com.aaravlabs.synapse.ftc.HardwareView(this);
com.aaravlabs.synapse.ftc.HardwareView view = hardwareView;
ScheduledFuture<?> f = hardwareThread.scheduleWithFixedDelay(() -> {
try {
reader.read(view);
Expand All @@ -446,9 +453,14 @@ public com.aaravlabs.synapse.ftc.HardwareActions.BulkReadHandle scheduleHardware
return new com.aaravlabs.synapse.ftc.HardwareActions.BulkReadHandle(f);
}

/**
* @return the shared {@link com.aaravlabs.synapse.ftc.HardwareActions} facade for
* this orchestrator. The same instance is returned on every call, so it is
* safe to hold onto.
*/
@Override
public com.aaravlabs.synapse.ftc.HardwareActions hardware() {
return new com.aaravlabs.synapse.ftc.HardwareActions(this);
return hardwareActions;
}

/** Package-private: used by the binder to schedule @RunPeriodically methods. */
Expand Down
5 changes: 3 additions & 2 deletions src/main/java/com/aaravlabs/synapse/ftc/HardwareView.java
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,9 @@
* Provides convenient {@code publish} and {@code getLatestValue} methods that
* the bulk-read callback uses to record hardware readings onto the bus.
*
* <p>Instances are created per-callback by {@code HardwareActions.bulkRead}; you
* should not construct them yourself.
* <p>A single instance is created per orchestrator and shared by every
* {@link HardwareActions#bulkRead(int, BulkReader)} registration; you should not
* construct them yourself.
*/
public final class HardwareView {

Expand Down
23 changes: 7 additions & 16 deletions src/main/java/com/aaravlabs/synapse/internal/AnnotationBinder.java
Original file line number Diff line number Diff line change
Expand Up @@ -107,31 +107,22 @@ private static void wireSubscribedTos(OrchestratorImpl orchestrator, Node node,
}

Class<?> paramType = m.getParameterCount() == 1 ? m.getParameterTypes()[0] : null;
Class<?> topicType = paramType != null ? boxed(paramType) : Object.class;
// Normalised once here rather than per delivered message. Equivalent to
// `paramType.isPrimitive() ? boxed(paramType) : paramType`, because boxed()
// returns non-primitives unchanged.
final Class<?> effectiveParam = paramType != null ? boxed(paramType) : null;
Class<?> topicType = effectiveParam != null ? effectiveParam : Object.class;

m.setAccessible(true);

boolean onHardware = m.isAnnotationPresent(OnHardwareThread.class);
Runnable handlerBody = () -> {
try {
// Body is invoked by the dispatcher (callback pool or hardware
// thread). We capture paramType via a closure.
// The actual subscription handler is built per-topic below.
} catch (Throwable t) {
orchestrator.error("@SubscribedTo handler threw: " + m, t);
}
};
// We don't actually use handlerBody above — each topic gets its own
// handler that captures its own message type. The annotation-present
// check is the only thing we need from here.

for (SubscribedTo sub : subs) {
Subscription s = orchestrator.subscribeRaw(sub.topic(), topicType, msg -> {
try {
if (paramType == null) {
if (effectiveParam == null) {
m.invoke(node);
} else if (paramType.isPrimitive() ? boxed(paramType).isInstance(msg)
: paramType.isInstance(msg)) {
} else if (effectiveParam.isInstance(msg)) {
m.invoke(node, msg);
}
// else: silently drop — message type didn't match the parameter.
Expand Down
146 changes: 146 additions & 0 deletions src/test/java/com/aaravlabs/synapse/CoreHotPathTest.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
package com.aaravlabs.synapse;

import com.aaravlabs.synapse.annotation.SubscribedTo;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;

import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

import static org.junit.jupiter.api.Assertions.*;

/**
* Covers the core hot-path cleanups: the shared {@code HardwareActions} facade and the
* per-message reflection work hoisted out of the annotation binder's dispatch lambda.
*/
class CoreHotPathTest {

private RecordingLogSink logs;
private Orchestrator orchestrator;

@BeforeEach void setUp() {
logs = new RecordingLogSink();
orchestrator = Orchestrator.create("core-hot-path", logs);
}
@AfterEach void tearDown() { orchestrator.close(); }

/**
* Records what the core logged so tests can assert on log output, not just on
* side effects. Mirrors everything to {@link LogSink#STDERR} so a failing run
* still shows the orchestrator's trace.
*/
private static final class RecordingLogSink implements LogSink {
private final List<String> errors = new CopyOnWriteArrayList<>();

/** @return every {@code error(...)} logged so far, in order. */
List<String> errors() { return errors; }

@Override public void info(String tag, String message) { LogSink.STDERR.info(tag, message); }
@Override public void warn(String tag, String message) { LogSink.STDERR.warn(tag, message); }
@Override public void error(String tag, String message) {
errors.add(message);
LogSink.STDERR.error(tag, message);
}
@Override public void error(String tag, String message, Throwable t) {
errors.add(message + " (" + t + ")");
LogSink.STDERR.error(tag, message, t);
}
}

@Test
void hardwareActionsFacadeIsSharedAcrossCalls() throws Exception {
assertSame(orchestrator.hardware(), orchestrator.hardware());

// Cached, not frozen: the facade still reaches the hardware thread.
AtomicInteger hits = new AtomicInteger();
CountDownLatch done = new CountDownLatch(1);
orchestrator.hardware().run(() -> { hits.incrementAndGet(); done.countDown(); });
assertTrue(done.await(2, TimeUnit.SECONDS), "hardware().run() should have run");
assertEquals(1, hits.get());
}

@Test
void annotatedSubscriberStillReceivesPrimitiveAndWrapperTypedMessages() throws Exception {
// Covers both sides of the hoisted normalisation: a primitive parameter and a
// reference-typed one.
List<Double> doubles = new CopyOnWriteArrayList<>();
List<String> strings = new CopyOnWriteArrayList<>();
CountDownLatch done = new CountDownLatch(2);

class N extends Node {
N(Orchestrator o) { super(o); }
@SubscribedTo(topic = "hot/double")
public void onDouble(double v) { doubles.add(v); done.countDown(); }
@SubscribedTo(topic = "hot/string")
public void onString(String v) { strings.add(v); done.countDown(); }
}
orchestrator.registerNode("n", new N(orchestrator));

orchestrator.publish("hot/double", 0.25);
orchestrator.publish("hot/string", "go");

assertTrue(done.await(2, TimeUnit.SECONDS), "both handlers should run");
assertEquals(List.of(0.25), doubles);
assertEquals(List.of("go"), strings);
}

@Test
void annotatedSubscriberSilentlyDropsMismatchedMessages() throws Exception {
AtomicInteger calls = new AtomicInteger();
CountDownLatch first = new CountDownLatch(1);

// The topic has to be typed more loosely than the handler parameter for the
// drop branch to be reachable: publish() rejects a String for a Double topic
// before any subscriber sees it. Topic types are first-writer-wins, so
// creating it as Object first is the realistic case — two handlers with
// different parameter types on one topic.
orchestrator.getOrCreateTopic("hot/drop", Object.class);

class N extends Node {
N(Orchestrator o) { super(o); }
@SubscribedTo(topic = "hot/drop")
public void onDouble(double v) { calls.incrementAndGet(); first.countDown(); }
}
orchestrator.registerNode("n", new N(orchestrator));

orchestrator.publish("hot/drop", "text");
assertFalse(first.await(150, TimeUnit.MILLISECONDS),
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
"mismatched message must be dropped, not delivered");
assertEquals(0, calls.get(), "mismatched message must be dropped, not delivered");
// "Not delivered" alone cannot tell a silent drop apart from a delivery
// that threw and was swallowed by the binder's catch (Throwable): if the
// isInstance guard were removed, m.invoke(node, "text") on a double
// parameter raises IllegalArgumentException, the handler never runs, and
// the two assertions above still hold. A dropped message leaves no trace;
// a failed delivery is logged as an error, so that is what we assert on.
assertEquals(List.of(), logs.errors(),
"mismatched message must be dropped silently, not delivered and failed");

// A real Double still gets through.
orchestrator.publish("hot/drop", 1.5);
assertTrue(first.await(2, TimeUnit.SECONDS));
assertEquals(1, calls.get());
assertEquals(List.of(), logs.errors(), "matching message must not log an error");
}

@Test
void bulkReadCallbackGetsAWorkingView() throws Exception {
CountDownLatch done = new CountDownLatch(1);
com.aaravlabs.synapse.ftc.HardwareActions.BulkReadHandle handle =
orchestrator.hardware().bulkRead(200, view -> {
view.publish("hot/bulk", 1.0f);
done.countDown();
});
try {
assertTrue(done.await(2, TimeUnit.SECONDS), "bulk read callback should have run");
assertEquals(Float.valueOf(1.0f),
orchestrator.getLatestValue("hot/bulk", Float.class).orElse(null));
} finally {
handle.cancel();
}
}
}
Loading