mirror of
https://github.com/ProtonMail/protoncore_android.git
synced 2026-06-14 09:54:49 +00:00
fix(observability): Use proper delays for scheduling sending the obserability events.
This commit is contained in:
@@ -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 <init> (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 <init> ()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;
|
||||
}
|
||||
|
||||
|
||||
+3
-3
@@ -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() })
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 <init> (Lkotlin/jvm/functions/Function0;Landroidx/work/WorkManager;)V
|
||||
public fun <init> (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 <init> (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 {
|
||||
|
||||
+2
-21
@@ -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<Long?>(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<ObservabilityWorker>()
|
||||
.setConstraints(
|
||||
@@ -66,14 +55,6 @@ public class ObservabilityWorkerManagerImpl constructor(
|
||||
workManager.beginUniqueWork(WORK_NAME, policy, request).enqueue()
|
||||
}
|
||||
|
||||
private class MutexValue<T>(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"
|
||||
}
|
||||
|
||||
+1
-24
@@ -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
|
||||
}
|
||||
}
|
||||
+5
-10
@@ -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<ListenableWorker.Result.Failure>(result)
|
||||
|
||||
coVerify(exactly = 0) { observabilityWorkerManager.setLastSentNow() }
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -178,8 +175,6 @@ class ObservabilityWorkerTest {
|
||||
|
||||
val result = makeAndRunWorker()
|
||||
assertIs<ListenableWorker.Result.Failure>(result)
|
||||
|
||||
coVerify(exactly = 0) { observabilityWorkerManager.setLastSentNow() }
|
||||
}
|
||||
|
||||
private fun makeWorker(): ObservabilityWorker = TestListenableWorkerBuilder<ObservabilityWorker>(context)
|
||||
|
||||
@@ -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 <init> (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 <init> (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 <init> (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;
|
||||
}
|
||||
|
||||
|
||||
+15
-4
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
+31
@@ -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<Long?>(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<T>(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 }
|
||||
}
|
||||
}
|
||||
-8
@@ -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.
|
||||
|
||||
+3
-3
@@ -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() {
|
||||
|
||||
+69
-39
@@ -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<ObservabilityData>(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<ObservabilityData>(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<Duration>()
|
||||
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<ObservabilityData>(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<ObservabilityData>(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<ObservabilityData>(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) }
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user