IJPL-233558 optimize rollout parsing

GitOrigin-RevId: 6a654db2863cc3f0d616d39acde43cf2a9fe9683
This commit is contained in:
Vladimir Krivosheev
2026-02-15 21:08:33 +00:00
committed by intellij-monorepo-bot
parent 72c947766f
commit d9fb8debee
24 changed files with 1247 additions and 215 deletions
+2 -1
View File
@@ -52,6 +52,7 @@
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/claude/sessions/intellij.agent.workbench.claude.sessions.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/claude/sessions/intellij.agent.workbench.claude.sessions.iml" />
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/codex/common/intellij.agent.workbench.codex.common.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/codex/common/intellij.agent.workbench.codex.common.iml" />
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/codex/sessions/intellij.agent.workbench.codex.sessions.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/codex/sessions/intellij.agent.workbench.codex.sessions.iml" />
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/json/intellij.agent.workbench.json.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/json/intellij.agent.workbench.json.iml" />
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/plugin/intellij.agent.workbench.plugin.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/plugin/intellij.agent.workbench.plugin.iml" />
<module fileurl="file://$PROJECT_DIR$/plugins/agent-workbench/sessions/intellij.agent.workbench.sessions.iml" filepath="$PROJECT_DIR$/plugins/agent-workbench/sessions/intellij.agent.workbench.sessions.iml" />
<module fileurl="file://$PROJECT_DIR$/android/android-adb/intellij.android.adb.iml" filepath="$PROJECT_DIR$/android/android-adb/intellij.android.adb.iml" />
@@ -1778,4 +1779,4 @@
<module fileurl="file://$PROJECT_DIR$/plugins/kotlin/util/test-generator-fir/kotlin.util.test-generator-fir.iml" filepath="$PROJECT_DIR$/plugins/kotlin/util/test-generator-fir/kotlin.util.test-generator-fir.iml" />
</modules>
</component>
</project>
</project>
+1
View File
@@ -879,6 +879,7 @@ plugins/agent-workbench/claude/common
plugins/agent-workbench/claude/sessions
plugins/agent-workbench/codex/common
plugins/agent-workbench/codex/sessions
plugins/agent-workbench/json
plugins/agent-workbench/plugin
plugins/agent-workbench/sessions
plugins/ant
@@ -15,6 +15,7 @@ jvm_library(
resources = [":common_resources"],
deps = [
"@lib//:kotlin-stdlib",
"//plugins/agent-workbench/json",
"//libraries/jackson/jackson",
]
)
@@ -9,6 +9,7 @@
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="intellij.agent.workbench.json" />
<orderEntry type="module" module-name="intellij.libraries.jackson" />
</component>
</module>
@@ -4,16 +4,16 @@ package com.intellij.agent.workbench.claude.common
import com.fasterxml.jackson.core.JsonFactory
import com.fasterxml.jackson.core.JsonParser
import com.fasterxml.jackson.core.JsonToken
import com.intellij.agent.workbench.json.WorkbenchJsonlScanner
import java.nio.file.Files
import java.nio.file.Path
import java.time.Instant
import kotlin.io.path.invariantSeparatorsPathString
import kotlin.io.path.nameWithoutExtension
private const val CLAUDE_PROJECTS_DIR = "projects"
private const val CLAUDE_INDEX_FILE = "sessions-index.json"
private const val CLAUDE_INDEX_VERSION = 1L
private const val MAX_JSONL_SCAN_LINES = 240
private const val MAX_JSONL_SCAN_OBJECTS = 240
private const val MAX_TITLE_LENGTH = 120
data class ClaudeSessionThread(
@@ -180,64 +180,54 @@ class ClaudeSessionsStore(
}
private fun parseJsonlMetadata(path: Path, targetProjectPath: String): ParsedJsonlMetadata? {
var firstPrompt: String? = null
var sessionId: String? = null
var isSidechain = false
var hasPathSignal = false
var pathMatches = false
var updatedAt: Long? = null
var hasConversationSignal = false
Files.newBufferedReader(path).use { reader ->
var scannedLines = 0
while (scannedLines < MAX_JSONL_SCAN_LINES) {
val line = reader.readLine() ?: break
scannedLines++
val trimmed = line.trim()
if (trimmed.isEmpty()) continue
parseJsonlLine(trimmed)?.let { lineData ->
if (lineData.isSidechain) {
isSidechain = true
return@use
}
if (sessionId == null && !lineData.sessionId.isNullOrBlank()) {
sessionId = lineData.sessionId
}
if (firstPrompt == null && !lineData.firstPrompt.isNullOrBlank()) {
firstPrompt = lineData.firstPrompt
}
if (lineData.hasConversationSignal) {
hasConversationSignal = true
}
val lineTimestamp = lineData.timestampMillis
if (lineTimestamp != null) {
updatedAt = maxOf(updatedAt ?: 0L, lineTimestamp)
}
if (!lineData.cwd.isNullOrBlank()) {
hasPathSignal = true
val normalizedCwd = normalizePath(lineData.cwd)
if (normalizedCwd == targetProjectPath) {
pathMatches = true
}
}
val state = WorkbenchJsonlScanner.scanJsonObjects(
path = path,
jsonFactory = jsonFactory,
maxObjects = MAX_JSONL_SCAN_OBJECTS,
newState = ::JsonlMetadataScanState,
) { parser, scanState ->
val lineData = parseJsonlLine(parser) ?: return@scanJsonObjects true
if (lineData.isSidechain) {
scanState.isSidechain = true
return@scanJsonObjects false
}
if (scanState.sessionId == null && !lineData.sessionId.isNullOrBlank()) {
scanState.sessionId = lineData.sessionId
}
if (scanState.firstPrompt == null && !lineData.firstPrompt.isNullOrBlank()) {
scanState.firstPrompt = lineData.firstPrompt
}
if (lineData.hasConversationSignal) {
scanState.hasConversationSignal = true
}
val lineTimestamp = lineData.timestampMillis
if (lineTimestamp != null) {
scanState.updatedAt = maxOf(scanState.updatedAt ?: 0L, lineTimestamp)
}
if (!lineData.cwd.isNullOrBlank()) {
scanState.hasPathSignal = true
val normalizedCwd = normalizePath(lineData.cwd)
if (normalizedCwd == targetProjectPath) {
scanState.pathMatches = true
}
}
true
}
if (hasPathSignal && !pathMatches) return null
if (!hasConversationSignal) return null
val normalizedSessionId = sessionId?.trim()?.takeIf { it.isNotEmpty() } ?: return null
if (state.hasPathSignal && !state.pathMatches) return null
if (!state.hasConversationSignal) return null
val normalizedSessionId = state.sessionId?.trim()?.takeIf { it.isNotEmpty() } ?: return null
return ParsedJsonlMetadata(
sessionId = normalizedSessionId,
firstPrompt = firstPrompt,
isSidechain = isSidechain,
updatedAt = updatedAt,
firstPrompt = state.firstPrompt,
isSidechain = state.isSidechain,
updatedAt = state.updatedAt,
)
}
private fun parseJsonlLine(line: String): ParsedJsonlLine? {
jsonFactory.createParser(line).use { parser ->
if (parser.nextToken() != JsonToken.START_OBJECT) return null
private fun parseJsonlLine(parser: JsonParser): ParsedJsonlLine? {
return try {
if (parser.currentToken != JsonToken.START_OBJECT) return null
var sessionId: String? = null
var cwd: String? = null
var isSidechain = false
@@ -281,6 +271,9 @@ class ClaudeSessionsStore(
hasConversationSignal = hasConversationSignal,
)
}
catch (_: Throwable) {
null
}
}
private fun readMessageObject(parser: JsonParser): ParsedMessageObject {
@@ -465,3 +458,13 @@ private data class ParsedJsonlMetadata(
val isSidechain: Boolean,
val updatedAt: Long?,
)
private data class JsonlMetadataScanState(
var firstPrompt: String? = null,
var sessionId: String? = null,
var isSidechain: Boolean = false,
var hasPathSignal: Boolean = false,
var pathMatches: Boolean = false,
var updatedAt: Long? = null,
var hasConversationSignal: Boolean = false,
)
@@ -8,6 +8,7 @@ import org.junit.jupiter.api.Test
import org.junit.jupiter.api.io.TempDir
import java.nio.file.Files
import java.nio.file.Path
import java.time.Instant
class ClaudeSessionsStoreTest {
@TempDir
@@ -173,5 +174,31 @@ class ClaudeSessionsStoreTest {
assertThat(threads).isEmpty()
}
}
@Test
fun skipsMalformedJsonlLineWhenReadingFallbackMetadata() {
val projectPath = "/work/project-malformed"
val encodedPath = "-work-project-malformed"
val projectDir = tempDir.resolve(".claude").resolve("projects").resolve(encodedPath)
Files.createDirectories(projectDir)
val transcript = projectDir.resolve("malformed-1111-2222-3333-444444444444.jsonl")
Files.writeString(
transcript,
"""
{"type":"user","sessionId":"malformed-1111-2222-3333-444444444444","cwd":"$projectPath","isSidechain":false,"timestamp":"2026-02-08T01:00:00.000Z","message":{"role":"user","content":"Recover after malformed line"}}
{"type":"assistant","sessionId":"malformed-1111-2222-3333-444444444444","cwd":"$projectPath","isSidechain":false,"timestamp":"2026-02-08T01:00:01.000Z","message":{"role":"assistant"
{"type":"assistant","sessionId":"malformed-1111-2222-3333-444444444444","cwd":"$projectPath","isSidechain":false,"timestamp":"2026-02-08T01:00:02.000Z","message":{"role":"assistant","content":[{"type":"text","text":"done"}]}}
""".trimIndent()
)
val store = ClaudeSessionsStore(claudeHomeProvider = { tempDir.resolve(".claude") })
val threads = runBlocking { store.listThreads(projectPath) }
assertThat(threads).hasSize(1)
val thread = threads.single()
assertThat(thread.id).isEqualTo("malformed-1111-2222-3333-444444444444")
assertThat(thread.title).contains("Recover after malformed line")
assertThat(thread.updatedAt).isEqualTo(Instant.parse("2026-02-08T01:00:02.000Z").toEpochMilli())
}
}
@@ -15,12 +15,14 @@ jvm_library(
resources = [":sessions_resources"],
deps = [
"@lib//:kotlin-stdlib",
"//libraries/fastutil",
"//platform/core-api:core",
"//platform/projectModel-api:projectModel",
"//platform/util",
"//platform/platform-api:ide",
"//platform/platform-util-io:ide-util-io",
"//plugins/agent-workbench/codex/common",
"//plugins/agent-workbench/json",
"//plugins/agent-workbench/sessions",
"//libraries/jackson/jackson",
]
@@ -33,12 +35,14 @@ jvm_library(
associates = [":sessions"],
deps = [
"@lib//:kotlin-stdlib",
"//libraries/fastutil",
"//platform/core-api:core",
"//platform/projectModel-api:projectModel",
"//platform/util",
"//platform/platform-api:ide",
"//platform/platform-util-io:ide-util-io",
"//plugins/agent-workbench/codex/common",
"//plugins/agent-workbench/json",
"//plugins/agent-workbench/sessions",
"//plugins/agent-workbench/sessions:sessions_test_lib",
"//libraries/jackson/jackson",
@@ -10,12 +10,14 @@
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="intellij.libraries.fastutil" />
<orderEntry type="module" module-name="intellij.platform.core" />
<orderEntry type="module" module-name="intellij.platform.projectModel" />
<orderEntry type="module" module-name="intellij.platform.util" />
<orderEntry type="module" module-name="intellij.platform.ide" />
<orderEntry type="module" module-name="intellij.platform.ide.util.io" />
<orderEntry type="module" module-name="intellij.agent.workbench.codex.common" />
<orderEntry type="module" module-name="intellij.agent.workbench.json" />
<orderEntry type="module" module-name="intellij.agent.workbench.sessions" />
<orderEntry type="module" module-name="intellij.libraries.jackson" />
<orderEntry type="module" module-name="intellij.libraries.junit5" scope="TEST" />
@@ -15,10 +15,14 @@ import com.intellij.agent.workbench.sessions.AgentSessionThread
import com.intellij.agent.workbench.sessions.AgentSubAgent
import com.intellij.agent.workbench.sessions.providers.BaseAgentSessionSource
import com.intellij.openapi.project.Project
import kotlinx.coroutines.flow.Flow
class CodexSessionSource(
private val backend: CodexSessionBackend = createDefaultCodexSessionBackend(),
) : BaseAgentSessionSource(provider = AgentSessionProvider.CODEX, canReportExactThreadCount = false) {
override val updates: Flow<Unit>
get() = backend.updates
override suspend fun listThreads(path: String, openProject: Project?): List<AgentSessionThread> {
return backend.listThreads(path = path, openProject = openProject).map { it.toAgentSessionThread() }
}
@@ -3,6 +3,8 @@ package com.intellij.agent.workbench.codex.sessions.backend
import com.intellij.agent.workbench.codex.common.CodexThread
import com.intellij.openapi.project.Project
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.emptyFlow
data class CodexBackendThread(
val thread: CodexThread,
@@ -19,6 +21,9 @@ enum class CodexSessionActivity {
interface CodexSessionBackend {
suspend fun listThreads(path: String, openProject: Project?): List<CodexBackendThread>
val updates: Flow<Unit>
get() = emptyFlow()
/**
* Prefetch threads for multiple paths in a single backend call.
* Returns a map of path to threads. Empty map means no prefetch (use per-path calls).
@@ -13,228 +13,388 @@ import com.intellij.agent.workbench.codex.sessions.backend.CodexBackendThread
import com.intellij.agent.workbench.codex.sessions.backend.CodexSessionActivity
import com.intellij.agent.workbench.codex.sessions.backend.CodexSessionBackend
import com.intellij.agent.workbench.codex.sessions.resolveProjectDirectoryFromPath
import com.intellij.agent.workbench.json.WorkbenchJsonlScanner
import com.intellij.openapi.diagnostic.logger
import com.intellij.openapi.project.Project
import it.unimi.dsi.fastutil.objects.Object2ObjectOpenHashMap
import it.unimi.dsi.fastutil.objects.ObjectArrayList
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.channels.awaitClose
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.callbackFlow
import kotlinx.coroutines.flow.conflate
import kotlinx.coroutines.withContext
import java.nio.file.ClosedWatchServiceException
import java.nio.file.FileSystems
import java.nio.file.FileVisitResult
import java.nio.file.Files
import java.nio.file.LinkOption
import java.nio.file.Path
import java.nio.file.SimpleFileVisitor
import java.nio.file.StandardWatchEventKinds
import java.nio.file.WatchKey
import java.nio.file.attribute.BasicFileAttributes
import java.time.Instant
import java.time.format.DateTimeParseException
import java.util.concurrent.atomic.AtomicBoolean
import kotlin.io.path.invariantSeparatorsPathString
private const val ROLLOUT_FILE_PREFIX = "rollout-"
private const val ROLLOUT_FILE_SUFFIX = ".jsonl"
private const val MAX_TITLE_LENGTH = 120
private const val USER_MESSAGE_BEGIN = "## My request for Codex:"
private const val ENVIRONMENT_CONTEXT_OPEN_TAG = "<environment_context>"
private const val TURN_ABORTED_OPEN_TAG = "<turn_aborted>"
private val LOG = logger<CodexRolloutSessionBackend>()
internal class CodexRolloutSessionBackend(
private val codexHomeProvider: () -> Path = { Path.of(System.getProperty("user.home"), ".codex") },
) : CodexSessionBackend {
private val jsonFactory = JsonFactory()
private val cacheLock = Any()
private val cachedFilesByPath = Object2ObjectOpenHashMap<String, CachedRolloutFile>()
private val threadsByCwd = Object2ObjectOpenHashMap<String, ObjectArrayList<CodexBackendThread>>()
override val updates: Flow<Unit> = callbackFlow {
val watcher = runCatching {
CodexRolloutSessionsWatcher(codexHomeProvider = codexHomeProvider) {
trySend(Unit)
}
}.onFailure { t ->
LOG.warn("Failed to initialize Codex rollout watcher", t)
}.getOrNull()
if (watcher == null) {
awaitClose { }
return@callbackFlow
}
awaitClose {
watcher.close()
}
}.conflate()
override suspend fun listThreads(path: String, @Suppress("UNUSED_PARAMETER") openProject: Project?): List<CodexBackendThread> {
return withContext(Dispatchers.IO) {
val workingDirectory = resolveProjectDirectoryFromPath(path)
?: return@withContext emptyList()
val cwdFilter = normalizeRootPath(workingDirectory.invariantSeparatorsPathString)
val sessionsDir = codexHomeProvider().resolve("sessions")
if (!Files.isDirectory(sessionsDir)) return@withContext emptyList()
val threads = mutableListOf<CodexBackendThread>()
try {
Files.walk(sessionsDir).use { stream ->
val iterator = stream.iterator()
while (iterator.hasNext()) {
val candidate = iterator.next()
if (!Files.isRegularFile(candidate)) continue
val fileName = candidate.fileName?.toString() ?: continue
if (!isRolloutFileName(fileName)) continue
parseRolloutFile(candidate, cwdFilter)?.let(threads::add)
}
}
}
catch (_: Throwable) {
return@withContext emptyList()
}
threads.sortedByDescending { it.thread.updatedAt }
collectThreadsByCwd(setOf(cwdFilter))[cwdFilter].orEmpty()
}
}
private fun parseRolloutFile(path: Path, cwdFilter: String): CodexBackendThread? {
var sessionId: String? = null
var sessionCwd: String? = null
var gitBranch: String? = null
var title: String? = null
var updatedAt = 0L
var processing = false
var reviewing = false
var latestUserMessageAt = Long.MIN_VALUE
var latestAgentMessageAt = Long.MIN_VALUE
var pendingUserInputAt: Long? = null
override suspend fun prefetchThreads(paths: List<String>): Map<String, List<CodexBackendThread>> {
return withContext(Dispatchers.IO) {
val pathFilters = resolvePathFilters(paths)
if (pathFilters.isEmpty()) return@withContext emptyMap()
try {
Files.newBufferedReader(path).use { reader ->
while (true) {
val line = reader.readLine() ?: break
if (line.isBlank()) continue
val event = parseEvent(line) ?: continue
val threadsByCwd = collectThreadsByCwd(pathFilters.map { (_, cwdFilter) -> cwdFilter }.toSet())
pathFilters.associate { (path, cwdFilter) ->
path to threadsByCwd[cwdFilter].orEmpty()
}
}
}
updatedAt = maxTimestamp(updatedAt, event.timestampMs)
updatedAt = maxTimestamp(updatedAt, event.sessionTimestampMs)
sessionId = sessionId ?: event.sessionId
sessionCwd = sessionCwd ?: event.sessionCwd
gitBranch = gitBranch ?: event.gitBranch
private fun collectThreadsByCwd(cwdFilters: Set<String>): Map<String, List<CodexBackendThread>> {
if (cwdFilters.isEmpty()) return emptyMap()
val eventTimestamp = event.timestampMs
when (event.topLevelType) {
"event_msg" -> {
when (event.payloadType) {
"task_started" -> processing = true
"task_complete", "turn_aborted" -> processing = false
"user_message" -> {
latestUserMessageAt = maxTimestamp(latestUserMessageAt, eventTimestamp)
title = title ?: extractTitle(event.payloadMessage)
val pendingInputAt = pendingUserInputAt
if (pendingInputAt != null && eventTimestamp != null && eventTimestamp >= pendingInputAt) {
pendingUserInputAt = null
}
}
val sessionsDir = codexHomeProvider().resolve("sessions")
if (!Files.isDirectory(sessionsDir)) {
synchronized(cacheLock) {
cachedFilesByPath.clear()
threadsByCwd.clear()
}
return emptyMap()
}
"agent_message" -> {
latestAgentMessageAt = maxTimestamp(latestAgentMessageAt, eventTimestamp)
val scannedFiles = try {
scanRolloutFiles(sessionsDir)
}
catch (_: Throwable) {
return emptyMap()
}
val filesToParse = ObjectArrayList<RolloutFileStat>()
var removedAny = false
synchronized(cacheLock) {
val iterator = cachedFilesByPath.object2ObjectEntrySet().iterator()
while (iterator.hasNext()) {
val entry = iterator.next()
if (!scannedFiles.containsKey(entry.key)) {
iterator.remove()
removedAny = true
}
}
for (entry in scannedFiles.object2ObjectEntrySet()) {
val stat = entry.value
val cached = cachedFilesByPath[entry.key]
if (cached == null || cached.lastModifiedMs != stat.lastModifiedMs || cached.sizeBytes != stat.sizeBytes) {
filesToParse.add(stat)
}
}
}
if (filesToParse.isNotEmpty()) {
val parsedUpdates = Object2ObjectOpenHashMap<String, CachedRolloutFile>(filesToParse.size)
for (stat in filesToParse) {
parsedUpdates[stat.pathKey] = CachedRolloutFile(
lastModifiedMs = stat.lastModifiedMs,
sizeBytes = stat.sizeBytes,
parsedThread = parseRolloutFile(stat.path),
)
}
synchronized(cacheLock) {
for (entry in parsedUpdates.object2ObjectEntrySet()) {
cachedFilesByPath[entry.key] = entry.value
}
}
}
if (removedAny || filesToParse.isNotEmpty()) {
synchronized(cacheLock) {
rebuildThreadsByCwd()
}
}
synchronized(cacheLock) {
val result = Object2ObjectOpenHashMap<String, List<CodexBackendThread>>(cwdFilters.size)
for (cwdFilter in cwdFilters) {
val threads = threadsByCwd[cwdFilter] ?: continue
result[cwdFilter] = ArrayList(threads)
}
return result
}
}
private fun rebuildThreadsByCwd() {
threadsByCwd.clear()
for (entry in cachedFilesByPath.object2ObjectEntrySet()) {
val parsedThread = entry.value.parsedThread ?: continue
var threads = threadsByCwd[parsedThread.normalizedCwd]
if (threads == null) {
threads = ObjectArrayList()
threadsByCwd[parsedThread.normalizedCwd] = threads
}
threads.add(parsedThread.thread)
}
for (threads in threadsByCwd.values) {
threads.sortWith(Comparator { left, right ->
right.thread.updatedAt.compareTo(left.thread.updatedAt)
})
}
}
private fun scanRolloutFiles(sessionsDir: Path): Object2ObjectOpenHashMap<String, RolloutFileStat> {
val scannedFiles = Object2ObjectOpenHashMap<String, RolloutFileStat>()
Files.walk(sessionsDir).use { stream ->
val iterator = stream.iterator()
while (iterator.hasNext()) {
val candidate = iterator.next()
if (!Files.isRegularFile(candidate)) continue
val fileName = candidate.fileName?.toString() ?: continue
if (!isRolloutFileName(fileName)) continue
val lastModifiedMs = try {
Files.getLastModifiedTime(candidate).toMillis()
}
catch (_: Throwable) {
continue
}
val sizeBytes = try {
Files.size(candidate)
}
catch (_: Throwable) {
continue
}
val pathKey = candidate.invariantSeparatorsPathString
scannedFiles[pathKey] = RolloutFileStat(
pathKey = pathKey,
path = candidate,
lastModifiedMs = lastModifiedMs,
sizeBytes = sizeBytes,
)
}
}
return scannedFiles
}
private fun resolvePathFilters(paths: List<String>): List<Pair<String, String>> {
return paths.mapNotNull { path ->
resolveProjectDirectoryFromPath(path)?.let { directory ->
path to normalizeRootPath(directory.invariantSeparatorsPathString)
}
}
}
private fun parseRolloutFile(path: Path): ParsedRolloutThread? {
val state = try {
WorkbenchJsonlScanner.scanJsonObjects(
path = path,
jsonFactory = jsonFactory,
newState = ::RolloutParseState,
) { parser, parseState ->
val event = parseEvent(parser) ?: return@scanJsonObjects true
parseState.updatedAt = maxTimestamp(parseState.updatedAt, event.timestampMs)
parseState.updatedAt = maxTimestamp(parseState.updatedAt, event.sessionTimestampMs)
parseState.sessionId = parseState.sessionId ?: event.sessionId
parseState.sessionCwd = parseState.sessionCwd ?: event.sessionCwd
parseState.gitBranch = parseState.gitBranch ?: event.gitBranch
val eventTimestamp = event.timestampMs
when (event.topLevelType) {
"event_msg" -> {
when (event.payloadType) {
"task_started" -> parseState.processing = true
"task_complete", "turn_aborted" -> parseState.processing = false
"user_message" -> {
parseState.latestUserMessageAt = maxTimestamp(parseState.latestUserMessageAt, eventTimestamp)
parseState.title = parseState.title ?: extractTitle(event.payloadMessage)
val pendingInputAt = parseState.pendingUserInputAt
if (pendingInputAt != null && eventTimestamp != null && eventTimestamp >= pendingInputAt) {
parseState.pendingUserInputAt = null
}
}
if (event.payloadType?.contains("requestUserInput", ignoreCase = true) == true) {
pendingUserInputAt = maxTimestamp(pendingUserInputAt ?: Long.MIN_VALUE, eventTimestamp)
}
when (event.itemType) {
"enteredReviewMode" -> reviewing = true
"exitedReviewMode" -> reviewing = false
"agent_message" -> {
parseState.latestAgentMessageAt = maxTimestamp(parseState.latestAgentMessageAt, eventTimestamp)
}
}
"response_item" -> {
if (event.payloadType == "message") {
when (event.payloadRole) {
"user" -> {
latestUserMessageAt = maxTimestamp(latestUserMessageAt, eventTimestamp)
val pendingInputAt = pendingUserInputAt
if (pendingInputAt != null && eventTimestamp != null && eventTimestamp >= pendingInputAt) {
pendingUserInputAt = null
}
}
if (event.payloadType?.contains("requestUserInput", ignoreCase = true) == true) {
parseState.pendingUserInputAt = maxTimestamp(parseState.pendingUserInputAt ?: Long.MIN_VALUE, eventTimestamp)
}
"assistant" -> {
latestAgentMessageAt = maxTimestamp(latestAgentMessageAt, eventTimestamp)
when (event.itemType) {
"enteredReviewMode" -> parseState.reviewing = true
"exitedReviewMode" -> parseState.reviewing = false
}
}
"response_item" -> {
if (event.payloadType == "message") {
when (event.payloadRole) {
"user" -> {
parseState.latestUserMessageAt = maxTimestamp(parseState.latestUserMessageAt, eventTimestamp)
val pendingInputAt = parseState.pendingUserInputAt
if (pendingInputAt != null && eventTimestamp != null && eventTimestamp >= pendingInputAt) {
parseState.pendingUserInputAt = null
}
}
"assistant" -> {
parseState.latestAgentMessageAt = maxTimestamp(parseState.latestAgentMessageAt, eventTimestamp)
}
}
}
}
}
true
}
}
catch (_: Throwable) {
return null
}
val normalizedCwd = normalizeRootPath(sessionCwd ?: return null)
if (normalizedCwd != cwdFilter) return null
val resolvedSessionId = sessionId ?: return null
val hasUnread = latestAgentMessageAt > latestUserMessageAt
val hasPendingUserInput = pendingUserInputAt != null
val normalizedCwd = normalizeRootPath(state.sessionCwd ?: return null)
val resolvedSessionId = state.sessionId ?: return null
val hasUnread = state.latestAgentMessageAt > state.latestUserMessageAt
val hasPendingUserInput = state.pendingUserInputAt != null
val activity = when {
hasPendingUserInput || hasUnread -> CodexSessionActivity.UNREAD
reviewing -> CodexSessionActivity.REVIEWING
processing -> CodexSessionActivity.PROCESSING
state.reviewing -> CodexSessionActivity.REVIEWING
state.processing -> CodexSessionActivity.PROCESSING
else -> CodexSessionActivity.READY
}
val fallbackUpdatedAt = runCatching { Files.getLastModifiedTime(path).toMillis() }.getOrDefault(0L)
val resolvedUpdatedAt = if (updatedAt > 0L) updatedAt else fallbackUpdatedAt
val resolvedTitle = title ?: "Thread ${resolvedSessionId.take(8)}"
val resolvedUpdatedAt = if (state.updatedAt > 0L) state.updatedAt else fallbackUpdatedAt
val resolvedTitle = state.title ?: "Thread ${resolvedSessionId.take(8)}"
return CodexBackendThread(
thread = CodexThread(
id = resolvedSessionId,
title = resolvedTitle,
updatedAt = resolvedUpdatedAt,
archived = false,
gitBranch = gitBranch,
return ParsedRolloutThread(
normalizedCwd = normalizedCwd,
thread = CodexBackendThread(
thread = CodexThread(
id = resolvedSessionId,
title = resolvedTitle,
updatedAt = resolvedUpdatedAt,
archived = false,
gitBranch = state.gitBranch,
),
activity = activity,
),
activity = activity,
)
}
private fun parseEvent(line: String): RolloutEvent? {
private fun parseEvent(parser: JsonParser): RolloutEvent? {
return try {
jsonFactory.createParser(line).use { parser ->
if (parser.nextToken() != JsonToken.START_OBJECT) return null
if (parser.currentToken != JsonToken.START_OBJECT) return null
var topLevelType: String? = null
var timestampMs: Long? = null
var payloadType: String? = null
var payloadRole: String? = null
var payloadMessage: String? = null
var sessionId: String? = null
var sessionCwd: String? = null
var sessionTimestampMs: Long? = null
var gitBranch: String? = null
var itemType: String? = null
var topLevelType: String? = null
var timestampMs: Long? = null
var payloadType: String? = null
var payloadRole: String? = null
var payloadMessage: String? = null
var sessionId: String? = null
var sessionCwd: String? = null
var sessionTimestampMs: Long? = null
var gitBranch: String? = null
var itemType: String? = null
forEachObjectField(parser) { fieldName ->
when (fieldName) {
"timestamp" -> timestampMs = parseIsoTimestamp(readStringOrNull(parser))
"type" -> topLevelType = readStringOrNull(parser)
"payload" -> {
if (parser.currentToken == JsonToken.START_OBJECT) {
forEachObjectField(parser) { payloadField ->
when (payloadField) {
"type" -> payloadType = readStringOrNull(parser)
"role" -> payloadRole = readStringOrNull(parser)
"message" -> payloadMessage = readStringOrNull(parser)
"id" -> sessionId = readStringOrNull(parser)
"cwd" -> sessionCwd = readStringOrNull(parser)
"timestamp" -> sessionTimestampMs = parseIsoTimestamp(readStringOrNull(parser))
"git" -> {
gitBranch = parseNestedStringField(parser, "branch")
}
"item" -> {
itemType = parseNestedStringField(parser, "type")
}
else -> parser.skipChildren()
forEachObjectField(parser) { fieldName ->
when (fieldName) {
"timestamp" -> timestampMs = parseIsoTimestamp(readStringOrNull(parser))
"type" -> topLevelType = readStringOrNull(parser)
"payload" -> {
if (parser.currentToken == JsonToken.START_OBJECT) {
forEachObjectField(parser) { payloadField ->
when (payloadField) {
"type" -> payloadType = readStringOrNull(parser)
"role" -> payloadRole = readStringOrNull(parser)
"message" -> payloadMessage = readStringOrNull(parser)
"id" -> sessionId = readStringOrNull(parser)
"cwd" -> sessionCwd = readStringOrNull(parser)
"timestamp" -> sessionTimestampMs = parseIsoTimestamp(readStringOrNull(parser))
"git" -> {
gitBranch = parseNestedStringField(parser, "branch")
}
true
"item" -> {
itemType = parseNestedStringField(parser, "type")
}
else -> parser.skipChildren()
}
}
else {
parser.skipChildren()
true
}
}
else -> parser.skipChildren()
else {
parser.skipChildren()
}
}
true
}
RolloutEvent(
topLevelType = topLevelType,
timestampMs = timestampMs,
payloadType = payloadType,
payloadRole = payloadRole,
payloadMessage = payloadMessage,
sessionId = sessionId,
sessionCwd = sessionCwd,
sessionTimestampMs = sessionTimestampMs,
gitBranch = gitBranch,
itemType = itemType,
)
else -> parser.skipChildren()
}
true
}
RolloutEvent(
topLevelType = topLevelType,
timestampMs = timestampMs,
payloadType = payloadType,
payloadRole = payloadRole,
payloadMessage = payloadMessage,
sessionId = sessionId,
sessionCwd = sessionCwd,
sessionTimestampMs = sessionTimestampMs,
gitBranch = gitBranch,
itemType = itemType,
)
}
catch (_: Throwable) {
null
@@ -255,6 +415,164 @@ private data class RolloutEvent(
@JvmField val itemType: String?,
)
private data class ParsedRolloutThread(
@JvmField val normalizedCwd: String,
@JvmField val thread: CodexBackendThread,
)
private data class RolloutParseState(
@JvmField var sessionId: String? = null,
@JvmField var sessionCwd: String? = null,
@JvmField var gitBranch: String? = null,
@JvmField var title: String? = null,
@JvmField var updatedAt: Long = 0L,
@JvmField var processing: Boolean = false,
@JvmField var reviewing: Boolean = false,
@JvmField var latestUserMessageAt: Long = Long.MIN_VALUE,
@JvmField var latestAgentMessageAt: Long = Long.MIN_VALUE,
@JvmField var pendingUserInputAt: Long? = null,
)
private data class RolloutFileStat(
@JvmField val pathKey: String,
@JvmField val path: Path,
@JvmField val lastModifiedMs: Long,
@JvmField val sizeBytes: Long,
)
private data class CachedRolloutFile(
@JvmField val lastModifiedMs: Long,
@JvmField val sizeBytes: Long,
@JvmField val parsedThread: ParsedRolloutThread?,
)
private class CodexRolloutSessionsWatcher(
private val codexHomeProvider: () -> Path,
private val onRolloutChange: () -> Unit,
) : AutoCloseable {
private val watchService = FileSystems.getDefault().newWatchService()
private val running = AtomicBoolean(true)
private val watchKeysByPath = Object2ObjectOpenHashMap<WatchKey, Path>()
private val watchKeysLock = Any()
private val sessionsRoot: Path
get() = codexHomeProvider().resolve("sessions")
private val thread = Thread(::runWatchLoop, "CodexRolloutSessionBackendWatcher").apply {
isDaemon = true
start()
}
init {
registerInitialPaths()
}
override fun close() {
if (!running.compareAndSet(true, false)) return
watchService.close()
thread.interrupt()
}
private fun registerInitialPaths() {
val codexHome = codexHomeProvider()
if (Files.isDirectory(codexHome)) {
registerDirectory(codexHome)
}
val sessions = sessionsRoot
if (Files.isDirectory(sessions)) {
registerDirectoryRecursively(sessions)
}
}
private fun runWatchLoop() {
while (running.get()) {
val watchKey = try {
watchService.take()
}
catch (_: InterruptedException) {
continue
}
catch (_: ClosedWatchServiceException) {
break
}
catch (_: Throwable) {
continue
}
val watchedPath = synchronized(watchKeysLock) { watchKeysByPath[watchKey] }
if (watchedPath == null) {
watchKey.reset()
continue
}
var hasRolloutChange = false
for (event in watchKey.pollEvents()) {
val kind = event.kind()
if (kind == StandardWatchEventKinds.OVERFLOW) {
hasRolloutChange = true
continue
}
val contextPath = event.context() as? Path ?: continue
val eventPath = watchedPath.resolve(contextPath)
if (kind == StandardWatchEventKinds.ENTRY_CREATE && Files.isDirectory(eventPath, LinkOption.NOFOLLOW_LINKS)) {
registerDirectoryRecursively(eventPath)
}
if (isRolloutPath(eventPath)) {
hasRolloutChange = true
}
}
if (!watchKey.reset()) {
synchronized(watchKeysLock) {
watchKeysByPath.remove(watchKey)
}
}
if (hasRolloutChange) {
onRolloutChange()
}
}
}
private fun isRolloutPath(path: Path): Boolean {
val fileName = path.fileName?.toString() ?: return false
return path.startsWith(sessionsRoot) && isRolloutFileName(fileName)
}
private fun registerDirectoryRecursively(root: Path) {
if (!Files.isDirectory(root, LinkOption.NOFOLLOW_LINKS)) return
try {
Files.walkFileTree(root, object : SimpleFileVisitor<Path>() {
override fun preVisitDirectory(dir: Path, attrs: BasicFileAttributes): FileVisitResult {
registerDirectory(dir)
return FileVisitResult.CONTINUE
}
})
}
catch (_: Throwable) {
}
}
private fun registerDirectory(path: Path) {
try {
val watchKey = path.register(
watchService,
StandardWatchEventKinds.ENTRY_CREATE,
StandardWatchEventKinds.ENTRY_DELETE,
StandardWatchEventKinds.ENTRY_MODIFY,
)
synchronized(watchKeysLock) {
watchKeysByPath[watchKey] = path
}
}
catch (_: Throwable) {
}
}
}
private fun parseIsoTimestamp(value: String?): Long? {
val text = value?.trim().takeIf { !it.isNullOrEmpty() } ?: return null
return try {
@@ -270,15 +588,30 @@ private fun isRolloutFileName(fileName: String): Boolean {
}
private fun extractTitle(message: String?): String? {
val candidate = message
?.lineSequence()
?.map(String::trim)
?.firstOrNull { it.isNotEmpty() }
val candidate = stripUserMessagePrefix(message ?: return null)
.lineSequence()
.map(String::trim)
.firstOrNull { it.isNotEmpty() }
?: return null
if (candidate.startsWith("<environment_context>")) return null
if (isSessionPrefix(candidate)) return null
return trimTitle(candidate.replace(Regex("\\s+"), " "))
}
private fun stripUserMessagePrefix(text: String): String {
val markerIndex = text.indexOf(USER_MESSAGE_BEGIN)
return if (markerIndex >= 0) {
text.substring(markerIndex + USER_MESSAGE_BEGIN.length).trim()
}
else {
text.trim()
}
}
private fun isSessionPrefix(text: String): Boolean {
val normalized = text.trimStart().lowercase()
return normalized.startsWith(ENVIRONMENT_CONTEXT_OPEN_TAG) || normalized.startsWith(TURN_ABORTED_OPEN_TAG)
}
private fun trimTitle(value: String): String {
if (value.length <= MAX_TITLE_LENGTH) return value
return value.take(MAX_TITLE_LENGTH - 3).trimEnd() + "..."
@@ -208,6 +208,200 @@ class CodexRolloutSessionBackendTest {
assertThat(threads.single().thread.gitBranch).isEqualTo("feature/codex-rollout")
}
}
@Test
fun prefetchThreadsMapsPerResolvedPath() {
runBlocking {
val projectA = tempDir.resolve("project-prefetch-a")
val projectB = tempDir.resolve("project-prefetch-b")
val projectC = tempDir.resolve("project-prefetch-c")
Files.createDirectories(projectA)
Files.createDirectories(projectB)
Files.createDirectories(projectC)
val sessionsRoot = tempDir.resolve("sessions").resolve("2026").resolve("02").resolve("14")
writeRollout(
file = sessionsRoot.resolve("rollout-prefetch-a-old.jsonl"),
lines = listOf(
sessionMetaLine(timestamp = "2026-02-14T09:00:00.000Z", id = "session-a-old", cwd = projectA),
"""{"timestamp":"2026-02-14T09:00:05.000Z","type":"event_msg","payload":{"type":"user_message","message":"A old"}}""",
),
)
writeRollout(
file = sessionsRoot.resolve("rollout-prefetch-a-new.jsonl"),
lines = listOf(
sessionMetaLine(timestamp = "2026-02-14T10:00:00.000Z", id = "session-a-new", cwd = projectA),
"""{"timestamp":"2026-02-14T10:00:05.000Z","type":"event_msg","payload":{"type":"user_message","message":"A new"}}""",
),
)
writeRollout(
file = sessionsRoot.resolve("rollout-prefetch-b.jsonl"),
lines = listOf(
sessionMetaLine(timestamp = "2026-02-14T11:00:00.000Z", id = "session-b", cwd = projectB),
),
)
val backend = CodexRolloutSessionBackend(codexHomeProvider = { tempDir })
val unresolvedPath = tempDir.resolve("project-prefetch-missing").toString()
val prefetched = backend.prefetchThreads(
listOf(projectA.toString(), projectB.toString(), projectC.toString(), unresolvedPath)
)
assertThat(prefetched.keys).containsExactlyInAnyOrder(projectA.toString(), projectB.toString(), projectC.toString())
assertThat(prefetched.getValue(projectA.toString()).map { it.thread.id }).containsExactly("session-a-new", "session-a-old")
assertThat(prefetched.getValue(projectB.toString()).map { it.thread.id }).containsExactly("session-b")
assertThat(prefetched.getValue(projectC.toString())).isEmpty()
}
}
@Test
fun refreshesCachedThreadsWhenRolloutFilesChangeAndDelete() {
runBlocking {
val projectDir = tempDir.resolve("project-cache")
Files.createDirectories(projectDir)
val sessionsRoot = tempDir.resolve("sessions").resolve("2026").resolve("02").resolve("15")
val rolloutA = sessionsRoot.resolve("rollout-cache-a.jsonl")
val rolloutB = sessionsRoot.resolve("rollout-cache-b.jsonl")
writeRollout(
file = rolloutA,
lines = listOf(
sessionMetaLine(timestamp = "2026-02-15T10:00:00.000Z", id = "session-a", cwd = projectDir),
"""{"timestamp":"2026-02-15T10:00:01.000Z","type":"event_msg","payload":{"type":"user_message","message":"Initial title"}}""",
),
)
val backend = CodexRolloutSessionBackend(codexHomeProvider = { tempDir })
val initialThreads = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(initialThreads.map { it.thread.id }).containsExactly("session-a")
assertThat(initialThreads.single().thread.title).isEqualTo("Initial title")
writeRollout(
file = rolloutA,
lines = listOf(
sessionMetaLine(timestamp = "2026-02-15T10:00:00.000Z", id = "session-a", cwd = projectDir),
"""{"timestamp":"2026-02-15T10:05:00.000Z","type":"event_msg","payload":{"type":"user_message","message":"Updated title with extra text"}}""",
"""{"timestamp":"2026-02-15T10:05:01.000Z","type":"event_msg","payload":{"type":"agent_message","message":"Done"}}""",
),
)
writeRollout(
file = rolloutB,
lines = listOf(
sessionMetaLine(timestamp = "2026-02-15T10:06:00.000Z", id = "session-b", cwd = projectDir),
"""{"timestamp":"2026-02-15T10:06:01.000Z","type":"event_msg","payload":{"type":"user_message","message":"Newest"}}""",
),
)
val afterRewriteAndAdd = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(afterRewriteAndAdd.map { it.thread.id }).containsExactly("session-b", "session-a")
assertThat(afterRewriteAndAdd.first { it.thread.id == "session-a" }.thread.title).isEqualTo("Updated title with extra text")
assertThat(afterRewriteAndAdd.first { it.thread.id == "session-a" }.thread.updatedAt)
.isEqualTo(Instant.parse("2026-02-15T10:05:01.000Z").toEpochMilli())
Files.delete(rolloutB)
val afterDelete = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(afterDelete.map { it.thread.id }).containsExactly("session-a")
}
}
@Test
fun retriesPreviouslyUnparseableRolloutAfterRewrite() {
runBlocking {
val projectDir = tempDir.resolve("project-retry")
Files.createDirectories(projectDir)
val rollout = tempDir.resolve("sessions").resolve("2026").resolve("02").resolve("15")
.resolve("rollout-retry.jsonl")
writeRollout(
file = rollout,
lines = listOf(
sessionMetaLineWithoutId(cwd = projectDir),
"""{"timestamp":"2026-02-15T11:00:01.000Z","type":"event_msg","payload":{"type":"user_message","message":"Ignored without id"}}""",
),
)
val backend = CodexRolloutSessionBackend(codexHomeProvider = { tempDir })
val beforeRewrite = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(beforeRewrite).isEmpty()
writeRollout(
file = rollout,
lines = listOf(
sessionMetaLine(timestamp = "2026-02-15T11:00:00.000Z", id = "session-retry", cwd = projectDir),
"""{"timestamp":"2026-02-15T11:00:02.000Z","type":"event_msg","payload":{"type":"user_message","message":"Recovered title"}}""",
),
)
val afterRewrite = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(afterRewrite).hasSize(1)
assertThat(afterRewrite.single().thread.id).isEqualTo("session-retry")
assertThat(afterRewrite.single().thread.title).isEqualTo("Recovered title")
}
}
@Test
fun usesFirstNonEnvironmentUserMessageAsTitle() {
runBlocking {
val projectDir = tempDir.resolve("project-title")
Files.createDirectories(projectDir)
writeRollout(
file = tempDir.resolve("sessions").resolve("2026").resolve("02").resolve("14")
.resolve("rollout-title.jsonl"),
lines = listOf(
sessionMetaLine(
timestamp = "2026-02-14T12:00:00.000Z",
id = "session-title",
cwd = projectDir,
),
"""{"timestamp":"2026-02-14T12:00:02.000Z","type":"event_msg","payload":{"type":"user_message","message":"<environment_context>\n<cwd>${projectDir.toString().replace("\\", "\\\\")}</cwd>"}}""",
"""{"timestamp":"2026-02-14T12:00:03.000Z","type":"event_msg","payload":{"type":"user_message","message":"<TURN_ABORTED>\nreason"}}""",
"""{"timestamp":"2026-02-14T12:00:04.000Z","type":"event_msg","payload":{"type":"user_message","message":"<prior context> ## My request for Codex: Real title line "}}""",
),
)
val backend = CodexRolloutSessionBackend(codexHomeProvider = { tempDir })
val threads = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(threads).hasSize(1)
assertThat(threads.single().thread.title).isEqualTo("Real title line")
}
}
@Test
fun skipsMalformedJsonLineAndKeepsParsingLaterEvents() {
runBlocking {
val projectDir = tempDir.resolve("project-malformed")
Files.createDirectories(projectDir)
writeRollout(
file = tempDir.resolve("sessions").resolve("2026").resolve("02").resolve("14")
.resolve("rollout-malformed.jsonl"),
lines = listOf(
sessionMetaLine(
timestamp = "2026-02-14T13:00:00.000Z",
id = "session-malformed",
cwd = projectDir,
),
"""{"timestamp":"2026-02-14T13:00:02.000Z","type":"event_msg","payload":{"type":"user_message","message":"Initial title"}}""",
"""{"timestamp":"2026-02-14T13:00:03.000Z","type":"event_msg","payload":{"type":"user_message"""",
"""{"timestamp":"2026-02-14T13:00:04.000Z","type":"event_msg","payload":{"type":"agent_message","message":"Still works"}}""",
),
)
val backend = CodexRolloutSessionBackend(codexHomeProvider = { tempDir })
val threads = backend.listThreads(path = projectDir.toString(), openProject = null)
assertThat(threads).hasSize(1)
val thread = threads.single()
assertThat(thread.thread.id).isEqualTo("session-malformed")
assertThat(thread.thread.title).isEqualTo("Initial title")
assertThat(thread.thread.updatedAt).isEqualTo(Instant.parse("2026-02-14T13:00:04.000Z").toEpochMilli())
assertThat(thread.activity).isEqualTo(CodexSessionActivity.UNREAD)
}
}
}
private fun sessionMetaLine(timestamp: String, id: String, cwd: Path): String {
+14
View File
@@ -0,0 +1,14 @@
### auto-generated section `build intellij.agent.workbench.json` start
load("@rules_jvm//:jvm.bzl", "jvm_library")
jvm_library(
name = "json",
module_name = "intellij.agent.workbench.json",
visibility = ["//visibility:public"],
srcs = glob(["src/**/*.kt", "src/**/*.java", "src/**/*.form"], allow_empty = True),
deps = [
"@lib//:kotlin-stdlib",
"//libraries/jackson/jackson",
]
)
### auto-generated section `build intellij.agent.workbench.json` end
@@ -0,0 +1,13 @@
<?xml version="1.0" encoding="UTF-8"?>
<module type="JAVA_MODULE" version="4">
<component name="NewModuleRootManager" inherit-compiler-output="true">
<exclude-output />
<content url="file://$MODULE_DIR$">
<sourceFolder url="file://$MODULE_DIR$/src" isTestSource="false" packagePrefix="com.intellij.agent.workbench.json" />
</content>
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="intellij.libraries.jackson" />
</component>
</module>
@@ -0,0 +1,110 @@
// Copyright 2000-2026 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
@file:Suppress("PackageDirectoryMismatch", "unused")
package com.intellij.agent.workbench.json
import com.fasterxml.jackson.core.JsonFactory
import com.fasterxml.jackson.core.JsonParser
import com.fasterxml.jackson.core.JsonToken
import java.nio.file.Files
import java.nio.file.Path
object WorkbenchJsonlScanner {
fun <S> scanJsonObjects(
path: Path,
jsonFactory: JsonFactory,
maxObjects: Int = Int.MAX_VALUE,
newState: () -> S,
onObject: (JsonParser, S) -> Boolean,
): S {
if (maxObjects <= 0) return newState()
val fastState = newState()
val parsedWithSingleParser = try {
parseWithSingleParser(path, jsonFactory, maxObjects, fastState, onObject)
true
}
catch (_: Throwable) {
false
}
if (parsedWithSingleParser) {
return fastState
}
val fallbackState = newState()
parseLineByLine(path, jsonFactory, maxObjects, fallbackState, onObject)
return fallbackState
}
private fun <S> parseWithSingleParser(
path: Path,
jsonFactory: JsonFactory,
maxObjects: Int,
state: S,
onObject: (JsonParser, S) -> Boolean,
) {
Files.newBufferedReader(path).use { reader ->
jsonFactory.createParser(reader).use { parser ->
var parsedObjects = 0
while (parsedObjects < maxObjects) {
val token = parser.nextToken() ?: return
if (token != JsonToken.START_OBJECT) {
parser.skipChildren()
continue
}
parsedObjects++
if (!onObject(parser, state)) {
return
}
}
}
}
}
private fun <S> parseLineByLine(
path: Path,
jsonFactory: JsonFactory,
maxObjects: Int,
state: S,
onObject: (JsonParser, S) -> Boolean,
) {
Files.newBufferedReader(path).use { reader ->
var parsedObjects = 0
while (parsedObjects < maxObjects) {
val line = reader.readLine() ?: return
val trimmed = line.trim()
if (trimmed.isEmpty()) continue
val shouldContinue = parseLineObject(trimmed, jsonFactory, state, onObject)
if (shouldContinue == null) {
continue
}
parsedObjects++
if (!shouldContinue) {
return
}
}
}
}
private fun <S> parseLineObject(
line: String,
jsonFactory: JsonFactory,
state: S,
onObject: (JsonParser, S) -> Boolean,
): Boolean? {
return try {
jsonFactory.createParser(line).use { parser ->
if (parser.nextToken() != JsonToken.START_OBJECT) {
return null
}
onObject(parser, state)
}
}
catch (_: Throwable) {
null
}
}
}
@@ -15,6 +15,7 @@ jvm_library(
resources = [":sessions_resources"],
deps = [
"@lib//:kotlin-stdlib",
"//libraries/fastutil",
"//libraries/kotlinx/serialization/core",
"//platform/compose",
"//platform/core-api:core",
@@ -40,6 +41,7 @@ jvm_library(
associates = [":sessions"],
deps = [
"@lib//:kotlin-stdlib",
"//libraries/fastutil",
"//libraries/kotlinx/serialization/core",
"//platform/compose",
"//platform/compose:compose_test_lib",
@@ -34,6 +34,7 @@
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="intellij.libraries.fastutil" />
<orderEntry type="module" module-name="intellij.libraries.kotlinx.serialization.core" />
<orderEntry type="module" module-name="intellij.platform.compose" />
<orderEntry type="module" module-name="intellij.platform.core" />
@@ -58,4 +59,4 @@
<orderEntry type="module" module-name="intellij.libraries.assertj.core" scope="TEST" />
<orderEntry type="module" module-name="intellij.platform.projectFrame" />
</component>
</module>
</module>
@@ -6,6 +6,7 @@ package com.intellij.agent.workbench.sessions
// @spec community/plugins/agent-workbench/spec/actions/new-thread.spec.md
import com.intellij.agent.workbench.chat.AgentChatEditorService
import com.intellij.agent.workbench.chat.AgentChatTabSelectionService
import com.intellij.agent.workbench.sessions.providers.AgentSessionProviderBridges
import com.intellij.agent.workbench.sessions.providers.AgentSessionSource
import com.intellij.ide.RecentProjectsManager
@@ -26,12 +27,17 @@ import com.intellij.openapi.ui.DoNotAskOption
import com.intellij.openapi.ui.MessageDialogBuilder
import com.intellij.openapi.util.NlsContexts
import com.intellij.openapi.util.io.FileUtilRt
import com.intellij.openapi.wm.ToolWindowManager
import it.unimi.dsi.fastutil.objects.Object2ObjectOpenHashMap
import it.unimi.dsi.fastutil.objects.ObjectOpenHashSet
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
@@ -44,6 +50,7 @@ import java.nio.file.InvalidPathException
import java.nio.file.Path
import kotlin.io.path.invariantSeparatorsPathString
import kotlin.io.path.name
import kotlin.time.Duration.Companion.milliseconds
private val LOG = logger<AgentSessionsService>()
@@ -52,6 +59,7 @@ private const val OPEN_PROJECT_ACTION_KEY_PREFIX = "project-open"
private const val CREATE_SESSION_ACTION_KEY_PREFIX = "session-create"
private const val OPEN_THREAD_ACTION_KEY_PREFIX = "thread-open"
private const val OPEN_SUB_AGENT_ACTION_KEY_PREFIX = "subagent-open"
private const val SOURCE_UPDATE_DEBOUNCE_MS = 350L
@Service(Service.Level.APP)
internal class AgentSessionsService private constructor(
@@ -85,11 +93,15 @@ internal class AgentSessionsService private constructor(
private val actionGate = SingleFlightActionGate()
private val onDemandLoading = LinkedHashSet<String>()
private val onDemandWorktreeLoading = LinkedHashSet<String>()
private val sourceRefreshJobs = Object2ObjectOpenHashMap<AgentSessionProvider, Job>()
private val sourceRefreshJobsLock = Any()
private val mutableState = MutableStateFlow(AgentSessionsState())
val state: StateFlow<AgentSessionsState> = mutableState.asStateFlow()
init {
observeSessionSourceUpdates()
if (subscribeToProjectLifecycle) {
ApplicationManager.getApplication().messageBus.connect(serviceScope)
.subscribe(ProjectManager.TOPIC, object : ProjectManagerListener {
@@ -106,6 +118,161 @@ internal class AgentSessionsService private constructor(
}
}
private fun observeSessionSourceUpdates() {
serviceScope.launch {
for (source in sessionSourcesProvider()) {
launch {
source.updates.collect {
scheduleSourceRefresh(source.provider)
}
}
}
}
}
private fun scheduleSourceRefresh(provider: AgentSessionProvider) {
synchronized(sourceRefreshJobsLock) {
sourceRefreshJobs.remove(provider)?.cancel()
val job = serviceScope.launch(Dispatchers.IO) {
delay(SOURCE_UPDATE_DEBOUNCE_MS.milliseconds)
if (!isSourceRefreshGateActive()) return@launch
refreshLoadedProviderThreads(provider)
}
sourceRefreshJobs[provider] = job
job.invokeOnCompletion {
synchronized(sourceRefreshJobsLock) {
if (sourceRefreshJobs[provider] === job) {
sourceRefreshJobs.remove(provider)
}
}
}
}
}
private suspend fun refreshLoadedProviderThreads(provider: AgentSessionProvider) {
if (!refreshMutex.tryLock()) return
try {
val source = sessionSourcesProvider().firstOrNull { it.provider == provider } ?: return
val stateSnapshot = mutableState.value
val targetPaths = collectLoadedPaths(stateSnapshot)
if (targetPaths.isEmpty()) return
val prefetched = try {
source.prefetchThreads(targetPaths)
}
catch (_: Throwable) {
emptyMap()
}
val outcomes = Object2ObjectOpenHashMap<String, ProviderRefreshOutcome>(targetPaths.size)
for (path in targetPaths) {
val prefetchedThreads = prefetched[path]
if (prefetchedThreads != null) {
outcomes[path] = ProviderRefreshOutcome(threads = prefetchedThreads)
continue
}
try {
outcomes[path] = ProviderRefreshOutcome(threads = source.listThreadsFromClosedProject(path))
}
catch (e: Throwable) {
if (e is CancellationException) throw e
LOG.warn("Failed to refresh ${provider.value} sessions for $path", e)
outcomes[path] = ProviderRefreshOutcome(
warningMessage = resolveProviderWarningMessage(provider, e),
)
}
}
mutableState.update { state ->
var changed = false
val nextProjects = state.projects.map { project ->
val updatedProject = if (project.hasLoaded) {
val outcome = outcomes[project.path]
if (outcome != null) {
changed = true
project.withProviderRefreshOutcome(provider, outcome)
}
else {
project
}
}
else {
project
}
val nextWorktrees = updatedProject.worktrees.map { worktree ->
if (!worktree.hasLoaded) return@map worktree
val outcome = outcomes[worktree.path] ?: return@map worktree
changed = true
worktree.withProviderRefreshOutcome(provider, outcome)
}
if (nextWorktrees == updatedProject.worktrees) {
updatedProject
}
else {
updatedProject.copy(worktrees = nextWorktrees)
}
}
if (!changed) {
state
}
else {
state.copy(
projects = nextProjects,
lastUpdatedAt = System.currentTimeMillis(),
)
}
}
}
finally {
refreshMutex.unlock()
}
}
private fun collectLoadedPaths(state: AgentSessionsState): List<String> {
val paths = ObjectOpenHashSet<String>()
for (project in state.projects) {
if (project.hasLoaded) {
paths.add(project.path)
}
for (worktree in project.worktrees) {
if (worktree.hasLoaded) {
paths.add(worktree.path)
}
}
}
return ArrayList(paths)
}
private suspend fun isSourceRefreshGateActive(): Boolean = withContext(Dispatchers.EDT) {
val openProjects = ProjectManager.getInstance().openProjects
if (openProjects.isEmpty()) {
val stateSnapshot = mutableState.value
return@withContext stateSnapshot.projects.any { project ->
project.isOpen || project.worktrees.any { it.isOpen }
}
}
openProjects.any { project ->
isSessionsToolWindowVisible(project) || isAgentChatActive(project)
}
}
private fun isSessionsToolWindowVisible(project: Project): Boolean {
return ToolWindowManager.getInstance(project)
.getToolWindow(AGENT_SESSIONS_TOOL_WINDOW_ID)
?.isVisible == true
}
private fun isAgentChatActive(project: Project): Boolean {
return runCatching {
project.service<AgentChatTabSelectionService>().selectedChatTab.value != null
}.getOrDefault(false)
}
fun refresh() {
serviceScope.launch(Dispatchers.IO) {
if (!refreshMutex.tryLock()) {
@@ -619,6 +786,58 @@ internal class AgentSessionsService private constructor(
return warnings + warning
}
private fun AgentProjectSessions.withProviderRefreshOutcome(
provider: AgentSessionProvider,
outcome: ProviderRefreshOutcome,
): AgentProjectSessions {
val mergedThreads = outcome.threads?.let { threads ->
mergeThreadsForProvider(this.threads, provider, threads)
} ?: this.threads
return copy(
threads = mergedThreads,
providerWarnings = replaceProviderWarning(this.providerWarnings, provider, outcome.warningMessage),
)
}
private fun AgentWorktree.withProviderRefreshOutcome(
provider: AgentSessionProvider,
outcome: ProviderRefreshOutcome,
): AgentWorktree {
val mergedThreads = outcome.threads?.let { threads ->
mergeThreadsForProvider(this.threads, provider, threads)
} ?: this.threads
return copy(
threads = mergedThreads,
providerWarnings = replaceProviderWarning(this.providerWarnings, provider, outcome.warningMessage),
)
}
private fun replaceProviderWarning(
warnings: List<AgentSessionProviderWarning>,
provider: AgentSessionProvider,
warningMessage: String?,
): List<AgentSessionProviderWarning> {
val withoutProvider = warnings.filterNot { it.provider == provider }
return if (warningMessage == null) {
withoutProvider
}
else {
withoutProvider + AgentSessionProviderWarning(provider = provider, message = warningMessage)
}
}
private fun mergeThreadsForProvider(
existingThreads: List<AgentSessionThread>,
provider: AgentSessionProvider,
newProviderThreads: List<AgentSessionThread>,
): List<AgentSessionThread> {
val mergedThreads = ArrayList<AgentSessionThread>(existingThreads.size + newProviderThreads.size)
existingThreads.filterTo(mergedThreads) { it.provider != provider }
mergedThreads.addAll(newProviderThreads)
mergedThreads.sortByDescending { it.updatedAt }
return mergedThreads
}
private suspend fun openOrFocusProjectInternal(path: String) {
val normalized = normalizePath(path)
val openProject = findOpenProject(normalized)
@@ -1167,6 +1386,11 @@ internal class AgentSessionsService private constructor(
return null
}
private data class ProviderRefreshOutcome(
val threads: List<AgentSessionThread>? = null,
val warningMessage: String? = null,
)
internal data class ProjectEntry(
val path: String,
val name: String,
@@ -0,0 +1,4 @@
// Copyright 2000-2026 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.agent.workbench.sessions.json
// Intentionally empty: scanner moved to com.intellij.agent.workbench.json.WorkbenchJsonlScanner.
@@ -4,12 +4,17 @@ package com.intellij.agent.workbench.sessions.providers
import com.intellij.agent.workbench.sessions.AgentSessionProvider
import com.intellij.agent.workbench.sessions.AgentSessionThread
import com.intellij.openapi.project.Project
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.emptyFlow
interface AgentSessionSource {
val provider: AgentSessionProvider
val canReportExactThreadCount: Boolean
get() = true
val updates: Flow<Unit>
get() = emptyFlow()
suspend fun listThreadsFromOpenProject(path: String, project: Project): List<AgentSessionThread>
suspend fun listThreadsFromClosedProject(path: String): List<AgentSessionThread>
@@ -8,6 +8,8 @@ import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.emptyFlow
import java.lang.reflect.InvocationHandler
import java.lang.reflect.Proxy
import kotlin.time.Duration.Companion.milliseconds
@@ -18,6 +20,7 @@ internal const val WORKTREE_PATH = "/work/project-feature"
internal class ScriptedSessionSource(
override val provider: AgentSessionProvider,
override val canReportExactThreadCount: Boolean = true,
override val updates: Flow<Unit> = emptyFlow(),
private val listFromOpenProject: suspend (path: String, project: Project) -> List<AgentSessionThread> = { _, _ -> emptyList() },
private val listFromClosedProject: suspend (path: String) -> List<AgentSessionThread> = { _ -> emptyList() },
) : AgentSessionSource {
@@ -2,6 +2,7 @@
package com.intellij.agent.workbench.sessions
import com.intellij.testFramework.junit5.TestApplication
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.runBlocking
import org.assertj.core.api.Assertions.assertThat
import org.junit.jupiter.api.Test
@@ -194,4 +195,71 @@ class AgentSessionsServiceRefreshIntegrationTest {
assertThat(project.threads.map { it.id }).containsExactly("claude-1")
}
}
@Test
fun providerUpdateRefreshesOnlyMatchingProviderThreads() = runBlocking {
val codexUpdates = MutableSharedFlow<Unit>(extraBufferCapacity = 1)
var codexUpdatedAt = 100L
withService(
sessionSourcesProvider = {
listOf(
ScriptedSessionSource(
provider = AgentSessionProvider.CODEX,
canReportExactThreadCount = false,
updates = codexUpdates,
listFromOpenProject = { path, _ ->
if (path == PROJECT_PATH) {
listOf(thread(id = "codex-1", updatedAt = codexUpdatedAt, provider = AgentSessionProvider.CODEX))
}
else {
emptyList()
}
},
listFromClosedProject = { path ->
if (path == PROJECT_PATH) {
listOf(thread(id = "codex-1", updatedAt = codexUpdatedAt, provider = AgentSessionProvider.CODEX))
}
else {
emptyList()
}
},
),
ScriptedSessionSource(
provider = AgentSessionProvider.CLAUDE,
listFromOpenProject = { path, _ ->
if (path == PROJECT_PATH) listOf(thread(id = "claude-1", updatedAt = 200, provider = AgentSessionProvider.CLAUDE))
else emptyList()
},
listFromClosedProject = { path ->
if (path == PROJECT_PATH) listOf(thread(id = "claude-1", updatedAt = 200, provider = AgentSessionProvider.CLAUDE))
else emptyList()
},
),
)
},
projectEntriesProvider = {
listOf(openProjectEntry(PROJECT_PATH, "Project A"))
},
) { service ->
service.refresh()
waitForCondition {
service.state.value.projects.firstOrNull { it.path == PROJECT_PATH }?.hasLoaded == true
}
codexUpdatedAt = 300L
codexUpdates.emit(Unit)
waitForCondition {
val project = service.state.value.projects.firstOrNull { it.path == PROJECT_PATH } ?: return@waitForCondition false
val codexThread = project.threads.firstOrNull { it.provider == AgentSessionProvider.CODEX } ?: return@waitForCondition false
codexThread.updatedAt == 300L
}
val project = service.state.value.projects.single { it.path == PROJECT_PATH }
assertThat(project.threads.map { it.id }).containsExactly("codex-1", "claude-1")
assertThat(project.threads.first { it.provider == AgentSessionProvider.CLAUDE }.updatedAt).isEqualTo(200L)
assertThat(project.threads.first { it.provider == AgentSessionProvider.CODEX }.updatedAt).isEqualTo(300L)
}
}
}
@@ -14,10 +14,10 @@ targets:
# Codex Sessions Rollout Source
Status: Draft
Date: 2026-02-13
Date: 2026-02-15
## Summary
Codex thread discovery for Agent Threads defaults to rollout files under `~/.codex/sessions` instead of app-server `thread/list`. The source computes thread activity (`unread`, `reviewing`, `processing`, `ready`) and maps it to indicator colors. Existing app-server loading remains available behind an alternative `SessionBackend` implementation.
Codex thread discovery for Agent Threads defaults to rollout files under `~/.codex/sessions` instead of app-server `thread/list`. The source computes thread activity (`unread`, `reviewing`, `processing`, `ready`) and maps it to indicator colors. Existing app-server loading remains available behind an alternative `SessionBackend` implementation. Title extraction semantics are aligned with Codex rollout parsing behavior.
## Goals
- Make Codex thread indicators reflect real activity based on rollout data.
@@ -40,8 +40,14 @@ Codex thread discovery for Agent Threads defaults to rollout files under `~/.cod
- Unknown backend override values must log a warning and fall back to rollout.
- Rollout backend must scan only `~/.codex/sessions/**/rollout-*.jsonl`.
- Rollout backend must filter sessions by normalized `cwd` matching project/worktree path.
- Rollout backend must support multi-path prefetch and return per-path filtered thread lists from a shared scan.
- Thread id must come from `session_meta.payload.id` (not rollout filename).
- Rollout backend must skip files missing `session_meta.payload.id` (no filename fallback).
- Title extraction must use the first qualifying `event_msg` with `payload.type=user_message`.
- Title extraction must strip `## My request for Codex:` when present and use the text after the marker.
- Title extraction must ignore session-prefix user messages starting with `<environment_context>` or `<turn_aborted>` (case-insensitive, leading whitespace ignored).
- Title extraction must trim and whitespace-normalize text, then apply bounded title trim.
- If no qualifying title is found, title must fall back to `Thread <id-prefix>`.
- Thread activity precedence must be: `unread` > `reviewing` > `processing` > `ready`.
- Session tree indicator colors must match CodexMonitor classes:
- `unread`: blue (`#4DA3FF`)
@@ -53,7 +59,9 @@ Codex thread discovery for Agent Threads defaults to rollout files under `~/.cod
## Data & Backend
- Rollout backend computes `updatedAt` from latest event timestamp with file mtime fallback.
- Rollout backend derives title from user-message content with bounded trim and fallback `Thread <id-prefix>`.
- Rollout backend derives title from the first qualifying `event_msg.user_message`, using Codex marker stripping and session-prefix filtering rules.
- `response_item` entries contribute to unread/activity timing but do not provide title source data.
- Rollout backend prefetch for multiple paths uses one filesystem scan and groups parsed threads by normalized `cwd`.
- Rollout backend carries branch from session meta when present; no branch fallback source is used.
- `CodexSessionSource` maps rollout backend data directly and does not use `CodexSessionBranchStore` fallback.
@@ -61,6 +69,8 @@ Codex thread discovery for Agent Threads defaults to rollout files under `~/.cod
- `./tests.cmd '-Dintellij.build.test.patterns=com.intellij.agent.workbench.codex.sessions.CodexRolloutSessionBackendTest'`
- `./tests.cmd '-Dintellij.build.test.patterns=com.intellij.agent.workbench.sessions.AgentSessionCliTest'`
[@test] ../codex/sessions/testSrc/CodexRolloutSessionBackendTest.kt
## References
- `spec/agent-sessions.spec.md`
- `spec/agent-chat-editor.spec.md`
@@ -14,7 +14,7 @@ targets:
# Agent Threads Tool Window
Status: Draft
Date: 2026-02-13
Date: 2026-02-15
## Summary
Define the Agent Threads tool window as a provider-agnostic, project-scoped session browser. Threads from supported providers are aggregated per project/worktree, rendered in one tree, and opened through a shared chat routing flow.
@@ -48,6 +48,7 @@ Define the Agent Threads tool window as a provider-agnostic, project-scoped sess
- Claude: `claude --resume <sessionId>`
- New-session action behavior (provider options, Codex/Claude command mapping, and Full Auto semantics) is defined in `spec/actions/new-thread.spec.md` and must be used by both project and worktree rows.
- Codex thread discovery must default to rollout session files; app-server thread discovery remains an explicit compatibility override path.
- Codex thread title normalization and filtering rules are defined in `spec/agent-sessions-codex-rollout-source.spec.md` and must be used for Codex thread rows.
- Branch mismatch between thread origin and current worktree branch must show a warning confirmation before opening chat.
[@test] ../sessions/testSrc/AgentSessionLoadAggregationTest.kt
@@ -88,3 +89,4 @@ Define the Agent Threads tool window as a provider-agnostic, project-scoped sess
- `spec/agent-dedicated-frame.spec.md`
- `spec/agent-chat-editor.spec.md`
- `spec/actions/new-thread.spec.md`
- `spec/agent-sessions-codex-rollout-source.spec.md`