fix(log): разблокировать спул лог-коллектора при устаревших таймстемпах
Раньше одна пачка в спуле с таймстемпом вне окна приёма Loki (≈7 суток) топила весь пуш: Loki отвечал 400 на весь batch, спул никогда не чистился, и новые логи (в т.ч. телефона) не доходили. Теперь flush() выгребает ВСЕ готовые пачки за тик (FIFO) и отбрасывает записи вне окна Loki, а при неудаче отправки спула буфер не трогается (ждёт следующий тик). Добавлены тесты на очистку устаревших записей и FIFO-ретрай спула.
This commit is contained in:
@@ -6,6 +6,9 @@ import kotlinx.coroutines.delay
|
|||||||
import kotlinx.coroutines.isActive
|
import kotlinx.coroutines.isActive
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
|
|
||||||
|
/** Окно приёма Loki: записи старше этого возраста (или из будущего) отбрасываем. */
|
||||||
|
private const val RETENTION_MS = 7L * 24 * 60 * 60 * 1000
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* «Абстрактный коллектор» лог-пачек: собирает записи в батчи и по времени
|
* «Абстрактный коллектор» лог-пачек: собирает записи в батчи и по времени
|
||||||
* шлёт их «вниз» через [sink] (абстракция следующего коллектора).
|
* шлёт их «вниз» через [sink] (абстракция следующего коллектора).
|
||||||
@@ -95,17 +98,40 @@ class BatchingLogCollector(
|
|||||||
* 2) новый батч из буфера; не слано — прихраняем в [spool].
|
* 2) новый батч из буфера; не слано — прихраняем в [spool].
|
||||||
*/
|
*/
|
||||||
suspend fun flush() {
|
suspend fun flush() {
|
||||||
val spooled = spool.peek()
|
// Ретраи спула: выгребаем ВСЕ готовые пачки за тик (FIFO), а не одну.
|
||||||
if (spooled != null) {
|
// Записи с «неприемлемым» таймстемпом (старше окна Loki ~7 дней или из
|
||||||
if (sendBatch(spooled)) {
|
// будущего) отбрасываем — иначе одна такая пачка топит весь пуш (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()
|
spool.take()
|
||||||
}
|
spooled = spool.peek()
|
||||||
|
} else {
|
||||||
|
// спул не ушёл — буфер не трогаем, ждём следующий тик
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
val batch = drainBuffer() ?: return
|
|
||||||
if (!sendBatch(batch)) {
|
|
||||||
spool.write(batch)
|
|
||||||
}
|
}
|
||||||
|
val batch = drainBuffer() ?: return
|
||||||
|
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 =
|
private suspend fun sendBatch(batch: LogBatch): Boolean =
|
||||||
|
|||||||
@@ -137,4 +137,39 @@ class BatchingLogCollectorTest {
|
|||||||
collector.flush()
|
collector.flush()
|
||||||
assertNull(sink.sent.firstOrNull())
|
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())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user