diff --git a/CHANGELOG.md b/CHANGELOG.md index df01ceb..7cded3b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java index 067cfc4..6aa2f54 100644 --- a/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java +++ b/src/main/java/com/aaravlabs/synapse/OrchestratorImpl.java @@ -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) { @@ -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; @@ -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); @@ -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. */ diff --git a/src/main/java/com/aaravlabs/synapse/ftc/HardwareView.java b/src/main/java/com/aaravlabs/synapse/ftc/HardwareView.java index 039f97a..f3cf09f 100644 --- a/src/main/java/com/aaravlabs/synapse/ftc/HardwareView.java +++ b/src/main/java/com/aaravlabs/synapse/ftc/HardwareView.java @@ -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. * - *

Instances are created per-callback by {@code HardwareActions.bulkRead}; you - * should not construct them yourself. + *

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 { diff --git a/src/main/java/com/aaravlabs/synapse/internal/AnnotationBinder.java b/src/main/java/com/aaravlabs/synapse/internal/AnnotationBinder.java index d2886cd..d8da54e 100644 --- a/src/main/java/com/aaravlabs/synapse/internal/AnnotationBinder.java +++ b/src/main/java/com/aaravlabs/synapse/internal/AnnotationBinder.java @@ -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. diff --git a/src/test/java/com/aaravlabs/synapse/CoreHotPathTest.java b/src/test/java/com/aaravlabs/synapse/CoreHotPathTest.java new file mode 100644 index 0000000..97fc67c --- /dev/null +++ b/src/test/java/com/aaravlabs/synapse/CoreHotPathTest.java @@ -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 errors = new CopyOnWriteArrayList<>(); + + /** @return every {@code error(...)} logged so far, in order. */ + List 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 doubles = new CopyOnWriteArrayList<>(); + List 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), + "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(); + } + } +}