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.runtime.common.utils;
020
021import org.apache.reef.annotations.audience.DriverSide;
022import org.apache.reef.annotations.audience.Private;
023import org.apache.reef.tang.util.MonotonicHashMap;
024import org.apache.reef.util.ExceptionHandlingEventHandler;
025import org.apache.reef.wake.EventHandler;
026import org.apache.reef.wake.impl.ThreadPoolStage;
027
028import java.util.Collections;
029import java.util.Map;
030import java.util.Set;
031
032/**
033 * Delayed event router that dispatches messages to the proper event handler by type.
034 * This class is used in EvaluatorManager to isolate user threads from REEF.
035 */
036@Private
037@DriverSide
038public final class DispatchingEStage implements AutoCloseable {
039
040  /**
041   * A map of event handlers, populated in the register() method.
042   */
043  private final Map<Class<?>, EventHandler<?>> handlers =
044      Collections.synchronizedMap(new MonotonicHashMap<Class<?>, EventHandler<?>>());
045  /**
046   * Exception handler, one for all event handlers.
047   */
048  private final EventHandler<Throwable> errorHandler;
049  /**
050   * Thread pool to process delayed event handler invocations.
051   */
052  private final ThreadPoolStage<DelayedOnNext> stage;
053
054  /**
055   * @param errorHandler used for exceptions thrown from the event handlers registered.
056   * @param numThreads   number of threads to allocate to dispatch events.
057   * @param stageName    the name to use for the underlying stage.
058   *                     It will be carried over to name the Thread(s) spawned.
059   */
060  public DispatchingEStage(final EventHandler<Throwable> errorHandler,
061                           final int numThreads,
062                           final String stageName) {
063    this.errorHandler = errorHandler;
064    this.stage = new ThreadPoolStage<>(stageName,
065        new EventHandler<DelayedOnNext>() {
066          @Override
067          public void onNext(final DelayedOnNext promise) {
068            promise.handler.onNext(promise.message);
069          }
070        }, numThreads
071    );
072
073  }
074
075  /**
076   * Constructs a DispatchingEStage that uses the Thread pool and ErrorHandler of another one.
077   *
078   * @param other
079   */
080  public DispatchingEStage(final DispatchingEStage other) {
081    this.errorHandler = other.errorHandler;
082    this.stage = other.stage;
083  }
084
085  /**
086   * Register a new event handler.
087   *
088   * @param type     Message type to process with this handler.
089   * @param handlers A set of handlers that process that type of message.
090   * @param <T>      Message type.
091   * @param <U>      Type of message that event handler supports. Must be a subclass of T.
092   */
093  @SuppressWarnings("checkstyle:hiddenfield")
094  public <T, U extends T> void register(final Class<T> type, final Set<EventHandler<U>> handlers) {
095    this.handlers.put(type, new ExceptionHandlingEventHandler<>(
096        new BroadCastEventHandler<>(handlers), this.errorHandler));
097  }
098
099  /**
100   * Dispatch a new message by type.
101   *
102   * @param type    Type of event handler - must match the register() call.
103   * @param message A message to process. Must be a subclass of T.
104   * @param <T>     Message type that event handler supports.
105   * @param <U>     input message type. Must be a subclass of T.
106   */
107  @SuppressWarnings("unchecked")
108  public <T, U extends T> void onNext(final Class<T> type, final U message) {
109    final EventHandler<T> handler = (EventHandler<T>) this.handlers.get(type);
110    this.stage.onNext(new DelayedOnNext(handler, message));
111  }
112
113  /**
114   * Return true if there are no messages queued or in processing, false otherwise.
115   */
116  public boolean isEmpty() {
117    return this.stage.getQueueLength() == 0;
118  }
119
120  /**
121   * Close the internal thread pool.
122   *
123   * @throws Exception forwarded from EStage.close() call.
124   */
125  @Override
126  public void close() throws Exception {
127    this.stage.close();
128  }
129
130  /**
131   * Delayed EventHandler.onNext() call.
132   * Contains a message object and EventHandler to process it.
133   */
134  private static final class DelayedOnNext {
135
136    private final EventHandler<Object> handler;
137    private final Object message;
138
139    @SuppressWarnings("unchecked")
140    <T, U extends T> DelayedOnNext(final EventHandler<T> handler, final U message) {
141      this.handler = (EventHandler<Object>) handler;
142      this.message = message;
143    }
144  }
145}