diff --git a/observability/dagger/api/observability-dagger.api b/observability/dagger/api/observability-dagger.api index bd7b57f1c..89421b043 100644 --- a/observability/dagger/api/observability-dagger.api +++ b/observability/dagger/api/observability-dagger.api @@ -18,14 +18,14 @@ public abstract interface class me/proton/core/observability/dagger/CoreObservab } public final class me/proton/core/observability/dagger/CoreObservabilityModule$Companion { - public final fun provideObservabilityWorkerManagerImpl (Landroidx/work/WorkManager;)Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl; + public final fun provideObservabilityTimeTracker ()Lme/proton/core/observability/domain/ObservabilityTimeTracker; } -public final class me/proton/core/observability/dagger/CoreObservabilityModule_Companion_ProvideObservabilityWorkerManagerImplFactory : dagger/internal/Factory { - public fun (Ljavax/inject/Provider;)V - public static fun create (Ljavax/inject/Provider;)Lme/proton/core/observability/dagger/CoreObservabilityModule_Companion_ProvideObservabilityWorkerManagerImplFactory; +public final class me/proton/core/observability/dagger/CoreObservabilityModule_Companion_ProvideObservabilityTimeTrackerFactory : dagger/internal/Factory { + public fun ()V + public static fun create ()Lme/proton/core/observability/dagger/CoreObservabilityModule_Companion_ProvideObservabilityTimeTrackerFactory; public synthetic fun get ()Ljava/lang/Object; - public fun get ()Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl; - public static fun provideObservabilityWorkerManagerImpl (Landroidx/work/WorkManager;)Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl; + public fun get ()Lme/proton/core/observability/domain/ObservabilityTimeTracker; + public static fun provideObservabilityTimeTracker ()Lme/proton/core/observability/domain/ObservabilityTimeTracker; } diff --git a/observability/dagger/src/main/kotlin/me/proton/core/observability/dagger/CoreObservabilityModule.kt b/observability/dagger/src/main/kotlin/me/proton/core/observability/dagger/CoreObservabilityModule.kt index 2ca0f0e2e..45e11b7ef 100644 --- a/observability/dagger/src/main/kotlin/me/proton/core/observability/dagger/CoreObservabilityModule.kt +++ b/observability/dagger/src/main/kotlin/me/proton/core/observability/dagger/CoreObservabilityModule.kt @@ -19,7 +19,6 @@ package me.proton.core.observability.dagger import android.os.SystemClock -import androidx.work.WorkManager import dagger.Binds import dagger.Module import dagger.Provides @@ -30,6 +29,7 @@ import me.proton.core.observability.data.ObservabilityRepositoryImpl import me.proton.core.observability.data.usecase.SendObservabilityEventsImpl import me.proton.core.observability.data.worker.ObservabilityWorkerManagerImpl import me.proton.core.observability.domain.ObservabilityRepository +import me.proton.core.observability.domain.ObservabilityTimeTracker import me.proton.core.observability.domain.ObservabilityWorkerManager import me.proton.core.observability.domain.usecase.IsObservabilityEnabled import me.proton.core.observability.domain.usecase.SendObservabilityEvents @@ -53,7 +53,7 @@ public interface CoreObservabilityModule { public companion object { @Provides @Singleton - public fun provideObservabilityWorkerManagerImpl(workManager: WorkManager): ObservabilityWorkerManagerImpl = - ObservabilityWorkerManagerImpl({ SystemClock.elapsedRealtime() }, workManager) + public fun provideObservabilityTimeTracker(): ObservabilityTimeTracker = + ObservabilityTimeTracker(clockMillis = { SystemClock.elapsedRealtime() }) } } diff --git a/observability/data/api/observability-data.api b/observability/data/api/observability-data.api index b772178bb..0b7b2213f 100644 --- a/observability/data/api/observability-data.api +++ b/observability/data/api/observability-data.api @@ -87,11 +87,17 @@ public final class me/proton/core/observability/data/usecase/SendObservabilityEv } public final class me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl : me/proton/core/observability/domain/ObservabilityWorkerManager { - public fun (Lkotlin/jvm/functions/Function0;Landroidx/work/WorkManager;)V + public fun (Landroidx/work/WorkManager;)V public fun cancel ()V - public fun getDurationSinceLastShipment-LV8wdWc (Lkotlin/coroutines/Continuation;)Ljava/lang/Object; public fun schedule-LRDsOJo (J)V - public fun setLastSentNow (Lkotlin/coroutines/Continuation;)Ljava/lang/Object; +} + +public final class me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl_Factory : dagger/internal/Factory { + public fun (Ljavax/inject/Provider;)V + public static fun create (Ljavax/inject/Provider;)Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl_Factory; + public synthetic fun get ()Ljava/lang/Object; + public fun get ()Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl; + public static fun newInstance (Landroidx/work/WorkManager;)Lme/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl; } public abstract interface class me/proton/core/observability/data/worker/ObservabilityWorker_AssistedFactory : androidx/hilt/work/WorkerAssistedFactory { diff --git a/observability/data/src/main/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl.kt b/observability/data/src/main/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl.kt index 83ded591d..12d6e9b39 100644 --- a/observability/data/src/main/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl.kt +++ b/observability/data/src/main/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImpl.kt @@ -23,30 +23,19 @@ import androidx.work.ExistingWorkPolicy import androidx.work.NetworkType import androidx.work.OneTimeWorkRequestBuilder import androidx.work.WorkManager -import kotlinx.coroutines.sync.Mutex -import kotlinx.coroutines.sync.withLock import me.proton.core.observability.domain.ObservabilityWorkerManager import java.util.concurrent.TimeUnit +import javax.inject.Inject import kotlin.time.Duration -import kotlin.time.Duration.Companion.milliseconds -public class ObservabilityWorkerManagerImpl constructor( - private val clockMillis: () -> Long, +public class ObservabilityWorkerManagerImpl @Inject constructor( private val workManager: WorkManager ) : ObservabilityWorkerManager { - private val lastSentAtMs = MutexValue(null) override fun cancel() { workManager.cancelUniqueWork(WORK_NAME) } - override suspend fun getDurationSinceLastShipment(): Duration? = - lastSentAtMs.getValue()?.let { clockMillis() - it }?.milliseconds - - override suspend fun setLastSentNow() { - lastSentAtMs.setValue(clockMillis()) - } - override fun schedule(delay: Duration) { val request = OneTimeWorkRequestBuilder() .setConstraints( @@ -66,14 +55,6 @@ public class ObservabilityWorkerManagerImpl constructor( workManager.beginUniqueWork(WORK_NAME, policy, request).enqueue() } - private class MutexValue(initialValue: T) { - private val mutex = Mutex() - private var value: T = initialValue - - suspend fun getValue(): T = mutex.withLock { value } - suspend fun setValue(newValue: T) = mutex.withLock { value = newValue } - } - private companion object { private const val WORK_NAME = "me.proton.core.observability.data.worker" } diff --git a/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImplTest.kt b/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImplTest.kt index 1a335cc46..f0917b42a 100644 --- a/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImplTest.kt +++ b/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerManagerImplTest.kt @@ -26,37 +26,20 @@ import io.mockk.every import io.mockk.mockk import io.mockk.slot import io.mockk.verify -import kotlinx.coroutines.test.runTest import kotlin.test.BeforeTest import kotlin.test.Test import kotlin.test.assertEquals -import kotlin.test.assertNull import kotlin.time.Duration.Companion.ZERO -import kotlin.time.Duration.Companion.milliseconds import kotlin.time.Duration.Companion.minutes class ObservabilityWorkerManagerImplTest { - private lateinit var clock: FakeClock private lateinit var tested: ObservabilityWorkerManagerImpl private lateinit var workManager: WorkManager @BeforeTest fun setUp() { - clock = FakeClock() workManager = mockk() - tested = ObservabilityWorkerManagerImpl(clock::now, workManager) - } - - @Test - fun durationSinceLastShipment() = runTest { - assertNull(tested.getDurationSinceLastShipment()) - tested.setLastSentNow() - clock.current = 1000 - assertEquals(1000.milliseconds, tested.getDurationSinceLastShipment()) - - tested.setLastSentNow() - clock.current = 1500 - assertEquals(500.milliseconds, tested.getDurationSinceLastShipment()) + tested = ObservabilityWorkerManagerImpl(workManager) } @Test @@ -95,10 +78,4 @@ class ObservabilityWorkerManagerImplTest { assertEquals(ExistingWorkPolicy.KEEP, workPolicySlot.captured) assertEquals(2.minutes.inWholeMilliseconds, requestSlot.captured.workSpec.initialDelay) } - - private class FakeClock { - var current: Long = 0 - - fun now(): Long = current - } } \ No newline at end of file diff --git a/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerTest.kt b/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerTest.kt index f507c0f4e..5efe2641e 100644 --- a/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerTest.kt +++ b/observability/data/src/test/kotlin/me/proton/core/observability/data/worker/ObservabilityWorkerTest.kt @@ -35,6 +35,7 @@ import me.proton.core.network.domain.ApiException import me.proton.core.network.domain.ApiResult import me.proton.core.network.domain.HttpResponseCodes import me.proton.core.observability.domain.ObservabilityRepository +import me.proton.core.observability.domain.ObservabilityTimeTracker import me.proton.core.observability.domain.ObservabilityWorkerManager import me.proton.core.observability.domain.entity.ObservabilityEvent import me.proton.core.observability.domain.usecase.IsObservabilityEnabled @@ -71,6 +72,9 @@ class ObservabilityWorkerTest { @BindValue internal lateinit var sendObservabilityEvents: SendObservabilityEvents + @BindValue + internal lateinit var timeTracker: ObservabilityTimeTracker + private lateinit var context: Context @Before @@ -81,6 +85,7 @@ class ObservabilityWorkerTest { observabilityWorkerManager = mockk(relaxUnitFun = true) repository = mockk(relaxUnitFun = true) sendObservabilityEvents = mockk(relaxUnitFun = true) + timeTracker = mockk(relaxUnitFun = true) } @@ -92,7 +97,6 @@ class ObservabilityWorkerTest { assertEquals(ListenableWorker.Result.success(), result) coVerify(exactly = 0) { sendObservabilityEvents.invoke(any()) } - coVerify { observabilityWorkerManager.setLastSentNow() } coVerify { repository.deleteAllEvents() } } @@ -119,7 +123,6 @@ class ObservabilityWorkerTest { coVerify(exactly = 1) { sendObservabilityEvents.invoke(events) } coVerify(exactly = 1) { repository.deleteEvents(events) } - coVerify(exactly = 1) { observabilityWorkerManager.setLastSentNow() } } @Test @@ -140,8 +143,6 @@ class ObservabilityWorkerTest { coVerify(exactly = 2) { sendObservabilityEvents.invoke(any()) } coVerify(exactly = 2) { repository.deleteEvents(any()) } - - coVerify(exactly = 1) { observabilityWorkerManager.setLastSentNow() } } @Test @@ -153,8 +154,6 @@ class ObservabilityWorkerTest { val result = makeAndRunWorker() assertEquals(ListenableWorker.Result.retry(), result) - - coVerify(exactly = 0) { observabilityWorkerManager.setLastSentNow() } } @Test @@ -166,8 +165,6 @@ class ObservabilityWorkerTest { val result = makeAndRunWorker() assertIs(result) - - coVerify(exactly = 0) { observabilityWorkerManager.setLastSentNow() } } @Test @@ -178,8 +175,6 @@ class ObservabilityWorkerTest { val result = makeAndRunWorker() assertIs(result) - - coVerify(exactly = 0) { observabilityWorkerManager.setLastSentNow() } } private fun makeWorker(): ObservabilityWorker = TestListenableWorkerBuilder(context) diff --git a/observability/domain/api/observability-domain.api b/observability/domain/api/observability-domain.api index da698aeb2..3acb6084a 100644 --- a/observability/domain/api/observability-domain.api +++ b/observability/domain/api/observability-domain.api @@ -26,11 +26,13 @@ public final class me/proton/core/observability/domain/ObservabilityRepository$D public static synthetic fun getEvents$default (Lme/proton/core/observability/domain/ObservabilityRepository;Ljava/lang/Integer;Lkotlin/coroutines/Continuation;ILjava/lang/Object;)Ljava/lang/Object; } +public final class me/proton/core/observability/domain/ObservabilityTimeTracker { + public fun (Lkotlin/jvm/functions/Function0;)V +} + public abstract interface class me/proton/core/observability/domain/ObservabilityWorkerManager { public abstract fun cancel ()V - public abstract fun getDurationSinceLastShipment-LV8wdWc (Lkotlin/coroutines/Continuation;)Ljava/lang/Object; public abstract fun schedule-LRDsOJo (J)V - public abstract fun setLastSentNow (Lkotlin/coroutines/Continuation;)Ljava/lang/Object; } public final class me/proton/core/observability/domain/entity/ObservabilityEvent { @@ -1455,7 +1457,7 @@ public abstract interface class me/proton/core/observability/domain/usecase/IsOb } public final class me/proton/core/observability/domain/usecase/ProcessObservabilityEvents { - public fun (Lme/proton/core/observability/domain/usecase/IsObservabilityEnabled;Lme/proton/core/observability/domain/ObservabilityWorkerManager;Lme/proton/core/observability/domain/ObservabilityRepository;Lme/proton/core/observability/domain/usecase/SendObservabilityEvents;)V + public fun (Lme/proton/core/observability/domain/usecase/IsObservabilityEnabled;Lme/proton/core/observability/domain/ObservabilityRepository;Lme/proton/core/observability/domain/ObservabilityTimeTracker;Lme/proton/core/observability/domain/usecase/SendObservabilityEvents;)V public final fun invoke (Lkotlin/coroutines/Continuation;)Ljava/lang/Object; } diff --git a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityManager.kt b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityManager.kt index 59957d10d..4725b1634 100644 --- a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityManager.kt +++ b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityManager.kt @@ -35,6 +35,7 @@ public class ObservabilityManager @Inject internal constructor( private val isObservabilityEnabled: IsObservabilityEnabled, private val repository: ObservabilityRepository, private val scopeProvider: CoroutineScopeProvider, + private val timeTracker: ObservabilityTimeTracker, private val workerManager: ObservabilityWorkerManager, ) { /** Enqueues an event with a given [data] and [timestamp] to be sent at some point in the future. @@ -63,18 +64,28 @@ public class ObservabilityManager @Inject internal constructor( if (isObservabilityEnabled()) { repository.addEvent(event) workerManager.schedule(getSendDelay()) + + if (timeTracker.getDurationSinceFirstEvent() == null) { + timeTracker.setFirstEventNow() + } } else { workerManager.cancel() repository.deleteAllEvents() + timeTracker.clear() } } private suspend fun getSendDelay(): Duration { - val eventCount = repository.getEventCount() + suspend fun isMaxDurationExceeded(): Boolean { + val duration = timeTracker.getDurationSinceFirstEvent() + return if (duration != null) { + duration >= MAX_DELAY_MS.milliseconds + } else false + } + return when { - eventCount <= 1L -> MAX_DELAY_MS.milliseconds - eventCount >= MAX_EVENT_COUNT -> ZERO - workerManager.getDurationSinceLastShipment()?.let { it >= MAX_DELAY_MS.milliseconds } ?: false -> ZERO + repository.getEventCount() >= MAX_EVENT_COUNT -> ZERO + isMaxDurationExceeded() -> ZERO else -> MAX_DELAY_MS.milliseconds } } diff --git a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityTimeTracker.kt b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityTimeTracker.kt new file mode 100644 index 000000000..8090ec4c5 --- /dev/null +++ b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityTimeTracker.kt @@ -0,0 +1,31 @@ +package me.proton.core.observability.domain + +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import javax.inject.Singleton +import kotlin.time.Duration +import kotlin.time.Duration.Companion.milliseconds + +@Singleton +public class ObservabilityTimeTracker constructor( + private val clockMillis: () -> Long, +) { + private val firstEnqueuedEventAtMs = MutexValue(null) + + internal suspend fun clear() = firstEnqueuedEventAtMs.setValue(null) + + internal suspend fun getDurationSinceFirstEvent(): Duration? = + firstEnqueuedEventAtMs.getValue()?.let { clockMillis() - it }?.milliseconds + + internal suspend fun setFirstEventNow() { + firstEnqueuedEventAtMs.setValue(clockMillis()) + } + + private class MutexValue(initialValue: T) { + private val mutex = Mutex() + private var value: T = initialValue + + suspend fun getValue(): T = mutex.withLock { value } + suspend fun setValue(newValue: T) = mutex.withLock { value = newValue } + } +} diff --git a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityWorkerManager.kt b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityWorkerManager.kt index be27972d5..947ec996e 100644 --- a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityWorkerManager.kt +++ b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/ObservabilityWorkerManager.kt @@ -26,14 +26,6 @@ public interface ObservabilityWorkerManager { */ public fun cancel() - /** Returns the duration since the last successful shipment of observability events. - * Returns `null` if the shipment wasn't recorder yet. - */ - public suspend fun getDurationSinceLastShipment(): Duration? - - /** Marks that observability events have been successfully shipped (sent). */ - public suspend fun setLastSentNow() - /** Schedules a worker to send the observability events. * If a worker has been previously scheduled but hasn't yet executed, * the existing scheduled worker will be kept. diff --git a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/usecase/ProcessObservabilityEvents.kt b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/usecase/ProcessObservabilityEvents.kt index 3fa3b2600..b9638e707 100644 --- a/observability/domain/src/main/kotlin/me/proton/core/observability/domain/usecase/ProcessObservabilityEvents.kt +++ b/observability/domain/src/main/kotlin/me/proton/core/observability/domain/usecase/ProcessObservabilityEvents.kt @@ -19,14 +19,14 @@ package me.proton.core.observability.domain.usecase import me.proton.core.observability.domain.ObservabilityRepository -import me.proton.core.observability.domain.ObservabilityWorkerManager +import me.proton.core.observability.domain.ObservabilityTimeTracker import javax.inject.Inject /** Processes observability events in batches and sends them to the server. */ public class ProcessObservabilityEvents @Inject constructor( private val isObservabilityEnabled: IsObservabilityEnabled, - private val observabilityWorkerManager: ObservabilityWorkerManager, private val repository: ObservabilityRepository, + private val timeTracker: ObservabilityTimeTracker, private val sendObservabilityEvents: SendObservabilityEvents ) { public suspend operator fun invoke() { @@ -34,7 +34,7 @@ public class ProcessObservabilityEvents @Inject constructor( true -> processSingleBatch() else -> repository.deleteAllEvents() } - observabilityWorkerManager.setLastSentNow() + timeTracker.clear() } private tailrec suspend fun processSingleBatch() { diff --git a/observability/domain/src/test/kotlin/me/proton/core/observability/domain/ObservabilityManagerTest.kt b/observability/domain/src/test/kotlin/me/proton/core/observability/domain/ObservabilityManagerTest.kt index ffe1c5fe7..b1c41f9a6 100644 --- a/observability/domain/src/test/kotlin/me/proton/core/observability/domain/ObservabilityManagerTest.kt +++ b/observability/domain/src/test/kotlin/me/proton/core/observability/domain/ObservabilityManagerTest.kt @@ -18,39 +18,50 @@ package me.proton.core.observability.domain +import io.mockk.clearMocks import io.mockk.coEvery import io.mockk.coVerify import io.mockk.mockk +import io.mockk.verify import kotlinx.coroutines.test.UnconfinedTestDispatcher import kotlinx.coroutines.test.runTest import me.proton.core.observability.domain.metrics.ObservabilityData import me.proton.core.observability.domain.usecase.IsObservabilityEnabled import me.proton.core.test.kotlin.TestCoroutineScopeProvider import me.proton.core.test.kotlin.TestDispatcherProvider -import me.proton.core.test.kotlin.assertTrue import me.proton.core.util.kotlin.CoroutineScopeProvider import kotlin.test.BeforeTest import kotlin.test.Test -import kotlin.test.assertContentEquals -import kotlin.test.assertEquals -import kotlin.time.Duration +import kotlin.time.Duration.Companion.ZERO import kotlin.time.Duration.Companion.milliseconds class ObservabilityManagerTest { + private var currentClockMillis: Long = 0L + private lateinit var isObservabilityEnabled: IsObservabilityEnabled private lateinit var observabilityRepository: ObservabilityRepository private lateinit var scopeProvider: CoroutineScopeProvider - private lateinit var workerManager: FakeWorkerManager + private lateinit var workerManager: ObservabilityWorkerManager + private lateinit var timeTracker: ObservabilityTimeTracker private lateinit var tested: ObservabilityManager @BeforeTest fun setUp() { + currentClockMillis = 0L isObservabilityEnabled = mockk() observabilityRepository = mockk(relaxUnitFun = true) - scopeProvider = TestCoroutineScopeProvider(TestDispatcherProvider(UnconfinedTestDispatcher())) - workerManager = FakeWorkerManager() + scopeProvider = + TestCoroutineScopeProvider(TestDispatcherProvider(UnconfinedTestDispatcher())) + timeTracker = ObservabilityTimeTracker { currentClockMillis } + workerManager = mockk(relaxed = true) - tested = ObservabilityManager(isObservabilityEnabled, observabilityRepository, scopeProvider, workerManager) + tested = ObservabilityManager( + isObservabilityEnabled, + observabilityRepository, + scopeProvider, + timeTracker, + workerManager + ) } @Test @@ -64,9 +75,7 @@ class ObservabilityManagerTest { coVerify(exactly = 0) { observabilityRepository.addEvent(any()) } coVerify(exactly = 1) { observabilityRepository.deleteAllEvents() } - assertTrue(workerManager.scheduledCalls.isEmpty()) { - "Unexpected call to `ObservabilityWorkerManager.schedule`." - } + verify(exactly = 0) { workerManager.schedule(any()) } } @Test @@ -79,10 +88,10 @@ class ObservabilityManagerTest { // THEN coVerify(exactly = 0) { observabilityRepository.deleteAllEvents() } - assertEquals(0, workerManager.cancelCount) + verify(exactly = 0) { workerManager.cancel() } coVerify(exactly = 1) { observabilityRepository.addEvent(any()) } - assertContentEquals(listOf(ObservabilityManager.MAX_DELAY_MS.milliseconds), workerManager.scheduledCalls) + verify(exactly = 1) { workerManager.schedule(ObservabilityManager.MAX_DELAY_MS.milliseconds) } } @Test @@ -95,66 +104,87 @@ class ObservabilityManagerTest { // THEN coVerify(exactly = 0) { observabilityRepository.deleteAllEvents() } - assertEquals(0, workerManager.cancelCount) + verify(exactly = 0) { workerManager.cancel() } coVerify(exactly = 1) { observabilityRepository.addEvent(any()) } - assertContentEquals(listOf(Duration.ZERO), workerManager.scheduledCalls) + verify(exactly = 1) { workerManager.schedule(ZERO) } } @Test fun durationSinceLastShipmentExceeded() = runTest { coEvery { isObservabilityEnabled.invoke() } returns true coEvery { observabilityRepository.getEventCount() } returns ObservabilityManager.MAX_EVENT_COUNT / 2 - workerManager.duration = ObservabilityManager.MAX_DELAY_MS.milliseconds + timeTracker.setFirstEventNow() + currentClockMillis = ObservabilityManager.MAX_DELAY_MS // WHEN tested.enqueue(mockk(relaxed = true)) // THEN coVerify(exactly = 0) { observabilityRepository.deleteAllEvents() } - assertEquals(0, workerManager.cancelCount) + verify(exactly = 0) { workerManager.cancel() } coVerify(exactly = 1) { observabilityRepository.addEvent(any()) } - assertContentEquals(listOf(Duration.ZERO), workerManager.scheduledCalls) + verify(exactly = 1) { workerManager.schedule(ZERO) } } @Test fun durationSinceLastShipmentNotExceeded() = runTest { coEvery { isObservabilityEnabled.invoke() } returns true coEvery { observabilityRepository.getEventCount() } returns ObservabilityManager.MAX_EVENT_COUNT / 2 - workerManager.duration = (ObservabilityManager.MAX_DELAY_MS - 1).milliseconds + currentClockMillis = ObservabilityManager.MAX_DELAY_MS - 1 // WHEN tested.enqueue(mockk(relaxed = true)) // THEN coVerify(exactly = 0) { observabilityRepository.deleteAllEvents() } - assertEquals(0, workerManager.cancelCount) + verify(exactly = 0) { workerManager.cancel() } coVerify(exactly = 1) { observabilityRepository.addEvent(any()) } - assertContentEquals(listOf(ObservabilityManager.MAX_DELAY_MS.milliseconds), workerManager.scheduledCalls) + verify(exactly = 1) { workerManager.schedule(ObservabilityManager.MAX_DELAY_MS.milliseconds) } } - /** - * The [getDurationSinceLastShipment] method returns inline class (Duration) which is not supported by mockk, - * so we need to use a fake class. - */ - private class FakeWorkerManager : ObservabilityWorkerManager { - var cancelCount = 0 - private set - var scheduledCalls = mutableListOf() - private set + @Test + fun schedulingEventsUsesProperDelays() = runTest { + // 1. ========================= + // GIVEN + coEvery { isObservabilityEnabled.invoke() } returns true + coEvery { observabilityRepository.getEventCount() } returns 1 + currentClockMillis = ObservabilityManager.MAX_DELAY_MS - var duration: Duration? = null + // WHEN + tested.enqueue(mockk(relaxed = true)) - override fun cancel() { - cancelCount += 1 - } + // THEN + // The first event should be enqueued with max delay. + verify(exactly = 1) { workerManager.schedule(ObservabilityManager.MAX_DELAY_MS.milliseconds) } - override suspend fun getDurationSinceLastShipment(): Duration? = duration - override suspend fun setLastSentNow() = Unit - override fun schedule(delay: Duration) { - scheduledCalls.add(delay) - } + // 2. ========================= + // GIVEN + clearMocks(workerManager) + coEvery { observabilityRepository.getEventCount() } returns 2 + currentClockMillis += 10 + + // WHEN + tested.enqueue(mockk(relaxed = true)) + + // THEN + // Subsequent events should be also enqueued with max delay. + verify(exactly = 1) { workerManager.schedule(ObservabilityManager.MAX_DELAY_MS.milliseconds) } + + // 3. ========================= + // GIVEN + clearMocks(workerManager) + coEvery { observabilityRepository.getEventCount() } returns 3 + currentClockMillis += ObservabilityManager.MAX_DELAY_MS + + // WHEN + tested.enqueue(mockk(relaxed = true)) + + // THEN + // Events enqueued after MAX_DELAY_MS since the first event, + // should be enqueued with no delay. + verify(exactly = 1) { workerManager.schedule(ZERO) } } }