feat: управление яркостью/WiFi очков через Mercury + накопленные доработки

- Mercury (app-phone): управление яркостью (с авто-выкл. перед
  применением) и WiFi (скан/подключение/список сетей с паролем);
  добавлено право BLUETOOTH_ADVERTISE в манифест и runtime-запрос
  (без него GATT-сервер не рекламируется → очки не подключаются)
- логирование (lib-core core.log): BatchingLogCollector + LogSink +
  LogSpool + LokiLogSink; цепочка очки→телефон→Loki; lib-core log()
  переадресован в коллектор
- STT: Qwen3-ASR через pw.binom.asr:asr-qwen3-android, убран in-house
  whisper-фоллбек
- audio-sync: P-регулятор с насыщением вместо фиксированной скорости
- media-split: очки только video, телефон только audio (UI-чипы)
- Components: форматирование размера через Locale.ROOT
This commit is contained in:
2026-08-29 23:06:52 +03:00
parent 62d2f06797
commit 6cd4a8de43
40 changed files with 1409 additions and 203 deletions
@@ -0,0 +1,68 @@
package pw.binom.viewmate.core.log
import java.io.File
class FileLogSpool(
dir: String,
maxFiles: Int,
) : LogSpool {
private val filesDir: File = File(dir).also { if (!it.exists()) it.mkdirs() }
private val limit = maxFiles
/** Сериал следующих файлов: продолжение после старейшего существующего. */
private var seq: Long = currentFiles().mapNotNull { it.seqNumber() }.maxOrNull() ?: 0L
override fun write(batch: LogBatch) {
seq += 1
File(filesDir, "$PREFIX$seq.json").writeText(LogBatch.encodeBatch(batch))
trim()
}
override fun peek(): LogBatch? {
var file = oldestFile()
while (file != null) {
val batch = decodeQuiet(file)
if (batch != null) return batch
file.delete() // битый файл не должен затыкать очередь
file = oldestFile()
}
return null
}
override fun take(): LogBatch? {
var file = oldestFile()
while (file != null) {
val batch = decodeQuiet(file)
file.delete()
if (batch != null) return batch
file = oldestFile()
}
return null
}
override fun pendingCount(): Int = currentFiles().size
private fun oldestFile(): File? = currentFiles().minByOrNull { it.seqNumber() }
private fun currentFiles(): List<File> = filesDir
.listFiles { f -> f.isFile && f.name.startsWith(PREFIX) && f.name.endsWith(".json") }
?.toList()
?: emptyList()
private fun File.seqNumber(): Long =
name.substringAfter(PREFIX).removeSuffix(".json").toLongOrNull() ?: 0L
private fun decodeQuiet(file: File): LogBatch? =
runCatching { LogBatch.decodeBatch(file.readText()) }.getOrNull()
/** Потолок: выбросить старейшие файлы сверх [limit]. */
private fun trim() {
val files = currentFiles().sortedBy { it.seqNumber() }
if (files.size <= limit) return
files.take(files.size - limit).forEach { it.delete() }
}
private companion object {
const val PREFIX = "logspool-"
}
}
@@ -0,0 +1,16 @@
package pw.binom.viewmate.core.log
import java.util.concurrent.locks.ReentrantLock
actual class LogLock {
private val delegate = ReentrantLock()
actual fun <T> withLock(action: () -> T): T {
delegate.lock()
return try {
action()
} finally {
delegate.unlock()
}
}
}
@@ -0,0 +1,5 @@
package pw.binom.viewmate.core.log
actual class SystemClock : LogClock {
override fun nowMillis(): Long = System.currentTimeMillis()
}
@@ -1,6 +1,17 @@
package pw.binom.viewmate.core
/** Единый лог общего слоя (KMP): без java.time, просто тег + сообщение. */
import pw.binom.viewmate.core.log.BatchingLogCollector
/**
* Единый лог общего слоя (KMP): без java.time, просто тег + сообщение.
*
* Приложение подключает свой коллектор через [logCollector] (в onCreate),
* и каждый вызов [log] попадает в него (дальше — в спул / на сервер / Loki).
* По умолчанию коллектор не задан — лог только в println.
*/
var logCollector: BatchingLogCollector? = null
fun log(tag: String, message: String) {
println("[view-mate] $tag $message")
logCollector?.record(tag, message)
}
@@ -0,0 +1,121 @@
package pw.binom.viewmate.core.log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
/**
* «Абстрактный коллектор» лог-пачек: собирает записи в батчи и по времени
* шлёт их «вниз» через [sink] (абстракция следующего коллектора).
*
* Смысл:
* - [record]/[recordBatch] — синхронно, с любой нитки: запись в буфер (лок),
* при переполнении [maxBatchEntries] — немедленный flush;
* - таймер (период [flushIntervalMs]) — собирает буфер в [LogBatch] и шлёт в [sink];
* - «не получилось послать» ([LogSink.send] = false / бросил) — пачка
* **прихраняется в [spool]** и на каждом следующем тике ретраится ДО новых
* пачек (FIFO: старое раньше);
* - «чужие» записи (например, логи очков, пришедшие на телефон) — [recordBatch]
* с сохранением их [LogEntry.source]/[LogEntry.ts].
*
* Поток: буфер — из любых ниток (под [LogLock]); спул и отправка — только
* таймер-корутина (однопоточно).
*/
class BatchingLogCollector(
/** Устройство-носитель коллектора: "glasses" / "phone". */
val device: String,
/** «Следующий» коллектор: куда шлём пачки. */
private val sink: LogSink,
/** Куда прихраним не отправленные пачки (retry позже). */
private val spool: LogSpool,
/** Часы для меток времени записей. */
private val clock: LogClock,
/** Период отправки пачек, мс. */
private val flushIntervalMs: Long = 30_000L,
/** Лимит записей в одной пачке — превышение триггерит немедленный flush. */
private val maxBatchEntries: Int = 500,
) {
private val buffer = ArrayList<LogEntry>()
private val lock = LogLock()
/** Потолок буфера в памяти (страховка, если sink долго не принимал). */
private val maxBuffer: Int = maxBatchEntries * 20
@Volatile
private var timerJob: Job? = null
@Volatile
private var scope: CoroutineScope? = null
/** Записать свою лог-строку (синхронно, с любой нитки). */
fun record(tag: String, message: String) {
recordBatch(listOf(LogEntry(clock.nowMillis(), tag, message, device)))
}
/**
* Прибавить «чужие» записи (готовые [LogEntry] — со своими ts/source):
* телефон так пересылает логи очков.
*/
fun recordBatch(entries: List<LogEntry>) {
if (entries.isEmpty()) return
val overflowed = lock.withLock {
buffer.addAll(entries)
val drop = buffer.size - maxBuffer
if (drop > 0) repeat(drop) { buffer.removeAt(0) }
buffer.size >= maxBatchEntries
}
if (overflowed) scope?.launch { flush() }
}
/** Запустить таймер: каждые [flushIntervalMs] — flush (спул-ретраи + новый батч). */
fun start(scope: CoroutineScope) {
if (timerJob != null) return
this.scope = scope
timerJob = scope.launch {
while (isActive) {
delay(flushIntervalMs)
runCatching { flush() }
}
}
}
/** Остановить таймер (буфер/спул не трогаются — дошлёт следующий запуск). */
fun stop() {
timerJob?.cancel()
timerJob = null
scope = null
}
/**
* Цикл отправки:
* 1) ретраи спула — только старейшая пачка; не слано — ждём следующий тик
* (новый батч не «перескакивает» через неё, порядок сохраняется);
* 2) новый батч из буфера; не слано — прихраняем в [spool].
*/
suspend fun flush() {
val spooled = spool.peek()
if (spooled != null) {
if (sendBatch(spooled)) {
spool.take()
}
return
}
val batch = drainBuffer() ?: return
if (!sendBatch(batch)) {
spool.write(batch)
}
}
private suspend fun sendBatch(batch: LogBatch): Boolean =
runCatching { sink.send(batch) }.getOrDefault(false)
/** Взять все накопленные записи в одну пачку (null — пусто). */
private fun drainBuffer(): LogBatch? = lock.withLock {
if (buffer.isEmpty()) return@withLock null
val batch = LogBatch(device, buffer.toList())
buffer.clear()
batch
}
}
@@ -0,0 +1,36 @@
package pw.binom.viewmate.core.log
import io.ktor.client.HttpClient
import io.ktor.client.request.header
import io.ktor.client.request.post
import io.ktor.client.request.setBody
import io.ktor.client.statement.HttpResponse
import io.ktor.http.ContentType
import io.ktor.http.contentType
import io.ktor.http.isSuccess
import pw.binom.viewmate.core.media.defaultHttpClient
/**
* [LogSink] по HTTP: POST JSON-пачки на [url] (заголовок X-API-Key, если [apiKey] задан).
* 2xx — принято (true); любой другой ответ/сетевая ошибка — false (пачка уйдёт в спул).
*/
class HttpLogSink(
private val client: HttpClient,
private val url: String,
private val apiKey: String? = null,
) : LogSink {
override suspend fun send(batch: LogBatch): Boolean {
val response: HttpResponse = runCatching {
client.post(url) {
if (apiKey != null) header("X-API-Key", apiKey)
contentType(ContentType.Application.Json)
setBody(LogBatch.encodeBatch(batch))
}
}.getOrNull() ?: return false
return response.status.isSuccess()
}
}
/** Готовый HTTP-лог-синк со своим клиентом: вызывающая сторона не знает про Ktor. */
fun httpLogSink(url: String, apiKey: String? = null): LogSink =
HttpLogSink(defaultHttpClient(), url, apiKey)
@@ -0,0 +1,44 @@
package pw.binom.viewmate.core.log
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json
/**
* Одна лог-запись.
* [ts] — часы устройства-источника, epoch-мс.
* [tag] — подсистема (stt, hub, mc, …).
* [source] — кто сгенерировал ("glasses" / "phone") — чтобы телефон мог пересылать
* логи очков на сервер с сохранением источника.
*/
@Serializable
data class LogEntry(
val ts: Long,
val tag: String,
val message: String,
val source: String = "unknown",
)
/**
* Пачка логов — единица передачи в «следующий» коллектор.
* [device] — кто шлёт эту пачку (устройство-коллектор).
*/
@Serializable
data class LogBatch(
val device: String,
val entries: List<LogEntry>,
) {
companion object {
private val json: Json = Json {
ignoreUnknownKeys = true
encodeDefaults = true
}
/** Пачка → JSON-строка (файл спула). */
fun encodeBatch(batch: LogBatch): String = json.encodeToString(serializer(), batch)
/** JSON-строка → пачка; некорректный/битый текст — null. */
fun decodeBatch(text: String): LogBatch? = runCatching {
json.decodeFromString(serializer(), text)
}.getOrNull()
}
}
@@ -0,0 +1,19 @@
package pw.binom.viewmate.core.log
/** Часы для лог-записей: у common-кода нет java.time — устройство-носитель подаёт свои. */
fun interface LogClock {
fun nowMillis(): Long
}
/** Стенговые (wall-clock) часы — JVM-реализация (Android/JVM: `System.currentTimeMillis`). */
expect class SystemClock() : LogClock
/** Зависанный (fixed) ход времени для юнит-тестов. */
class FixedClock(private var now: Long = 0L) : LogClock {
override fun nowMillis(): Long = now
/** Продвинуть «время» на [ms]. */
fun advance(ms: Long) {
now += ms
}
}
@@ -0,0 +1,10 @@
package pw.binom.viewmate.core.log
/**
* Лока (thread-safe) для буфера коллектора: KMP-common не даёт нативных лок-ов —
* каждая платформа подставляет свою (Android/JVM: `java.util.concurrent.locks.ReentrantLock`).
*/
expect class LogLock() {
/** Выполнить [action] под локом; вернуть результат. */
fun <T> withLock(action: () -> T): T
}
@@ -0,0 +1,20 @@
package pw.binom.viewmate.core.log
/**
* Абстракция «следующего» получателя лог-пачек (сервер у телефона, сам телефон
* у очков, чёрная дыра в тестах). Коллектор не знает про транспорт/протокол —
* только про результат отправки.
*/
fun interface LogSink {
/**
* Отправить [batch] вниз по цепочке.
* @return `true` — принято (пачка потреблена); `false` — не принято
* (коллектор положит пачку в спул и попробует позже).
*/
suspend fun send(batch: LogBatch): Boolean
}
/** LogSink, который пачки проглатывает (логирование отключено / нет цели). */
object NoopLogSink : LogSink {
override suspend fun send(batch: LogBatch): Boolean = true
}
@@ -0,0 +1,41 @@
package pw.binom.viewmate.core.log
/**
* Хранилище не отправленных пачек (спул): «не удалось послать → прихраним у себя
* и позже попробуем ещё». FIFO: первым ретраится самая старая пачка.
*
* Доступ — из одного потока (таймер-корутина коллектора): peek/take на каждом
* тике, write — при неудачной отправке.
*/
interface LogSpool {
/** Добавить [batch] в конец очереди. */
fun write(batch: LogBatch)
/** Самая старая пачка без удаления (null — пусто). */
fun peek(): LogBatch?
/** Самая старая пачка с удалением (null — пусто). */
fun take(): LogBatch?
/** Сколько пачек ждёт отправки. */
fun pendingCount(): Int
}
/**
* In-память спул без персистентности (тесты, «прихранил в процессе»):
* очередь с потолком [maxBatches] — при переполнении выбрасывается старейшая.
*/
class InMemoryLogSpool(private val maxBatches: Int = 32) : LogSpool {
private val queue = ArrayDeque<LogBatch>()
override fun write(batch: LogBatch) {
if (queue.size >= maxBatches) queue.removeFirst()
queue.addLast(batch)
}
override fun peek(): LogBatch? = queue.firstOrNull()
override fun take(): LogBatch? = queue.removeFirstOrNull()
override fun pendingCount(): Int = queue.size
}
@@ -0,0 +1,94 @@
package pw.binom.viewmate.core.log
import io.ktor.client.HttpClient
import io.ktor.client.request.header
import io.ktor.client.request.post
import io.ktor.client.request.setBody
import io.ktor.client.statement.HttpResponse
import io.ktor.http.ContentType
import io.ktor.http.contentType
import io.ktor.http.isSuccess
import pw.binom.viewmate.core.media.defaultHttpClient
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json
/**
* Коллектор в Loki (Grafana): POST JSON в [LokiConfig.url] (Basic-auth).
*
* Важное: instance каждого stream'а = [LogEntry.source] записи ("phone"/"glasses"),
* то есть в Loki сразу видно, откуда лог — с телефона или с очков. Пачка на
* телефоне может содержать и то, и другое (телефон пересылает логи очков),
* поэтому группируем по source и шлём несколько stream'ов одним запросом.
*/
class LokiLogSink(
private val client: HttpClient,
private val config: LokiConfig = LokiConfig(),
) : LogSink {
private val authHeader = "Basic " + base64("${config.user}:${config.password}")
private val json = Json { ignoreUnknownKeys = true }
override suspend fun send(batch: LogBatch): Boolean {
if (batch.entries.isEmpty()) return true
val payload = buildLokiPush(batch.entries, config.job)
val response: HttpResponse = runCatching {
client.post(config.url) {
header("Authorization", authHeader)
contentType(ContentType.Application.Json)
setBody(json.encodeToString(LokiPush.serializer(), payload))
}
}.getOrNull() ?: return false
return response.status.isSuccess()
}
}
/** Параметры Loki. Логин/пароль/адрес зашиты по умолчанию (см. запрос пользователя). */
data class LokiConfig(
val url: String = "https://loki.binom.pw/loki/api/v1/push",
val user: String = "loki-ingest",
val password: String = "q0OH8na60bE_KhGc6c1YOsUuFPrLjw2K",
val job: String = "view-mate",
)
/** Готовый Loki-синк со своим клиентом: вызывающая сторона не знает про Ktor/Basic-auth. */
fun lokiLogSink(config: LokiConfig = LokiConfig()): LogSink =
LokiLogSink(defaultHttpClient(), config)
@Serializable
internal data class LokiStream(val stream: Map<String, String>, val values: List<List<String>>)
@Serializable
internal data class LokiPush(val streams: List<LokiStream>)
internal fun buildLokiPush(entries: List<LogEntry>, job: String): LokiPush {
val streams = entries
.groupBy { it.source.ifBlank { "phone" } }
.map { (source, es) ->
val values = es
.map { e -> listOf((e.ts * 1_000_000).toString(), "${e.tag}: ${e.message}") }
.sortedBy { it[0] }
LokiStream(stream = mapOf("job" to job, "instance" to source), values = values)
}
return LokiPush(streams = streams)
}
private val B64 = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"
/** Минимальный base64 (только кодирование) — для Basic-auth, без новых зависимостей. */
private fun base64(input: String): String {
val bytes = input.encodeToByteArray()
val out = StringBuilder()
var i = 0
while (i < bytes.size) {
val b0 = bytes[i].toInt() and 0xFF
val b1 = if (i + 1 < bytes.size) bytes[i + 1].toInt() and 0xFF else 0
val b2 = if (i + 2 < bytes.size) bytes[i + 2].toInt() and 0xFF else 0
val triple = (b0 shl 16) or (b1 shl 8) or b2
out.append(B64[(triple ushr 18) and 0x3F])
out.append(B64[(triple ushr 12) and 0x3F])
out.append(if (i + 1 < bytes.size) B64[(triple ushr 6) and 0x3F] else '=')
out.append(if (i + 2 < bytes.size) B64[triple and 0x3F] else '=')
i += 3
}
return out.toString()
}
@@ -128,8 +128,9 @@ class PhoneActions(
}
/**
* Скачать контент в память очков: DownloadFiles с presigned URL всех файлов
* (video.mkv + audio-N.ogg) и известными размерами дорожек [sizes] (key → байт).
* Скачать контент в память очков: DownloadFiles с presigned URL только видео (video.mkv).
* Звук на очки не идёт — аудио-дорожки скачиваются только на телефон.
* [sizes] — известные размеры дорожек (key → байт).
*/
suspend fun downloadToGlasses(
itemId: String,
@@ -3,6 +3,8 @@ package pw.binom.viewmate.core.protocol
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import pw.binom.viewmate.core.log.LogEntry
/**
* Сообщения от очков к хосту (телефону).
* Полиморфная сериализация: дискриминатор "type" со значением @SerialName.
@@ -102,3 +104,12 @@ data class SttAudio(val data: ByteArray) : GlassesToHost {
@Serializable
@SerialName("stop_stt")
data class StopStt(val cancel: Boolean) : GlassesToHost
/**
* Лог-пачка очков (очки → телефон). Телефон прибавляет эти записи к своему
* коллектору и пересылает на сервер вместе со своими лог-записями;
* источник — [LogEntry.source]. Старые очки такое сообщение не шлют.
*/
@Serializable
@SerialName("log_batch")
data class LogBatchMsg(val entries: List<LogEntry>) : GlassesToHost
@@ -0,0 +1,140 @@
package pw.binom.viewmate.core.log
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.cancel
import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runTest
class BatchingLogCollectorTest {
/** Фейк «следующего» коллектора: фиксирует пачки, [fail] — имитация сбоя. */
class FakeSink(var fail: Boolean = false) : LogSink {
val sent = mutableListOf<LogBatch>()
override suspend fun send(batch: LogBatch): Boolean {
if (fail) return false
sent += batch
return true
}
}
@Test
fun `timer flush sends buffered entries`() = runTest {
val sink = FakeSink()
val collector = BatchingLogCollector(
device = "glasses",
sink = sink,
spool = InMemoryLogSpool(),
clock = FixedClock(1000),
flushIntervalMs = 100,
maxBatchEntries = 10,
)
@OptIn(ExperimentalCoroutinesApi::class)
val scope = CoroutineScope(UnconfinedTestDispatcher(testScheduler))
collector.start(scope)
collector.record("tag1", "hello")
collector.record("tag2", "world")
testScheduler.advanceTimeBy(150)
val batch = sink.sent.single()
assertEquals(2, batch.entries.size)
assertEquals("glasses", batch.device)
assertEquals(listOf("tag1", "tag2"), batch.entries.map { it.tag })
assertEquals("glasses", batch.entries[0].source)
scope.cancel()
}
@Test
fun `failed sink spools batch, later flush retries`() = runTest {
val sink = FakeSink(fail = true)
val spool = InMemoryLogSpool()
val collector = BatchingLogCollector("phone", sink, spool, FixedClock())
collector.record("a", "1")
collector.flush()
assertEquals(0, sink.sent.size)
assertEquals(1, spool.pendingCount())
sink.fail = false
collector.flush() // ретраи спула
assertEquals(1, sink.sent.size)
assertEquals("a", sink.sent[0].entries.single().tag)
assertEquals(0, spool.pendingCount())
}
@Test
fun `spool backlog blocks new batches, FIFO order kept`() = runTest {
val sink = FakeSink(fail = true)
val spool = InMemoryLogSpool()
val collector = BatchingLogCollector("phone", sink, spool, FixedClock())
collector.record("a", "1")
collector.flush() // batch1 → спул
collector.record("b", "2")
collector.flush() // спул затык — batch2 остаётся в буфере
assertEquals(1, spool.pendingCount())
sink.fail = false
collector.flush() // доразгрузили спул (batch1)
collector.flush() // теперь batch2
assertEquals(2, sink.sent.size)
assertEquals(listOf("a", "b"), sink.sent.flatMap { it.entries }.map { it.tag })
assertEquals(0, spool.pendingCount())
}
@Test
fun `foreign entries keep their source and ts`() = runTest {
val sink = FakeSink()
val collector = BatchingLogCollector("phone", sink, InMemoryLogSpool(), FixedClock(999))
collector.recordBatch(listOf(LogEntry(42, "mc", "glasses event", "glasses")))
collector.record("phone", "own")
collector.flush()
val batch = sink.sent.single()
assertEquals("phone", batch.device)
assertEquals(2, batch.entries.size)
assertEquals("glasses", batch.entries[0].source)
assertEquals(42, batch.entries[0].ts)
assertEquals("phone", batch.entries[1].source)
}
@OptIn(ExperimentalCoroutinesApi::class)
@Test
fun `overflowing batch flushes immediately`() = runTest {
val sink = FakeSink()
val collector = BatchingLogCollector(
device = "glasses",
sink = sink,
spool = InMemoryLogSpool(),
clock = FixedClock(),
flushIntervalMs = 60_000,
maxBatchEntries = 2,
)
val scope = CoroutineScope(UnconfinedTestDispatcher(testScheduler))
collector.start(scope)
collector.record("a", "1")
collector.record("b", "2") // буфер == лимит → немедленный flush
testScheduler.advanceTimeBy(0)
val batch = sink.sent.single()
assertEquals(2, batch.entries.size)
scope.cancel()
}
@Test
fun `empty recordBatch is no-op`() = runTest {
val sink = FakeSink()
val collector = BatchingLogCollector("phone", sink, InMemoryLogSpool(), FixedClock())
collector.recordBatch(emptyList())
collector.flush()
assertNull(sink.sent.firstOrNull())
}
}
@@ -0,0 +1,50 @@
package pw.binom.viewmate.core.log
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
class LogSpoolTest {
private fun batch(tag: String) = LogBatch("phone", listOf(LogEntry(0, tag, tag, "phone")))
@Test
fun `in-memory spool is FIFO with cap dropping oldest`() {
val spool = InMemoryLogSpool(maxBatches = 2)
spool.write(batch("1"))
spool.write(batch("2"))
spool.write(batch("3")) // «1» выброшена
assertEquals(2, spool.pendingCount())
assertEquals("2", spool.peek()?.entries?.single()?.tag)
assertEquals("2", spool.take()?.entries?.single()?.tag)
assertEquals("3", spool.take()?.entries?.single()?.tag)
assertNull(spool.take())
assertEquals(0, spool.pendingCount())
}
@Test
fun `in-memory spool peek does not consume`() {
val spool = InMemoryLogSpool()
spool.write(batch("a"))
assertEquals("a", spool.peek()?.entries?.single()?.tag)
assertEquals("a", spool.peek()?.entries?.single()?.tag)
assertEquals(1, spool.pendingCount())
}
@Test
fun `batch json roundtrip`() {
val batch = LogBatch(
"phone",
listOf(LogEntry(1, "a", "x", "phone"), LogEntry(2, "b", "y", "glasses")),
)
val text = LogBatch.encodeBatch(batch)
assertEquals(batch, LogBatch.decodeBatch(text))
}
@Test
fun `batch json decode garbage is null`() {
assertNull(LogBatch.decodeBatch("garbage"))
assertNull(LogBatch.decodeBatch(""))
}
}
@@ -0,0 +1,40 @@
package pw.binom.viewmate.core.log
import kotlin.test.Test
import kotlin.test.assertEquals
class LokiLogSinkTest {
@Test
fun `payload groups by source and labels instance`() {
val entries = listOf(
LogEntry(1_000, "hub", "phone line", "phone"),
LogEntry(2_000, "mc", "glasses line", "glasses"),
LogEntry(3_000, "hub", "another phone", "phone"),
)
val push = buildLokiPush(entries, "view-mate")
assertEquals(2, push.streams.size)
val phoneStream = push.streams.first { it.stream["instance"] == "phone" }
val glassesStream = push.streams.first { it.stream["instance"] == "glasses" }
assertEquals("view-mate", phoneStream.stream["job"])
assertEquals(2, phoneStream.values.size)
// timestamp: ms * 1e6, sorted ascending
assertEquals("1000000000", phoneStream.values[0][0])
assertEquals("hub: phone line", phoneStream.values[0][1])
assertEquals("3000000000", phoneStream.values[1][0])
assertEquals(1, glassesStream.values.size)
assertEquals("2000000000", glassesStream.values[0][0])
assertEquals("mc: glasses line", glassesStream.values[0][1])
}
@Test
fun `blank source falls back to phone`() {
val push = buildLokiPush(listOf(LogEntry(0, "t", "m", "")), "job")
assertEquals(1, push.streams.size)
assertEquals("phone", push.streams[0].stream["instance"])
}
}
@@ -0,0 +1,68 @@
package pw.binom.viewmate.core.log
import java.io.File
class FileLogSpool(
dir: String,
maxFiles: Int,
) : LogSpool {
private val filesDir: File = File(dir).also { if (!it.exists()) it.mkdirs() }
private val limit = maxFiles
/** Сериал следующих файлов: продолжение после старейшего существующего. */
private var seq: Long = currentFiles().mapNotNull { it.seqNumber() }.maxOrNull() ?: 0L
override fun write(batch: LogBatch) {
seq += 1
File(filesDir, "$PREFIX$seq.json").writeText(LogBatch.encodeBatch(batch))
trim()
}
override fun peek(): LogBatch? {
var file = oldestFile()
while (file != null) {
val batch = decodeQuiet(file)
if (batch != null) return batch
file.delete() // битый файл не должен затыкать очередь
file = oldestFile()
}
return null
}
override fun take(): LogBatch? {
var file = oldestFile()
while (file != null) {
val batch = decodeQuiet(file)
file.delete()
if (batch != null) return batch
file = oldestFile()
}
return null
}
override fun pendingCount(): Int = currentFiles().size
private fun oldestFile(): File? = currentFiles().minByOrNull { it.seqNumber() }
private fun currentFiles(): List<File> = filesDir
.listFiles { f -> f.isFile && f.name.startsWith(PREFIX) && f.name.endsWith(".json") }
?.toList()
?: emptyList()
private fun File.seqNumber(): Long =
name.substringAfter(PREFIX).removeSuffix(".json").toLongOrNull() ?: 0L
private fun decodeQuiet(file: File): LogBatch? =
runCatching { LogBatch.decodeBatch(file.readText()) }.getOrNull()
/** Потолок: выбросить старейшие файлы сверх [limit]. */
private fun trim() {
val files = currentFiles().sortedBy { it.seqNumber() }
if (files.size <= limit) return
files.take(files.size - limit).forEach { it.delete() }
}
private companion object {
const val PREFIX = "logspool-"
}
}
@@ -0,0 +1,16 @@
package pw.binom.viewmate.core.log
import java.util.concurrent.locks.ReentrantLock
actual class LogLock {
private val delegate = ReentrantLock()
actual fun <T> withLock(action: () -> T): T {
delegate.lock()
return try {
action()
} finally {
delegate.unlock()
}
}
}
@@ -0,0 +1,5 @@
package pw.binom.viewmate.core.log
actual class SystemClock : LogClock {
override fun nowMillis(): Long = System.currentTimeMillis()
}
@@ -0,0 +1,66 @@
package pw.binom.viewmate.core.log
import java.io.File
import java.nio.file.Files
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
class FileLogSpoolTest {
private var tempDir: File? = null
private fun batch(tag: String) = LogBatch("phone", listOf(LogEntry(0, tag, tag, "phone")))
private fun freshDir(name: String): File {
val dir = Files.createTempDirectory("logspool-$name").toFile()
tempDir = dir
return dir
}
@AfterTest
fun cleanup() {
tempDir?.deleteRecursively()
}
@Test
fun `writes and drains FIFO, survives restart, trims cap`() {
val dir = freshDir("fif")
var spool = FileLogSpool(dir.absolutePath, maxFiles = 2)
spool.write(batch("1"))
spool.write(batch("2"))
spool.write(batch("3")) // «1» выбрасывается по потолку
assertEquals(2, spool.pendingCount())
assertEquals("2", spool.peek()?.entries?.single()?.tag)
assertEquals("2", spool.take()?.entries?.single()?.tag)
// «перезапуск»: новая инстанция видит прихранианные файлы и продолжает сериал
spool = FileLogSpool(dir.absolutePath, maxFiles = 2)
assertEquals(1, spool.pendingCount())
assertEquals("3", spool.take()?.entries?.single()?.tag)
assertNull(spool.take())
}
@Test
fun `corrupt file is skipped without blocking queue`() {
val dir = freshDir("corrupt")
File(dir, "logspool-1.json").writeText("not a json")
val spool = FileLogSpool(dir.absolutePath, 64)
spool.write(batch("a"))
assertEquals("a", spool.take()?.entries?.single()?.tag) // битый удалён, «a» дошёл
assertNull(spool.take())
assertEquals(0, spool.pendingCount())
}
@Test
fun `empty directory is empty spool`() {
val dir = freshDir("empty")
val spool = FileLogSpool(dir.absolutePath, 64)
assertEquals(0, spool.pendingCount())
assertNull(spool.peek())
assertNull(spool.take())
}
}