diff options
Diffstat (limited to 'container-messagebus')
27 files changed, 267 insertions, 172 deletions
diff --git a/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerHolder.java b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerHolder.java new file mode 100644 index 00000000000..9020bfaffe3 --- /dev/null +++ b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerHolder.java @@ -0,0 +1,45 @@ +// Copyright Verizon Media. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. +package com.yahoo.container.jdisc.messagebus; + +import com.yahoo.component.AbstractComponent; +import com.yahoo.messagebus.network.Network; +import com.yahoo.messagebus.network.NetworkMultiplexer; +import com.yahoo.messagebus.network.rpc.RPCNetwork; +import com.yahoo.messagebus.network.rpc.RPCNetworkParams; +import com.yahoo.messagebus.shared.NullNetwork; + +/** + * Holds a reference to a singleton {@link NetworkMultiplexer}. + * + * @author jonmv + */ +public class NetworkMultiplexerHolder extends AbstractComponent { + + private final Object monitor = new Object(); + private boolean destroyed = false; + private NetworkMultiplexer net; + + /** Get the singleton RPCNetworkAdapter, creating it if this hasn't yet been done. */ + public NetworkMultiplexer get(RPCNetworkParams params) { + synchronized (monitor) { + if (destroyed) + throw new IllegalStateException("Component already destroyed"); + + return net = net != null ? net : NetworkMultiplexer.shared(newNetwork(params)); + } + } + + private Network newNetwork(RPCNetworkParams params) { + return params.getSlobroksConfig().slobrok().isEmpty() ? new NullNetwork() : new RPCNetwork(params); + } + + @Override + public void deconstruct() { + synchronized (monitor) { + net.destroy(); + net = null; + destroyed = true; + } + } + +} diff --git a/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerProvider.java b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerProvider.java new file mode 100644 index 00000000000..650ac92e779 --- /dev/null +++ b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/NetworkMultiplexerProvider.java @@ -0,0 +1,44 @@ +// Copyright Verizon Media. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. +package com.yahoo.container.jdisc.messagebus; + +import com.google.inject.Inject; +import com.yahoo.cloud.config.SlobroksConfig; +import com.yahoo.container.jdisc.ContainerMbusConfig; +import com.yahoo.messagebus.network.Identity; +import com.yahoo.messagebus.network.NetworkMultiplexer; +import com.yahoo.messagebus.network.rpc.RPCNetworkParams; + +/** + * Injectable component which provides an {@link NetworkMultiplexer}, creating one if needed, + * i.e., the first time this is created in a container--subsequent creations of this will reuse + * the underlying network that was created initially. This breaks the DI pattern, but must be done + * because the network is a unique resource which cannot exist in several versions simultaneously. + * + * @author jonmv + */ +public class NetworkMultiplexerProvider { + + private final NetworkMultiplexer net; + + @Inject + public NetworkMultiplexerProvider(NetworkMultiplexerHolder net, ContainerMbusConfig mbusConfig, SlobroksConfig slobroksConfig) { + this(net, mbusConfig, slobroksConfig, System.getProperty("config.id")); //: + } + + public NetworkMultiplexerProvider(NetworkMultiplexerHolder net, ContainerMbusConfig mbusConfig, SlobroksConfig slobroksConfig, String identity) { + this.net = net.get(asParameters(mbusConfig, slobroksConfig, identity)); + } + + public static RPCNetworkParams asParameters(ContainerMbusConfig mbusConfig, SlobroksConfig slobroksConfig, String identity) { + return new RPCNetworkParams().setSlobroksConfig(slobroksConfig) + .setIdentity(new Identity(identity)) + .setListenPort(mbusConfig.port()) + .setNumTargetsPerSpec(mbusConfig.numconnectionspertarget()) + .setNumNetworkThreads(mbusConfig.numthreads()) + .setTransportEventsBeforeWakeup(mbusConfig.transport_events_before_wakeup()) + .setOptimization(RPCNetworkParams.Optimization.valueOf(mbusConfig.optimize_for().name())); + } + + public NetworkMultiplexer net() { return net; } + +} diff --git a/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/SessionCache.java b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/SessionCache.java index d65b2f7cc12..3e3902b35aa 100644 --- a/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/SessionCache.java +++ b/container-messagebus/src/main/java/com/yahoo/container/jdisc/messagebus/SessionCache.java @@ -1,34 +1,40 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.container.jdisc.messagebus; +import com.google.inject.Inject; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.component.AbstractComponent; -import com.yahoo.config.subscription.ConfigGetter; import com.yahoo.container.jdisc.ContainerMbusConfig; import com.yahoo.document.DocumentTypeManager; import com.yahoo.document.DocumentUtil; +import com.yahoo.document.config.DocumentmanagerConfig; import com.yahoo.documentapi.messagebus.loadtypes.LoadTypeSet; import com.yahoo.documentapi.messagebus.protocol.DocumentProtocol; +import com.yahoo.documentapi.messagebus.protocol.DocumentProtocolPoliciesConfig; import com.yahoo.jdisc.ReferencedResource; import com.yahoo.jdisc.References; import com.yahoo.jdisc.ResourceReference; import com.yahoo.jdisc.SharedResource; -import java.util.logging.Level; import com.yahoo.messagebus.ConfigAgent; import com.yahoo.messagebus.DynamicThrottlePolicy; import com.yahoo.messagebus.IntermediateSessionParams; +import com.yahoo.messagebus.MessageBus; import com.yahoo.messagebus.MessageBusParams; +import com.yahoo.messagebus.MessagebusConfig; import com.yahoo.messagebus.Protocol; import com.yahoo.messagebus.SourceSessionParams; import com.yahoo.messagebus.StaticThrottlePolicy; import com.yahoo.messagebus.ThrottlePolicy; -import com.yahoo.messagebus.network.Identity; -import com.yahoo.messagebus.network.rpc.RPCNetworkParams; +import com.yahoo.messagebus.network.NetworkMultiplexer; import com.yahoo.messagebus.shared.SharedIntermediateSession; import com.yahoo.messagebus.shared.SharedMessageBus; import com.yahoo.messagebus.shared.SharedSourceSession; +import com.yahoo.vespa.config.content.DistributionConfig; +import com.yahoo.vespa.config.content.LoadTypeConfig; import java.util.HashMap; import java.util.Map; +import java.util.logging.Level; import java.util.logging.Logger; /** @@ -37,24 +43,19 @@ import java.util.logging.Logger; * @author Steinar Knutsen * @author Einar Rosenvinge */ -// TODO jonmv: Remove this? Only used sensibly by FeedHandlerV3, where only timeout varies. +// TODO jonmv: Remove this: only used with more than one entry by FeedHandlerV3, where only timeout varies. +// rant: This whole construct is because DI at one point didn't exist, so getting hold of a shared resource +// or session was hard(?), and one resorted to routing through the Container, using URIs, to the correct +// MbusClient, with or without throttling. This introduced the problem of ownership during shutdown, +// which was solved with manual reference counting. This is all much better solved with DI, which (now) +// owns everything, and does component shutdown in reverse construction order, which is always right. +// So for the sake of everyone's mental health, this should all just be removed now! I suspect this is +// even the case for Request; we can track in handlers, and warn when requests have been misplaced. public final class SessionCache extends AbstractComponent { private static final Logger log = Logger.getLogger(SessionCache.class.getName()); - //config - private final String messagebusConfigId; - private final String slobrokConfigId; - private final String identity; - private final String containerMbusConfigId; - private final String documentManagerConfigId; - private final String loadTypeConfigId; - private final DocumentTypeManager documentTypeManager; - - // initialized in start() - private ConfigAgent configAgent; - private SharedMessageBus messageBus; - + private final SharedMessageBus messageBus; private final Object intermediateLock = new Object(); private final Map<String, SharedIntermediateSession> intermediates = new HashMap<>(); private final IntermediateSessionCreator intermediatesCreator = new IntermediateSessionCreator(); @@ -63,50 +64,43 @@ public final class SessionCache extends AbstractComponent { private final Map<SourceSessionKey, SharedSourceSession> sources = new HashMap<>(); private final SourceSessionCreator sourcesCreator = new SourceSessionCreator(); - public SessionCache(String messagebusConfigId, String slobrokConfigId, String identity, - String containerMbusConfigId, String documentManagerConfigId, - String loadTypeConfigId, - DocumentTypeManager documentTypeManager) { - this.messagebusConfigId = messagebusConfigId; - this.slobrokConfigId = slobrokConfigId; - this.identity = identity; - this.containerMbusConfigId = containerMbusConfigId; - this.documentManagerConfigId = documentManagerConfigId; - this.loadTypeConfigId = loadTypeConfigId; - this.documentTypeManager = documentTypeManager; - } + @Inject + public SessionCache(NetworkMultiplexerProvider nets, ContainerMbusConfig containerMbusConfig, + DocumentmanagerConfig documentmanagerConfig, + LoadTypeConfig loadTypeConfig, MessagebusConfig messagebusConfig, + DocumentProtocolPoliciesConfig policiesConfig, + DistributionConfig distributionConfig) { + this(nets.net(), containerMbusConfig, documentmanagerConfig, + loadTypeConfig, messagebusConfig, policiesConfig, distributionConfig); - public SessionCache(final String identity) { - this(identity, identity, identity, identity, identity, identity, new DocumentTypeManager()); } - public void deconstruct() { - if (configAgent != null) { - configAgent.shutdown(); - } + public SessionCache(NetworkMultiplexer net, ContainerMbusConfig containerMbusConfig, + DocumentmanagerConfig documentmanagerConfig, + LoadTypeConfig loadTypeConfig, MessagebusConfig messagebusConfig, + DocumentProtocolPoliciesConfig policiesConfig, + DistributionConfig distributionConfig) { + this(net, + containerMbusConfig, + messagebusConfig, + new DocumentProtocol(new DocumentTypeManager(documentmanagerConfig), + new LoadTypeSet(loadTypeConfig), + policiesConfig, + distributionConfig)); } - - private void start() { - ContainerMbusConfig mbusConfig = ConfigGetter.getConfig(ContainerMbusConfig.class, containerMbusConfigId); - if (documentManagerConfigId != null) { - documentTypeManager.configure(documentManagerConfigId); - } - LoadTypeSet loadTypeSet = new LoadTypeSet(loadTypeConfigId); - DocumentProtocol protocol = new DocumentProtocol(documentTypeManager, identity, loadTypeSet); - messageBus = createSharedMessageBus(mbusConfig, slobrokConfigId, identity, protocol); - // TODO: stop doing subscriptions to config when that is to be solved in slobrok as well - configAgent = new ConfigAgent(messagebusConfigId, messageBus.messageBus()); - configAgent.subscribe(); + public SessionCache(NetworkMultiplexer net, ContainerMbusConfig containerMbusConfig, + MessagebusConfig messagebusConfig, Protocol protocol) { + this.messageBus = createSharedMessageBus(net, containerMbusConfig, messagebusConfig, protocol); } - - private boolean isStarted() { - return messageBus != null; + public void deconstruct() { + messageBus.release(); } - private static SharedMessageBus createSharedMessageBus(ContainerMbusConfig mbusConfig, - String slobrokConfigId, String identity, + private static SharedMessageBus createSharedMessageBus(NetworkMultiplexer net, + ContainerMbusConfig mbusConfig, + MessagebusConfig messagebusConfig, Protocol protocol) { MessageBusParams mbusParams = new MessageBusParams().addProtocol(protocol); @@ -118,15 +112,9 @@ public final class SessionCache extends AbstractComponent { mbusParams.setMaxPendingCount(mbusConfig.maxpendingcount()); mbusParams.setMaxPendingSize(maxPendingSize); - RPCNetworkParams netParams = new RPCNetworkParams() - .setSlobrokConfigId(slobrokConfigId) - .setIdentity(new Identity(identity)) - .setListenPort(mbusConfig.port()) - .setNumTargetsPerSpec(mbusConfig.numconnectionspertarget()) - .setNumNetworkThreads(mbusConfig.numthreads()) - .setTransportEventsBeforeWakeup(mbusConfig.transport_events_before_wakeup()) - .setOptimization(RPCNetworkParams.Optimization.valueOf(mbusConfig.optimize_for().name())); - return SharedMessageBus.newInstance(mbusParams, netParams); + MessageBus bus = new MessageBus(net, mbusParams); + new ConfigAgent(messagebusConfig, bus); // Configure the wrapped MessageBus with a routing table. + return new SharedMessageBus(bus); } private static void logSystemInfo(ContainerMbusConfig containerMbusConfig, long maxPendingSize) { @@ -144,20 +132,10 @@ public final class SessionCache extends AbstractComponent { } ReferencedResource<SharedIntermediateSession> retainIntermediate(final IntermediateSessionParams p) { - synchronized (this) { - if (!isStarted()) { - start(); - } - } return intermediatesCreator.retain(intermediateLock, intermediates, p); } public ReferencedResource<SharedSourceSession> retainSource(final SourceSessionParams p) { - synchronized (this) { - if (!isStarted()) { - start(); - } - } return sourcesCreator.retain(sourceLock, sources, p); } diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/MbusServer.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/MbusServer.java index e26e1e7e134..a2131f22dc4 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/MbusServer.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/MbusServer.java @@ -50,6 +50,7 @@ public final class MbusServer extends AbstractResource implements ServerProvider public void close() { log.log(Level.FINE, "Closing message bus server."); running.set(false); + session.close(); } @Override diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ClientTestDriver.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ClientTestDriver.java index 111805d61b0..ea0ed7eadc8 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ClientTestDriver.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ClientTestDriver.java @@ -33,7 +33,7 @@ public class ClientTestDriver { this.server = server; MessageBusParams mbusParams = new MessageBusParams().addProtocol(protocol); - RPCNetworkParams netParams = new RPCNetworkParams().setSlobrokConfigId(server.slobrokId()); + RPCNetworkParams netParams = new RPCNetworkParams().setSlobroksConfig(server.slobroksConfig()); SharedMessageBus mbus = SharedMessageBus.newInstance(mbusParams, netParams); session = mbus.newSourceSession(new SourceSessionParams()); client = new MbusClient(session); @@ -128,7 +128,4 @@ public class ClientTestDriver { return new ClientTestDriver(RemoteServer.newInstanceWithInternSlobrok(), protocol); } - public static ClientTestDriver newInstanceWithExternSlobrok(String slobrokId) { - return new ClientTestDriver(RemoteServer.newInstanceWithExternSlobrok(slobrokId), new SimpleProtocol()); - } } diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteClient.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteClient.java index 57d0abd980b..6cd8fb8f34d 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteClient.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteClient.java @@ -1,9 +1,17 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.messagebus.jdisc.test; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.jrt.ListenFailedException; import com.yahoo.jrt.slobrok.server.Slobrok; -import com.yahoo.messagebus.*; +import com.yahoo.messagebus.Message; +import com.yahoo.messagebus.MessageBus; +import com.yahoo.messagebus.MessageBusParams; +import com.yahoo.messagebus.Protocol; +import com.yahoo.messagebus.Reply; +import com.yahoo.messagebus.Result; +import com.yahoo.messagebus.SourceSession; +import com.yahoo.messagebus.SourceSessionParams; import com.yahoo.messagebus.network.local.LocalNetwork; import com.yahoo.messagebus.network.rpc.RPCNetwork; import com.yahoo.messagebus.network.rpc.RPCNetworkParams; @@ -17,16 +25,14 @@ import java.util.concurrent.TimeUnit; public class RemoteClient { private final Slobrok slobrok; - private final String slobrokId; private final MessageBus mbus; private final ReplyQueue queue = new ReplyQueue(); private final SourceSession session; - private RemoteClient(Slobrok slobrok, String slobrokId, Protocol protocol, boolean network) { - this.slobrok = slobrok; - this.slobrokId = slobrok != null ? slobrok.configId() : slobrokId; + private RemoteClient(Protocol protocol, boolean network) { + this.slobrok = newSlobrok(); mbus = network - ? new MessageBus(new RPCNetwork(new RPCNetworkParams().setSlobrokConfigId(this.slobrokId)), + ? new MessageBus(new RPCNetwork(new RPCNetworkParams().setSlobroksConfig(slobroksConfig())), new MessageBusParams().addProtocol(protocol)) : new MessageBus(new LocalNetwork(), new MessageBusParams().addProtocol(protocol)); session = mbus.createSourceSession(new SourceSessionParams().setThrottlePolicy(null).setReplyHandler(queue)); @@ -40,28 +46,22 @@ public class RemoteClient { return queue.awaitReply(timeout, unit); } - public String slobrokId() { - return slobrokId; + public SlobroksConfig slobroksConfig() { + return TestUtils.configFor(slobrok); } public void close() { session.destroy(); mbus.destroy(); - if (slobrok != null) { - slobrok.stop(); - } + slobrok.stop(); } public static RemoteClient newInstanceWithInternSlobrok(boolean network) { - return new RemoteClient(newSlobrok(), null, new SimpleProtocol(), network); - } - - public static RemoteClient newInstanceWithExternSlobrok(String slobrokId, boolean network) { - return new RemoteClient(null, slobrokId, new SimpleProtocol(), network); + return new RemoteClient(new SimpleProtocol(), network); } public static RemoteClient newInstanceWithProtocolAndInternSlobrok(Protocol protocol, boolean network) { - return new RemoteClient(newSlobrok(), null, protocol, network); + return new RemoteClient(protocol, network); } private static Slobrok newSlobrok() { diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteServer.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteServer.java index 1f0f82c4903..0fb23e33709 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteServer.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/RemoteServer.java @@ -1,9 +1,16 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.messagebus.jdisc.test; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.jrt.ListenFailedException; import com.yahoo.jrt.slobrok.server.Slobrok; -import com.yahoo.messagebus.*; +import com.yahoo.messagebus.DestinationSession; +import com.yahoo.messagebus.DestinationSessionParams; +import com.yahoo.messagebus.Message; +import com.yahoo.messagebus.MessageBus; +import com.yahoo.messagebus.MessageBusParams; +import com.yahoo.messagebus.Protocol; +import com.yahoo.messagebus.Reply; import com.yahoo.messagebus.network.Identity; import com.yahoo.messagebus.network.rpc.RPCNetwork; import com.yahoo.messagebus.network.rpc.RPCNetworkParams; @@ -17,16 +24,14 @@ import java.util.concurrent.TimeUnit; public class RemoteServer { private final Slobrok slobrok; - private final String slobrokId; private final MessageBus mbus; private final MessageQueue queue = new MessageQueue(); private final DestinationSession session; - private RemoteServer(Slobrok slobrok, String slobrokId, Protocol protocol, String identity) { - this.slobrok = slobrok; - this.slobrokId = slobrok != null ? slobrok.configId() : slobrokId; + private RemoteServer(Protocol protocol, String identity) { + this.slobrok = newSlobrok(); mbus = new MessageBus(new RPCNetwork(new RPCNetworkParams() - .setSlobrokConfigId(this.slobrokId) + .setSlobroksConfig(slobroksConfig()) .setIdentity(new Identity(identity))), new MessageBusParams().addProtocol(protocol)); session = mbus.createDestinationSession(new DestinationSessionParams().setMessageHandler(queue)); @@ -48,32 +53,22 @@ public class RemoteServer { session.reply(reply); } - public String slobrokId() { - return slobrokId; + public SlobroksConfig slobroksConfig() { + return TestUtils.configFor(slobrok); } public void close() { session.destroy(); mbus.destroy(); - if (slobrok != null) { - slobrok.stop(); - } + slobrok.stop(); } public static RemoteServer newInstanceWithInternSlobrok() { - return new RemoteServer(newSlobrok(), null, new SimpleProtocol(), "remote"); - } - - public static RemoteServer newInstanceWithExternSlobrok(String slobrokId) { - return new RemoteServer(null, slobrokId, new SimpleProtocol(), "remote"); - } - - public static RemoteServer newInstance(String slobrokId, String identity, Protocol protocol) { - return new RemoteServer(null, slobrokId, protocol, identity); + return new RemoteServer(new SimpleProtocol(), "remote"); } - public static RemoteServer newInstanceWithProtocol(Protocol protocol) { - return new RemoteServer(newSlobrok(), null, protocol, "remote"); + public static RemoteServer newInstance(String identity, Protocol protocol) { + return new RemoteServer(protocol, identity); } private static Slobrok newSlobrok() { diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ServerTestDriver.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ServerTestDriver.java index e59db28e886..fa0dd37ed13 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ServerTestDriver.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/ServerTestDriver.java @@ -42,7 +42,7 @@ public class ServerTestDriver { } MessageBusParams mbusParams = new MessageBusParams().addProtocol(protocol); - RPCNetworkParams netParams = new RPCNetworkParams().setSlobrokConfigId(client.slobrokId()); + RPCNetworkParams netParams = new RPCNetworkParams().setSlobroksConfig(client.slobroksConfig()); SharedMessageBus mbus = SharedMessageBus.newInstance(mbusParams, netParams); ServerSession session = mbus.newDestinationSession(new DestinationSessionParams()); server = new MbusServer(driver, session); @@ -130,13 +130,6 @@ public class ServerTestDriver { guiceModules); } - public static ServerTestDriver newInstanceWithExternSlobrok(String slobrokId, RequestHandler requestHandler, - boolean network, Module... guiceModules) - { - return new ServerTestDriver(RemoteClient.newInstanceWithExternSlobrok(slobrokId, network), - true, requestHandler, new SimpleProtocol(), guiceModules); - } - public static ServerTestDriver newInactiveInstance(boolean network, Module... guiceModules) { return new ServerTestDriver(RemoteClient.newInstanceWithInternSlobrok(network), false, null, new SimpleProtocol(), guiceModules); diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/TestUtils.java b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/TestUtils.java new file mode 100644 index 00000000000..85ad241259f --- /dev/null +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/jdisc/test/TestUtils.java @@ -0,0 +1,21 @@ +// Copyright Verizon Media. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. +package com.yahoo.messagebus.jdisc.test; + +import com.yahoo.cloud.config.SlobroksConfig; +import com.yahoo.jrt.Spec; +import com.yahoo.jrt.slobrok.server.Slobrok; + +/** + * @author jonmv + */ +public class TestUtils { + + private TestUtils() { } + + public static SlobroksConfig configFor(Slobrok slobrok) { + return new SlobroksConfig.Builder().slobrok(new SlobroksConfig.Slobrok.Builder() + .connectionspec(new Spec("localhost", slobrok.port()).toString())) + .build(); + } + +} diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/NullNetwork.java b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/NullNetwork.java index ad58d6b9a5e..a7cf5fec7d1 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/NullNetwork.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/NullNetwork.java @@ -12,9 +12,9 @@ import java.util.List; /** * <p>Used by SharedMessageBus as a network when the container runs in LocalApplication with no network services.</p> * - * @author <a href="mailto:vegardh@yahoo-inc.com">Vegard Havdal</a> + * @author Vegard Havdal */ -class NullNetwork implements Network { +public class NullNetwork implements Network { @Override public boolean waitUntilReady(double seconds) { diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/ServerSession.java b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/ServerSession.java index 56713815c7a..4cd4a776292 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/ServerSession.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/ServerSession.java @@ -10,13 +10,16 @@ import com.yahoo.messagebus.Reply; */ public interface ServerSession extends SharedResource { - public MessageHandler getMessageHandler(); + MessageHandler getMessageHandler(); - public void setMessageHandler(MessageHandler msgHandler); + void setMessageHandler(MessageHandler msgHandler); - public void sendReply(Reply reply); + void sendReply(Reply reply); - public String connectionSpec(); + String connectionSpec(); + + String name(); + + void close(); - public String name(); } diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedDestinationSession.java b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedDestinationSession.java index 7da164757cd..5a9cd39c5b4 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedDestinationSession.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedDestinationSession.java @@ -77,6 +77,11 @@ public class SharedDestinationSession extends AbstractResource implements Messag } @Override + public void close() { + session.destroy(); + } + + @Override protected void destroy() { log.log(Level.FINE, "Destroying shared destination session."); session.destroy(); diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedIntermediateSession.java b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedIntermediateSession.java index 5c9fab46e34..64cc1aaf510 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedIntermediateSession.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedIntermediateSession.java @@ -96,6 +96,11 @@ public class SharedIntermediateSession extends AbstractResource } @Override + public void close() { + session.destroy(); + } + + @Override protected void destroy() { log.log(Level.FINE, "Destroying shared intermediate session."); session.destroy(); diff --git a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedMessageBus.java b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedMessageBus.java index dd135a51378..39bd6e86aa2 100644 --- a/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedMessageBus.java +++ b/container-messagebus/src/main/java/com/yahoo/messagebus/shared/SharedMessageBus.java @@ -3,7 +3,10 @@ package com.yahoo.messagebus.shared; import com.yahoo.config.subscription.ConfigGetter; import com.yahoo.jdisc.AbstractResource; + +import java.util.Objects; import java.util.logging.Level; + import com.yahoo.messagebus.DestinationSessionParams; import com.yahoo.messagebus.IntermediateSessionParams; import com.yahoo.messagebus.MessageBus; @@ -25,8 +28,7 @@ public class SharedMessageBus extends AbstractResource { private final MessageBus mbus; public SharedMessageBus(MessageBus mbus) { - mbus.getClass(); // throws NullPointerException - this.mbus = mbus; + this.mbus = Objects.requireNonNull(mbus); } public MessageBus messageBus() { @@ -65,4 +67,6 @@ public class SharedMessageBus extends AbstractResource { } return new RPCNetwork(params); } + + } diff --git a/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusClientProviderTest.java b/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusClientProviderTest.java index 6335cf01d8c..2a1c4e9a15a 100644 --- a/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusClientProviderTest.java +++ b/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusClientProviderTest.java @@ -1,9 +1,15 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.container.jdisc.messagebus; +import com.yahoo.container.jdisc.ContainerMbusConfig; import com.yahoo.container.jdisc.config.SessionConfig; -import com.yahoo.container.jdisc.messagebus.MbusClientProvider; -import com.yahoo.container.jdisc.messagebus.SessionCache; +import com.yahoo.document.config.DocumentmanagerConfig; +import com.yahoo.documentapi.messagebus.protocol.DocumentProtocolPoliciesConfig; +import com.yahoo.messagebus.MessagebusConfig; +import com.yahoo.messagebus.network.NetworkMultiplexer; +import com.yahoo.messagebus.shared.NullNetwork; +import com.yahoo.vespa.config.content.DistributionConfig; +import com.yahoo.vespa.config.content.LoadTypeConfig; import org.junit.Test; import static org.junit.Assert.assertNotNull; @@ -30,7 +36,14 @@ public class MbusClientProviderTest { } private void testClient(SessionConfig config) { - MbusClientProvider p = new MbusClientProvider(new SessionCache("dir:src/test/resources/config/clientprovider"), config); + SessionCache cache = new SessionCache(NetworkMultiplexer.dedicated(new NullNetwork()), + new ContainerMbusConfig.Builder().build(), + new DocumentmanagerConfig.Builder().build(), + new LoadTypeConfig.Builder().build(), + new MessagebusConfig.Builder().build(), + new DocumentProtocolPoliciesConfig.Builder().build(), + new DistributionConfig.Builder().build()); + MbusClientProvider p = new MbusClientProvider(cache, config); assertNotNull(p.get()); p.deconstruct(); } diff --git a/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusSessionKeyTestCase.java b/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusSessionKeyTestCase.java index 89bc5b9cecd..709f78f5cf0 100644 --- a/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusSessionKeyTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/container/jdisc/messagebus/MbusSessionKeyTestCase.java @@ -17,7 +17,7 @@ import static org.junit.Assert.assertTrue; /** * Check the completeness of the mbus session key classes. * - * @author <a href="mailto:steinar@yahoo-inc.com">Steinar Knutsen</a> + * @author Steinar Knutsen */ public class MbusSessionKeyTestCase { diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/MbusServerTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/MbusServerTestCase.java index 9d45d2e7abf..6ebb41c4ab7 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/MbusServerTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/MbusServerTestCase.java @@ -331,9 +331,7 @@ public class MbusServerTestCase { int refCount = 1; @Override - public void sendReply(Reply reply) { - - } + public void sendReply(Reply reply) { } @Override public MessageHandler getMessageHandler() { @@ -341,9 +339,7 @@ public class MbusServerTestCase { } @Override - public void setMessageHandler(MessageHandler msgHandler) { - - } + public void setMessageHandler(MessageHandler msgHandler) { } @Override public String connectionSpec() { @@ -356,14 +352,12 @@ public class MbusServerTestCase { } @Override + public void close() { } + + @Override public ResourceReference refer() { ++refCount; - return new ResourceReference() { - @Override - public void close() { - --refCount; - } - }; + return () -> --refCount; } @Override diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ClientTestDriverTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ClientTestDriverTestCase.java index ef290a070cb..d24acb9d0d8 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ClientTestDriverTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ClientTestDriverTestCase.java @@ -23,10 +23,6 @@ public class ClientTestDriverTestCase { driver = ClientTestDriver.newInstanceWithProtocol(new SimpleProtocol()); assertNotNull(driver); assertTrue(driver.close()); - - Slobrok slobrok = new Slobrok(); - driver = ClientTestDriver.newInstanceWithExternSlobrok(slobrok.configId()); - assertNotNull(driver); - assertTrue(driver.close()); } + } diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ServerTestDriverTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ServerTestDriverTestCase.java index f6ae2335d12..7b2c1b68549 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ServerTestDriverTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/jdisc/test/ServerTestDriverTestCase.java @@ -24,11 +24,6 @@ public class ServerTestDriverTestCase { driver = ServerTestDriver.newInstanceWithProtocol(new SimpleProtocol(), new NonWorkingRequestHandler(), false); assertNotNull(driver); assertTrue(driver.close()); - - Slobrok slobrok = new Slobrok(); - driver = ServerTestDriver.newInstanceWithExternSlobrok(slobrok.configId(), new NonWorkingRequestHandler(), false); - assertNotNull(driver); - assertTrue(driver.close()); } } diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedDestinationSessionTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedDestinationSessionTestCase.java index 78e79da4b9f..759509946ce 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedDestinationSessionTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedDestinationSessionTestCase.java @@ -1,12 +1,14 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.messagebus.shared; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.jrt.ListenFailedException; import com.yahoo.jrt.slobrok.server.Slobrok; import com.yahoo.messagebus.*; import com.yahoo.messagebus.jdisc.test.MessageQueue; import com.yahoo.messagebus.jdisc.test.RemoteClient; import com.yahoo.messagebus.jdisc.test.ReplyQueue; +import com.yahoo.messagebus.jdisc.test.TestUtils; import com.yahoo.messagebus.network.rpc.RPCNetworkParams; import com.yahoo.messagebus.routing.Route; import com.yahoo.messagebus.test.SimpleMessage; @@ -98,7 +100,7 @@ public class SharedDestinationSessionTestCase { RemoteClient client = RemoteClient.newInstanceWithInternSlobrok(true); MessageQueue queue = new MessageQueue(); DestinationSessionParams params = new DestinationSessionParams().setMessageHandler(queue); - SharedDestinationSession session = newDestinationSession(client.slobrokId(), params); + SharedDestinationSession session = newDestinationSession(client.slobroksConfig(), params); Route route = Route.parse(session.connectionSpec()); assertTrue(client.sendMessage(new SimpleMessage("foo").setRoute(route)).isAccepted()); @@ -120,11 +122,11 @@ public class SharedDestinationSessionTestCase { } catch (ListenFailedException e) { fail(); } - return newDestinationSession(slobrok.configId(), new DestinationSessionParams()); + return newDestinationSession(TestUtils.configFor(slobrok), new DestinationSessionParams()); } - private static SharedDestinationSession newDestinationSession(String slobrokId, DestinationSessionParams params) { - RPCNetworkParams netParams = new RPCNetworkParams().setSlobrokConfigId(slobrokId); + private static SharedDestinationSession newDestinationSession(SlobroksConfig slobroksConfig, DestinationSessionParams params) { + RPCNetworkParams netParams = new RPCNetworkParams().setSlobroksConfig(slobroksConfig); MessageBusParams mbusParams = new MessageBusParams().addProtocol(new SimpleProtocol()); SharedMessageBus mbus = SharedMessageBus.newInstance(mbusParams, netParams); SharedDestinationSession session = mbus.newDestinationSession(params); diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedIntermediateSessionTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedIntermediateSessionTestCase.java index 87958415149..ca8301d1ea9 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedIntermediateSessionTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedIntermediateSessionTestCase.java @@ -1,6 +1,7 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.messagebus.shared; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.jrt.ListenFailedException; import com.yahoo.jrt.slobrok.server.Slobrok; import com.yahoo.messagebus.*; @@ -8,6 +9,7 @@ import com.yahoo.messagebus.jdisc.test.MessageQueue; import com.yahoo.messagebus.jdisc.test.RemoteClient; import com.yahoo.messagebus.jdisc.test.RemoteServer; import com.yahoo.messagebus.jdisc.test.ReplyQueue; +import com.yahoo.messagebus.jdisc.test.TestUtils; import com.yahoo.messagebus.network.local.LocalNetwork; import com.yahoo.messagebus.network.local.LocalWire; import com.yahoo.messagebus.network.rpc.RPCNetworkParams; @@ -77,7 +79,7 @@ public class SharedIntermediateSessionTestCase { public void requireThatReplyHandlerCanNotBeSet() throws ListenFailedException { Slobrok slobrok = new Slobrok(); try { - newIntermediateSession(slobrok.configId(), + newIntermediateSession(TestUtils.configFor(slobrok), new IntermediateSessionParams().setReplyHandler(new ReplyQueue()), false); fail(); @@ -111,7 +113,7 @@ public class SharedIntermediateSessionTestCase { @Test public void requireThatSessionCanSendMessage() throws InterruptedException { RemoteServer server = RemoteServer.newInstanceWithInternSlobrok(); - SharedIntermediateSession session = newIntermediateSession(server.slobrokId(), + SharedIntermediateSession session = newIntermediateSession(server.slobroksConfig(), new IntermediateSessionParams(), true); ReplyQueue queue = new ReplyQueue(); @@ -134,7 +136,7 @@ public class SharedIntermediateSessionTestCase { RemoteClient client = RemoteClient.newInstanceWithInternSlobrok(true); MessageQueue queue = new MessageQueue(); IntermediateSessionParams params = new IntermediateSessionParams().setMessageHandler(queue); - SharedIntermediateSession session = newIntermediateSession(client.slobrokId(), params, true); + SharedIntermediateSession session = newIntermediateSession(client.slobroksConfig(), params, true); Route route = Route.parse(session.connectionSpec()); assertTrue(client.sendMessage(new SimpleMessage("foo").setRoute(route)).isAccepted()); @@ -156,13 +158,13 @@ public class SharedIntermediateSessionTestCase { } catch (ListenFailedException e) { fail(); } - return newIntermediateSession(slobrok.configId(), new IntermediateSessionParams(), network); + return newIntermediateSession(TestUtils.configFor(slobrok), new IntermediateSessionParams(), network); } - private static SharedIntermediateSession newIntermediateSession(String slobrokId, + private static SharedIntermediateSession newIntermediateSession(SlobroksConfig slobroksConfig, IntermediateSessionParams params, boolean network) { - RPCNetworkParams netParams = new RPCNetworkParams().setSlobrokConfigId(slobrokId); + RPCNetworkParams netParams = new RPCNetworkParams().setSlobroksConfig(slobroksConfig); MessageBusParams mbusParams = new MessageBusParams().addProtocol(new SimpleProtocol()); SharedMessageBus mbus = network ? SharedMessageBus.newInstance(mbusParams, netParams) diff --git a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedSourceSessionTestCase.java b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedSourceSessionTestCase.java index 1f0966fc961..7a4a3d17170 100644 --- a/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedSourceSessionTestCase.java +++ b/container-messagebus/src/test/java/com/yahoo/messagebus/shared/SharedSourceSessionTestCase.java @@ -1,6 +1,7 @@ // Copyright 2017 Yahoo Holdings. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root. package com.yahoo.messagebus.shared; +import com.yahoo.cloud.config.SlobroksConfig; import com.yahoo.jrt.ListenFailedException; import com.yahoo.jrt.slobrok.server.Slobrok; import com.yahoo.messagebus.Message; @@ -8,6 +9,7 @@ import com.yahoo.messagebus.MessageBusParams; import com.yahoo.messagebus.SourceSessionParams; import com.yahoo.messagebus.jdisc.test.RemoteServer; import com.yahoo.messagebus.jdisc.test.ReplyQueue; +import com.yahoo.messagebus.jdisc.test.TestUtils; import com.yahoo.messagebus.network.rpc.RPCNetworkParams; import com.yahoo.messagebus.routing.Route; import com.yahoo.messagebus.test.SimpleMessage; @@ -59,7 +61,7 @@ public class SharedSourceSessionTestCase { @Test public void requireThatSessionCanSendMessage() throws InterruptedException { RemoteServer server = RemoteServer.newInstanceWithInternSlobrok(); - SharedSourceSession session = newSourceSession(server.slobrokId(), + SharedSourceSession session = newSourceSession(server.slobroksConfig(), new SourceSessionParams()); ReplyQueue queue = new ReplyQueue(); Message msg = new SimpleMessage("foo").setRoute(Route.parse(server.connectionSpec())); @@ -80,11 +82,11 @@ public class SharedSourceSessionTestCase { } catch (ListenFailedException e) { fail(); } - return newSourceSession(slobrok.configId(), params); + return newSourceSession(TestUtils.configFor(slobrok), params); } - private static SharedSourceSession newSourceSession(String slobrokId, SourceSessionParams params) { - RPCNetworkParams netParams = new RPCNetworkParams().setSlobrokConfigId(slobrokId); + private static SharedSourceSession newSourceSession(SlobroksConfig slobroksConfig, SourceSessionParams params) { + RPCNetworkParams netParams = new RPCNetworkParams().setSlobroksConfig(slobroksConfig); MessageBusParams mbusParams = new MessageBusParams().addProtocol(new SimpleProtocol()); SharedMessageBus mbus = SharedMessageBus.newInstance(mbusParams, netParams); SharedSourceSession session = mbus.newSourceSession(params); diff --git a/container-messagebus/src/test/resources/config/clientprovider/container-mbus.cfg b/container-messagebus/src/test/resources/config/clientprovider/container-mbus.cfg deleted file mode 100644 index e69de29bb2d..00000000000 --- a/container-messagebus/src/test/resources/config/clientprovider/container-mbus.cfg +++ /dev/null diff --git a/container-messagebus/src/test/resources/config/clientprovider/documentmanager.cfg b/container-messagebus/src/test/resources/config/clientprovider/documentmanager.cfg deleted file mode 100644 index e69de29bb2d..00000000000 --- a/container-messagebus/src/test/resources/config/clientprovider/documentmanager.cfg +++ /dev/null diff --git a/container-messagebus/src/test/resources/config/clientprovider/load-type.cfg b/container-messagebus/src/test/resources/config/clientprovider/load-type.cfg deleted file mode 100644 index e69de29bb2d..00000000000 --- a/container-messagebus/src/test/resources/config/clientprovider/load-type.cfg +++ /dev/null diff --git a/container-messagebus/src/test/resources/config/clientprovider/messagebus.cfg b/container-messagebus/src/test/resources/config/clientprovider/messagebus.cfg deleted file mode 100644 index e69de29bb2d..00000000000 --- a/container-messagebus/src/test/resources/config/clientprovider/messagebus.cfg +++ /dev/null diff --git a/container-messagebus/src/test/resources/config/clientprovider/slobroks.cfg b/container-messagebus/src/test/resources/config/clientprovider/slobroks.cfg deleted file mode 100644 index e69de29bb2d..00000000000 --- a/container-messagebus/src/test/resources/config/clientprovider/slobroks.cfg +++ /dev/null |