chore(Android): migrate MessageQueueThread and it's implementation to Kotlin (#48652)

Summary:
Continuing our usual journey, this time migrating MessageQueueThread. Was not expecting to see that many assertions in ReactContext.
One important thing to note on this PR: I had to add an extra Throw RuntimeException to `startNewBackgroundThread`. It already had one from the`dataFuture.getOrThrow()`, but if your thread is not associated with a Looper, calling `myLooper` (Line 203 of MessageQueueThreadImpl) Can return null. Until now, this would have been a `NPE`, i just decided to make it a `RuntimeException` with the message `Looper not found for thread`. Let me know if you want me to make that function return a nullable, or throw another message or exception, or if you want me to treat Looper as Nullable in the whole class.

## Changelog:

[INTERNAL] [FIXED] - Migrate MessageQueueThread and MessageQueueThreadImpl to Kotlin

Pull Request resolved: https://github.com/facebook/react-native/pull/48652

Test Plan: <img width="1318" alt="Screenshot 2025-01-13 at 21 03 00" src="https://github.com/user-attachments/assets/462cc8af-4648-4437-9260-5ffa6c69e763" />

Reviewed By: tdn120

Differential Revision: D68155120

Pulled By: rshest

fbshipit-source-id: 1082a832df5b1d8ee64ca7be26e4b85e79152f88
This commit is contained in:
Parsa Nasirimehr
2025-01-16 07:34:03 -08:00
committed by Facebook GitHub Bot
parent 9f1236eff2
commit 166347ead9
5 changed files with 223 additions and 262 deletions
@@ -1664,13 +1664,15 @@ public class com/facebook/react/bridge/queue/MessageQueueThreadHandler : android
public fun dispatchMessage (Landroid/os/Message;)V
}
public class com/facebook/react/bridge/queue/MessageQueueThreadImpl : com/facebook/react/bridge/queue/MessageQueueThread {
public final class com/facebook/react/bridge/queue/MessageQueueThreadImpl : com/facebook/react/bridge/queue/MessageQueueThread {
public static final field Companion Lcom/facebook/react/bridge/queue/MessageQueueThreadImpl$Companion;
public synthetic fun <init> (Ljava/lang/String;Landroid/os/Looper;Lcom/facebook/react/bridge/queue/QueueThreadExceptionHandler;Lcom/facebook/react/bridge/queue/MessageQueueThreadPerfStats;Lkotlin/jvm/internal/DefaultConstructorMarker;)V
public fun assertIsOnThread ()V
public fun assertIsOnThread (Ljava/lang/String;)V
public fun callOnQueue (Ljava/util/concurrent/Callable;)Ljava/util/concurrent/Future;
public static fun create (Lcom/facebook/react/bridge/queue/MessageQueueThreadSpec;Lcom/facebook/react/bridge/queue/QueueThreadExceptionHandler;)Lcom/facebook/react/bridge/queue/MessageQueueThreadImpl;
public fun getLooper ()Landroid/os/Looper;
public fun getName ()Ljava/lang/String;
public static final fun create (Lcom/facebook/react/bridge/queue/MessageQueueThreadSpec;Lcom/facebook/react/bridge/queue/QueueThreadExceptionHandler;)Lcom/facebook/react/bridge/queue/MessageQueueThreadImpl;
public final fun getLooper ()Landroid/os/Looper;
public final fun getName ()Ljava/lang/String;
public fun getPerfStats ()Lcom/facebook/react/bridge/queue/MessageQueueThreadPerfStats;
public fun isIdle ()Z
public fun isOnThread ()Z
@@ -1679,6 +1681,10 @@ public class com/facebook/react/bridge/queue/MessageQueueThreadImpl : com/facebo
public fun runOnQueue (Ljava/lang/Runnable;)Z
}
public final class com/facebook/react/bridge/queue/MessageQueueThreadImpl$Companion {
public final fun create (Lcom/facebook/react/bridge/queue/MessageQueueThreadSpec;Lcom/facebook/react/bridge/queue/QueueThreadExceptionHandler;)Lcom/facebook/react/bridge/queue/MessageQueueThreadImpl;
}
public class com/facebook/react/bridge/queue/MessageQueueThreadPerfStats {
public field cpuTime J
public field wallTime J
@@ -5,74 +5,67 @@
* LICENSE file in the root directory of this source tree.
*/
package com.facebook.react.bridge.queue;
package com.facebook.react.bridge.queue
import androidx.annotation.Nullable;
import com.facebook.proguard.annotations.DoNotStrip;
import java.util.concurrent.Callable;
import java.util.concurrent.Future;
import com.facebook.proguard.annotations.DoNotStripAny
import com.facebook.react.bridge.AssertionException
import java.util.concurrent.Callable
import java.util.concurrent.Future
/** Encapsulates a Thread that can accept Runnables. */
@DoNotStrip
@DoNotStripAny
public interface MessageQueueThread {
/**
* Runs the given Runnable on this Thread. It will be submitted to the end of the event queue even
* if it is being submitted from the same queue Thread.
*/
@DoNotStrip
boolean runOnQueue(Runnable runnable);
public fun runOnQueue(runnable: Runnable): Boolean
/**
* Runs the given Callable on this Thread. It will be submitted to the end of the event queue even
* if it is being submitted from the same queue Thread.
*/
@DoNotStrip
<T> Future<T> callOnQueue(final Callable<T> callable);
public fun <T> callOnQueue(callable: Callable<T>): Future<T>
/**
* @return whether the current Thread is also the Thread associated with this MessageQueueThread.
* Tells whether the current Thread is also the Thread associated with this MessageQueueThread.
*/
@DoNotStrip
boolean isOnThread();
public fun isOnThread(): Boolean
/**
* Asserts {@link #isOnThread()}, throwing a {@link AssertionException} (NOT an {@link
* AssertionError}) if the assertion fails.
* Asserts [isOnThread], throwing a [AssertionException] (NOT an [AssertionError]) if the
* assertion fails.
*/
@DoNotStrip
void assertIsOnThread();
public fun assertIsOnThread()
/**
* Asserts {@link #isOnThread()}, throwing a {@link AssertionException} (NOT an {@link
* AssertionError}) if the assertion fails. The given message is appended to the error.
* Asserts [isOnThread], throwing a [AssertionException] (NOT an [AssertionError]) if the
* assertion fails. The given message is appended to the error.
*/
@DoNotStrip
void assertIsOnThread(String message);
public fun assertIsOnThread(message: String)
/**
* Quits this MessageQueueThread. If called from this MessageQueueThread, this will be the last
* thing the thread runs. If called from a separate thread, this will block until the thread can
* be quit and joined.
*/
@DoNotStrip
void quitSynchronous();
public fun quitSynchronous()
/**
* Returns the perf counters taken when the framework was started. This method is intended to be
* used for instrumentation purposes.
*/
@Nullable
@DoNotStrip
MessageQueueThreadPerfStats getPerfStats();
public fun getPerfStats(): MessageQueueThreadPerfStats?
/**
* Resets the perf counters. This is useful if the RN threads are being re-used. This method is
* intended to be used for instrumentation purposes.
*/
@DoNotStrip
void resetPerfStats();
public fun resetPerfStats()
/** Returns true if the message queue is idle */
@DoNotStrip
boolean isIdle();
/**
* Resets the perf counters. This is useful if the RN threads are being re-used. This method is
* intended to be used for instrumentation purposes.
*/
public fun isIdle(): Boolean
}
@@ -1,226 +0,0 @@
/*
* Copyright (c) Meta Platforms, Inc. and affiliates.
*
* This source code is licensed under the MIT license found in the
* LICENSE file in the root directory of this source tree.
*/
package com.facebook.react.bridge.queue;
import android.os.Looper;
import android.os.Process;
import android.os.SystemClock;
import android.util.Pair;
import androidx.annotation.Nullable;
import com.facebook.common.logging.FLog;
import com.facebook.proguard.annotations.DoNotStrip;
import com.facebook.react.bridge.AssertionException;
import com.facebook.react.bridge.SoftAssertions;
import com.facebook.react.common.ReactConstants;
import com.facebook.react.common.futures.SimpleSettableFuture;
import java.util.concurrent.Callable;
import java.util.concurrent.Future;
/** Encapsulates a Thread that has a {@link Looper} running on it that can accept Runnables. */
@DoNotStrip
public class MessageQueueThreadImpl implements MessageQueueThread {
private final String mName;
private final Looper mLooper;
private final MessageQueueThreadHandler mHandler;
private final String mAssertionErrorMessage;
private final @Nullable MessageQueueThreadPerfStats mPerfStats;
private volatile boolean mIsFinished = false;
private MessageQueueThreadImpl(
String name, Looper looper, QueueThreadExceptionHandler exceptionHandler) {
this(name, looper, exceptionHandler, null);
}
private MessageQueueThreadImpl(
String name,
Looper looper,
QueueThreadExceptionHandler exceptionHandler,
@Nullable MessageQueueThreadPerfStats stats) {
mName = name;
mLooper = looper;
mHandler = new MessageQueueThreadHandler(looper, exceptionHandler);
mPerfStats = stats;
mAssertionErrorMessage = "Expected to be called from the '" + getName() + "' thread!";
}
/**
* Runs the given Runnable on this Thread. It will be submitted to the end of the event queue even
* if it is being submitted from the same queue Thread.
*/
@DoNotStrip
@Override
public boolean runOnQueue(Runnable runnable) {
if (mIsFinished) {
FLog.w(
ReactConstants.TAG,
"Tried to enqueue runnable on already finished thread: '"
+ getName()
+ "... dropping Runnable.");
return false;
}
mHandler.post(runnable);
return true;
}
@DoNotStrip
@Override
public <T> Future<T> callOnQueue(final Callable<T> callable) {
final SimpleSettableFuture<T> future = new SimpleSettableFuture<>();
runOnQueue(
() -> {
try {
future.set(callable.call());
} catch (Exception e) {
future.setException(e);
}
});
return future;
}
/**
* @return whether the current Thread is also the Thread associated with this MessageQueueThread.
*/
@DoNotStrip
@Override
public boolean isOnThread() {
return mLooper.getThread() == Thread.currentThread();
}
/**
* Asserts {@link #isOnThread()}, throwing a {@link AssertionException} (NOT an {@link
* AssertionError}) if the assertion fails.
*/
@DoNotStrip
@Override
public void assertIsOnThread() {
SoftAssertions.assertCondition(isOnThread(), mAssertionErrorMessage);
}
/**
* Asserts {@link #isOnThread()}, throwing a {@link AssertionException} (NOT an {@link
* AssertionError}) if the assertion fails.
*/
@DoNotStrip
@Override
public void assertIsOnThread(String message) {
SoftAssertions.assertCondition(
isOnThread(),
new StringBuilder().append(mAssertionErrorMessage).append(" ").append(message).toString());
}
/**
* Quits this queue's Looper. If that Looper was running on a different Thread than the current
* Thread, also waits for the last message being processed to finish and the Thread to die.
*/
@DoNotStrip
@Override
public void quitSynchronous() {
mIsFinished = true;
mLooper.quit();
if (mLooper.getThread() != Thread.currentThread()) {
try {
mLooper.getThread().join();
} catch (InterruptedException e) {
throw new RuntimeException("Got interrupted waiting to join thread " + mName);
}
}
}
@Nullable
@DoNotStrip
@Override
public MessageQueueThreadPerfStats getPerfStats() {
return mPerfStats;
}
@DoNotStrip
@Override
public void resetPerfStats() {
assignToPerfStats(mPerfStats, -1, -1);
runOnQueue(
() -> {
long wallTime = SystemClock.uptimeMillis();
long cpuTime = SystemClock.currentThreadTimeMillis();
assignToPerfStats(mPerfStats, wallTime, cpuTime);
});
}
@DoNotStrip
@Override
public boolean isIdle() {
return mLooper.getQueue().isIdle();
}
private static void assignToPerfStats(
@Nullable MessageQueueThreadPerfStats stats, long wall, long cpu) {
if (stats == null) {
return;
}
stats.wallTime = wall;
stats.cpuTime = cpu;
}
public Looper getLooper() {
return mLooper;
}
public String getName() {
return mName;
}
public static MessageQueueThreadImpl create(
MessageQueueThreadSpec spec, QueueThreadExceptionHandler exceptionHandler) {
switch (spec.getThreadType()) {
case MAIN_UI:
return createForMainThread(spec.getName(), exceptionHandler);
case NEW_BACKGROUND:
return startNewBackgroundThread(spec.getName(), spec.getStackSize(), exceptionHandler);
default:
throw new RuntimeException("Unknown thread type: " + spec.getThreadType());
}
}
/**
* @return a MessageQueueThreadImpl corresponding to Android's main UI thread.
*/
private static MessageQueueThreadImpl createForMainThread(
String name, QueueThreadExceptionHandler exceptionHandler) {
return new MessageQueueThreadImpl(name, Looper.getMainLooper(), exceptionHandler);
}
/**
* Creates and starts a new MessageQueueThreadImpl encapsulating a new Thread with a new Looper
* running on it. Give it a name for easier debugging and optionally a suggested stack size. When
* this method exits, the new MessageQueueThreadImpl is ready to receive events.
*/
private static MessageQueueThreadImpl startNewBackgroundThread(
final String name, long stackSize, QueueThreadExceptionHandler exceptionHandler) {
final SimpleSettableFuture<Pair<Looper, MessageQueueThreadPerfStats>> dataFuture =
new SimpleSettableFuture<>();
Thread bgThread =
new Thread(
null,
() -> {
Process.setThreadPriority(Process.THREAD_PRIORITY_DISPLAY);
Looper.prepare();
MessageQueueThreadPerfStats stats = new MessageQueueThreadPerfStats();
long wallTime = SystemClock.uptimeMillis();
long cpuTime = SystemClock.currentThreadTimeMillis();
assignToPerfStats(stats, wallTime, cpuTime);
dataFuture.set(new Pair<>(Looper.myLooper(), stats));
Looper.loop();
},
"mqt_" + name,
stackSize);
bgThread.start();
Pair<Looper, MessageQueueThreadPerfStats> pair = dataFuture.getOrThrow();
return new MessageQueueThreadImpl(name, pair.first, exceptionHandler, pair.second);
}
}
@@ -0,0 +1,188 @@
/*
* Copyright (c) Meta Platforms, Inc. and affiliates.
*
* This source code is licensed under the MIT license found in the
* LICENSE file in the root directory of this source tree.
*/
package com.facebook.react.bridge.queue
import android.os.Looper
import android.os.Process
import android.os.SystemClock
import android.util.Pair
import com.facebook.common.logging.FLog
import com.facebook.proguard.annotations.DoNotStripAny
import com.facebook.react.bridge.AssertionException
import com.facebook.react.bridge.SoftAssertions
import com.facebook.react.bridge.queue.MessageQueueThreadSpec.ThreadType
import com.facebook.react.common.ReactConstants
import com.facebook.react.common.futures.SimpleSettableFuture
import java.util.concurrent.Callable
import java.util.concurrent.Future
import kotlin.concurrent.Volatile
/** Encapsulates a Thread that has a [Looper] running on it that can accept Runnables. */
@DoNotStripAny
public class MessageQueueThreadImpl
private constructor(
public val name: String,
public val looper: Looper,
exceptionHandler: QueueThreadExceptionHandler,
private val stats: MessageQueueThreadPerfStats? = null
) : MessageQueueThread {
private val handler = MessageQueueThreadHandler(looper, exceptionHandler)
private val assertionErrorMessage = "Expected to be called from the '$name' thread!"
@Volatile private var isFinished = false
/**
* Runs the given Runnable on this Thread. It will be submitted to the end of the event queue even
* if it is being submitted from the same queue Thread.
*/
public override fun runOnQueue(runnable: Runnable): Boolean {
if (isFinished) {
FLog.w(
ReactConstants.TAG,
"Tried to enqueue runnable on already finished thread: '$name... dropping Runnable.")
return false
}
handler.post(runnable)
return true
}
public override fun <T> callOnQueue(callable: Callable<T>): Future<T> {
val future = SimpleSettableFuture<T>()
runOnQueue {
try {
future.set(callable.call())
} catch (e: Exception) {
future.setException(e)
}
}
return future
}
/**
* @return whether the current Thread is also the Thread associated with this MessageQueueThread.
*/
override fun isOnThread(): Boolean = looper.thread === Thread.currentThread()
/**
* Asserts [isOnThread], throwing a [AssertionException] (NOT an [AssertionError]) if the
* assertion fails.
*/
@Throws(AssertionException::class)
override fun assertIsOnThread() {
SoftAssertions.assertCondition(isOnThread(), assertionErrorMessage)
}
/**
* Asserts [isOnThread], throwing a [AssertionException] (NOT an [AssertionError]) if the
* assertion fails.
*/
@Throws(AssertionException::class)
public override fun assertIsOnThread(message: String) {
SoftAssertions.assertCondition(
isOnThread(),
StringBuilder().append(assertionErrorMessage).append(" ").append(message).toString())
}
/**
* Quits this queue's Looper. If that Looper was running on a different Thread than the current
* Thread, also waits for the last message being processed to finish and the Thread to die.
*/
@Throws(RuntimeException::class)
override fun quitSynchronous() {
isFinished = true
looper.quit()
if (looper.thread !== Thread.currentThread()) {
try {
looper.thread.join()
} catch (e: InterruptedException) {
throw RuntimeException("Got interrupted waiting to join thread $name")
}
}
}
override fun getPerfStats(): MessageQueueThreadPerfStats? = stats
override fun resetPerfStats() {
assignToPerfStats(stats, -1, -1)
runOnQueue {
val wallTime = SystemClock.uptimeMillis()
val cpuTime = SystemClock.currentThreadTimeMillis()
assignToPerfStats(stats, wallTime, cpuTime)
}
}
public override fun isIdle(): Boolean = looper.queue.isIdle
public companion object {
private fun assignToPerfStats(stats: MessageQueueThreadPerfStats?, wall: Long, cpu: Long) {
stats?.let { s ->
s.wallTime = wall
s.cpuTime = cpu
}
}
@JvmStatic
@Throws(RuntimeException::class)
public fun create(
spec: MessageQueueThreadSpec,
exceptionHandler: QueueThreadExceptionHandler
): MessageQueueThreadImpl {
return when (spec.threadType) {
ThreadType.MAIN_UI -> createForMainThread(spec.name, exceptionHandler)
ThreadType.NEW_BACKGROUND ->
startNewBackgroundThread(spec.name, spec.stackSize, exceptionHandler)
else -> throw RuntimeException("Unknown thread type: " + spec.threadType)
}
}
/** Returns a MessageQueueThreadImpl corresponding to Android's main UI thread. */
private fun createForMainThread(
name: String,
exceptionHandler: QueueThreadExceptionHandler
): MessageQueueThreadImpl =
MessageQueueThreadImpl(name, Looper.getMainLooper(), exceptionHandler)
/**
* Creates and starts a new MessageQueueThreadImpl encapsulating a new Thread with a new Looper
* running on it. Give it a name for easier debugging and optionally a suggested stack size.
* When this method exits, the new MessageQueueThreadImpl is ready to receive events. throws a
* Runtime exception if there was no looper for current thread or looper and stats couldn't be
* made
*/
@Throws(RuntimeException::class)
private fun startNewBackgroundThread(
name: String,
stackSize: Long,
exceptionHandler: QueueThreadExceptionHandler
): MessageQueueThreadImpl {
val dataFuture = SimpleSettableFuture<Pair<Looper?, MessageQueueThreadPerfStats>>()
val bgThread =
Thread(
null,
{
Process.setThreadPriority(Process.THREAD_PRIORITY_DISPLAY)
Looper.prepare()
val stats = MessageQueueThreadPerfStats()
val wallTime = SystemClock.uptimeMillis()
val cpuTime = SystemClock.currentThreadTimeMillis()
assignToPerfStats(stats, wallTime, cpuTime)
dataFuture.set(Pair(Looper.myLooper(), stats))
Looper.loop()
},
"mqt_$name",
stackSize)
bgThread.start()
val pair = dataFuture.getOrThrow()
val looper = pair?.first ?: throw RuntimeException("Looper not found for thread")
return MessageQueueThreadImpl(name, looper, exceptionHandler, pair.second)
}
}
}
@@ -17,7 +17,7 @@ import java.util.concurrent.TimeoutException
* A super simple Future-like class that can safely notify another Thread when a value is ready.
* Does not support canceling.
*/
internal class SimpleSettableFuture<T> : Future<T?> {
internal class SimpleSettableFuture<T> : Future<T> {
private val readyLatch = CountDownLatch(1)
private var result: T? = null