diff --git a/.changeset/durable-queue-lifecycle.md b/.changeset/durable-queue-lifecycle.md new file mode 100644 index 000000000..06439f14a --- /dev/null +++ b/.changeset/durable-queue-lifecycle.md @@ -0,0 +1,5 @@ +--- +"posthog": patch +--- + +Retain bounded durable queue entries across retryable transport and HTTP failures, pause while offline, and acknowledge successful batches by unique queue-entry identity. diff --git a/posthog/src/main/java/com/posthog/PostHogConfig.kt b/posthog/src/main/java/com/posthog/PostHogConfig.kt index fb9cd492b..716274b68 100644 --- a/posthog/src/main/java/com/posthog/PostHogConfig.kt +++ b/posthog/src/main/java/com/posthog/PostHogConfig.kt @@ -125,7 +125,8 @@ public open class PostHogConfig( */ public var maxBatchSize: Int = DEFAULT_MAX_BATCH_SIZE, /** - * Maximum number of retries for failed flush attempts before events are dropped + * Maximum number of retries for push subscription registration failures. + * Durable ingestion queues retain retryable records and are bounded by their queue size. * Defaults to 3 */ public var maxRetries: Int = 3, diff --git a/posthog/src/main/java/com/posthog/internal/EndpointSpec.kt b/posthog/src/main/java/com/posthog/internal/EndpointSpec.kt index 6377c8733..b6460c158 100644 --- a/posthog/src/main/java/com/posthog/internal/EndpointSpec.kt +++ b/posthog/src/main/java/com/posthog/internal/EndpointSpec.kt @@ -6,7 +6,6 @@ import com.posthog.PostHogInternal import com.posthog.logs.PostHogLogRecord import java.io.InputStream import java.io.OutputStream -import java.util.UUID /** * Per-endpoint specification consumed by [PostHogQueue]. Carries @@ -31,7 +30,6 @@ public class EndpointSpec internal constructor( internal val send: (List) -> Unit, internal val isRetriableStatusCode: (Int) -> Boolean, internal val isFatalRecord: (Record) -> Boolean = { false }, - internal val recordUuid: (Record) -> UUID? = { null }, ) { public companion object { @JvmStatic @@ -57,7 +55,6 @@ public class EndpointSpec internal constructor( send = { events -> api.batch(events) }, isRetriableStatusCode = ::isEventsRetriableStatusCode, isFatalRecord = { it.isFatalExceptionEvent() }, - recordUuid = { it.uuid }, ) @JvmStatic @@ -83,7 +80,6 @@ public class EndpointSpec internal constructor( send = { events -> api.snapshot(events) }, isRetriableStatusCode = ::isEventsRetriableStatusCode, isFatalRecord = { it.isFatalExceptionEvent() }, - recordUuid = { it.uuid }, ) /** diff --git a/posthog/src/main/java/com/posthog/internal/PostHogQueue.kt b/posthog/src/main/java/com/posthog/internal/PostHogQueue.kt index a4f9d68f7..b9abacca2 100644 --- a/posthog/src/main/java/com/posthog/internal/PostHogQueue.kt +++ b/posthog/src/main/java/com/posthog/internal/PostHogQueue.kt @@ -3,12 +3,12 @@ package com.posthog.internal import com.posthog.PostHogConfig import com.posthog.PostHogInternal import com.posthog.PostHogVisibleForTesting -import com.posthog.vendor.uuid.TimeBasedEpochGenerator import java.io.File import java.io.IOException import java.util.Date import java.util.Timer import java.util.TimerTask +import java.util.UUID import java.util.concurrent.ExecutorService import java.util.concurrent.atomic.AtomicBoolean import kotlin.concurrent.schedule @@ -63,8 +63,7 @@ public class PostHogQueue( dirCreated = true } - val uuid = spec.recordUuid(record) ?: TimeBasedEpochGenerator.generate() - val file = File(dir, "$uuid.event") + val file = File(dir, "${UUID.randomUUID()}.event") synchronized(dequeLock) { deque.add(file) } @@ -81,6 +80,9 @@ public class PostHogQueue( config.logger.log("${spec.describe(record)}: ${file.name} failed to parse: $e.") // if for some reason the file failed to serialize, lets delete it + synchronized(dequeLock) { + deque.remove(file) + } file.deleteSafely(config) } @@ -203,19 +205,11 @@ public class PostHogQueue( } catch (e: Throwable) { config.logger.log("Flushing failed: $e.") - retryCount++ - - if (retryCount > config.maxRetries) { - config.logger.log("Max retries (${config.maxRetries}) exceeded, dropping ${spec.recordsLabel}.") - retryCount = 0 - pausedUntil = null - dropAllRecords() - } else { - retry = true + retryCount = (retryCount + 1).coerceAtMost(maxRetryDelaySeconds) + retry = true - if (e is PostHogApiError) { - retryAfterSeconds = e.retryAfterSeconds - } + if (e is PostHogApiError) { + retryAfterSeconds = e.retryAfterSeconds } } finally { calculateDelay(retry, retryAfterSeconds) @@ -285,13 +279,8 @@ public class PostHogQueue( throw e } } catch (e: IOException) { - // no connection should try again - if (e.isNetworkingError()) { - deleteFiles = false - config.logger.log("Flushing failed because of a network error, let's try again soon.") - } else { - config.logger.log("Flushing failed: $e") - } + deleteFiles = false + config.logger.log("Flushing failed because of a network error, let's try again soon.") throw e } finally { if (deleteFiles) { @@ -432,7 +421,14 @@ public class PostHogQueue( // sort by last modified date ascending so records are sent in order files.sortBy { file -> file.lastModified() } - return files + + val maxQueueSize = spec.maxQueueSize(config).coerceAtLeast(1) + val overflow = (files.size - maxQueueSize).coerceAtLeast(0) + if (overflow > 0) { + files.take(overflow).forEach { it.deleteSafely(config) } + config.logger.log("Dropped $overflow oldest cached ${spec.recordsLabel} to enforce queue capacity.") + } + return files.drop(overflow) } private fun reloadFromDiskSync() { @@ -490,6 +486,10 @@ public class PostHogQueue( internal val currentFlushAtForTesting: Int @PostHogVisibleForTesting get() = batchLimits.flushAt + + internal val currentRetryCountForTesting: Int + @PostHogVisibleForTesting + get() = retryCount } internal class BatchLimits( diff --git a/posthog/src/test/java/com/posthog/internal/PostHogQueueTest.kt b/posthog/src/test/java/com/posthog/internal/PostHogQueueTest.kt index cedd223b6..b0cc7a8e2 100644 --- a/posthog/src/test/java/com/posthog/internal/PostHogQueueTest.kt +++ b/posthog/src/test/java/com/posthog/internal/PostHogQueueTest.kt @@ -10,14 +10,21 @@ import com.posthog.internal.errortracking.ThrowableCoercer import com.posthog.mockHttp import com.posthog.shutdownAndAwaitTermination import com.posthog.vendor.uuid.TimeBasedEpochGenerator +import okhttp3.OkHttpClient import okhttp3.mockwebserver.MockResponse import okhttp3.mockwebserver.SocketPolicy import org.junit.Assert.assertFalse import org.junit.Rule import org.junit.rules.TemporaryFolder import java.io.File +import java.io.IOException +import java.util.Collections +import java.util.Date import java.util.UUID +import java.util.concurrent.CountDownLatch import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicInteger import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertTrue @@ -36,6 +43,8 @@ internal class PostHogQueueTest { dateProvider: PostHogDateProvider = PostHogDeviceDateProvider(), maxBatchSize: Int = 50, networkStatus: PostHogNetworkStatus? = null, + maxRetries: Int = 3, + httpClient: OkHttpClient? = null, ): PostHogQueue { val config = PostHogConfig(API_KEY, host).apply { @@ -45,11 +54,25 @@ internal class PostHogQueueTest { this.networkStatus = networkStatus this.maxBatchSize = maxBatchSize this.dateProvider = dateProvider + this.maxRetries = maxRetries + this.httpClient = httpClient } val api = PostHogApi(config) return PostHogQueue(config, EndpointSpec.batch(config, api, config.storagePrefix), executor) } + private fun waitUntil( + timeoutMillis: Long = 5_000, + condition: () -> Boolean, + ): Boolean { + val deadline = System.currentTimeMillis() + timeoutMillis + while (System.currentTimeMillis() < deadline) { + if (condition()) return true + Thread.sleep(10) + } + return condition() + } + @Test fun `respect maxQueueSize and deletes the first if full`() { val http = mockHttp() @@ -247,6 +270,307 @@ internal class PostHogQueueTest { assertEquals(0, sut.dequeList.size) } + @Test + fun `known offline flushes do not consume retries and availability drains the queue`() { + val http = mockHttp() + var connected = false + var onAvailableCallback: (() -> Unit)? = null + val path = tmpDir.newFolder().absolutePath + val sut = + getSut( + host = http.url("/").toString(), + storagePrefix = path, + flushAt = 1, + maxRetries = 0, + networkStatus = + object : PostHogNetworkStatus { + override fun isConnected() = connected + + override fun register(callback: () -> Unit) { + onAvailableCallback = callback + } + }, + ) + + try { + sut.start() + sut.add(generateEvent()) + executor.awaitExecution() + + repeat(3) { + sut.flush() + executor.awaitExecution() + } + + assertEquals(0, http.requestCount) + assertEquals(0, sut.currentRetryCountForTesting) + assertEquals(1, sut.dequeList.size) + assertEquals(1, File(path, API_KEY).listFiles()!!.size) + + connected = true + onAvailableCallback?.invoke() + executor.awaitExecution() + + assertEquals(1, http.requestCount) + assertEquals(0, sut.currentRetryCountForTesting) + assertEquals(0, sut.dequeList.size) + assertEquals(0, File(path, API_KEY).listFiles()!!.size) + } finally { + sut.stop() + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + + @Test + fun `retryable HTTP exhaustion retains a bounded queue and later success drains it`() { + val http = mockHttp(response = MockResponse().setResponseCode(503).setBody("error")) + http.enqueue(MockResponse().setResponseCode(503).setBody("error")) + http.enqueue(MockResponse().setBody("")) + val fakeCurrentTime = FakePostHogDateProvider() + fakeCurrentTime.setAddSecondsToCurrentDate(parseISO8601Date("1970-09-20T11:58:49.000Z")!!) + val path = tmpDir.newFolder().absolutePath + val sut = + getSut( + host = http.url("/").toString(), + storagePrefix = path, + flushAt = 100, + maxQueueSize = 2, + maxRetries = 0, + dateProvider = fakeCurrentTime, + ) + + try { + sut.add(generateEvent("first", givenUuuid = UUID.randomUUID())) + sut.add(generateEvent("second", givenUuuid = UUID.randomUUID())) + executor.awaitExecution() + val firstFile = sut.dequeList.first() + + repeat(2) { + sut.flush() + executor.awaitExecution() + + assertEquals(2, sut.dequeList.size) + assertEquals(2, File(path, API_KEY).listFiles()!!.size) + } + + sut.add(generateEvent("replacement", givenUuuid = UUID.randomUUID())) + executor.awaitExecution() + + assertEquals(2, sut.dequeList.size) + assertFalse(sut.dequeList.contains(firstFile)) + assertFalse(firstFile.exists()) + assertEquals(2, File(path, API_KEY).listFiles()!!.size) + + sut.flush() + executor.awaitExecution() + + assertEquals(3, http.requestCount) + assertEquals(0, sut.dequeList.size) + assertEquals(0, File(path, API_KEY).listFiles()!!.size) + } finally { + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + + @Test + fun `retry after remains authoritative beyond the exponential backoff cap`() { + val http = mockHttp(response = MockResponse().setResponseCode(429).setHeader("Retry-After", "120").setBody("error")) + var scheduledDelay = 0 + val dateProvider = + object : PostHogDateProvider { + override fun currentDate() = Date() + + override fun addSecondsToCurrentDate(seconds: Int): Date { + scheduledDelay = seconds + return Date(System.currentTimeMillis() + seconds * 1000L) + } + + override fun currentTimeMillis() = System.currentTimeMillis() + + override fun nanoTime() = System.nanoTime() + } + val path = tmpDir.newFolder().absolutePath + val sut = + getSut( + host = http.url("/").toString(), + storagePrefix = path, + flushAt = 1, + dateProvider = dateProvider, + ) + + try { + sut.add(generateEvent()) + executor.awaitExecution() + + assertEquals(120, scheduledDelay) + assertEquals(1, sut.dequeList.size) + } finally { + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + + @Test + fun `generic transport IO failures retain files beyond max retries and recover`() { + val http = mockHttp() + val failedAttempts = 3 + val attempts = AtomicInteger() + val httpClient = + OkHttpClient.Builder() + .addInterceptor { chain -> + if (attempts.incrementAndGet() <= failedAttempts) { + throw IOException("connection reset") + } + chain.proceed(chain.request()) + }.build() + val fakeCurrentTime = FakePostHogDateProvider() + fakeCurrentTime.setAddSecondsToCurrentDate(parseISO8601Date("1970-09-20T11:58:49.000Z")!!) + val path = tmpDir.newFolder().absolutePath + val sut = + getSut( + host = http.url("/").toString(), + storagePrefix = path, + flushAt = 100, + maxRetries = 0, + dateProvider = fakeCurrentTime, + httpClient = httpClient, + ) + + try { + sut.add(generateEvent()) + executor.awaitExecution() + + repeat(failedAttempts) { + sut.flush() + executor.awaitExecution() + + assertEquals(1, sut.dequeList.size) + assertEquals(1, File(path, API_KEY).listFiles()!!.size) + } + assertEquals(0, http.requestCount) + + sut.flush() + executor.awaitExecution() + + assertEquals(failedAttempts + 1, attempts.get()) + assertEquals(1, http.requestCount) + assertEquals(0, sut.dequeList.size) + assertEquals(0, File(path, API_KEY).listFiles()!!.size) + } finally { + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + + @Test + fun `duplicate payload UUIDs create distinct durable queue entries`() { + val http = mockHttp() + val path = tmpDir.newFolder().absolutePath + val sut = getSut(host = http.url("/").toString(), storagePrefix = path) + val payloadUuid = UUID.randomUUID() + val event = generateEvent("same", givenUuuid = payloadUuid) + + try { + sut.add(event) + sut.add(event) + executor.awaitExecution() + + assertEquals(2, sut.dequeList.size) + assertEquals(2, sut.dequeList.map { it.name }.toSet().size) + assertEquals(2, File(path, API_KEY).listFiles()!!.size) + } finally { + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + + @Test + fun `successful in-flight batch removes exact entries after full queue replacement`() { + val queueExecutor = Executors.newFixedThreadPool(2, PostHogThreadFactory("ConcurrentQueueTest")) + val path = tmpDir.newFolder().absolutePath + val config = + PostHogConfig(API_KEY).apply { + storagePrefix = path + maxQueueSize = 2 + maxBatchSize = 2 + flushAt = 100 + } + val initialRecordsWritten = CountDownLatch(2) + val replacementWritten = CountDownLatch(1) + val sendStarted = CountDownLatch(1) + val releaseSend = CountDownLatch(1) + val sendAttempts = AtomicInteger() + val sentRecords = Collections.synchronizedList(mutableListOf()) + val spec = + EndpointSpec( + recordsLabel = "records", + storagePrefix = path, + initialCap = { it.maxBatchSize }, + initialFlushAt = { it.flushAt }, + maxQueueSize = { it.maxQueueSize }, + flushIntervalSeconds = { it.flushIntervalSeconds }, + encode = { record, stream -> + stream.write(record.toByteArray()) + if (record == "replacement") { + replacementWritten.countDown() + } else { + initialRecordsWritten.countDown() + } + }, + decode = { stream -> String(stream.readBytes()) }, + describe = { it }, + send = { records -> + if (sendAttempts.incrementAndGet() == 1) { + sentRecords.addAll(records) + sendStarted.countDown() + check(releaseSend.await(5, TimeUnit.SECONDS)) + } else { + throw IOException("retain replacement after its later send attempt") + } + }, + isRetriableStatusCode = { false }, + ) + val sut = PostHogQueue(config, spec, queueExecutor) + + try { + sut.add("first") + sut.add("second") + assertTrue(initialRecordsWritten.await(5, TimeUnit.SECONDS)) + val initialFiles = sut.dequeList.toSet() + assertEquals(2, initialFiles.size) + + sut.flush() + assertTrue(sendStarted.await(5, TimeUnit.SECONDS)) + + sut.add("replacement") + assertTrue(replacementWritten.await(5, TimeUnit.SECONDS)) + val replacementFile = sut.dequeList.single { it !in initialFiles } + + releaseSend.countDown() + assertTrue { + waitUntil { + sut.currentRetryCountForTesting == 1 && sut.dequeList == listOf(replacementFile) + } + } + + assertEquals(2, sendAttempts.get()) + assertEquals(setOf("first", "second"), sentRecords.toSet()) + assertEquals("replacement", replacementFile.readText()) + assertTrue(replacementFile.exists()) + } finally { + releaseSend.countDown() + sut.clear() + queueExecutor.shutdownAndAwaitTermination() + } + } + @Test fun `does not delete file if API is 3xx`() { val http = mockHttp(response = MockResponse().setResponseCode(300).setBody("error")) @@ -617,6 +941,36 @@ internal class PostHogQueueTest { assertEquals(4, sut.dequeList.size) } + @Test + fun `reload evicts oldest cached files beyond queue capacity`() { + val http = mockHttp() + val path = tmpDir.newFolder().absolutePath + val dir = File(path, API_KEY) + dir.mkdirs() + val eventContent = File("src/test/resources/json/basic-event.json").readText() + val cachedFiles = + (1..3).map { index -> + File(dir, "${UUID.randomUUID()}.event").apply { + writeText(eventContent) + setLastModified(System.currentTimeMillis() - (4 - index) * 1000L) + } + } + val sut = getSut(host = http.url("/").toString(), storagePrefix = path, maxQueueSize = 2) + + try { + sut.reloadFromDisk() + + assertEquals(2, sut.dequeList.size) + assertFalse(cachedFiles.first().exists()) + assertEquals(cachedFiles.drop(1), sut.dequeList) + assertEquals(2, dir.listFiles()!!.size) + } finally { + sut.clear() + executor.shutdownAndAwaitTermination() + http.shutdown() + } + } + @Test fun `loads cached events and flushes them when add triggers threshold`() { val http = mockHttp()