From 411948109e85fb847f1d200d374a2d5e5efeef65 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 23 Aug 2026 18:23:33 +0000 Subject: [PATCH] feat(engine): defer RouterContext.attach to run after registry bootstrap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit RouterContext.attach(RouterConfig) was called as the first statement of EngineWorker.doInit(), before this worker's own EngineRegistry exists and before EngineManager has processed any namespace, yet its contract already lets an implementation attach a synthesized composite namespace via RouteableContext.attachComposite() as part of that same call — a namespace registration that reads current per-worker registry state. Every existing RouterContext implementation already builds its stream factory once, at construction time in Router.supply(), and attach() just returns that value unchanged; RouterConfig itself is unused. So the stream factory never depended on attach() running early — only bindings' own construction-time capture of it did. Split the two: add RouterContext.streamFactory(), sourced by EngineWorker right after RouterContext is constructed, and change attach() to void, called once per worker only after Engine.start() has run the engine's own manager.start() bootstrap to completion, dispatched onto that worker's own thread the same way ordinary namespace attachment already is. --- .../aklivity/zilla/runtime/engine/Engine.java | 7 ++++++ .../registry/EngineRouterContext.java | 9 +++++-- .../internal/registry/EngineWorker.java | 25 +++++++++++++++++-- .../runtime/engine/router/RouterContext.java | 22 +++++++++++++--- .../internal/router/TestRouterContext.java | 9 +++++-- 5 files changed, 62 insertions(+), 10 deletions(-) diff --git a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/Engine.java b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/Engine.java index 6b34abb88c..815a8add02 100644 --- a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/Engine.java +++ b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/Engine.java @@ -44,6 +44,7 @@ import java.util.Objects; import java.util.ServiceLoader; import java.util.ServiceLoader.Provider; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ThreadFactory; @@ -372,6 +373,12 @@ public void start() throws Exception if (!readonly) { manager.start(); + + workers.stream() + .map(w -> w.attachRouter(routerConfig)) + .reduce(CompletableFuture::allOf) + .ifPresent(CompletableFuture::join); + Ready.markReady(config.directory()); } } diff --git a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineRouterContext.java b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineRouterContext.java index 60211ebe23..ebfcc16648 100644 --- a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineRouterContext.java +++ b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineRouterContext.java @@ -33,12 +33,17 @@ final class EngineRouterContext implements RouterContext } @Override - public BindingHandler attach( - RouterConfig config) + public BindingHandler streamFactory() { return streamFactory; } + @Override + public void attach( + RouterConfig config) + { + } + @Override public void detach( long routerId) diff --git a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineWorker.java b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineWorker.java index 5ed783e288..a82f6ca31a 100644 --- a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineWorker.java +++ b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/internal/registry/EngineWorker.java @@ -430,6 +430,7 @@ public EngineWorker( EngineRouteable routeable = new EngineRouteable(config, this::newStream, this::attachComposite, this::detachComposite, this::supplyStore); this.router = router.supply(routeable); + this.streamFactory = this.router.streamFactory(); this.bindings = bindings; this.exporters = exporters; @@ -993,8 +994,6 @@ public void onStart() private void doInit() { - this.streamFactory = router.attach(routerConfig); - Map bindingsByType = bindings.stream() .collect(toMap(Binding::name, b -> b.supply(this), (a, b) -> a, LinkedHashMap::new)); @@ -1194,6 +1193,28 @@ public CompletableFuture detach( return detachTask.future(); } + public CompletableFuture attachRouter( + RouterConfig config) + { + assert thread != Thread.currentThread(); + + CompletableFuture future = new CompletableFuture<>(); + dispatch(() -> + { + try + { + router.attach(config); + future.complete(null); + } + catch (Throwable ex) + { + future.completeExceptionally(ex); + } + }); + + return future; + } + public AgentRunner runner() { return runner; diff --git a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/router/RouterContext.java b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/router/RouterContext.java index 82309a6e02..e80141cee5 100644 --- a/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/router/RouterContext.java +++ b/runtime/engine/src/main/java/io/aklivity/zilla/runtime/engine/router/RouterContext.java @@ -31,13 +31,27 @@ public interface RouterContext { /** - * Attaches a router configuration to this thread's context, returning the composed - * {@link BindingHandler} that will be installed as the engine's stream factory. + * Returns this thread's composed {@link BindingHandler}, installed as the engine's + * stream factory. + *

+ * Available immediately once this context is constructed, independent of {@link #attach}, + * so the engine can wire it into bindings before any registry-dependent setup has run. + *

+ * + * @return the composed stream factory + */ + BindingHandler streamFactory(); + + /** + * Attaches a router configuration to this thread's context. + *

+ * Called once this thread's namespace registry reflects the engine's bootstrap + * configuration, so implementations may safely attach synthesized namespaces here. + *

* * @param config the router configuration to activate - * @return a {@link BindingHandler} that the engine installs as its stream factory */ - BindingHandler attach( + void attach( RouterConfig config); /** diff --git a/runtime/engine/src/test/java/io/aklivity/zilla/runtime/engine/test/internal/router/TestRouterContext.java b/runtime/engine/src/test/java/io/aklivity/zilla/runtime/engine/test/internal/router/TestRouterContext.java index fa170442d8..92fd52d21e 100644 --- a/runtime/engine/src/test/java/io/aklivity/zilla/runtime/engine/test/internal/router/TestRouterContext.java +++ b/runtime/engine/src/test/java/io/aklivity/zilla/runtime/engine/test/internal/router/TestRouterContext.java @@ -34,12 +34,17 @@ public TestRouterContext( } @Override - public BindingHandler attach( - RouterConfig config) + public BindingHandler streamFactory() { return streamFactory; } + @Override + public void attach( + RouterConfig config) + { + } + @Override public void detach( long routerId)