From bbff73108e51fec18775f09a43d965b5e69afe3c Mon Sep 17 00:00:00 2001 From: kami Date: Sun, 17 May 2026 14:56:36 +0400 Subject: [PATCH] feat(server): wire log4j2 + log every event emission Add log4j2/SLF4J deps and log4j2.xml, install Ktor CallLogging, replace manual System.err.println sites with leveled loggers, and wrap the EventStore with a LoggingEventStore decorator so every append/appendAll is logged in one place. --- apps/server/build.gradle | 5 ++ .../com/correx/apps/server/Application.kt | 4 ++ .../kotlin/com/correx/apps/server/Main.kt | 3 +- .../apps/server/logging/LoggingEventStore.kt | 61 +++++++++++++++++++ .../registry/FileSystemWorkflowRegistry.kt | 5 +- .../apps/server/routes/SessionRoutes.kt | 5 +- .../apps/server/ws/GlobalStreamHandler.kt | 20 +++--- apps/server/src/main/resources/log4j2.xml | 25 ++++++++ 8 files changed, 116 insertions(+), 12 deletions(-) create mode 100644 apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt create mode 100644 apps/server/src/main/resources/log4j2.xml diff --git a/apps/server/build.gradle b/apps/server/build.gradle index 0164233b..b5ef6485 100644 --- a/apps/server/build.gradle +++ b/apps/server/build.gradle @@ -36,4 +36,9 @@ dependencies { implementation "io.ktor:ktor-server-content-negotiation:$ktor_version" implementation "io.ktor:ktor-serialization-kotlinx-json:$ktor_version" implementation "io.ktor:ktor-server-status-pages:$ktor_version" + implementation "io.ktor:ktor-server-call-logging:$ktor_version" + + implementation "org.apache.logging.log4j:log4j-core:2.24.1" + implementation "org.apache.logging.log4j:log4j-slf4j2-impl:2.24.1" + implementation "org.slf4j:slf4j-api:2.0.16" } 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 db9b08bd..df5991b4 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 @@ -8,6 +8,8 @@ import io.ktor.http.HttpStatusCode import io.ktor.serialization.kotlinx.json.json import io.ktor.server.application.Application import io.ktor.server.application.install +import io.ktor.server.plugins.calllogging.CallLogging +import org.slf4j.event.Level import io.ktor.server.plugins.contentnegotiation.ContentNegotiation import io.ktor.server.plugins.statuspages.StatusPages import io.ktor.server.response.respond @@ -18,6 +20,8 @@ import io.ktor.server.websocket.webSocket import kotlinx.serialization.json.Json fun Application.configureServer(module: ServerModule) { + install(CallLogging) { level = Level.INFO } + install(WebSockets) install(ContentNegotiation) { 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 7bd14349..f4fdd183 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 @@ -23,6 +23,7 @@ import com.correx.core.sessions.projections.replay.DefaultEventReplayer import com.correx.core.transitions.evaluation.PromptResolver import com.correx.core.transitions.resolution.DefaultTransitionResolver import com.correx.core.validation.pipeline.ValidationPipeline +import com.correx.apps.server.logging.LoggingEventStore import com.correx.infrastructure.InfrastructureModule import com.correx.infrastructure.inference.DefaultProviderRegistry import com.correx.infrastructure.inference.FirstAvailableRoutingStrategy @@ -30,7 +31,7 @@ import io.ktor.server.engine.embeddedServer import io.ktor.server.netty.Netty fun main() { - val eventStore = InfrastructureModule.createEventStore() + val eventStore = LoggingEventStore(InfrastructureModule.createEventStore()) val llamaProvider = InfrastructureModule.createLlamaCppProvider( modelId = System.getenv("CORREX_MODEL_ID") ?: "default", diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt b/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt new file mode 100644 index 00000000..8b46a940 --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/logging/LoggingEventStore.kt @@ -0,0 +1,61 @@ +package com.correx.apps.server.logging + +import com.correx.core.events.events.NewEvent +import com.correx.core.events.events.StoredEvent +import com.correx.core.events.stores.EventStore +import com.correx.core.events.types.SessionId +import kotlinx.coroutines.flow.Flow +import org.slf4j.LoggerFactory + +private val log = LoggerFactory.getLogger(LoggingEventStore::class.java) + +class LoggingEventStore(private val delegate: EventStore) : EventStore { + + override fun append(event: NewEvent): StoredEvent = + runCatching { delegate.append(event) } + .onSuccess { stored -> + log.info( + "event emitted session={} type={} id={}", + event.metadata.sessionId.value, + event.payload::class.simpleName, + stored.metadata.eventId.value, + ) + } + .onFailure { ex -> + log.error( + "append failed session={} type={}", + event.metadata.sessionId.value, + event.payload::class.simpleName, + ex, + ) + } + .getOrThrow() + + override fun appendAll(events: List): List = + runCatching { delegate.appendAll(events) } + .onSuccess { stored -> + val breakdown = events.groupingBy { it.payload::class.simpleName }.eachCount() + log.info("events emitted count={} types={}", stored.size, breakdown) + stored.forEach { s -> + log.debug( + "event emitted session={} type={} id={}", + s.metadata.sessionId.value, + s.payload::class.simpleName, + s.metadata.eventId.value, + ) + } + } + .onFailure { ex -> + log.error("appendAll failed count={}", events.size, ex) + } + .getOrThrow() + + override fun read(sessionId: SessionId): List = delegate.read(sessionId) + + override fun readFrom(sessionId: SessionId, fromSequence: Long): List = + delegate.readFrom(sessionId, fromSequence) + + override fun lastSequence(sessionId: SessionId): Long? = delegate.lastSequence(sessionId) + + override fun subscribe(sessionId: SessionId): Flow = delegate.subscribe(sessionId) +} diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/registry/FileSystemWorkflowRegistry.kt b/apps/server/src/main/kotlin/com/correx/apps/server/registry/FileSystemWorkflowRegistry.kt index 6f953dcd..a47d9cc8 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/registry/FileSystemWorkflowRegistry.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/registry/FileSystemWorkflowRegistry.kt @@ -2,10 +2,13 @@ package com.correx.apps.server.registry import com.correx.core.transitions.graph.WorkflowGraph import com.correx.infrastructure.workflow.WorkflowLoader +import org.slf4j.LoggerFactory import java.nio.file.Path import kotlin.io.path.exists import kotlin.io.path.listDirectoryEntries +private val log = LoggerFactory.getLogger(FileSystemWorkflowRegistry::class.java) + class FileSystemWorkflowRegistry( private val loader: WorkflowLoader, private val workflowsDir: Path = Path.of( @@ -34,6 +37,6 @@ class FileSystemWorkflowRegistry( } private fun logWarning(file: Path, ex: Throwable) { - System.err.println("[WorkflowRegistry] Skipping $file: ${ex.message}") + log.warn("Skipping {}: {}", file, ex.message) } } diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt b/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt index af841d88..ad246d4e 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt @@ -13,6 +13,7 @@ import com.correx.core.utils.TypeId import io.ktor.http.HttpStatusCode import io.ktor.server.application.call import io.ktor.server.application.application +import org.slf4j.LoggerFactory import io.ktor.server.request.receive import io.ktor.server.response.respond import io.ktor.server.routing.Route @@ -40,6 +41,8 @@ data class EventResponse(val eventId: String, val sequence: Long, val sessionId: @Serializable data class StartSessionResponse(val sessionId: String) +private val log = LoggerFactory.getLogger("com.correx.apps.server.routes.SessionRoutes") + fun Route.sessionRoutes(module: ServerModule) { val streamHandler = SessionStreamHandler(module) route("/sessions") { @@ -67,7 +70,7 @@ private fun Route.startSessionRoute(module: ServerModule) { call.application.launch { runCatching { module.orchestrator.run(sessionId, graph, OrchestrationConfig()) } .onFailure { ex -> - System.err.println("[SessionRoutes] run failed for $sessionId: ${ex.message}") + log.error("run failed for session={}: {}", sessionId.value, ex.message, ex) module.eventStore.append( NewEvent( metadata = EventMetadata( diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt index 8219b51a..d5283add 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt @@ -21,14 +21,16 @@ import kotlinx.coroutines.delay import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.datetime.Clock +import org.slf4j.LoggerFactory import java.util.* private const val HEARTBEAT_INTERVAL_MS = 30_000L +private val log = LoggerFactory.getLogger(GlobalStreamHandler::class.java) class GlobalStreamHandler(private val module: ServerModule) { suspend fun handle(session: DefaultWebSocketServerSession) { - System.err.println("[GlobalStream] client connected") + log.info("client connected") sendInitialSnapshot(session) val heartbeatJob = session.launch { @@ -45,18 +47,18 @@ class GlobalStreamHandler(private val module: ServerModule) { val text = frame.readText() runCatching { ProtocolSerializer.decodeClientMessage(text) } .onSuccess { msg -> - System.err.println("[GlobalStream] recv: ${msg::class.simpleName}") + log.debug("recv: {}", msg::class.simpleName) handleClientMessage(session, msg) } .onFailure { - System.err.println("[GlobalStream] decode error: ${it.message}") + log.warn("decode error: {}", it.message) val error = ServerMessage.ProtocolError("Unknown message: ${it.message}") session.send(Frame.Text(ProtocolSerializer.encodeServerMessage(error))) } } } } catch (_: ClosedReceiveChannelException) { - System.err.println("[GlobalStream] client disconnected") + log.info("client disconnected") } finally { heartbeatJob.cancel() } @@ -64,16 +66,16 @@ class GlobalStreamHandler(private val module: ServerModule) { private suspend fun sendInitialSnapshot(session: DefaultWebSocketServerSession) { val providerHealth = runCatching { module.providerRegistry.healthCheckAll() } - .onFailure { System.err.println("[GlobalStream] healthCheckAll failed: ${it.message}") } + .onFailure { log.warn("healthCheckAll failed: {}", it.message) } .getOrDefault(emptyMap()) - System.err.println("[GlobalStream] snapshot: ${providerHealth.size} provider(s)") + log.info("snapshot: {} provider(s)", providerHealth.size) providerHealth.forEach { (providerId, health) -> val status = when (health) { is ProviderHealth.Healthy -> "healthy" is ProviderHealth.Degraded -> "degraded" is ProviderHealth.Unavailable -> "unavailable" } - System.err.println("[GlobalStream] provider=${providerId.value} status=$status") + log.debug("provider={} status={}", providerId.value, status) val msg = ServerMessage.ProviderStatusChanged( providerId = providerId.value, status = ProviderHealthDto(providerId.value, status, null), @@ -94,11 +96,11 @@ class GlobalStreamHandler(private val module: ServerModule) { return } val sessionId: SessionId = TypeId(UUID.randomUUID().toString()) - System.err.println("[GlobalStream] starting session=${sessionId.value} workflow=${msg.workflowId}") + log.info("starting session={} workflow={}", sessionId.value, msg.workflowId) session.launch { runCatching { module.orchestrator.run(sessionId, graph, OrchestrationConfig()) } .onFailure { ex -> - System.err.println("[GlobalStreamHandler] run failed for $sessionId: ${ex.message}") + log.error("run failed for session={}: {}", sessionId.value, ex.message, ex) module.eventStore.append( NewEvent( metadata = EventMetadata( diff --git a/apps/server/src/main/resources/log4j2.xml b/apps/server/src/main/resources/log4j2.xml new file mode 100644 index 00000000..e5f7d81c --- /dev/null +++ b/apps/server/src/main/resources/log4j2.xml @@ -0,0 +1,25 @@ + + + + + + + + + + + + + + + + + + + + + + + + +