Skip to content

Commit 5511d97

Browse files
committed
refactor(messaging): use daemon threads and handle rejected executions
1 parent db87137 commit 5511d97

2 files changed

Lines changed: 101 additions & 4 deletions

File tree

‎src/main/java/com/google/firebase/messaging/FirebaseMessagingClientImpl.java‎

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
import com.google.common.base.Strings;
3939
import com.google.common.collect.ImmutableList;
4040
import com.google.common.collect.ImmutableMap;
41+
import com.google.common.util.concurrent.ThreadFactoryBuilder;
4142
import com.google.firebase.ErrorCode;
4243
import com.google.firebase.FirebaseApp;
4344
import com.google.firebase.FirebaseException;
@@ -61,6 +62,7 @@
6162
import java.util.concurrent.ExecutorService;
6263
import java.util.concurrent.Executors;
6364
import java.util.concurrent.LinkedBlockingQueue;
65+
import java.util.concurrent.RejectedExecutionException;
6466
import java.util.concurrent.ThreadFactory;
6567
import java.util.concurrent.ThreadPoolExecutor;
6668
import java.util.concurrent.TimeUnit;
@@ -111,11 +113,18 @@ private FirebaseMessagingClientImpl(Builder builder) {
111113
}
112114

113115
private static ExecutorService createDefaultExecutor(ThreadFactory threadFactory) {
116+
ThreadFactory baseFactory =
117+
threadFactory != null ? threadFactory : Executors.defaultThreadFactory();
118+
ThreadFactory factory = new ThreadFactoryBuilder()
119+
.setThreadFactory(baseFactory)
120+
.setNameFormat("firebase-messaging-topics-%d")
121+
.setDaemon(true)
122+
.build();
114123
ThreadPoolExecutor pool = new ThreadPoolExecutor(
115124
100, 100,
116125
60L, TimeUnit.SECONDS,
117126
new LinkedBlockingQueue<Runnable>(),
118-
threadFactory != null ? threadFactory : Executors.defaultThreadFactory());
127+
factory);
119128
pool.allowCoreThreadTimeOut(true);
120129
return pool;
121130
}
@@ -244,9 +253,14 @@ private TopicManagementResponse sendTopicManagementRequest(
244253
for (int i = 0; i < registrationTokens.size(); i++) {
245254
final int index = i;
246255
final String token = registrationTokens.get(i);
247-
futures.add(CompletableFuture.supplyAsync(
248-
() -> sendSingleTopicRequest(token, topicName, isSubscribe, index),
249-
this.executor));
256+
try {
257+
futures.add(CompletableFuture.supplyAsync(
258+
() -> sendSingleTopicRequest(token, topicName, isSubscribe, index),
259+
this.executor));
260+
} catch (RejectedExecutionException e) {
261+
futures.add(CompletableFuture.completedFuture(
262+
TopicResult.error(index, "REJECTED_BY_EXECUTOR")));
263+
}
250264
}
251265

252266
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();

‎src/test/java/com/google/firebase/messaging/FirebaseMessagingClientImplTest.java‎

Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,10 +53,15 @@
5353
import java.util.HashMap;
5454
import java.util.List;
5555
import java.util.Map;
56+
import java.util.concurrent.AbstractExecutorService;
57+
import java.util.concurrent.CountDownLatch;
5658
import java.util.concurrent.ExecutorService;
5759
import java.util.concurrent.Executors;
60+
import java.util.concurrent.RejectedExecutionException;
5861
import java.util.concurrent.ThreadFactory;
5962
import java.util.concurrent.TimeUnit;
63+
import java.util.concurrent.atomic.AtomicBoolean;
64+
import java.util.concurrent.atomic.AtomicReference;
6065
import org.junit.Before;
6166
import org.junit.Test;
6267

@@ -776,4 +781,82 @@ public void testCustomThreadFactory() {
776781

777782
assertSame(customThreadFactory, clientWithThreadFactory.getThreadFactory());
778783
}
784+
785+
@Test
786+
public void testDefaultExecutorUsesDaemonThreads() throws Exception {
787+
FirebaseMessagingClientImpl clientWithDefaultExecutor =
788+
FirebaseMessagingClientImpl.builder()
789+
.setProjectId("test-project")
790+
.setJsonFactory(ApiClientUtils.getDefaultJsonFactory())
791+
.setRequestFactory(new MockHttpTransport().createRequestFactory())
792+
.setChildRequestFactory(ApiClientUtils.getDefaultTransport().createRequestFactory())
793+
.build();
794+
795+
final AtomicBoolean isDaemon = new AtomicBoolean();
796+
final AtomicReference<String> threadName = new AtomicReference<>();
797+
final CountDownLatch latch = new CountDownLatch(1);
798+
clientWithDefaultExecutor.getExecutor().execute(() -> {
799+
Thread current = Thread.currentThread();
800+
isDaemon.set(current.isDaemon());
801+
threadName.set(current.getName());
802+
latch.countDown();
803+
});
804+
805+
assertTrue(latch.await(5, TimeUnit.SECONDS));
806+
assertTrue(isDaemon.get());
807+
assertNotNull(threadName.get());
808+
assertTrue(threadName.get().startsWith("firebase-messaging-topics-"));
809+
}
810+
811+
@Test
812+
public void testTopicManagementRejectedExecution() throws Exception {
813+
ExecutorService rejectingExecutor = new AbstractExecutorService() {
814+
@Override
815+
public void shutdown() {}
816+
817+
@Override
818+
public List<Runnable> shutdownNow() {
819+
return ImmutableList.of();
820+
}
821+
822+
@Override
823+
public boolean isShutdown() {
824+
return false;
825+
}
826+
827+
@Override
828+
public boolean isTerminated() {
829+
return false;
830+
}
831+
832+
@Override
833+
public boolean awaitTermination(long timeout, TimeUnit unit) {
834+
return false;
835+
}
836+
837+
@Override
838+
public void execute(Runnable command) {
839+
throw new RejectedExecutionException("Task rejected");
840+
}
841+
};
842+
843+
FirebaseMessagingClientImpl clientWithRejection = FirebaseMessagingClientImpl.builder()
844+
.setProjectId("test-project")
845+
.setJsonFactory(ApiClientUtils.getDefaultJsonFactory())
846+
.setRequestFactory(new MockHttpTransport().createRequestFactory())
847+
.setChildRequestFactory(ApiClientUtils.getDefaultTransport().createRequestFactory())
848+
.setExecutor(rejectingExecutor)
849+
.build();
850+
851+
TopicManagementResponse result = clientWithRejection.subscribeToTopic(
852+
"test-topic", ImmutableList.of("id1", "id2"));
853+
854+
assertEquals(0, result.getSuccessCount());
855+
assertEquals(2, result.getFailureCount());
856+
assertEquals(2, result.getErrors().size());
857+
assertEquals(0, result.getErrors().get(0).getIndex());
858+
assertEquals("rejected-by-executor", result.getErrors().get(0).getReason());
859+
assertEquals(1, result.getErrors().get(1).getIndex());
860+
assertEquals("rejected-by-executor", result.getErrors().get(1).getReason());
861+
}
779862
}

0 commit comments

Comments
 (0)