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}