diff --git a/apps/cli/src/main/kotlin/com/correx/apps/cli/CorrexCli.kt b/apps/cli/src/main/kotlin/com/correx/apps/cli/CorrexCli.kt index 60a3a7a6..8d7fd0e6 100644 --- a/apps/cli/src/main/kotlin/com/correx/apps/cli/CorrexCli.kt +++ b/apps/cli/src/main/kotlin/com/correx/apps/cli/CorrexCli.kt @@ -2,6 +2,7 @@ package com.correx.apps.cli import com.correx.apps.cli.commands.ApproveCommand import com.correx.apps.cli.commands.EventsCommand +import com.correx.apps.cli.commands.HealthCommand import com.correx.apps.cli.commands.ProviderCommand import com.correx.apps.cli.commands.ReplayCommand import com.correx.apps.cli.commands.RunCommand @@ -31,4 +32,5 @@ fun buildCli(): CorrexCli = CorrexCli().subcommands( EventsCommand(), ReplayCommand(), StatsCommand(), + HealthCommand(), ) diff --git a/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/HealthCommand.kt b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/HealthCommand.kt new file mode 100644 index 00000000..1b2a5860 --- /dev/null +++ b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/HealthCommand.kt @@ -0,0 +1,79 @@ +package com.correx.apps.cli.commands + +import com.correx.apps.cli.CorrexCli +import com.correx.apps.cli.DEFAULT_PORT +import com.github.ajalt.clikt.core.CliktCommand +import com.github.ajalt.clikt.parameters.options.default +import com.github.ajalt.clikt.parameters.options.option +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.engine.cio.CIO +import io.ktor.client.plugins.contentnegotiation.ContentNegotiation +import io.ktor.client.request.get +import io.ktor.client.statement.bodyAsText +import io.ktor.serialization.kotlinx.json.json +import kotlinx.coroutines.runBlocking +import kotlinx.serialization.Serializable +import kotlinx.serialization.json.Json + +@Serializable +data class HealthSubjectReportDto( + val subject: String, + val status: String, + val metric: String, + val observedValue: Long = 0, + val detail: String = "", + val since: String = "", +) + +@Serializable +data class HealthReportDto( + val overall: String = "HEALTHY", + val subjects: List = emptyList(), + val checkedAt: String = "", +) + +fun renderHealth(report: HealthReportDto): String { + val lines = mutableListOf() + lines += "Overall: ${report.overall}" + if (report.checkedAt.isNotEmpty()) lines += "Checked: ${report.checkedAt}" + lines += "" + if (report.subjects.isEmpty()) { + lines += "No health checks recorded yet." + return lines.joinToString("\n") + } + for (s in report.subjects) { + lines += "%-14s %-9s %s".format(s.subject, s.status, s.detail) + if (s.since.isNotEmpty()) lines += " since ${s.since}" + } + return lines.joinToString("\n") +} + +class HealthCommand : CliktCommand(name = "health") { + private val host by option("--host").default("localhost") + private val port by option("--port").default("$DEFAULT_PORT") + + private val healthJson = Json { ignoreUnknownKeys = true } + + override fun run(): Unit = runBlocking { + val portInt = port.toIntOrNull() ?: DEFAULT_PORT + val outputJson = (currentContext.findRoot().command as? CorrexCli)?.json ?: false + val url = "http://$host:$portInt/health/checks" + + val client = HttpClient(CIO) { + install(ContentNegotiation) { json(healthJson) } + } + + runCatching { + if (outputJson) { + println(client.get(url).bodyAsText()) + } else { + println(renderHealth(client.get(url).body())) + } + }.onFailure { e -> + System.err.println("Error fetching health: ${e.message}") + } + + client.close() + } +} diff --git a/apps/cli/src/test/kotlin/com/correx/apps/cli/commands/HealthRenderTest.kt b/apps/cli/src/test/kotlin/com/correx/apps/cli/commands/HealthRenderTest.kt new file mode 100644 index 00000000..2e8083f3 --- /dev/null +++ b/apps/cli/src/test/kotlin/com/correx/apps/cli/commands/HealthRenderTest.kt @@ -0,0 +1,37 @@ +package com.correx.apps.cli.commands + +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class HealthRenderTest { + + @Test + fun `render shows overall and each subject row`() { + val report = HealthReportDto( + overall = "DEGRADED", + subjects = listOf( + HealthSubjectReportDto( + subject = "DISK", + status = "DEGRADED", + metric = "total_bytes", + observedValue = 6_000_000_000, + detail = "event log + CAS at 5.6 GiB (warn ≥ 4.7 GiB)", + since = "2026-06-13T01:00:00Z", + ), + ), + checkedAt = "2026-06-13T02:00:00Z", + ) + val out = renderHealth(report) + assertTrue(out.contains("Overall: DEGRADED"), "overall missing") + assertTrue(out.contains("DISK"), "subject row missing") + assertTrue(out.contains("5.6 GiB"), "detail missing") + assertTrue(out.contains("since 2026-06-13T01:00:00Z"), "since line missing") + } + + @Test + fun `render handles the empty no-checks case`() { + val out = renderHealth(HealthReportDto(overall = "HEALTHY")) + assertTrue(out.contains("Overall: HEALTHY")) + assertTrue(out.contains("No health checks recorded yet."), "empty hint missing") + } +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/Application.kt b/apps/server/src/main/kotlin/com/correx/apps/server/Application.kt index ecd4824e..108b8acb 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/Application.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/Application.kt @@ -1,5 +1,6 @@ package com.correx.apps.server +import com.correx.apps.server.health.HealthInspectionService import com.correx.apps.server.routes.providerRoutes import com.correx.apps.server.routes.sessionRoutes import com.correx.apps.server.routes.workflowRoutes @@ -35,6 +36,7 @@ fun Application.configureServer(module: ServerModule) { } val globalStreamHandler = GlobalStreamHandler(module) + val healthInspection = HealthInspectionService(module.eventStore) routing { get("/health") { @@ -42,6 +44,12 @@ fun Application.configureServer(module: ServerModule) { call.respond(mapOf("status" to "ok", "providers" to providerHealth.size.toString())) } + // observability-spec §4: the recorded health-check report (disk watermark, etc.) derived + // by replaying the system session's degraded/restored events. Backs `correx health`. + get("/health/checks") { + call.respond(healthInspection.inspect()) + } + webSocket("/stream") { globalStreamHandler.handle(this) } diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt index ce5f0fb6..4f1dcd66 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt @@ -404,6 +404,30 @@ fun main() { runPhase2 = { sid, graph, cfg -> orchestrator.run(sid, graph, cfg) }, requestPlanApproval = { sid, planJson -> orchestrator.requestPlanApproval(sid, planJson) }, ) + // observability-spec §4: continuous health watch. Seed the monitor's last-status from the + // recorded system-session events so a restart doesn't re-emit a degraded already in the log. + val healthMonitor = correxConfig.health.let { hc -> + if (!hc.enabled) { + null + } else { + val dataDir = Paths.get(System.getProperty("user.home"), ".config", "correx") + val seeded = DefaultEventReplayer(eventStore, com.correx.apps.server.health.HealthProjection()) + .rebuild(com.correx.apps.server.health.SYSTEM_SESSION) + .subjects.mapValues { it.value.status } + com.correx.apps.server.health.HealthMonitor( + eventStore = eventStore, + probes = listOf( + com.correx.apps.server.health.DiskWatermarkProbe( + monitoredPaths = listOf(dataDir), + warnBytes = hc.diskWarnBytes, + warnGrowthBytesPerMin = hc.diskWarnGrowthBytesPerMin, + ), + ), + intervalMs = hc.intervalMs, + initialStatuses = seeded, + ) + } + } val module = ServerModule( orchestrator = orchestrator, eventStore = eventStore, @@ -426,6 +450,7 @@ fun main() { freestyleDriver = freestyleDriver, operatorProfile = operatorProfile, profileAdaptationService = profileAdaptationService, + healthMonitor = healthMonitor, ) // Wire live config editing: persist to TOML, swap the holder, and rebuild config-derived // services. Built after the module so the rebuild hook can swap them in place. diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/ServerModule.kt b/apps/server/src/main/kotlin/com/correx/apps/server/ServerModule.kt index ccbf756e..c5ea422f 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/ServerModule.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/ServerModule.kt @@ -99,6 +99,9 @@ class ServerModule( // Propose learned profile updates at session end. Null when learn=false or no profile. // Rebuilt by ConfigService when personalization.enabled/learn is toggled live. profileAdaptationService: com.correx.apps.server.personalization.ProfileAdaptationService? = null, + // Continuous health watch (observability-spec §4). Null disables it (tests / health.enabled=false); + // records degraded/restored edge events on the system session, started in [start]. + private val healthMonitor: com.correx.apps.server.health.HealthMonitor? = null, ) { val approvalCoordinator: ApprovalCoordinator = approvalCoordinator ?: ApprovalCoordinator( orchestrator = orchestrator, @@ -180,6 +183,10 @@ class ServerModule( // Live-only: subscribeAll() replays nothing and ServerModule is never built under // ReplayOrchestrator, so narration never re-fires on restart/replay (invariant #8). NarrationSubscriber(eventStore = eventStore, routerFacade = routerFacade, scope = moduleScope, maxPerRun = narrationMaxPerRun).start() + + // Continuous, edge-triggered health watch (observability-spec §4). Live-only like narration: + // probes read the environment and record degraded/restored events; replay reads those facts. + healthMonitor?.start(moduleScope) } /** @@ -410,9 +417,12 @@ class ServerModule( private val sessionSummaryProjector = SessionSummaryProjector() fun listSessionSummaries(): List = - eventStore.allSessionIds().map { sessionId -> - sessionSummaryProjector.project(sessionId, eventStore.readFrom(sessionId, fromSequence = 0L)) - } + eventStore.allSessionIds() + // The system session carries global health events, not a user workflow — hide it. + .filter { it != com.correx.apps.server.health.SYSTEM_SESSION } + .map { sessionId -> + sessionSummaryProjector.project(sessionId, eventStore.readFrom(sessionId, fromSequence = 0L)) + } private fun preRegisterPendingApprovals() { val projector = ApprovalProjector(DefaultApprovalReducer()) diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/DiskWatermarkProbe.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/DiskWatermarkProbe.kt new file mode 100644 index 00000000..63a9d8fc --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/DiskWatermarkProbe.kt @@ -0,0 +1,104 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.HealthSubject +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.withContext +import kotlinx.datetime.Clock +import kotlinx.datetime.Instant +import java.nio.file.Files +import java.nio.file.Path + +private const val MILLIS_PER_MINUTE = 60_000.0 +private const val BYTES_PER_KIB = 1024.0 + +/** + * Watches the on-disk footprint of the event log + CAS (observability-spec §4: "disk watermark"). + * Append-only storage grows monotonically, so two thresholds matter: absolute size (the warning + * that arrives before the 2am surprise) and growth rate (when log compaction stops being + * theoretical). Size takes priority over growth when both cross. + * + * The directory walk is the environment observation; it runs only live in the monitor loop. The + * prior-sample state used to compute the growth rate is per-run and never affects replay. + */ +class DiskWatermarkProbe( + private val monitoredPaths: List, + private val warnBytes: Long, + private val warnGrowthBytesPerMin: Long, + private val clock: Clock = Clock.System, +) : HealthProbe { + + override val subject: HealthSubject = HealthSubject.DISK + + private data class Sample(val bytes: Long, val at: Instant) + + private var lastSample: Sample? = null + + override suspend fun observe(): HealthObservation = withContext(Dispatchers.IO) { + val now = clock.now() + val totalBytes = monitoredPaths.sumOf { sizeOf(it) } + val growthPerMin = growthBytesPerMin(totalBytes, now) + lastSample = Sample(totalBytes, now) + + when { + totalBytes >= warnBytes -> degraded( + metric = "total_bytes", + value = totalBytes, + threshold = warnBytes, + detail = "event log + CAS at ${humanBytes(totalBytes)} (warn ≥ ${humanBytes(warnBytes)})", + ) + growthPerMin >= warnGrowthBytesPerMin -> degraded( + metric = "growth_bytes_per_min", + value = growthPerMin, + threshold = warnGrowthBytesPerMin, + detail = "storage growing ${humanBytes(growthPerMin)}/min " + + "(warn ≥ ${humanBytes(warnGrowthBytesPerMin)}/min)", + ) + else -> HealthObservation( + subject = subject, + status = HealthStatus.HEALTHY, + metric = "total_bytes", + observedValue = totalBytes, + threshold = warnBytes, + detail = "event log + CAS at ${humanBytes(totalBytes)}", + ) + } + } + + private fun degraded(metric: String, value: Long, threshold: Long, detail: String) = + HealthObservation(subject, HealthStatus.DEGRADED, metric, value, threshold, detail) + + /** Bytes-per-minute since the previous sample; 0 on the first tick (no baseline yet). */ + private fun growthBytesPerMin(totalBytes: Long, now: Instant): Long { + val prev = lastSample ?: return 0L + val elapsedMin = (now - prev.at).inWholeMilliseconds.toDouble() / MILLIS_PER_MINUTE + return if (elapsedMin <= 0.0) { + 0L + } else { + ((totalBytes - prev.bytes).toDouble() / elapsedMin).toLong().coerceAtLeast(0L) + } + } + + private fun sizeOf(path: Path): Long { + if (!Files.exists(path)) return 0L + return runCatching { + Files.walk(path).use { stream -> + stream.filter { Files.isRegularFile(it) } + .mapToLong { runCatching { Files.size(it) }.getOrDefault(0L) } + .sum() + } + }.getOrDefault(0L) + } +} + +private val UNITS = listOf("B", "KiB", "MiB", "GiB", "TiB") + +private fun humanBytes(bytes: Long): String { + if (bytes < BYTES_PER_KIB) return "$bytes B" + var value = bytes.toDouble() + var unit = 0 + while (value >= BYTES_PER_KIB && unit < UNITS.size - 1) { + value /= BYTES_PER_KIB + unit++ + } + return "%.1f %s".format(value, UNITS[unit]) +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthInspectionService.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthInspectionService.kt new file mode 100644 index 00000000..cad33f47 --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthInspectionService.kt @@ -0,0 +1,63 @@ +package com.correx.apps.server.health + +import com.correx.core.events.stores.EventStore +import com.correx.core.sessions.projections.replay.DefaultEventReplayer +import kotlinx.datetime.Clock +import kotlinx.serialization.Serializable + +@Serializable +data class HealthSubjectReport( + val subject: String, + val status: String, + val metric: String, + val observedValue: Long, + val detail: String, + val since: String, +) + +/** + * Render-friendly snapshot of system health. [overall] is DEGRADED if any subject is degraded, + * else HEALTHY (including the "nothing recorded yet" case). [checkedAt] is when this report was + * rendered, distinct from each subject's [HealthSubjectReport.since] edge timestamp. + */ +@Serializable +data class HealthReport( + val overall: String, + val subjects: List, + val checkedAt: String, +) + +/** + * Rebuilds [HealthState] by replaying the system session through [HealthProjection], then renders + * a [HealthReport]. Replay-only: no probing, no network (Hard Invariant #8). Backs `GET /health/checks` + * and the `correx health` CLI, the same way [com.correx.apps.server.metrics.MetricsInspectionService] + * backs `correx stats`. + */ +class HealthInspectionService( + eventStore: EventStore, + private val clock: Clock = Clock.System, +) { + private val replayer = DefaultEventReplayer(eventStore, HealthProjection()) + + fun inspect(): HealthReport { + val state = replayer.rebuild(SYSTEM_SESSION) + val subjects = state.subjects.entries + .sortedBy { it.key.name } + .map { (subject, health) -> + HealthSubjectReport( + subject = subject.name, + status = health.status.name, + metric = health.metric, + observedValue = health.observedValue, + detail = health.detail, + since = health.since.toString(), + ) + } + val overall = if (subjects.any { it.status == HealthStatus.DEGRADED.name }) { + HealthStatus.DEGRADED.name + } else { + HealthStatus.HEALTHY.name + } + return HealthReport(overall = overall, subjects = subjects, checkedAt = clock.now().toString()) + } +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthMonitor.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthMonitor.kt new file mode 100644 index 00000000..50616dfd --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthMonitor.kt @@ -0,0 +1,103 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.EventMetadata +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent +import com.correx.core.events.events.HealthSubject +import com.correx.core.events.events.NewEvent +import com.correx.core.events.stores.EventStore +import com.correx.core.events.types.EventId +import com.correx.core.events.types.SessionId +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import kotlinx.datetime.Clock +import org.slf4j.LoggerFactory +import java.util.UUID + +/** + * The reserved session that carries system-global facts (health checks) that belong to no user + * workflow. Filtered out of operator-facing session listings; it exists only so global events + * have a home in the per-session event log and flow on the global stream. + */ +val SYSTEM_SESSION: SessionId = SessionId("__system__") + +/** + * Small, paranoid, continuous health watch (observability-spec §4). Each tick polls every + * [HealthProbe] and, comparing against the last recorded status per subject, appends a + * [HealthDegradedEvent] / [HealthRestoredEvent] only on a transition (hysteresis — a probe that + * keeps reporting the same status writes nothing). Health is thus event-sourced: observable, + * replayable (Invariant #8/#9, the probe runs only live), and renderable via [HealthProjection]. + * + * [initialStatuses] seeds the last-status map from the projection at start so a server restart + * does not re-emit a Degraded that the log already records. + */ +class HealthMonitor( + private val eventStore: EventStore, + private val probes: List, + private val intervalMs: Long, + private val clock: Clock = Clock.System, + private val sessionId: SessionId = SYSTEM_SESSION, + initialStatuses: Map = emptyMap(), +) { + private val log = LoggerFactory.getLogger(HealthMonitor::class.java) + private val lastStatus: MutableMap = HashMap(initialStatuses) + + /** Launches the polling loop on [scope]; cancelling the scope stops the monitor. */ + fun start(scope: CoroutineScope): Job = scope.launch { + log.info("HealthMonitor: starting ({} probe(s), every {} ms)", probes.size, intervalMs) + while (isActive) { + runCatching { tick() }.onFailure { log.warn("HealthMonitor tick failed: {}", it.message) } + delay(intervalMs) + } + } + + /** One poll of every probe. Visible for deterministic testing without the loop/delay. */ + suspend fun tick() { + for (probe in probes) { + val observation = runCatching { probe.observe() } + .getOrElse { log.warn("Health probe {} failed: {}", probe.subject, it.message); null } + ?: continue + val previous = lastStatus[probe.subject] ?: HealthStatus.HEALTHY + if (observation.status == previous) continue + append(observation) + lastStatus[probe.subject] = observation.status + } + } + + private suspend fun append(observation: HealthObservation) { + val payload = when (observation.status) { + HealthStatus.DEGRADED -> HealthDegradedEvent( + sessionId = sessionId, + subject = observation.subject, + metric = observation.metric, + observedValue = observation.observedValue, + threshold = observation.threshold, + detail = observation.detail, + ) + HealthStatus.HEALTHY -> HealthRestoredEvent( + sessionId = sessionId, + subject = observation.subject, + metric = observation.metric, + observedValue = observation.observedValue, + detail = observation.detail, + ) + } + log.info("Health {} -> {}: {}", observation.subject, observation.status, observation.detail) + eventStore.append( + NewEvent( + metadata = EventMetadata( + eventId = EventId(UUID.randomUUID().toString()), + sessionId = sessionId, + timestamp = clock.now(), + schemaVersion = 1, + causationId = null, + correlationId = null, + ), + payload = payload, + ), + ) + } +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProbe.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProbe.kt new file mode 100644 index 00000000..f2d90636 --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProbe.kt @@ -0,0 +1,31 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.HealthSubject + +/** + * A single point-in-time reading of one [subject]'s environment state. The [HealthMonitor] + * compares [status] against the last recorded status and emits an edge event only on change, + * so a probe is free to report the same status every tick without polluting the log. + * + * [metric]/[observedValue]/[threshold] describe the dimension that drove the [status] decision + * (e.g. `total_bytes` = 5.2 GB vs a 5.0 GB threshold), and are carried verbatim into the event. + */ +data class HealthObservation( + val subject: HealthSubject, + val status: HealthStatus, + val metric: String, + val observedValue: Long, + val threshold: Long, + val detail: String, +) + +/** + * Reads one subject's nondeterministic environment state (filesystem, process, store). Probes + * run only live, inside the monitor loop — never during replay; the recorded edge events are the + * sole replay input (Invariant #8/#9). Implementations may hold per-run state (e.g. the prior + * sample for a growth rate); that state never affects replay. + */ +interface HealthProbe { + val subject: HealthSubject + suspend fun observe(): HealthObservation +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProjection.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProjection.kt new file mode 100644 index 00000000..348a9d9f --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthProjection.kt @@ -0,0 +1,45 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent +import com.correx.core.events.events.StoredEvent +import com.correx.core.sessions.projections.Projection + +/** + * Pure fold deriving [HealthState] from the health events recorded on the system session + * (observability-spec §4: "health states are events ... renderable like everything else"). + * + * Records no new events and reads only the edge events the [HealthMonitor] already emits; + * a disposable read-model (Hard Invariant #2). Because health is event-sourced, the monitor + * can seed its in-memory last-status from this projection across a restart, and replay + * reconstructs the exact health timeline without re-probing the environment (Invariant #8/#9). + */ +class HealthProjection : Projection { + + override fun initial(): HealthState = HealthState() + + override fun apply(state: HealthState, event: StoredEvent): HealthState = + when (val payload = event.payload) { + is HealthDegradedEvent -> state.with( + payload.subject, + SubjectHealth( + status = HealthStatus.DEGRADED, + metric = payload.metric, + observedValue = payload.observedValue, + detail = payload.detail, + since = event.metadata.timestamp, + ), + ) + is HealthRestoredEvent -> state.with( + payload.subject, + SubjectHealth( + status = HealthStatus.HEALTHY, + metric = payload.metric, + observedValue = payload.observedValue, + detail = payload.detail, + since = event.metadata.timestamp, + ), + ) + else -> state + } +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthState.kt b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthState.kt new file mode 100644 index 00000000..c485d950 --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/health/HealthState.kt @@ -0,0 +1,32 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.HealthSubject +import kotlinx.datetime.Instant +import kotlinx.serialization.Serializable + +/** Current health of a monitored subject (observability-spec §4). */ +enum class HealthStatus { HEALTHY, DEGRADED } + +/** + * Read-model accumulated by [HealthProjection] — a pure fold over the system session's + * health events. Holds the latest known status per subject; rebuilt from events on demand, + * never persisted (Hard Invariant #1). A subject absent from [subjects] has never reported, + * which the render layer treats as healthy ("no degradation recorded"). + */ +@Serializable +data class HealthState( + val subjects: Map = emptyMap(), +) { + fun with(subject: HealthSubject, health: SubjectHealth): HealthState = + copy(subjects = subjects + (subject to health)) +} + +@Serializable +data class SubjectHealth( + val status: HealthStatus, + val metric: String, + val observedValue: Long, + val detail: String, + /** Timestamp of the edge event that put the subject into [status]. */ + val since: Instant, +) diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/health/DiskWatermarkProbeTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/health/DiskWatermarkProbeTest.kt new file mode 100644 index 00000000..1f7a59b9 --- /dev/null +++ b/apps/server/src/test/kotlin/com/correx/apps/server/health/DiskWatermarkProbeTest.kt @@ -0,0 +1,87 @@ +package com.correx.apps.server.health + +import kotlinx.coroutines.runBlocking +import kotlinx.datetime.Clock +import kotlinx.datetime.Instant +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.io.TempDir +import org.junit.jupiter.api.Test +import java.nio.file.Files +import java.nio.file.Path + +class DiskWatermarkProbeTest { + + /** Hand-advanced clock so growth-rate math is deterministic. */ + private class FakeClock(var current: Instant) : Clock { + override fun now(): Instant = current + } + + private fun writeFile(dir: Path, name: String, bytes: Int) { + Files.write(dir.resolve(name), ByteArray(bytes)) + } + + @Test + fun `healthy when under both thresholds`(@TempDir dir: Path): Unit = runBlocking { + writeFile(dir, "a.bin", 100) + val probe = DiskWatermarkProbe(listOf(dir), warnBytes = 10_000, warnGrowthBytesPerMin = 10_000) + val obs = probe.observe() + assertEquals(HealthStatus.HEALTHY, obs.status) + assertEquals(100, obs.observedValue) + } + + @Test + fun `degraded when total size crosses the watermark`(@TempDir dir: Path): Unit = runBlocking { + writeFile(dir, "big.bin", 2_000) + val probe = DiskWatermarkProbe(listOf(dir), warnBytes = 1_000, warnGrowthBytesPerMin = Long.MAX_VALUE) + val obs = probe.observe() + assertEquals(HealthStatus.DEGRADED, obs.status) + assertEquals("total_bytes", obs.metric) + assertEquals(2_000, obs.observedValue) + assertEquals(1_000, obs.threshold) + } + + @Test + fun `sums nested files recursively`(@TempDir dir: Path): Unit = runBlocking { + val sub = Files.createDirectories(dir.resolve("artifacts/seg")) + writeFile(dir, "log.db", 500) + writeFile(sub, "seg-0.dat", 700) + val probe = DiskWatermarkProbe(listOf(dir), warnBytes = Long.MAX_VALUE, warnGrowthBytesPerMin = Long.MAX_VALUE) + assertEquals(1_200, probe.observe().observedValue) + } + + @Test + fun `first tick reports no growth then degrades on a fast-growing second tick`(@TempDir dir: Path): Unit = + runBlocking { + val clock = FakeClock(Instant.parse("2026-06-13T00:00:00Z")) + writeFile(dir, "a.bin", 1_000) + val probe = DiskWatermarkProbe( + monitoredPaths = listOf(dir), + warnBytes = Long.MAX_VALUE, // isolate the growth dimension + warnGrowthBytesPerMin = 5_000, + clock = clock, + ) + // First tick: no baseline -> growth treated as 0 -> healthy. + assertEquals(HealthStatus.HEALTHY, probe.observe().status) + + // One minute later, +10_000 bytes => 10_000 B/min, over the 5_000 threshold. + clock.current = clock.current.plus(kotlin.time.Duration.parse("60s")) + writeFile(dir, "b.bin", 10_000) + val obs = probe.observe() + assertEquals(HealthStatus.DEGRADED, obs.status) + assertEquals("growth_bytes_per_min", obs.metric) + assertTrue(obs.observedValue >= 5_000, "growth ${obs.observedValue} should be >= threshold") + } + + @Test + fun `missing path contributes zero`(@TempDir dir: Path): Unit = runBlocking { + val probe = DiskWatermarkProbe( + listOf(dir.resolve("does-not-exist")), + warnBytes = 1, + warnGrowthBytesPerMin = Long.MAX_VALUE, + ) + val obs = probe.observe() + assertEquals(0, obs.observedValue) + assertEquals(HealthStatus.HEALTHY, obs.status) + } +} diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthMonitorTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthMonitorTest.kt new file mode 100644 index 00000000..deb2c95b --- /dev/null +++ b/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthMonitorTest.kt @@ -0,0 +1,96 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent +import com.correx.core.events.events.HealthSubject +import com.correx.infrastructure.persistence.InMemoryEventStore +import kotlinx.coroutines.runBlocking +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class HealthMonitorTest { + + /** Probe that returns a scripted sequence of statuses, one per tick (last value repeats). */ + private class ScriptedProbe( + override val subject: HealthSubject, + private val statuses: List, + ) : HealthProbe { + private var index = 0 + override suspend fun observe(): HealthObservation { + val status = statuses[index.coerceAtMost(statuses.size - 1)] + index++ + return HealthObservation(subject, status, "total_bytes", index.toLong(), 5_000, "tick $index") + } + } + + private fun monitor( + store: InMemoryEventStore, + probe: HealthProbe, + initial: Map = emptyMap(), + ) = HealthMonitor( + eventStore = store, + probes = listOf(probe), + intervalMs = 1_000, + initialStatuses = initial, + ) + + @Test + fun `emits only on status transitions`(): Unit = runBlocking { + val store = InMemoryEventStore() + val probe = ScriptedProbe( + HealthSubject.DISK, + listOf( + HealthStatus.HEALTHY, // tick 1: no edge (default healthy) + HealthStatus.DEGRADED, // tick 2: edge -> degraded + HealthStatus.DEGRADED, // tick 3: no edge + HealthStatus.HEALTHY, // tick 4: edge -> restored + ), + ) + val m = monitor(store, probe) + repeat(4) { m.tick() } + + val events = store.read(SYSTEM_SESSION).map { it.payload } + assertEquals(2, events.size, "exactly two edges expected, got $events") + assertTrue(events[0] is HealthDegradedEvent) + assertTrue(events[1] is HealthRestoredEvent) + } + + @Test + fun `seeded degraded status suppresses duplicate degraded on restart`(): Unit = runBlocking { + val store = InMemoryEventStore() + val probe = ScriptedProbe(HealthSubject.DISK, listOf(HealthStatus.DEGRADED)) + // Monitor restarts with last-known status already DEGRADED (rebuilt from the log). + val m = monitor(store, probe, initial = mapOf(HealthSubject.DISK to HealthStatus.DEGRADED)) + m.tick() + m.tick() + assertTrue(store.read(SYSTEM_SESSION).isEmpty(), "no edge — already degraded before restart") + } + + @Test + fun `degraded event carries probe metric and threshold`(): Unit = runBlocking { + val store = InMemoryEventStore() + val probe = ScriptedProbe(HealthSubject.DISK, listOf(HealthStatus.DEGRADED)) + monitor(store, probe).tick() + + val event = store.read(SYSTEM_SESSION).single().payload as HealthDegradedEvent + assertEquals(HealthSubject.DISK, event.subject) + assertEquals("total_bytes", event.metric) + assertEquals(5_000, event.threshold) + } + + @Test + fun `a failing probe does not abort the tick or emit`(): Unit = runBlocking { + val store = InMemoryEventStore() + val throwing = object : HealthProbe { + override val subject = HealthSubject.EVENT_STORE + override suspend fun observe(): HealthObservation = error("probe boom") + } + val ok = ScriptedProbe(HealthSubject.DISK, listOf(HealthStatus.DEGRADED)) + HealthMonitor(store, probes = listOf(throwing, ok), intervalMs = 1_000).tick() + + val events = store.read(SYSTEM_SESSION).map { it.payload } + assertEquals(1, events.size) + assertTrue(events.single() is HealthDegradedEvent) + } +} diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthProjectionTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthProjectionTest.kt new file mode 100644 index 00000000..b9192622 --- /dev/null +++ b/apps/server/src/test/kotlin/com/correx/apps/server/health/HealthProjectionTest.kt @@ -0,0 +1,87 @@ +package com.correx.apps.server.health + +import com.correx.core.events.events.EventMetadata +import com.correx.core.events.events.EventPayload +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent +import com.correx.core.events.events.HealthSubject +import com.correx.core.events.events.StoredEvent +import com.correx.core.events.types.EventId +import com.correx.core.utils.TypeId +import kotlinx.datetime.Instant +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertNull +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class HealthProjectionTest { + + private val t0 = Instant.parse("2026-06-13T00:00:00Z") + private fun at(seconds: Long): Instant = t0.plus(kotlin.time.Duration.parse("${seconds}s")) + + private fun stored(payload: EventPayload, seq: Long, at: Instant): StoredEvent = + StoredEvent( + metadata = EventMetadata( + eventId = EventId("e$seq"), + sessionId = SYSTEM_SESSION, + timestamp = at, + schemaVersion = 1, + causationId = null, + correlationId = null, + ), + sequence = seq, + sessionSequence = seq, + payload = payload, + ) + + private fun degraded(metric: String, value: Long, threshold: Long) = + HealthDegradedEvent(SYSTEM_SESSION, HealthSubject.DISK, metric, value, threshold, "d") + + private fun restored(metric: String, value: Long) = + HealthRestoredEvent(SYSTEM_SESSION, HealthSubject.DISK, metric, value, "r") + + private fun fold(events: List): HealthState { + val projection = HealthProjection() + return events.fold(projection.initial()) { state, e -> projection.apply(state, e) } + } + + @Test + fun `empty log yields no subjects`() { + assertTrue(fold(emptyList()).subjects.isEmpty()) + } + + @Test + fun `degraded then restored leaves subject healthy with latest observation`() { + val state = fold( + listOf( + stored(degraded("total_bytes", 6_000, 5_000), 1, at(0)), + stored(restored("total_bytes", 1_000), 2, at(60)), + ), + ) + val disk = state.subjects.getValue(HealthSubject.DISK) + assertEquals(HealthStatus.HEALTHY, disk.status) + assertEquals(1_000, disk.observedValue) + assertEquals(at(60), disk.since) + } + + @Test + fun `last edge wins`() { + val state = fold( + listOf( + stored(restored("total_bytes", 1_000), 1, at(0)), + stored(degraded("growth_bytes_per_min", 80, 50), 2, at(30)), + ), + ) + val disk = state.subjects.getValue(HealthSubject.DISK) + assertEquals(HealthStatus.DEGRADED, disk.status) + assertEquals("growth_bytes_per_min", disk.metric) + assertEquals(at(30), disk.since) + } + + @Test + fun `unrelated events are ignored`() { + val before = fold(listOf(stored(degraded("total_bytes", 6_000, 5_000), 1, at(0)))) + assertNull(before.subjects[HealthSubject.LLAMA_SERVER]) + assertEquals(1, before.subjects.size) + } +} diff --git a/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt b/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt index 6eeb01d3..57b787d4 100644 --- a/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt +++ b/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt @@ -16,6 +16,21 @@ data class CorrexConfig( val project: ProjectConfig = ProjectConfig(), val personalization: PersonalizationConfig = PersonalizationConfig(), val orchestration: OrchestrationKnobs = OrchestrationKnobs(), + val health: HealthConfig = HealthConfig(), +) + +/** + * Continuous health watch (observability-spec §4). When [enabled], a background monitor polls + * the on-disk footprint every [intervalMs] and records a degraded/restored event on the edge. + * [diskWarnBytes] is the absolute event-log + CAS size that trips a warning; [diskWarnGrowthBytesPerMin] + * is the growth rate that warns log compaction is becoming non-theoretical. Defaults suit a homelab. + */ +@Serializable +data class HealthConfig( + val enabled: Boolean = true, + val intervalMs: Long = 30_000, + val diskWarnBytes: Long = 5_000_000_000, + val diskWarnGrowthBytesPerMin: Long = 50_000_000, ) /** diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/HealthEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/HealthEvents.kt new file mode 100644 index 00000000..4f2ee132 --- /dev/null +++ b/core/events/src/main/kotlin/com/correx/core/events/events/HealthEvents.kt @@ -0,0 +1,56 @@ +package com.correx.core.events.events + +import com.correx.core.events.types.SessionId +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable + +/** + * The organs the health monitor watches (observability-spec §4). Each is small, paranoid, and + * continuous; a state change for any one is recorded as a [HealthDegradedEvent] / [HealthRestoredEvent] + * so health is observable, replayable, and renderable like every other fact in the log. + */ +@Serializable +enum class HealthSubject { + /** Local inference process: liveness + tokens/sec degradation (KV-cache pressure, thermal throttle). */ + LLAMA_SERVER, + + /** The append-only event log: append latency + fsync health. The organ that must never lie. */ + EVENT_STORE, + + /** CAS + log size and growth rate. Append-only means monotonic growth; warn before the 2am surprise. */ + DISK, +} + +/** + * A monitored [subject] crossed from healthy into degraded. Emitted only on the edge (the monitor + * holds the last status per subject), so the log records transitions, not every poll. The + * environment observation that triggered it is captured here at observe-time (Hard Invariant #9): + * replay reads this recorded fact and never re-probes the filesystem/process/store. + * + * [metric] names the dimension that crossed (e.g. `total_bytes`, `growth_bytes_per_min`); + * [observedValue] is its value at the edge; [threshold] is the watermark it crossed. + */ +@Serializable +@SerialName("HealthDegraded") +data class HealthDegradedEvent( + val sessionId: SessionId, + val subject: HealthSubject, + val metric: String, + val observedValue: Long, + val threshold: Long, + val detail: String, +) : EventPayload + +/** + * A monitored [subject] recovered from degraded back to healthy. Edge-triggered counterpart to + * [HealthDegradedEvent]; [observedValue] is the [metric] value that brought it back under the line. + */ +@Serializable +@SerialName("HealthRestored") +data class HealthRestoredEvent( + val sessionId: SessionId, + val subject: HealthSubject, + val metric: String, + val observedValue: Long, + val detail: String, +) : EventPayload diff --git a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt index 575dcb7d..6a0e4353 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt @@ -19,6 +19,8 @@ import com.correx.core.events.events.ContextTruncatedEvent import com.correx.core.events.events.EventPayload import com.correx.core.events.events.ExecutionPlanLockedEvent import com.correx.core.events.events.ExecutionPlanRejectedEvent +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent import com.correx.core.events.events.JournalCompactedEvent import com.correx.core.events.events.FileWrittenEvent import com.correx.core.events.events.L3MemoryRetrievedEvent @@ -113,6 +115,8 @@ val eventModule = SerializersModule { subclass(ExecutionPlanLockedEvent::class) subclass(ExecutionPlanRejectedEvent::class) subclass(JournalCompactedEvent::class) + subclass(HealthDegradedEvent::class) + subclass(HealthRestoredEvent::class) } } diff --git a/core/events/src/test/kotlin/com/correx/core/events/serialization/HealthEventsSerializationTest.kt b/core/events/src/test/kotlin/com/correx/core/events/serialization/HealthEventsSerializationTest.kt new file mode 100644 index 00000000..8d14aa55 --- /dev/null +++ b/core/events/src/test/kotlin/com/correx/core/events/serialization/HealthEventsSerializationTest.kt @@ -0,0 +1,62 @@ +package com.correx.core.events.serialization + +import com.correx.core.events.events.EventPayload +import com.correx.core.events.events.HealthDegradedEvent +import com.correx.core.events.events.HealthRestoredEvent +import com.correx.core.events.events.HealthSubject +import com.correx.core.events.types.SessionId +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class HealthEventsSerializationTest { + + @Test + fun `HealthDegradedEvent round-trips as polymorphic EventPayload`() { + val sample: EventPayload = HealthDegradedEvent( + sessionId = SessionId("__system__"), + subject = HealthSubject.DISK, + metric = "total_bytes", + observedValue = 6_000_000_000, + threshold = 5_000_000_000, + detail = "event log + CAS at 5.6 GiB (warn ≥ 4.7 GiB)", + ) + val encoded = eventJson.encodeToString(EventPayload.serializer(), sample) + assertTrue(encoded.contains("\"type\":\"HealthDegraded\""), "SerialName must be present: $encoded") + assertEquals(sample, eventJson.decodeFromString(EventPayload.serializer(), encoded)) + } + + @Test + fun `HealthRestoredEvent round-trips as polymorphic EventPayload`() { + val sample: EventPayload = HealthRestoredEvent( + sessionId = SessionId("__system__"), + subject = HealthSubject.DISK, + metric = "total_bytes", + observedValue = 1_000_000_000, + detail = "event log + CAS at 953.7 MiB", + ) + val encoded = eventJson.encodeToString(EventPayload.serializer(), sample) + assertTrue(encoded.contains("\"type\":\"HealthRestored\""), "SerialName must be present: $encoded") + assertEquals(sample, eventJson.decodeFromString(EventPayload.serializer(), encoded)) + } + + @Test + fun `every HealthSubject round-trips`() { + for (subject in HealthSubject.entries) { + val sample: EventPayload = HealthRestoredEvent( + sessionId = SessionId("__system__"), + subject = subject, + metric = "m", + observedValue = 0, + detail = "ok", + ) + assertEquals( + sample, + eventJson.decodeFromString( + EventPayload.serializer(), + eventJson.encodeToString(EventPayload.serializer(), sample), + ), + ) + } + } +}