From 746282f58c85f38a9fd909dbfe584f6f5bbe1ed7 Mon Sep 17 00:00:00 2001 From: Fabrizio Demaria Date: Fri, 7 Aug 2026 10:56:35 +0200 Subject: [PATCH 1/3] feat: improve event delivery reliability for tracked events Upload pending batches at startup, flush and upload on stop(), and add an optional periodic flush interval so low-volume events are not left on disk indefinitely. --- .../java/com/spotify/confidence/Confidence.kt | 7 +- .../spotify/confidence/EventSenderEngine.kt | 100 ++++++++++---- .../EventSenderEngineReliabilityTest.kt | 129 ++++++++++++++++++ .../openfeature/ConfidenceFeatureProvider.kt | 1 + 4 files changed, 205 insertions(+), 32 deletions(-) create mode 100644 Confidence/src/test/java/com/spotify/confidence/EventSenderEngineReliabilityTest.kt diff --git a/Confidence/src/main/java/com/spotify/confidence/Confidence.kt b/Confidence/src/main/java/com/spotify/confidence/Confidence.kt index 431be4c4..75dc3730 100644 --- a/Confidence/src/main/java/com/spotify/confidence/Confidence.kt +++ b/Confidence/src/main/java/com/spotify/confidence/Confidence.kt @@ -376,6 +376,7 @@ object ConfidenceFactory { * @param loggingLevel allows to print warnings or debugging information to the local console. * @param timeoutMillis sets a timeout for completing an HTTP call. Defaults to 10 seconds * @param visitorIdContextKey key to use for the visitor id in the context. Defaults to "visitor_id". + * @param eventFlushIntervalMillis optional periodic flush interval in milliseconds. Disabled by default. */ fun create( context: Context, @@ -385,7 +386,8 @@ object ConfidenceFactory { dispatcher: CoroutineDispatcher = Dispatchers.IO, loggingLevel: LoggingLevel = LoggingLevel.WARN, timeoutMillis: Long = 10000, - visitorIdContextKey: String = VISITOR_ID_CONTEXT_KEY + visitorIdContextKey: String = VISITOR_ID_CONTEXT_KEY, + eventFlushIntervalMillis: Long? = null ): Confidence { val debugLogger: DebugLogger? = if (loggingLevel == LoggingLevel.NONE) { null @@ -400,7 +402,8 @@ object ConfidenceFactory { flushPolicies = listOf(minBatchSizeFlushPolicy), sdkMetadata = sdkMetadata, dispatcher = dispatcher, - debugLogger = debugLogger + debugLogger = debugLogger, + flushIntervalMillis = eventFlushIntervalMillis ) val flagApplierClient = FlagApplierClientImpl( clientSecret, diff --git a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt index 877729ec..d665b1da 100644 --- a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt +++ b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt @@ -8,10 +8,14 @@ import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.CoroutineExceptionHandler import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking import okhttp3.OkHttpClient import java.io.File @@ -30,7 +34,8 @@ internal class EventSenderEngineImpl( private val clock: Clock = Clock.CalendarBacked.systemUTC(), private val dispatcher: CoroutineDispatcher = Dispatchers.IO, private val sdkMetadata: SdkMetadata, - private val debugLogger: DebugLogger? + private val debugLogger: DebugLogger?, + private val flushIntervalMillis: Long? = null ) : EventSenderEngine { private val writeReqChannel: Channel = Channel() private val sendChannel: Channel = Channel() @@ -43,13 +48,15 @@ internal class EventSenderEngineImpl( debugLogger?.logMessage(message = "EventSenderEngine error: $e", isWarning = true) } } + private var flushIntervalJob: Job? = null + @Volatile + private var isStopped = false init { flushPolicies.add(ManualFlushPolicy) coroutineScope.launch(exceptionHandler) { for (event in writeReqChannel) { if (event.eventDefinition != manualFlushEvent.eventDefinition) { - // skip storing manual flush event eventStorage.writeEvent(event) debugLogger?.logEvent(action = "DiskWrite ", event = event) } @@ -69,36 +76,24 @@ internal class EventSenderEngineImpl( } } - // upload might throw exceptions coroutineScope.launch(exceptionHandler) { for (flush in sendChannel) { - eventStorage.rollover() - val readyFiles = eventStorage.batchReadyFiles() - for (readyFile in readyFiles) { - val events = eventStorage.eventsFor(readyFile) - .map { e -> - EngineEvent( - "eventDefinitions/${e.eventDefinition}", - e.eventTime, - e.payload - ) - } - val batch = EventBatchRequest( - clientSecret = clientSecret, - events = events, - sendTime = clock.currentTime(), - sdk = Sdk(sdkMetadata.sdkId, sdkMetadata.sdkVersion) - ) - runCatching { - val shouldCleanup = uploader.upload(batch) - debugLogger?.logMessage(message = "Uploading events") - if (shouldCleanup) { - readyFile.delete() - } - } + uploadReadyBatches(sealCurrentBatch = true) + } + } + + if (flushIntervalMillis != null && flushIntervalMillis > 0) { + flushIntervalJob = coroutineScope.launch(exceptionHandler) { + while (isActive) { + delay(flushIntervalMillis) + flush() } } } + + coroutineScope.launch(exceptionHandler) { + uploadReadyBatches(sealCurrentBatch = false) + } } override fun onLowMemoryChannel(): Channel> { @@ -111,6 +106,9 @@ internal class EventSenderEngineImpl( data: ConfidenceFieldsType, context: Map ) { + if (isStopped) { + return + } coroutineScope.launch { val payload = payloadMerger(context, data) val event = EngineEvent( @@ -124,6 +122,9 @@ internal class EventSenderEngineImpl( } override fun flush() { + if (isStopped) { + return + } coroutineScope.launch { writeReqChannel.send(manualFlushEvent) debugLogger?.logEvent(action = "Flush ", event = manualFlushEvent) @@ -131,11 +132,46 @@ internal class EventSenderEngineImpl( } override fun stop() { + isStopped = true + flushIntervalJob?.cancel() + runBlocking(dispatcher) { + uploadReadyBatches(sealCurrentBatch = true) + } coroutineScope.cancel() eventStorage.stop() debugLogger?.logMessage(message = "EventSenderEngine closed ") } + private suspend fun uploadReadyBatches(sealCurrentBatch: Boolean) { + if (sealCurrentBatch) { + eventStorage.rollover() + } + val readyFiles = eventStorage.batchReadyFiles() + for (readyFile in readyFiles) { + val events = eventStorage.eventsFor(readyFile) + .map { e -> + EngineEvent( + "eventDefinitions/${e.eventDefinition}", + e.eventTime, + e.payload + ) + } + val batch = EventBatchRequest( + clientSecret = clientSecret, + events = events, + sendTime = clock.currentTime(), + sdk = Sdk(sdkMetadata.sdkId, sdkMetadata.sdkVersion) + ) + runCatching { + val shouldCleanup = uploader.upload(batch) + debugLogger?.logMessage(message = "Uploading events") + if (shouldCleanup) { + readyFile.delete() + } + } + } + } + companion object { private const val SEND_SIG = "FLUSH" private var Instance: EventSenderEngine? = null @@ -145,7 +181,8 @@ internal class EventSenderEngineImpl( sdkMetadata: SdkMetadata, flushPolicies: List = listOf(), dispatcher: CoroutineDispatcher = Dispatchers.IO, - debugLogger: DebugLogger? + debugLogger: DebugLogger?, + flushIntervalMillis: Long? = null ): EventSenderEngine { return Instance ?: run { EventSenderEngineImpl( @@ -155,8 +192,11 @@ internal class EventSenderEngineImpl( flushPolicies = flushPolicies.toMutableList(), dispatcher = dispatcher, sdkMetadata = sdkMetadata, - debugLogger = debugLogger - ) + debugLogger = debugLogger, + flushIntervalMillis = flushIntervalMillis + ).also { + Instance = it + } } } } diff --git a/Confidence/src/test/java/com/spotify/confidence/EventSenderEngineReliabilityTest.kt b/Confidence/src/test/java/com/spotify/confidence/EventSenderEngineReliabilityTest.kt new file mode 100644 index 00000000..d67d9668 --- /dev/null +++ b/Confidence/src/test/java/com/spotify/confidence/EventSenderEngineReliabilityTest.kt @@ -0,0 +1,129 @@ +package com.spotify.confidence + +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.test.UnconfinedTestDispatcher +import kotlinx.coroutines.test.advanceUntilIdle +import kotlinx.coroutines.test.runTest +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Before +import org.junit.Test +import java.util.Date + +@OptIn(ExperimentalCoroutinesApi::class) +class EventSenderEngineReliabilityTest { + private lateinit var testDispatcher: UnconfinedTestDispatcher + private lateinit var uploader: RecordingEventUploader + private lateinit var storage: RecordingEventStorage + + @Before + fun setUp() { + testDispatcher = UnconfinedTestDispatcher() + uploader = RecordingEventUploader() + storage = RecordingEventStorage() + } + + @Test + fun startupUploadsPendingReadyBatchesWithoutSealingCurrentBatch() = runTest(testDispatcher) { + storage.readyEvents["pending.batch"] = listOf( + EngineEvent("pending", Date(), mapOf()) + ) + storage.currentEvents.add( + EngineEvent("current", Date(), mapOf()) + ) + + EventSenderEngineImpl( + eventStorage = storage, + clientSecret = "secret", + uploader = uploader, + flushPolicies = mutableListOf(), + dispatcher = testDispatcher, + sdkMetadata = com.spotify.confidence.client.SdkMetadata("id", "1.0"), + debugLogger = null + ) + + advanceUntilIdle() + + assertEquals(listOf("pending"), uploader.uploadedEventNames) + assertEquals(listOf("current"), storage.currentEvents.map { it.eventDefinition }) + } + + @Test + fun stopUploadsCurrentBatch() = runTest(testDispatcher) { + val engine = EventSenderEngineImpl( + eventStorage = storage, + clientSecret = "secret", + uploader = uploader, + flushPolicies = mutableListOf(), + dispatcher = testDispatcher, + sdkMetadata = com.spotify.confidence.client.SdkMetadata("id", "1.0"), + debugLogger = null + ) + + engine.emit("session-end", mapOf(), mapOf()) + advanceUntilIdle() + engine.stop() + + assertTrue(uploader.uploadedEventNames.contains("session-end")) + } + + @Test + fun periodicFlushIntervalUploadsEvents() = runTest(testDispatcher) { + val engine = EventSenderEngineImpl( + eventStorage = storage, + clientSecret = "secret", + uploader = uploader, + flushPolicies = mutableListOf(), + dispatcher = testDispatcher, + sdkMetadata = com.spotify.confidence.client.SdkMetadata("id", "1.0"), + debugLogger = null, + flushIntervalMillis = 100 + ) + + engine.emit("interval-event", mapOf(), mapOf()) + advanceUntilIdle() + testScheduler.advanceTimeBy(150) + advanceUntilIdle() + engine.stop() + + assertTrue(uploader.uploadedEventNames.contains("interval-event")) + } + + private class RecordingEventUploader : EventSenderUploader { + val uploadedEventNames = mutableListOf() + + override suspend fun upload(events: EventBatchRequest): Boolean { + uploadedEventNames.addAll(events.events.map { it.eventDefinition.removePrefix("eventDefinitions/") }) + return true + } + } + + private class RecordingEventStorage : EventStorage { + val currentEvents = mutableListOf() + val readyEvents = mutableMapOf>() + + override suspend fun rollover() { + if (currentEvents.isNotEmpty()) { + readyEvents["batch-${readyEvents.size}"] = currentEvents.toList() + currentEvents.clear() + } + } + + override suspend fun writeEvent(event: EngineEvent) { + currentEvents.add(event) + } + + override suspend fun batchReadyFiles(): List { + return readyEvents.keys.map { java.io.File(it) } + } + + override suspend fun eventsFor(file: java.io.File): List { + return readyEvents[file.name].orEmpty() + } + + override fun onLowMemoryChannel() = kotlinx.coroutines.channels.Channel>() + + override fun stop() { + } + } +} diff --git a/Provider/src/main/java/com/spotify/confidence/openfeature/ConfidenceFeatureProvider.kt b/Provider/src/main/java/com/spotify/confidence/openfeature/ConfidenceFeatureProvider.kt index 73d2438f..0bd04ead 100644 --- a/Provider/src/main/java/com/spotify/confidence/openfeature/ConfidenceFeatureProvider.kt +++ b/Provider/src/main/java/com/spotify/confidence/openfeature/ConfidenceFeatureProvider.kt @@ -53,6 +53,7 @@ class ConfidenceFeatureProvider private constructor( } override fun shutdown() { + confidence.flush() } override suspend fun onContextSet( From b18370787be84f164fa8d66be55a68912fffc599 Mon Sep 17 00:00:00 2001 From: Fabrizio Demaria Date: Fri, 7 Aug 2026 11:03:02 +0200 Subject: [PATCH 2/3] Fix ktlint spacing around @Volatile in EventSenderEngine. --- .../src/main/java/com/spotify/confidence/EventSenderEngine.kt | 1 + 1 file changed, 1 insertion(+) diff --git a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt index d665b1da..62e9a365 100644 --- a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt +++ b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt @@ -49,6 +49,7 @@ internal class EventSenderEngineImpl( } } private var flushIntervalJob: Job? = null + @Volatile private var isStopped = false From 6d2ba34430ac855ec33bd53562fffd8abea0c6f9 Mon Sep 17 00:00:00 2001 From: Fabrizio Demaria Date: Fri, 7 Aug 2026 11:09:29 +0200 Subject: [PATCH 3/3] Fix EventSenderEngine instance caching regression. Do not assign the companion Instance field; main never cached instances and tests rely on a fresh engine per ConfidenceFactory.create() call. --- .../src/main/java/com/spotify/confidence/EventSenderEngine.kt | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt index 62e9a365..87baed4c 100644 --- a/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt +++ b/Confidence/src/main/java/com/spotify/confidence/EventSenderEngine.kt @@ -195,9 +195,7 @@ internal class EventSenderEngineImpl( sdkMetadata = sdkMetadata, debugLogger = debugLogger, flushIntervalMillis = flushIntervalMillis - ).also { - Instance = it - } + ) } } }