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)