diff --git a/core/java/com/android/internal/util/ConcurrentUtils.java b/core/java/com/android/internal/util/ConcurrentUtils.java new file mode 100644 index 0000000000000..e35f9f45acfe7 --- /dev/null +++ b/core/java/com/android/internal/util/ConcurrentUtils.java @@ -0,0 +1,89 @@ +/* + * Copyright (C) 2016 The Android Open Source Project + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License + */ + +package com.android.internal.util; + +import android.os.Process; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Utility methods for common functionality using java.util.concurrent package + * + * @hide + */ +public class ConcurrentUtils { + + private ConcurrentUtils() { + } + + /** + * Creates a thread pool using + * {@link java.util.concurrent.Executors#newFixedThreadPool(int, ThreadFactory)} + * + * @param nThreads the number of threads in the pool + * @param poolName base name of the threads in the pool + * @param linuxThreadPriority a Linux priority level. see {@link Process#setThreadPriority(int)} + * @return the newly created thread pool + */ + public static ExecutorService newFixedThreadPool(int nThreads, String poolName, + int linuxThreadPriority) { + return Executors.newFixedThreadPool(nThreads, + new ThreadFactory() { + private final AtomicInteger threadNum = new AtomicInteger(0); + + @Override + public Thread newThread(final Runnable r) { + return new Thread(poolName + threadNum.incrementAndGet()) { + @Override + public void run() { + Process.setThreadPriority(linuxThreadPriority); + r.run(); + } + }; + } + }); + } + + /** + * Waits if necessary for the computation to complete, and then retrieves its result. + *
If {@code InterruptedException} occurs, this method will interrupt the current thread + * and throw {@code IllegalStateException}
+ * + * @param future future to wait for result + * @param description short description of the operation + * @return the computed result + * @throws IllegalStateException if interrupted during wait + * @throws RuntimeException if an error occurs while waiting for {@link Future#get()} + * @see Future#get() + */ + public staticSystem services can {@link #submit(Runnable)} tasks for execution during boot.
+ * The pool will be shut down after {@link SystemService#PHASE_BOOT_COMPLETED}.
+ * New tasks should not be submitted afterwards.
+ *
+ * @hide
+ */
+public class SystemServerInitThreadPool {
+ private static final String TAG = SystemServerInitThreadPool.class.getSimpleName();
+ private static final int SHUTDOWN_TIMEOUT_MILLIS = 20000;
+ private static final boolean IS_DEBUGGABLE = Build.IS_DEBUGGABLE;
+
+ private static SystemServerInitThreadPool sInstance;
+
+ private ExecutorService mService = ConcurrentUtils.newFixedThreadPool(2,
+ "system-server-init-thread", Process.THREAD_PRIORITY_FOREGROUND);
+
+ public static synchronized SystemServerInitThreadPool get() {
+ if (sInstance == null) {
+ sInstance = new SystemServerInitThreadPool();
+ }
+ Preconditions.checkState(sInstance.mService != null, "Cannot get " + TAG
+ + " - it has been shut down");
+ return sInstance;
+ }
+
+ public Future> submit(Runnable runnable, String description) {
+ if (IS_DEBUGGABLE) {
+ return mService.submit(() -> {
+ Slog.d(TAG, "Started executing " + description);
+ try {
+ runnable.run();
+ } catch (RuntimeException e) {
+ Slog.e(TAG, "Failure in " + description + ": " + e, e);
+ throw e;
+ }
+ Slog.d(TAG, "Finished executing " + description);
+ });
+ }
+ return mService.submit(runnable);
+ }
+
+ static synchronized void shutdown() {
+ if (sInstance != null && sInstance.mService != null) {
+ sInstance.mService.shutdown();
+ boolean terminated;
+ try {
+ terminated = sInstance.mService.awaitTermination(SHUTDOWN_TIMEOUT_MILLIS,
+ TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException(TAG + " init interrupted");
+ }
+ List