View Javadoc
1   /*
2    * ====================================================================
3    * Licensed to the Apache Software Foundation (ASF) under one
4    * or more contributor license agreements.  See the NOTICE file
5    * distributed with this work for additional information
6    * regarding copyright ownership.  The ASF licenses this file
7    * to you under the Apache License, Version 2.0 (the
8    * "License"); you may not use this file except in compliance
9    * with the License.  You may obtain a copy of the License at
10   *
11   *   http://www.apache.org/licenses/LICENSE-2.0
12   *
13   * Unless required by applicable law or agreed to in writing,
14   * software distributed under the License is distributed on an
15   * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
16   * KIND, either express or implied.  See the License for the
17   * specific language governing permissions and limitations
18   * under the License.
19   * ====================================================================
20   *
21   * This software consists of voluntary contributions made by many
22   * individuals on behalf of the Apache Software Foundation.  For more
23   * information on the Apache Software Foundation, please see
24   * <http://www.apache.org/>.
25   *
26   */
27  
28  package org.apache.hc.core5.reactor;
29  
30  import java.io.IOException;
31  import java.net.SocketAddress;
32  import java.nio.channels.SocketChannel;
33  import java.util.ArrayList;
34  import java.util.Collections;
35  import java.util.Deque;
36  import java.util.List;
37  import java.util.Set;
38  import java.util.concurrent.ConcurrentLinkedDeque;
39  import java.util.concurrent.Future;
40  import java.util.concurrent.ThreadFactory;
41  
42  import org.apache.hc.core5.concurrent.DefaultThreadFactory;
43  import org.apache.hc.core5.concurrent.FutureCallback;
44  import org.apache.hc.core5.function.Callback;
45  import org.apache.hc.core5.function.Decorator;
46  import org.apache.hc.core5.io.CloseMode;
47  import org.apache.hc.core5.util.Args;
48  import org.apache.hc.core5.util.TimeValue;
49  
50  /**
51   * Multi-core I/O reactor that can ask as both {@link ConnectionInitiator}
52   * and {@link ConnectionAcceptor}. Internally this I/O reactor distributes newly created
53   * I/O session equally across multiple I/O worker threads for a more optimal resource
54   * utilization and a better I/O performance. Usually it is recommended to have
55   * one worker I/O reactor per physical CPU core.
56   *
57   * @since 4.0
58   */
59  public class DefaultListeningIOReactor extends AbstractIOReactorBase implements ConnectionAcceptor {
60  
61      private final static ThreadFactory DISPATCH_THREAD_FACTORY = new DefaultThreadFactory("I/O server dispatch", true);
62      private final static ThreadFactory LISTENER_THREAD_FACTORY = new DefaultThreadFactory("I/O listener", true);
63  
64      private final Deque<ExceptionEvent> auditLog;
65      private final int workerCount;
66      private final SingleCoreIOReactor[] workers;
67      private final SingleCoreListeningIOReactor listener;
68      private final MultiCoreIOReactor ioReactor;
69      private final IOWorkers.Selector workerSelector;
70  
71      /**
72       * Creates an instance of DefaultListeningIOReactor with the given configuration.
73       *
74       * @param eventHandlerFactory the factory to create I/O event handlers.
75       * @param ioReactorConfig I/O reactor configuration.
76       * @param listenerThreadFactory the factory to create listener thread.
77       *   Can be {@code null}.
78       *
79       * @since 5.0
80       */
81      public DefaultListeningIOReactor(
82              final IOEventHandlerFactory eventHandlerFactory,
83              final IOReactorConfig ioReactorConfig,
84              final ThreadFactory dispatchThreadFactory,
85              final ThreadFactory listenerThreadFactory,
86              final Decorator<IOSession> ioSessionDecorator,
87              final IOSessionListener sessionListener,
88              final Callback<IOSession> sessionShutdownCallback) {
89          Args.notNull(eventHandlerFactory, "Event handler factory");
90          this.auditLog = new ConcurrentLinkedDeque<>();
91          this.workerCount = ioReactorConfig != null ? ioReactorConfig.getIoThreadCount() : IOReactorConfig.DEFAULT.getIoThreadCount();
92          this.workers = new SingleCoreIOReactor[workerCount];
93          final Thread[] threads = new Thread[workerCount + 1];
94          for (int i = 0; i < this.workers.length; i++) {
95              final SingleCoreIOReactor dispatcher = new SingleCoreIOReactor(
96                      auditLog,
97                      eventHandlerFactory,
98                      ioReactorConfig != null ? ioReactorConfig : IOReactorConfig.DEFAULT,
99                      ioSessionDecorator,
100                     sessionListener,
101                     sessionShutdownCallback);
102             this.workers[i] = dispatcher;
103             threads[i + 1] = (dispatchThreadFactory != null ? dispatchThreadFactory : DISPATCH_THREAD_FACTORY).newThread(new IOReactorWorker(dispatcher));
104         }
105         final IOReactor[] ioReactors = new IOReactor[this.workerCount + 1];
106         System.arraycopy(this.workers, 0, ioReactors, 1, this.workerCount);
107         this.listener = new SingleCoreListeningIOReactor(auditLog, ioReactorConfig, new Callback<SocketChannel>() {
108 
109             @Override
110             public void execute(final SocketChannel channel) {
111                 enqueueChannel(channel);
112             }
113 
114         });
115         ioReactors[0] = this.listener;
116         threads[0] = (listenerThreadFactory != null ? listenerThreadFactory : LISTENER_THREAD_FACTORY).newThread(new IOReactorWorker(listener));
117 
118         this.ioReactor = new MultiCoreIOReactor(ioReactors, threads);
119 
120         workerSelector = IOWorkers.newSelector(workers);
121     }
122 
123     /**
124      * Creates an instance of DefaultListeningIOReactor with the given configuration.
125      *
126      * @param eventHandlerFactory the factory to create I/O event handlers.
127      * @param config I/O reactor configuration.
128      *   Can be {@code null}.
129      *
130      * @since 5.0
131      */
132     public DefaultListeningIOReactor(
133             final IOEventHandlerFactory eventHandlerFactory,
134             final IOReactorConfig config,
135             final Callback<IOSession> sessionShutdownCallback) {
136         this(eventHandlerFactory, config, null, null, null, null, sessionShutdownCallback);
137     }
138 
139     /**
140      * Creates an instance of DefaultListeningIOReactor with default configuration.
141      *
142      * @param eventHandlerFactory the factory to create I/O event handlers.
143      *
144      * @since 5.0
145      */
146     public DefaultListeningIOReactor(final IOEventHandlerFactory eventHandlerFactory) {
147         this(eventHandlerFactory, null, null);
148     }
149 
150     @Override
151     public void start() {
152         ioReactor.start();
153     }
154 
155     @Override
156     public Future<ListenerEndpoint> listen(final SocketAddress address, final FutureCallback<ListenerEndpoint> callback) {
157         return listener.listen(address, callback);
158     }
159 
160     public Future<ListenerEndpoint> listen(final SocketAddress address) {
161         return listen(address, null);
162     }
163 
164     @Override
165     public Set<ListenerEndpoint> getEndpoints() {
166         return listener.getEndpoints();
167     }
168 
169     @Override
170     public void pause() throws IOException {
171         listener.pause();
172     }
173 
174     @Override
175     public void resume() throws IOException {
176         listener.resume();
177     }
178 
179     @Override
180     public IOReactorStatus getStatus() {
181         return ioReactor.getStatus();
182     }
183 
184     @Override
185     IOWorkers.Selector getWorkerSelector() {
186         return workerSelector;
187     }
188 
189     @Override
190     public List<ExceptionEvent> getExceptionLog() {
191         return auditLog.isEmpty() ? Collections.<ExceptionEvent>emptyList() : new ArrayList<>(auditLog);
192     }
193 
194     private void enqueueChannel(final SocketChannel socketChannel) {
195         try {
196             workerSelector.next().enqueueChannel(socketChannel);
197         } catch (final IOReactorShutdownException ex) {
198             initiateShutdown();
199         }
200     }
201 
202 
203     @Override
204     public void initiateShutdown() {
205         ioReactor.initiateShutdown();
206     }
207 
208     @Override
209     public void awaitShutdown(final TimeValue waitTime) throws InterruptedException {
210         ioReactor.awaitShutdown(waitTime);
211     }
212 
213     @Override
214     public void close(final CloseMode closeMode) {
215         ioReactor.close(closeMode);
216     }
217 
218     @Override
219     public void close() throws IOException {
220         ioReactor.close();
221     }
222 
223 }