Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.launch
import kotlinx.coroutines.suspendCancellableCoroutine
Expand Down Expand Up @@ -128,6 +129,17 @@ class KiloBackendAppService private constructor(
private val _appState = MutableStateFlow<KiloAppState>(KiloAppState.Disconnected)
val appState: StateFlow<KiloAppState> = _appState.asStateFlow()

/**
* Whether the connected CLI allows background subagents, from `GET /experimental/capabilities`.
*
* Deliberately not part of [AppData]: it is a property of the connected CLI rather than loaded
* app data, it is probed off the load's critical path, and folding it into [AppData] would churn
* that object's identity — which `updateConfig` and [refreshConfigState] use to detect a
* concurrent reload, so a late probe would silently cancel a config write.
*/
private val _capabilities = MutableStateFlow(false)
val capabilities: StateFlow<Boolean> = _capabilities.asStateFlow()

val events: SharedFlow<SseEvent> get() = connection.events
val api: DefaultApi? get() = connection.api
val http: OkHttpClient? get() = connection.apiClient
Expand Down Expand Up @@ -340,7 +352,10 @@ class KiloBackendAppService private constructor(
return@collect
}
when (next) {
ConnectionState.Disconnected -> _appState.value = KiloAppState.Disconnected
ConnectionState.Disconnected -> {
_capabilities.value = false
_appState.value = KiloAppState.Disconnected
}
is ConnectionState.Downloading -> _appState.value = KiloAppState.Downloading(next.percent, next.version, next.platform)
ConnectionState.Connecting -> _appState.value = KiloAppState.Connecting
is ConnectionState.Connected -> {
Expand Down Expand Up @@ -467,6 +482,12 @@ class KiloBackendAppService private constructor(
"notifications=${notifs.size} ${configSummary(cfg)}",
)
log.info("Application started — config, profile, notifications loaded")
// Off the critical path on purpose: this is an optional probe, and the generated
// client's call is blocking, so a timeout inside the load's coroutineScope could
// not actually release it — structured concurrency would still wait for the
// socket and fail an otherwise-successful load. Ready therefore starts with the
// capability off and flips once the probe answers.
cs.launch { refreshCapabilities() }
} catch (e: TimeoutCancellationException) {
val err = LoadError(
resource = "app",
Expand Down Expand Up @@ -659,6 +680,37 @@ class KiloBackendAppService private constructor(
}
}

/**
* Reads the CLI's background-subagent capability. Returns false when it cannot be read — an older
* CLI has no `/experimental/capabilities` route at all — which matches VS Code's
* `data?.backgroundSubagents === true` fallback and hides the affordance rather than offering one
* that would fail.
*
* Called off the load's critical path — see the call site in the start flow.
*/
private suspend fun fetchCapabilities(): Boolean {
val client = connection.appLoadApi ?: return false
return try {
withContext(Dispatchers.IO) { client.experimentalCapabilitiesGet().backgroundSubagents }
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
log.warn("Experimental capabilities fetch failed: ${e.message}", e)
false
}
}

/**
* Publishes the probe result into [capabilities], ignoring a late answer that arrives after the
* connection changed under us.
*/
private suspend fun refreshCapabilities() {
val conn = connection.state.value as? ConnectionState.Connected ?: return
val enabled = fetchCapabilities()
if (connection.state.value != conn) return
_capabilities.value = enabled
}

private suspend fun refreshConfigState() {
val current = _appState.value as? KiloAppState.Ready ?: return
val connection = connection.state.value as? ConnectionState.Connected ?: return
Expand Down Expand Up @@ -836,6 +888,7 @@ class KiloBackendAppService private constructor(
profile = null
config = null
notifications = emptyList()
_capabilities.value = false
_appState.value = KiloAppState.Disconnected
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import ai.kilocode.log.KiloLog
import ai.kilocode.jetbrains.api.client.DefaultApi
import ai.kilocode.jetbrains.api.model.GlobalSession
import ai.kilocode.jetbrains.api.model.SessionStatus
import ai.kilocode.rpc.dto.BackgroundJobDto
import ai.kilocode.rpc.dto.CloudSessionListDto
import ai.kilocode.rpc.dto.SessionBoardDto
import ai.kilocode.rpc.dto.SessionChangeDto
Expand All @@ -19,14 +20,23 @@ import ai.kilocode.rpc.dto.SessionSummaryDto
import ai.kilocode.rpc.dto.SessionTimeDto
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.SharingStarted
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.onCompletion
import kotlinx.coroutines.flow.shareIn
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
import kotlinx.coroutines.withContext
import kotlinx.coroutines.Dispatchers
import kotlinx.serialization.json.JsonPrimitive
import okhttp3.HttpUrl.Companion.toHttpUrl
import okhttp3.MediaType.Companion.toMediaType
Expand All @@ -53,6 +63,11 @@ class KiloBackendSessionManager(
private val cs: CoroutineScope,
private val log: KiloLog,
) {
companion object {
private const val FAST_POLL_MS = 1_000L
private const val SLOW_POLL_MS = 5_000L
}

/** Per-session directory overrides (sessionId → worktree path). */
private val directories = ConcurrentHashMap<String, String>()

Expand All @@ -71,6 +86,9 @@ class KiloBackendSessionManager(
private var base: String? = null
private var watcher: Job? = null

/** One shared, poll-backed flow per (directory, parent session) key — see [backgroundJobs]. */
private val jobFlows = ConcurrentHashMap<String, Flow<List<BackgroundJobDto>>>()

fun start(api: DefaultApi, httpClient: OkHttpClient, port: Int, events: SharedFlow<SseEvent>) {
client = api
http = httpClient
Expand Down Expand Up @@ -110,6 +128,7 @@ class KiloBackendSessionManager(
http = null
base = null
owned.clear()
jobFlows.clear()
_statuses.value = emptyMap()
log.info("Session manager stopped")
}
Expand Down Expand Up @@ -326,6 +345,112 @@ class KiloBackendSessionManager(
}
}

// ------ background subagents ------

/**
* Observe background subagent jobs owned by root session [id] in [dir].
*
* Backed by a single poller per (directory, id) pair, shared across every subscriber via
* [kotlinx.coroutines.flow.shareIn] with [SharingStarted.WhileSubscribed] — the underlying
* `flow{}` coroutine starts on first collector and stops automatically once the last one leaves,
* so an open session's own subscription and a re-opened editor tab reuse the same poll loop
* instead of hitting the CLI twice. [http]/[base] are read fresh every iteration, so the poller
* pauses while disconnected (both go null in [stop]) and resumes on the next [start] without
* needing to be recreated. Cadence adapts: 1 s while any job is running, 5 s once the list is
* empty or every job is terminal, matching the VS Code webview's poll rate for the fast case.
*/
fun backgroundJobs(id: String, dir: String): Flow<List<BackgroundJobDto>> {
val key = jobKey(dir, id)
return jobFlows.getOrPut(key) {
pollBackgroundJobs(id, dir)
// [SharingStarted.WhileSubscribed] stops the poll loop when the last collector
// leaves, which cancels this upstream and runs onCompletion — the point to drop the
// cache entry. Without it the SharedFlow and its replayed jobs list would be retained
// for every session opened during a long IDE run. A collector arriving in the same
// instant keeps working (sharing simply restarts); the next caller shares a fresh
// flow, so the race costs at most one extra poller, never correctness.
.onCompletion { jobFlows.remove(key) }
.shareIn(cs, SharingStarted.WhileSubscribed(), replay = 1)
}
}

private fun jobKey(dir: String, id: String) = "$dir\u0000$id"

private fun pollBackgroundJobs(id: String, dir: String): Flow<List<BackgroundJobDto>> = flow {
var cadence = FAST_POLL_MS
while (true) {
val h = http
val url = base
if (h != null && url != null) {
val jobs = withContext(Dispatchers.IO) {
runCatching { fetchBackgroundJobs(h, url, id, dir) }
.onFailure { log.warn("${ChatLogSummary.sid(id)} kind=background-jobs poll=true failed message=${it.message}", it) }
.getOrNull()
}
if (jobs != null) {
emit(jobs)
cadence = if (jobs.any { it.status == "running" }) FAST_POLL_MS else SLOW_POLL_MS
}
}
delay(cadence)
}
}.distinctUntilChanged()

private fun fetchBackgroundJobs(h: OkHttpClient, url: String, id: String, dir: String): List<BackgroundJobDto> {
val target = url.toHttpUrl().newBuilder()
.addPathSegment("kilocode")
.addPathSegment("background-jobs")
.addQueryParameter("directory", dir)
.addQueryParameter("sessionID", id)
.build()
val request = Request.Builder().url(target).get().build()
h.newCall(request).execute().use { response ->
val raw = response.body?.string()
if (!response.isSuccessful) {
throw RuntimeException("Background jobs list failed: HTTP ${response.code} — $raw")
}
return KiloCliDataParser.parseBackgroundJobs(raw!!)
}
}

/**
* Cancel background job [id] via `POST /kilocode/background-jobs/{id}/cancel?directory={dir}`.
* Cancels the job's child session tree, not just the job entry. Raw HTTP — these routes are
* newer than the generated client built from the pinned CLI release, matching [sessionBoard].
*/
fun cancelBackgroundJob(id: String, dir: String): Boolean =
postBackgroundJobAction(id, dir, "cancel")

/**
* Continue background job [id] in the background via
* `POST /kilocode/background-jobs/{id}/promote?directory={dir}`. Returns `false` when the CLI's
* `KILO_EXPERIMENTAL_BACKGROUND_SUBAGENTS` kill switch is off — callers must not assume success.
*/
fun promoteBackgroundJob(id: String, dir: String): Boolean =
postBackgroundJobAction(id, dir, "promote")

private fun postBackgroundJobAction(id: String, dir: String, action: String): Boolean {
val h = http ?: throw IllegalStateException("Session manager not started")
val url = base ?: throw IllegalStateException("Session manager not started")
val target = url.toHttpUrl().newBuilder()
.addPathSegment("kilocode")
.addPathSegment("background-jobs")
.addPathSegment(id)
.addPathSegment(action)
.addQueryParameter("directory", dir)
.build()
log.info("Background job $action: POST $target")
val request = Request.Builder().url(target).post(ByteArray(0).toRequestBody(null)).build()
h.newCall(request).execute().use { response ->
val raw = response.body?.string()
if (!response.isSuccessful) {
log.warn("Background job $action failed: HTTP ${response.code}, body=$raw")
throw RuntimeException("Background job $action failed: HTTP ${response.code} — $raw")
}
return raw?.trim() == "true"
}
}

fun get(id: String, dir: String): SessionDto {
val all = requireClient().sessionList(directory = dir)
val raw = all.firstOrNull { it.id == id }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import ai.kilocode.backend.workspace.ModelTerminalBenchInfo
import ai.kilocode.backend.workspace.ProviderData
import ai.kilocode.backend.workspace.ProviderInfo
import ai.kilocode.rpc.dto.AgentConfigDto
import ai.kilocode.rpc.dto.BackgroundJobDto
import ai.kilocode.rpc.dto.BoardMessageDto
import ai.kilocode.rpc.dto.ChatEventDto
import ai.kilocode.rpc.dto.CloudSessionDto
Expand Down Expand Up @@ -450,6 +451,40 @@ object KiloCliDataParser {
)
}

/**
* Parse a background-jobs list response (`GET /kilocode/background-jobs`) into
* [BackgroundJobDto]s, flattening the open `metadata` map's `sessionId`/`parentSessionId`/
* `background` keys into typed fields. A row missing `id`/`type`/`status` is dropped rather
* than failing the whole list.
*
* `started_at`/`completed_at` normally arrive as numbers, but the generated SDK types widen them
* to `"NaN"`/`"Infinity"`/`"-Infinity"` strings — [JsonObject.num] parses both via
* [String.toDoubleOrNull], and a non-finite result falls back to the given default rather than a
* garbage [Long].
*/
fun parseBackgroundJobs(raw: String): List<BackgroundJobDto> {
val arr = tryParseArray(raw) ?: return emptyList()
return arr.mapNotNull { elem ->
val row = elem.obj() ?: return@mapNotNull null
val id = row.str("id") ?: return@mapNotNull null
val type = row.str("type") ?: return@mapNotNull null
val status = row.str("status") ?: return@mapNotNull null
val metadata = row["metadata"].obj()
BackgroundJobDto(
id = id,
type = type,
status = status,
title = row.str("title"),
startedAt = row.num("started_at")?.takeIf { it.isFinite() }?.toLong() ?: 0L,
completedAt = row.num("completed_at")?.takeIf { it.isFinite() }?.toLong(),
error = row.str("error"),
sessionId = metadata?.str("sessionId"),
parentSessionId = metadata?.str("parentSessionId"),
background = metadata.bool("background"),
)
}
}

/**
* Parse message history response (`GET /session/{id}/message`)
* into a list of messages with their parts.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ import com.intellij.openapi.project.RootsChangeRescanningInfo
import com.intellij.openapi.roots.ex.ProjectRootManagerEx
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.withContext
Expand All @@ -61,7 +62,8 @@ class KiloAppRpcApiImpl : KiloAppRpcApi {
override suspend fun connect() = app.connect()

override suspend fun state(): Flow<KiloAppStateDto> =
app.appState.map(::dto).distinctUntilChanged()
combine(app.appState, app.capabilities) { state, caps -> appStateDto(state, caps) }
.distinctUntilChanged()

override suspend fun health(): HealthDto = app.health()

Expand Down Expand Up @@ -104,7 +106,7 @@ class KiloAppRpcApiImpl : KiloAppRpcApi {

override suspend fun updateConfig(patch: ConfigPatchDto): KiloAppStateDto {
app.requireReady()
return appStateDto(app.updateConfig(patch))
return appStateDto(app.updateConfig(patch), app.capabilities.value)
}

override suspend fun applyLogConfig(config: LogConfigDto) {
Expand Down Expand Up @@ -147,11 +149,9 @@ class KiloAppRpcApiImpl : KiloAppRpcApi {
service<KiloBackendTelemetry>().capture(app.http, app.port, capture.event, capture.properties)
}

private fun dto(state: KiloAppState): KiloAppStateDto =
appStateDto(state)
}

internal fun appStateDto(state: KiloAppState): KiloAppStateDto =
internal fun appStateDto(state: KiloAppState, backgroundSubagents: Boolean = false): KiloAppStateDto =
when (state) {
KiloAppState.Disconnected -> KiloAppStateDto(KiloAppStatusDto.DISCONNECTED)
is KiloAppState.Downloading -> KiloAppStateDto(
Expand Down Expand Up @@ -179,6 +179,7 @@ internal fun appStateDto(state: KiloAppState): KiloAppStateDto =
),
config = state.data.config,
profile = state.data.profile?.let(::profileDto),
backgroundSubagents = backgroundSubagents,
)
is KiloAppState.Error -> KiloAppStateDto(
status = KiloAppStatusDto.ERROR,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import ai.kilocode.backend.app.KiloBackendSessionManager
import ai.kilocode.backend.workspace.KiloBackendWorkspaceManager
import ai.kilocode.log.ChatLogSummary
import ai.kilocode.rpc.KiloSessionRpcApi
import ai.kilocode.rpc.dto.BackgroundJobDto
import ai.kilocode.rpc.dto.ChatEventDto
import ai.kilocode.rpc.dto.CloudSessionListDto
import ai.kilocode.rpc.dto.DiffFileDto
Expand Down Expand Up @@ -335,6 +336,17 @@ class KiloSessionRpcApiImpl internal constructor(
override suspend fun resetSessionBoard(sessionID: String, directory: String, revision: Int): SessionBoardDto? =
ready { withContext(Dispatchers.IO) { sessions.resetSessionBoard(sessionID, directory, revision) } }

// ------ background subagents ------

override suspend fun backgroundJobs(id: String, directory: String): Flow<List<BackgroundJobDto>> =
sessions.backgroundJobs(id, directory)

override suspend fun cancelBackgroundJob(id: String, directory: String): Boolean =
ready { withContext(Dispatchers.IO) { sessions.cancelBackgroundJob(id, directory) } }

override suspend fun promoteBackgroundJob(id: String, directory: String): Boolean =
ready { withContext(Dispatchers.IO) { sessions.promoteBackgroundJob(id, directory) } }

private suspend fun <T> ready(block: suspend () -> T): T {
app.requireReady()
return block()
Expand Down
Loading
Loading