From 0bc3ef6218fd70caa37de1aef5360e82f861a20a Mon Sep 17 00:00:00 2001 From: Austin Brooks Date: Wed, 19 Aug 2026 15:38:49 -0700 Subject: [PATCH] RUM-18168: Fix DataFlusher racing the upload scheduler and duplicating events DataFlusher.flush() (the explicit Flush() call) can race the SDK's own periodic upload scheduler: both can independently claim and upload the same on-disk batch file, producing duplicate RUM events. drainAndShutdownExecutors() shuts down the persistence/upload executors and waits for them to terminate before DataFlusher runs, but it waited only DRAIN_WAIT_SECONDS (10s) - shorter than the upload call's own timeout (NETWORK_TIMEOUT_MS, 45s). An in-flight upload could still be running when that wait gave up, so drainAndShutdownExecutors() returned anyway while the upload kept racing DataFlusher for the same batch file. Fix: uploadExecutorService now waits up to NETWORK_TIMEOUT_MS instead of DRAIN_WAIT_SECONDS. Also log a warning if an executor doesn't terminate in time, so this doesn't silently regress. --- .../android/core/internal/CoreFeature.kt | 30 ++++++++++++++++-- .../android/core/internal/CoreFeatureTest.kt | 31 ++++++++++++++++++- 2 files changed, 57 insertions(+), 4 deletions(-) diff --git a/dd-sdk-android-core/src/main/kotlin/com/datadog/android/core/internal/CoreFeature.kt b/dd-sdk-android-core/src/main/kotlin/com/datadog/android/core/internal/CoreFeature.kt index 37ae8375ec..a08285fc50 100644 --- a/dd-sdk-android-core/src/main/kotlin/com/datadog/android/core/internal/CoreFeature.kt +++ b/dd-sdk-android-core/src/main/kotlin/com/datadog/android/core/internal/CoreFeature.kt @@ -359,7 +359,7 @@ internal class CoreFeature( contextExecutorService.queue.drainTo(contextTasks) contextExecutorService.shutdown() - contextExecutorService.awaitTermination(DRAIN_WAIT_SECONDS, TimeUnit.SECONDS) + awaitTerminationLogged(contextExecutorService, "contextExecutorService", DRAIN_WAIT_SECONDS, TimeUnit.SECONDS) contextTasks.forEach { it.run() } @@ -376,14 +376,38 @@ internal class CoreFeature( persistenceExecutorService.shutdown() uploadExecutorService.shutdown() - persistenceExecutorService.awaitTermination(DRAIN_WAIT_SECONDS, TimeUnit.SECONDS) - uploadExecutorService.awaitTermination(DRAIN_WAIT_SECONDS, TimeUnit.SECONDS) + awaitTerminationLogged( + persistenceExecutorService, + "persistenceExecutorService", + DRAIN_WAIT_SECONDS, + TimeUnit.SECONDS + ) + // uploadExecutorService can be mid-upload when this runs, and only NETWORK_TIMEOUT_MS + // bounds how long that upload can take. Failing to wait long enough here can lead to + // a DataFlusher race where it uploads the same batch twice. See RUM-18168. + awaitTerminationLogged( + uploadExecutorService, + "uploadExecutorService", + NETWORK_TIMEOUT_MS, + TimeUnit.MILLISECONDS + ) ioTasks.forEach { it.run() } } + @Suppress("UnsafeThirdPartyFunctionCall") // Used in Nightly tests only + private fun awaitTerminationLogged(executor: ExecutorService, executorName: String, wait: Long, unit: TimeUnit) { + if (!executor.awaitTermination(wait, unit)) { + internalLogger.log( + InternalLogger.Level.WARN, + InternalLogger.Target.MAINTAINER, + { "drainAndShutdownExecutors: $executorName did not terminate within $wait $unit" } + ) + } + } + // region Internal @WorkerThread diff --git a/dd-sdk-android-core/src/test/kotlin/com/datadog/android/core/internal/CoreFeatureTest.kt b/dd-sdk-android-core/src/test/kotlin/com/datadog/android/core/internal/CoreFeatureTest.kt index f199d5674d..10caffc66e 100644 --- a/dd-sdk-android-core/src/test/kotlin/com/datadog/android/core/internal/CoreFeatureTest.kt +++ b/dd-sdk-android-core/src/test/kotlin/com/datadog/android/core/internal/CoreFeatureTest.kt @@ -1555,10 +1555,39 @@ internal class CoreFeatureTest { // Then inOrder(mockUploadService) { verify(mockUploadService).shutdown() - verify(mockUploadService).awaitTermination(10, TimeUnit.SECONDS) + verify(mockUploadService).awaitTermination(CoreFeature.NETWORK_TIMEOUT_MS, TimeUnit.MILLISECONDS) } } + @Test + fun `M log a warning W drainAndShutdownExecutors() { upload executor doesn't terminate in time }`() { + // Given + testedFeature.initialize( + appContext.mockInstance, + fakeSdkInstanceId, + fakeConfig, + fakeConsent + ) + + val blockingQueue = LinkedBlockingQueue() + val mockUploadService: ScheduledThreadPoolExecutor = mock() + whenever(mockUploadService.queue).thenReturn(blockingQueue) + whenever(mockUploadService.awaitTermination(any(), any())) doReturn false + testedFeature.uploadExecutorService = mockUploadService + + // When + testedFeature.drainAndShutdownExecutors() + + // Then + mockInternalLogger.verifyLog( + InternalLogger.Level.WARN, + InternalLogger.Target.MAINTAINER, + "drainAndShutdownExecutors: uploadExecutorService did not terminate " + + "within ${CoreFeature.NETWORK_TIMEOUT_MS} ${TimeUnit.MILLISECONDS}", + mode = atLeastOnce() + ) + } + @Test fun `M shutdown with wait the context executor W drainAndShutdownExecutors()`() { // Given