diff --git a/lib-core/src/commonMain/kotlin/pw/binom/viewmate/core/log/BatchingLogCollector.kt b/lib-core/src/commonMain/kotlin/pw/binom/viewmate/core/log/BatchingLogCollector.kt index 0d36be8..090fbe4 100644 --- a/lib-core/src/commonMain/kotlin/pw/binom/viewmate/core/log/BatchingLogCollector.kt +++ b/lib-core/src/commonMain/kotlin/pw/binom/viewmate/core/log/BatchingLogCollector.kt @@ -6,6 +6,9 @@ import kotlinx.coroutines.delay import kotlinx.coroutines.isActive import kotlinx.coroutines.launch +/** Окно приёма Loki: записи старше этого возраста (или из будущего) отбрасываем. */ +private const val RETENTION_MS = 7L * 24 * 60 * 60 * 1000 + /** * «Абстрактный коллектор» лог-пачек: собирает записи в батчи и по времени * шлёт их «вниз» через [sink] (абстракция следующего коллектора). @@ -95,19 +98,42 @@ class BatchingLogCollector( * 2) новый батч из буфера; не слано — прихраняем в [spool]. */ suspend fun flush() { - val spooled = spool.peek() - if (spooled != null) { - if (sendBatch(spooled)) { + // Ретраи спула: выгребаем ВСЕ готовые пачки за тик (FIFO), а не одну. + // Записи с «неприемлемым» таймстемпом (старше окна Loki ~7 дней или из + // будущего) отбрасываем — иначе одна такая пачка топит весь пуш (Loki + // отвечает 400 на весь batch) и спул никогда не чистится. + var spooled = spool.peek() + while (spooled != null) { + val cleaned = clean(spooled) + val ok = if (cleaned.entries.isEmpty()) true else sendBatch(cleaned) + if (ok) { spool.take() + spooled = spool.peek() + } else { + // спул не ушёл — буфер не трогаем, ждём следующий тик + return } - return } val batch = drainBuffer() ?: return - if (!sendBatch(batch)) { - spool.write(batch) + val cleaned = clean(batch) + if (cleaned.entries.isNotEmpty()) { + if (!sendBatch(cleaned)) spool.write(cleaned) } } + /** + * Отбросить записи вне окна приёма Loki: старше [RETENTION_MS] или из будущего. + * Такие записи Loki всё равно отверг бы всем батчем (HTTP 400), поэтому их + * безопасно выкинуть — они не могут попасть в хранилище. + */ + private fun clean(batch: LogBatch): LogBatch { + val now = clock.nowMillis() + val min = now - RETENTION_MS + val max = now + 3_600_000L + val kept = batch.entries.filter { it.ts in min..max } + return if (kept.size == batch.entries.size) batch else batch.copy(entries = kept) + } + private suspend fun sendBatch(batch: LogBatch): Boolean = runCatching { sink.send(batch) }.getOrDefault(false) diff --git a/lib-core/src/commonTest/kotlin/pw/binom/viewmate/core/log/BatchingLogCollectorTest.kt b/lib-core/src/commonTest/kotlin/pw/binom/viewmate/core/log/BatchingLogCollectorTest.kt index 610edd6..931302d 100644 --- a/lib-core/src/commonTest/kotlin/pw/binom/viewmate/core/log/BatchingLogCollectorTest.kt +++ b/lib-core/src/commonTest/kotlin/pw/binom/viewmate/core/log/BatchingLogCollectorTest.kt @@ -137,4 +137,39 @@ class BatchingLogCollectorTest { collector.flush() assertNull(sink.sent.firstOrNull()) } + + @Test + fun `out-of-range spooled entries are dropped, buffer still flushed`() = runTest { + val sink = FakeSink() + val spool = InMemoryLogSpool() + val now = 1_000_000_000_000L + val collector = BatchingLogCollector("phone", sink, spool, FixedClock(now)) + + // «старые» пачки в спуле (метки из прошлого, как при загрузке без NTP) + spool.write(LogBatch("phone", listOf(LogEntry(0, "old", "stale", "phone")))) + spool.write(LogBatch("phone", listOf(LogEntry(123, "older", "stale2", "phone")))) + + collector.record("fresh", "current") // ts = now → валидно + collector.flush() + + // старые пачки выкинуты (не топят пуш), свежая ушла в sink + assertEquals(0, spool.pendingCount()) + assertEquals(1, sink.sent.size) + assertEquals(listOf("fresh"), sink.sent[0].entries.map { it.tag }) + } + + @Test + fun `valid spooled batch retried and sent before buffer`() = runTest { + val sink = FakeSink(fail = true) + val spool = InMemoryLogSpool() + val collector = BatchingLogCollector("phone", sink, spool, FixedClock(1_000_000_000_000L)) + collector.record("a", "1") + collector.flush() // → спул (fail) + sink.fail = false + collector.record("b", "2") + collector.flush() + assertEquals(2, sink.sent.size) + assertEquals(listOf("a", "b"), sink.sent.flatMap { it.entries }.map { it.tag }) + assertEquals(0, spool.pendingCount()) + } }