001/*
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *   http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing,
013 * software distributed under the License is distributed on an
014 * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
015 * KIND, either express or implied.  See the License for the
016 * specific language governing permissions and limitations
017 * under the License.
018 */
019package org.apache.reef.wake.remote.transport.netty;
020
021import org.apache.reef.tang.Injector;
022import org.apache.reef.tang.Tang;
023import org.apache.reef.tang.exceptions.InjectionException;
024import org.apache.reef.wake.EStage;
025import org.apache.reef.wake.EventHandler;
026import org.apache.reef.wake.impl.SyncStage;
027import org.apache.reef.wake.remote.RemoteConfiguration;
028import org.apache.reef.wake.remote.address.LocalAddressProvider;
029import org.apache.reef.wake.remote.impl.TransportEvent;
030import org.apache.reef.wake.remote.ports.TcpPortProvider;
031import org.apache.reef.wake.remote.transport.Transport;
032import org.apache.reef.wake.remote.transport.TransportFactory;
033
034import javax.inject.Inject;
035
036/**
037 * Factory that creates a messaging transport.
038 */
039public class MessagingTransportFactory implements TransportFactory {
040
041  private final String localAddress;
042
043  /**
044   * @deprecated Have an instance injected instead.
045   */
046  @Deprecated
047  @Inject
048  // TODO[JIRA REEF-703]: change constructor to private
049  public MessagingTransportFactory(final LocalAddressProvider localAddressProvider) {
050    this.localAddress = localAddressProvider.getLocalAddress();
051  }
052
053  /**
054   * Creates a transport.
055   *
056   * @param port          a listening port
057   * @param clientHandler a transport client side handler
058   * @param serverHandler a transport server side handler
059   * @param exHandler     a exception handler
060   */
061  @Override
062  public Transport newInstance(final int port,
063                               final EventHandler<TransportEvent> clientHandler,
064                               final EventHandler<TransportEvent> serverHandler,
065                               final EventHandler<Exception> exHandler) {
066
067    final Injector injector = Tang.Factory.getTang().newInjector();
068    injector.bindVolatileParameter(RemoteConfiguration.HostAddress.class, this.localAddress);
069    injector.bindVolatileParameter(RemoteConfiguration.Port.class, port);
070    injector.bindVolatileParameter(RemoteConfiguration.RemoteClientStage.class, new SyncStage<>(clientHandler));
071    injector.bindVolatileParameter(RemoteConfiguration.RemoteServerStage.class, new SyncStage<>(serverHandler));
072
073    final Transport transport;
074    try {
075      transport = injector.getInstance(NettyMessagingTransport.class);
076      transport.registerErrorHandler(exHandler);
077      return transport;
078    } catch (final InjectionException e) {
079      throw new RuntimeException(e);
080    }
081  }
082
083  /**
084   * Creates a transport.
085   *
086   * @param hostAddress   a host address
087   * @param port          a listening port
088   * @param clientStage   a client stage
089   * @param serverStage   a server stage
090   * @param numberOfTries a number of tries
091   * @param retryTimeout  a timeout for retry
092   */
093  @Override
094  public Transport newInstance(final String hostAddress,
095                               final int port,
096                               final EStage<TransportEvent> clientStage,
097                               final EStage<TransportEvent> serverStage,
098                               final int numberOfTries,
099                               final int retryTimeout) {
100    try {
101      TcpPortProvider tcpPortProvider = Tang.Factory.getTang().newInjector().getInstance(TcpPortProvider.class);
102      return newInstance(hostAddress, port, clientStage,
103              serverStage, numberOfTries, retryTimeout, tcpPortProvider);
104    } catch (final InjectionException e) {
105      throw new RuntimeException(e);
106    }
107  }
108
109  /**
110   * Creates a transport.
111   *
112   * @param hostAddress     a host address
113   * @param port            a listening port
114   * @param clientStage     a client stage
115   * @param serverStage     a server stage
116   * @param numberOfTries   a number of tries
117   * @param retryTimeout    a timeout for retry
118   * @param tcpPortProvider a provider for TCP port
119   */
120  @Override
121  public Transport newInstance(final String hostAddress,
122                               final int port,
123                               final EStage<TransportEvent> clientStage,
124                               final EStage<TransportEvent> serverStage,
125                               final int numberOfTries,
126                               final int retryTimeout,
127                               final TcpPortProvider tcpPortProvider) {
128
129    final Injector injector = Tang.Factory.getTang().newInjector();
130    injector.bindVolatileParameter(RemoteConfiguration.HostAddress.class, hostAddress);
131    injector.bindVolatileParameter(RemoteConfiguration.Port.class, port);
132    injector.bindVolatileParameter(RemoteConfiguration.RemoteClientStage.class, clientStage);
133    injector.bindVolatileParameter(RemoteConfiguration.RemoteServerStage.class, serverStage);
134    injector.bindVolatileParameter(RemoteConfiguration.NumberOfTries.class, numberOfTries);
135    injector.bindVolatileParameter(RemoteConfiguration.RetryTimeout.class, retryTimeout);
136    injector.bindVolatileInstance(TcpPortProvider.class, tcpPortProvider);
137    try {
138      return injector.getInstance(NettyMessagingTransport.class);
139    } catch (final InjectionException e) {
140      throw new RuntimeException(e);
141    }
142  }
143}