fix: complete all P2 audit findings — dead code, NPE guard, hygiene
P2-1: Guard SessionId NPE with safe-call in SessionsReducer
P2-2: Remove 16 dead event classes + reducer branches; replace with live orchestration events in tests
P2-3: Standardize ChatInput divergences — guard CancelSession, single-arg ProtocolError
P2-4: Remove dead TuiToolRecord.diff and ToolDisplayStatus.REQUESTED
P2-5: Remove dead streamLive from SessionEventBridge
P2-6: Extract SessionsReducerContext data class for 6+ param method
P2-7: Replace StageId("none") sentinel with TypeId.NONE
P2-9: Add TypeId.random() factory on all type-alias IDs
P2-10: Add KDoc on schemaVersion documenting reserved-for-migration
P2-8: Verified RouterReducer string-template already correct (TypeId value class toString)
This commit is contained in:
@@ -244,12 +244,4 @@ class SessionEventBridge(
|
||||
return orderedIds.mapNotNull { byInvocation[it] }
|
||||
}
|
||||
|
||||
suspend fun streamLive(sessionId: SessionId) {
|
||||
var sessionSequence = 0L
|
||||
eventStore.subscribe(sessionId).collect { event ->
|
||||
sessionSequence++
|
||||
domainEventToServerMessage(event, artifactStore, sessionSequence = sessionSequence)
|
||||
?.let { send(it) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -88,11 +88,7 @@ class GlobalStreamHandler(private val module: ServerModule) {
|
||||
}
|
||||
.onFailure {
|
||||
log.warn("decode error: {}", it.message)
|
||||
val error = ServerMessage.ProtocolError(
|
||||
message = "Unknown message: ${it.message}",
|
||||
sequence = null,
|
||||
sessionSequence = null,
|
||||
)
|
||||
val error = ServerMessage.ProtocolError("Unknown message: ${it.message}")
|
||||
session.send(Frame.Text(ProtocolSerializer.encodeServerMessage(error)))
|
||||
}
|
||||
}
|
||||
@@ -144,7 +140,10 @@ class GlobalStreamHandler(private val module: ServerModule) {
|
||||
is ClientMessage.Ping -> Unit
|
||||
is ClientMessage.StartSession -> handleStartSession(session, msg, sendFrame)
|
||||
is ClientMessage.StartChatSession -> handleStartChatSession(session, msg, sendFrame)
|
||||
is ClientMessage.CancelSession -> module.orchestrator.cancel(msg.sessionId)
|
||||
is ClientMessage.CancelSession -> {
|
||||
runCatching { module.orchestrator.cancel(msg.sessionId) }
|
||||
.onFailure { log.error("cancel failed for session={}: {}", msg.sessionId.value, it.message, it) }
|
||||
}
|
||||
is ClientMessage.ResumeSession -> session.send(encodeError("ResumeSession not supported"))
|
||||
is ClientMessage.ApprovalResponse -> handleApprovalResponse(msg, sendFrame)
|
||||
is ClientMessage.CreateGrant -> handleCreateGrant(msg, sendFrame)
|
||||
@@ -159,11 +158,7 @@ class GlobalStreamHandler(private val module: ServerModule) {
|
||||
))
|
||||
}.onFailure {
|
||||
log.error("routerFacade.onUserInput failed: {}", it.message)
|
||||
sendFrame(ServerMessage.ProtocolError(
|
||||
message = "Router error: ${it.message}",
|
||||
sequence = null,
|
||||
sessionSequence = null,
|
||||
))
|
||||
sendFrame(ServerMessage.ProtocolError("Router error: ${it.message}"))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -176,34 +171,16 @@ class GlobalStreamHandler(private val module: ServerModule) {
|
||||
val scopeSessionId = module.approvalCoordinator.lookupSession(msg.requestId)
|
||||
if (scopeSessionId == null) {
|
||||
log.warn("handleApprovalResponse: no session found for requestId={}", msg.requestId.value)
|
||||
sendFrame(
|
||||
ServerMessage.ProtocolError(
|
||||
message = "Unknown approval request",
|
||||
sequence = null,
|
||||
sessionSequence = null,
|
||||
),
|
||||
)
|
||||
sendFrame(ServerMessage.ProtocolError("Unknown approval request"))
|
||||
return
|
||||
}
|
||||
module.approvalCoordinator.handleResponse(msg, scopeSessionId)?.let { sendFrame(it) }
|
||||
}
|
||||
|
||||
private fun errorResponse(message: String) = ServerMessage.ProtocolError(
|
||||
message = message,
|
||||
sequence = null,
|
||||
sessionSequence = null,
|
||||
)
|
||||
private fun errorResponse(message: String) = ServerMessage.ProtocolError(message)
|
||||
|
||||
private fun encodeError(message: String): Frame.Text =
|
||||
Frame.Text(
|
||||
ProtocolSerializer.encodeServerMessage(
|
||||
ServerMessage.ProtocolError(
|
||||
message = message,
|
||||
sequence = null,
|
||||
sessionSequence = null,
|
||||
),
|
||||
),
|
||||
)
|
||||
Frame.Text(ProtocolSerializer.encodeServerMessage(ServerMessage.ProtocolError(message)))
|
||||
|
||||
private suspend fun handleCreateGrant(
|
||||
msg: ClientMessage.CreateGrant,
|
||||
|
||||
@@ -10,7 +10,7 @@ import com.correx.core.events.events.InferenceCompletedEvent
|
||||
import com.correx.core.events.events.InferenceStartedEvent
|
||||
import com.correx.core.events.events.InferenceTimeoutEvent
|
||||
import com.correx.core.events.events.OrchestrationPausedEvent
|
||||
import com.correx.core.events.events.SessionStartedEvent
|
||||
import com.correx.core.events.events.SteeringNoteAddedEvent
|
||||
import com.correx.core.events.events.StageCompletedEvent
|
||||
import com.correx.core.events.events.StageFailedEvent
|
||||
import com.correx.core.events.events.StoredEvent
|
||||
@@ -371,7 +371,7 @@ class DomainEventMapperTest {
|
||||
|
||||
@Test
|
||||
fun `unmapped event returns null`(): Unit = runTest {
|
||||
val event = storedEvent(SessionStartedEvent(sessionId = sessionId))
|
||||
val event = storedEvent(SteeringNoteAddedEvent(sessionId = sessionId, content = "test"))
|
||||
val result = domainEventToServerMessage(event, noopStore, sessionSequence = 0L)
|
||||
assertNull(result)
|
||||
}
|
||||
@@ -392,12 +392,12 @@ class DomainEventMapperTest {
|
||||
val origLevel = logger.level
|
||||
logger.level = Level.DEBUG
|
||||
try {
|
||||
val event = storedEvent(SessionStartedEvent(sessionId = sessionId))
|
||||
val event = storedEvent(SteeringNoteAddedEvent(sessionId = sessionId, content = "test"))
|
||||
val result = domainEventToServerMessage(event, noopStore, sessionSequence = 0L)
|
||||
assertNull(result)
|
||||
assertEquals(1, appender.events.size)
|
||||
assertEquals(Level.DEBUG, appender.events[0].level)
|
||||
assertTrue(appender.events[0].message.formattedMessage.contains("SessionStartedEvent"))
|
||||
assertTrue(appender.events[0].message.formattedMessage.contains("SteeringNoteAddedEvent"))
|
||||
} finally {
|
||||
logger.removeAppender(appender)
|
||||
appender.stop()
|
||||
|
||||
@@ -25,10 +25,8 @@ import com.correx.core.kernel.orchestration.OrchestrationRepository
|
||||
import com.correx.core.tools.contract.Tool
|
||||
import com.correx.core.tools.registry.ToolRegistry
|
||||
import com.correx.core.sessions.projections.replay.EventReplayer
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.datetime.Instant
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
@@ -186,36 +184,6 @@ class SessionEventBridgeTest {
|
||||
assertEquals(ServerMessage.SnapshotComplete, sent[1])
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `streamLive sends mapped messages from live flow`() = runTest {
|
||||
val events = listOf(
|
||||
storedEvent(WorkflowStartedEvent(sessionId, workflowId, stageId), seq = 1L),
|
||||
storedEvent(WorkflowCompletedEvent(sessionId, stageId, 1), seq = 2L),
|
||||
)
|
||||
val liveFlow: Flow<StoredEvent> = flow { events.forEach { emit(it) } }
|
||||
val store = fakeEventStore(liveFlow = liveFlow)
|
||||
val sent = mutableListOf<ServerMessage>()
|
||||
val bridge = SessionEventBridge(
|
||||
store,
|
||||
noopArtifactStore,
|
||||
activeOrchestrationRepository(),
|
||||
noopWorkflowRegistry,
|
||||
noopToolRegistry,
|
||||
) { sent.add(it) }
|
||||
|
||||
bridge.streamLive(sessionId)
|
||||
|
||||
assertEquals(2, sent.size)
|
||||
assertEquals(
|
||||
ServerMessage.SessionStarted(sessionId, workflowId, 1L, 1L),
|
||||
sent[0],
|
||||
)
|
||||
assertEquals(
|
||||
ServerMessage.SessionCompleted(sessionId, 2L, 2L),
|
||||
sent[1],
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `replaySnapshot with no sessions emits only SnapshotComplete`() = runTest {
|
||||
val store = fakeEventStore(allEventsList = emptyList())
|
||||
@@ -266,29 +234,4 @@ class SessionEventBridgeTest {
|
||||
assertEquals(3L, snapshot.lastSessionSequence)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `streamLive cancellation stops collection`() = runTest {
|
||||
val channel = Channel<StoredEvent>(Channel.UNLIMITED)
|
||||
val liveFlow: Flow<StoredEvent> = flow {
|
||||
for (event in channel) emit(event)
|
||||
}
|
||||
val store = fakeEventStore(liveFlow = liveFlow)
|
||||
val sent = mutableListOf<ServerMessage>()
|
||||
val bridge = SessionEventBridge(
|
||||
store,
|
||||
noopArtifactStore,
|
||||
activeOrchestrationRepository(),
|
||||
noopWorkflowRegistry,
|
||||
noopToolRegistry,
|
||||
) { sent.add(it) }
|
||||
|
||||
channel.send(storedEvent(WorkflowStartedEvent(sessionId, workflowId, stageId), seq = 1L))
|
||||
|
||||
val job = launch { bridge.streamLive(sessionId) }
|
||||
channel.send(storedEvent(WorkflowCompletedEvent(sessionId, stageId, 1), seq = 2L))
|
||||
channel.close()
|
||||
job.join()
|
||||
|
||||
assertEquals(2, sent.size)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user