Effect Migration for Kilo callsites (follow-up) (#10587)

* refactor(opencode): migrate ModelCache and Config to effect-native services

Remove legacy async wrapper functions from Config module and convert
ModelCache from a stateful namespace with module-level Maps into a
proper Effect service with Context/Layer semantics.

Key changes:
- Delete Config's `makeRuntime`-based async wrappers (get, getGlobal,
  update, warnings, etc.) — all callsites now use
  `Config.Service.use(...)` through AppRuntime
- Rewrite ModelCache as an Effect service with HttpClient dependency
  injection, replacing imperative Map-based caching with Effect-native
  Ref cells and TTL logic
- Convert KiloSessions.init and KilocodeBootstrap.init into proper
  Effect services with Layer-based dependency injection
- Wire ModelCache.Service into AppLayer, ProviderAuth, ModelsDev, and
  HTTP API handler layers
- Update Permission.layer to depend on Config.Service directly instead
  of calling Config async wrappers
- Add new test files for KiloSessions and ModelCache Effect integration
- Remove stale Config.get spyOn mocks from tests that no longer need
  them (experimental-session-list, recall)
- Fix indexing-auth to use typed IndexingConfig parameter instead of
  untyped record access

* fix(model-cache): resolve race conditions in concurrent fetch and cache invalidation

Introduce versioned cache cells with proper key derivation to prevent
stale responses from overwriting fresher data during concurrent fetches.

- Add version tracking to detect and discard outdated fetch results
- Derive cache keys from provider-specific options (baseURL, token, apiKey)
  to isolate concurrent requests with different credentials
- Make ModelCache.clear async to properly await invalidation across layers
- Update OrganizationDeps.clear signature to allow Promise<void> return
- Add concurrency and ordering tests for fetch/refresh race scenarios
- Rename local variable from `state` to `entry` in kilo-sessions sync loop

* chore(opencode): remove duplicate imports and fix test layer composition

Remove duplicate `AppRuntime` imports introduced during merge and update
kilo-sessions tests to use Effect-native Auth service instead of static
module calls.

- Remove duplicate `AppRuntime` import in index.ts and instance.ts
- Add Auth.defaultLayer to test layer helper
- Refactor test to yield Auth.Service and use instance methods
- Reorder Effect.provide/Effect.ensuring for correct resource cleanup

* style(opencode): normalize kilocode_change marker comments to block format

Standardize inline `// kilocode_change` annotations across source and
test files to use consistent `// kilocode_change start` / `// kilocode_change end`
block delimiters, improving readability and grep-ability of custom
modifications.
This commit is contained in:
Imanol Maiztegui
2026-05-27 11:43:32 +02:00
committed by GitHub
parent cb933859ff
commit 39a7305c97
43 changed files with 977 additions and 664 deletions
+2 -1
View File
@@ -1,6 +1,7 @@
// kilocode_change - new file
import { EOL } from "os"
import { Config } from "../../config/config"
import { AppRuntime } from "../../effect/app-runtime"
import { bootstrap } from "../bootstrap"
import { cmd } from "./cmd"
import { UI } from "../ui"
@@ -15,7 +16,7 @@ export const ConfigCommand = cmd({
describe: "check configuration for warnings and errors",
async handler() {
await bootstrap(process.cwd(), async () => {
const list = await Config.warnings()
const list = await AppRuntime.runPromise(Config.Service.use((svc) => svc.warnings()))
if (list.length === 0) {
process.stdout.write("No config warnings." + EOL)
return
-41
View File
@@ -53,7 +53,6 @@ import {
IndexingConfig as KiloIndexingConfig,
IndexingSchema as KiloIndexingSchema,
} from "@kilocode/kilo-indexing/config"
import { makeRuntime } from "@/effect/run-service"
import { unique } from "remeda"
// kilocode_change end
@@ -1100,44 +1099,4 @@ export const defaultLayer = layer.pipe(
Layer.provide(Npm.defaultLayer),
)
// kilocode_change start - keep async wrappers for Kilo callsites during Effect migration
const { runPromise } = makeRuntime(Service, defaultLayer)
export async function get() {
return runPromise((svc) => svc.get())
}
export async function getGlobal() {
return runPromise((svc) => svc.getGlobal())
}
export async function getConsoleState() {
return runPromise((svc) => svc.getConsoleState())
}
export async function update(config: Info) {
return runPromise((svc) => svc.update(config))
}
export async function updateGlobal(config: Info, options?: { dispose?: boolean }) {
return runPromise((svc) => svc.updateGlobal(config, options))
}
export async function invalidate() {
return runPromise((svc) => svc.invalidate())
}
export async function directories() {
return runPromise((svc) => svc.directories())
}
export async function waitForDependencies() {
return runPromise((svc) => svc.waitForDependencies())
}
export async function warnings() {
return runPromise((svc) => svc.warnings())
}
// kilocode_change end
export * as Config from "./config"
@@ -15,6 +15,7 @@ import { Storage } from "@/storage/storage"
import { Snapshot } from "@/snapshot"
import { Plugin } from "@/plugin"
import { ModelsDev } from "@/provider/models"
import { ModelCache } from "@/provider/model-cache" // kilocode_change
import { Provider } from "@/provider/provider"
import { ProviderAuth } from "@/provider/auth"
import { Agent } from "@/agent/agent"
@@ -68,6 +69,7 @@ export const AppLayer = Layer.mergeAll(
Storage.defaultLayer,
Snapshot.defaultLayer,
Plugin.defaultLayer,
ModelCache.defaultLayer, // kilocode_change
ModelsDev.defaultLayer,
Provider.defaultLayer,
ProviderAuth.defaultLayer,
+2 -2
View File
@@ -56,8 +56,8 @@ if (!process.env[ENV_VERSION]) {
process.env[ENV_VERSION] = InstallationVersion
}
import { Config } from "./config/config"
import { Auth } from "./auth"
import { AppRuntime } from "./effect/app-runtime"
import { Auth } from "./auth"
// kilocode_change end
import { DbCommand } from "./cli/cmd/db"
import path from "path"
@@ -148,7 +148,7 @@ let cli = yargs(args) // kilocode_change
})
// kilocode_change start - Initialize telemetry
const globalCfg = await Config.getGlobal()
const globalCfg = await AppRuntime.runPromise(Config.Service.use((cfg) => cfg.getGlobal()))
await Telemetry.init({
dataPath: Global.Path.data,
version: InstallationVersion,
@@ -15,9 +15,11 @@ import { IngestQueue } from "@/kilo-sessions/ingest-queue"
import { clearInFlightCache, withInFlightCache } from "@/kilo-sessions/inflight-cache"
import type * as SDK from "@kilocode/sdk/v2"
import z from "zod"
import { Schema } from "effect"
import { Context, Effect, Layer, Schema, Stream } from "effect"
import { KILO_API_BASE } from "@kilocode/kilo-gateway"
import { Config } from "@/config/config"
import { EffectBridge } from "@/effect/bridge"
import { InstanceState } from "@/effect/instance-state"
import { Instance } from "@/project/instance"
import { Vcs } from "@/project/vcs"
import simpleGit from "simple-git"
@@ -45,6 +47,12 @@ export namespace KiloSessions {
),
}
export interface Interface {
readonly init: () => Effect.Effect<void, unknown>
}
export class Service extends Context.Service<Service, Interface>()("@kilocode/KiloSessions") {}
const log = Log.create({ service: "kilo-sessions" })
const runtime = makeRuntime(Auth.Service, Auth.defaultLayer)
@@ -184,131 +192,127 @@ export namespace KiloSessions {
await ingest.sync(sessionID, [{ type: "session_status", data: { status } }])
}
export async function init() {
if (ingestDisabled) return
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const bus = yield* Bus.Service
const config = yield* Config.Service
const state = yield* InstanceState.make(
Effect.fn("KiloSessions.state")(function* () {
if (ingestDisabled) return
// kilocode_change start - Same type-erasure upstream uses in share/share-next.ts:165-169.
const watch = <D extends { type: string }>(def: D, fn: (evt: { properties: any }) => unknown) =>
Bus.subscribe(def as never, fn as never)
// kilocode_change end
Bus.subscribe(Session.Event.Created, (evt) => {
const sessionId = evt.properties.info.id
void create(sessionId).catch((error) => log.error("share init create failed", { sessionId, error }))
})
Bus.subscribe(Session.Event.Updated, async (evt) => {
const sessionID = evt.properties.sessionID // kilocode_change
const session = await Session.get(sessionID).catch(() => null) // kilocode_change
if (!session) return
await ingest.sync(sessionID, [
{
type: "kilo_meta",
data: await meta(sessionID),
},
{
type: "session",
data: session,
},
])
})
watch(MessageV2.Event.Updated, async (evt) => {
await ingest.sync(evt.properties.info.sessionID, [
{
type: "message",
data: evt.properties.info,
},
])
if (evt.properties.info.role === "user") {
await ingest.sync(evt.properties.info.sessionID, [
{
type: "model",
data: [
await Provider.getModel(evt.properties.info.model.providerID, evt.properties.info.model.modelID).then(
(m) => m,
const watch = <D extends { type: string }>(
def: D,
fn: (evt: { properties: any }) => unknown | Promise<unknown>,
) =>
bus.subscribe(def as never).pipe(
Stream.runForEach((evt) =>
EffectBridge.fromPromise(() => fn(evt as { properties: any })).pipe(
Effect.catchCause((cause) =>
Effect.sync(() => log.error("subscriber failed", { type: def.type, cause })),
),
),
),
],
},
])
}
})
Effect.forkScoped,
)
watch(MessageV2.Event.PartUpdated, async (evt) => {
await ingest.sync(evt.properties.part.sessionID, [
{
type: "part",
data: evt.properties.part,
},
])
})
yield* watch(Session.Event.Created, (evt) => {
const sessionID = evt.properties.info.id
return create(sessionID).catch((error) => log.error("share init create failed", { sessionID, error }))
})
yield* watch(Session.Event.Updated, async (evt) => {
const sessionID = evt.properties.sessionID
const session = await Session.get(sessionID).catch(() => null)
if (!session) return
await ingest.sync(sessionID, [
{ type: "kilo_meta", data: await meta(sessionID) },
{ type: "session", data: session },
])
})
yield* watch(MessageV2.Event.Updated, async (evt) => {
await ingest.sync(evt.properties.info.sessionID, [{ type: "message", data: evt.properties.info }])
if (evt.properties.info.role !== "user") return
const model = await Provider.getModel(
evt.properties.info.model.providerID,
evt.properties.info.model.modelID,
)
await ingest.sync(evt.properties.info.sessionID, [{ type: "model", data: [model] }])
})
yield* watch(MessageV2.Event.PartUpdated, (evt) =>
ingest.sync(evt.properties.part.sessionID, [{ type: "part", data: evt.properties.part }]),
)
yield* watch(Session.Event.Diff, (evt) =>
ingest.sync(evt.properties.sessionID, [{ type: "session_diff", data: evt.properties.diff }]),
)
yield* watch(Session.Event.TurnOpen, (evt) =>
ingest.sync(evt.properties.sessionID, [{ type: "session_open", data: {} }]),
)
yield* watch(Session.Event.TurnClose, (evt) =>
ingest.sync(evt.properties.sessionID, [{ type: "session_close", data: { reason: evt.properties.reason } }]),
)
watch(Session.Event.Diff, async (evt) => {
await ingest.sync(evt.properties.sessionID, [
{
type: "session_diff",
data: evt.properties.diff,
},
])
})
const sync = (evt: { properties: { sessionID: string } }) => {
const sessionID = evt.properties.sessionID
const current = statusSyncs.get(sessionID)
if (current?.running) {
current.dirty = true
return
}
Bus.subscribe(Session.Event.TurnOpen, async (evt) => {
await ingest.sync(evt.properties.sessionID, [{ type: "session_open", data: {} }])
})
const entry = current ?? { running: false, dirty: false }
statusSyncs.set(sessionID, entry)
Bus.subscribe(Session.Event.TurnClose, async (evt) => {
await ingest.sync(evt.properties.sessionID, [{ type: "session_close", data: { reason: evt.properties.reason } }])
})
const fail = (error: unknown) => {
const dirty = entry.dirty
statusSyncs.delete(sessionID)
log.error("status sync failed", { sessionID, error: String(error) })
if (dirty) sync(evt)
}
// Session status changes (busy/idle/retry), question lifecycle, permission lifecycle.
const syncStatus = (evt: { properties: { sessionID: string } }) => {
const sessionID = evt.properties.sessionID
const current = statusSyncs.get(sessionID)
if (current?.running) {
current.dirty = true
return
}
const loop = async () => {
entry.running = true
entry.dirty = false
await deriveAndSyncStatus(sessionID)
if (entry.dirty) {
void loop().catch(fail)
return
}
statusSyncs.delete(sessionID)
}
const state = current ?? { running: false, dirty: false }
statusSyncs.set(sessionID, state)
void loop().catch(fail)
}
yield* watch(SessionStatus.Event.Status, sync)
yield* watch(Question.Event.Asked, sync)
yield* watch(Question.Event.Replied, sync)
yield* watch(Question.Event.Rejected, sync)
yield* watch(Permission.Event.Asked, sync)
yield* watch(Permission.Event.Replied, sync)
const fail = (error: unknown) => {
const dirty = state.dirty
statusSyncs.delete(sessionID)
log.error("status sync failed", { sessionID, error: String(error) })
if (dirty) syncStatus(evt)
}
const cfg = yield* config.getGlobal()
if (remoteEnabled || cfg.remote_control) {
yield* Effect.sync(
() => void enableRemote().catch((err) => log.warn("remote not enabled", { error: String(err) })),
)
}
yield* Effect.addFinalizer(() =>
Effect.sync(() => {
statusSyncs.clear()
disableRemote()
}),
)
}),
)
const loop = async () => {
state.running = true
state.dirty = false
await deriveAndSyncStatus(sessionID)
if (state.dirty) {
void loop().catch(fail)
return
}
statusSyncs.delete(sessionID)
}
const init = Effect.fn("KiloSessions.init")(function* () {
yield* InstanceState.get(state)
})
void loop().catch(fail)
}
Bus.subscribe(SessionStatus.Event.Status, syncStatus)
Bus.subscribe(Question.Event.Asked, syncStatus)
Bus.subscribe(Question.Event.Replied, syncStatus)
Bus.subscribe(Question.Event.Rejected, syncStatus)
Bus.subscribe(Permission.Event.Asked, syncStatus)
Bus.subscribe(Permission.Event.Replied, syncStatus)
return Service.of({ init })
}),
)
const cfg = await Config.getGlobal()
if (remoteEnabled || cfg.remote_control)
enableRemote().catch((err) => log.warn("remote not enabled", { error: String(err) }))
// Use wildcard subscription so the dispose handler actually fires —
// Bus.state dispose only notifies "*" subscribers, not event-type ones.
Bus.subscribeAll((evt) => {
if (evt.type === Bus.InstanceDisposed.type) disableRemote()
})
}
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer), Layer.provide(Config.defaultLayer))
export async function enableRemote() {
if (remote) return
@@ -465,8 +465,8 @@ export async function remove(name: string) {
let found = false
// 1. Delete .md files from config directories
const { Config } = await import("../../config/config")
const dirs = await Config.directories()
const { AppRuntime } = await import("@/effect/app-runtime")
const dirs = await AppRuntime.runPromise(Config.Service.use((svc) => svc.directories()))
const patterns = ["{agent,agents}/**/" + name + ".md", "{mode,modes}/" + name + ".md"]
for (const dir of dirs) {
for (const pattern of patterns) {
+29 -5
View File
@@ -1,13 +1,37 @@
import { Cause, Context, Effect, Layer } from "effect"
import { EffectBridge } from "@/effect/bridge"
import { KiloSessions } from "@/kilo-sessions/kilo-sessions"
import * as Log from "@opencode-ai/core/util/log"
const log = Log.create({ service: "kilocode-bootstrap" })
export namespace KilocodeBootstrap {
export async function init() {
await KiloSessions.init()
void import("@/kilocode/indexing")
.then((mod) => mod.KiloIndexing.init())
.catch((err) => log.warn("indexing bootstrap failed", { err }))
export interface Interface {
readonly init: () => Effect.Effect<void, unknown>
}
export class Service extends Context.Service<Service, Interface>()("@kilocode/Bootstrap") {}
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const sessions = yield* KiloSessions.Service
const init = Effect.fn("KilocodeBootstrap.init")(function* () {
yield* sessions.init()
yield* EffectBridge.fromPromise(() =>
import("@/kilocode/indexing").then((mod) => mod.KiloIndexing.init()),
).pipe(
Effect.catchCause((cause) =>
Effect.sync(() => log.warn("indexing bootstrap failed", { err: Cause.squash(cause) })),
),
Effect.forkDetach,
)
})
return Service.of({ init })
}),
)
export const defaultLayer = layer.pipe(Layer.provide(KiloSessions.defaultLayer))
}
@@ -128,8 +128,9 @@ export namespace ConfigValidation {
async function existing(): Promise<string> {
try {
const warns = await Config.warnings()
if (!warns || warns.length === 0) return ""
const { AppRuntime } = await import("@/effect/app-runtime")
const warns = await AppRuntime.runPromise(Config.Service.use((svc) => svc.warnings()))
if (warns.length === 0) return ""
const items = warns.map((w: Config.Warning) => ` ${label(w.path)}: ${w.message}`).join("\n")
return `Pre-existing config issues (from session start):\n${items}\n\n`
} catch {
@@ -155,7 +156,6 @@ export namespace ConfigValidation {
}
const prefix = await existing()
const validation = JSONC_EXT.has(ext) ? await jsonc(filepath) : ext === ".md" ? await markdown(filepath) : ""
if (!validation) return ""
@@ -1,3 +1,5 @@
import type { IndexingConfig } from "@kilocode/kilo-indexing/config"
type Auth = unknown
type Env = {
@@ -106,9 +108,10 @@ export function shouldDefaultIndexingToKilo(indexing: unknown, auth: KiloIndexin
return !hasOtherProvider(cfg)
}
export function indexingWithKiloDefault(config: unknown, auth: KiloIndexingAuth) {
const cfg = record(config)
const indexing = cfg.indexing
export function indexingWithKiloDefault(
indexing: IndexingConfig | undefined,
auth: KiloIndexingAuth,
): IndexingConfig | undefined {
if (!shouldDefaultIndexingToKilo(indexing, auth)) return indexing
return { ...record(indexing), provider: "kilo" }
return { ...indexing, provider: "kilo" }
}
+7 -8
View File
@@ -13,6 +13,7 @@ import { fetchKiloEmbeddingModelCatalog } from "@kilocode/kilo-gateway"
import { Instance } from "@/project/instance"
import { Bus } from "@/bus"
import { Config } from "@/config/config"
import { AppRuntime } from "@/effect/app-runtime"
import { Auth } from "@/auth"
import { makeRuntime } from "@/effect/run-service"
import { registerDisposer } from "@/effect/instance-registry"
@@ -65,7 +66,7 @@ function pending(): z.infer<typeof IndexingStatus> {
}
}
async function kiloAuth(cfg: Awaited<ReturnType<typeof Config.get>>): Promise<KiloIndexingAuth> {
async function kiloAuth(cfg: Config.Info): Promise<KiloIndexingAuth> {
const info = await auth.runPromise((svc) => svc.get("kilo"))
return resolveKiloIndexingAuth({ config: cfg, auth: info })
}
@@ -221,7 +222,7 @@ export namespace KiloIndexing {
const boot = async (hit: Cache): Promise<Entry> => {
const dir = Instance.directory
const cfg = await Config.get()
const cfg = await AppRuntime.runPromise(Config.Service.use((svc) => svc.get()))
if (process.env["KILO_DISABLE_CODEBASE_INDEXING"] === "vscode-no-workspace") {
return track(hit, await inert(() => noWorkspace()))
}
@@ -246,12 +247,10 @@ export namespace KiloIndexing {
const root = path.join(Global.Path.state, "indexing")
const manager = new CodeIndexManager(dir, root)
const auth = await kiloAuth(cfg)
const globalConfig = await Config.getGlobal()
const merged = indexingWithKiloDefault(
{ ...cfg, indexing: { ...globalConfig.indexing, ...cfg.indexing } },
auth,
) as Config.Indexing | undefined
const cfgInput = await model(enrichKilo(input(merged, globalConfig.indexing), auth), auth)
const globalConfig = await AppRuntime.runPromise(Config.Service.use((svc) => svc.getGlobal()))
const global = globalConfig.indexing
const merged = indexingWithKiloDefault({ ...global, ...cfg.indexing }, auth)
const cfgInput = await model(enrichKilo(input(merged, global), auth), auth)
const box = { status: pending() as Status | undefined }
const current = () => box.status ?? normalizeIndexingStatus(manager)
let disposed = false
@@ -44,6 +44,7 @@ export const kiloGatewayHandlers = HttpApiBuilder.group(InstanceHttpApi, "kilo",
Effect.gen(function* () {
const auth = yield* Auth.Service
const store = yield* InstanceStore.Service
const cache = yield* ModelCache.Service
const profile = Effect.fn("KiloGatewayHttpApi.profile")(function* () {
const info = yield* auth.get("kilo").pipe(Effect.mapError(() => new HttpApiError.BadRequest({})))
@@ -210,7 +211,7 @@ export const kiloGatewayHandlers = HttpApiBuilder.group(InstanceHttpApi, "kilo",
})
.pipe(Effect.mapError(() => new HttpApiError.Unauthorized({})))
ModelCache.clear("kilo")
yield* cache.clear("kilo")
clearModesCache()
yield* store.disposeAll().pipe(Effect.mapError(() => new HttpApiError.Unauthorized({})))
return true
@@ -63,7 +63,10 @@ export function register(app: Hono): Hono {
Bus,
SessionCreatedEvent: Session.Event.Created,
Identifier,
ModelCache,
ModelCache: {
clear: (providerID: string) =>
AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(providerID))),
},
}),
)
}
@@ -4,6 +4,7 @@ import { describeRoute, resolver, validator } from "hono-openapi"
import z from "zod"
import { generateCommitMessage } from "../../commit-message"
import { Config } from "../../../config/config"
import { AppRuntime } from "../../../effect/app-runtime"
import { lazy } from "../../../util/lazy"
import { errors } from "../../../server/error"
@@ -39,7 +40,7 @@ export const CommitMessageRoutes = lazy(() =>
),
async (c) => {
const body = c.req.valid("json")
const config = await Config.get()
const config = await AppRuntime.runPromise(Config.Service.use((svc) => svc.get()))
const prompt = config.commit_message?.prompt || undefined
const result = await generateCommitMessage({ ...body, prompt })
return c.json({ message: result.message })
@@ -3,6 +3,7 @@
// Imported by ../../server/server.ts with minimal kilocode_change markers.
import { ModelCache } from "../../provider/model-cache"
import { AppRuntime } from "../../effect/app-runtime"
import { InstanceRuntime } from "../../project/instance-runtime"
/** Extra paths to skip request logging for */
@@ -20,7 +21,7 @@ export function corsOrigin(input: string): string | undefined {
/** Invalidate model cache and provider state after auth change */
export async function authChanged(providerID: string) {
ModelCache.clear(providerID)
await AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(providerID)))
await InstanceRuntime.disposeAllInstances()
}
+4 -3
View File
@@ -218,6 +218,7 @@ export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const bus = yield* Bus.Service
const config = yield* Config.Service // kilocode_change
const state = yield* InstanceState.make<State>(
Effect.fn("Permission.state")(function* (ctx) {
const row = Database.use((db) =>
@@ -360,7 +361,7 @@ export const layer = Layer.effect(
action: "allow" as const,
}))
if (alwaysRules.length > 0) {
yield* Effect.promise(() => Config.updateGlobal({ permission: toConfig(alwaysRules) }, { dispose: false }))
yield* config.updateGlobal({ permission: toConfig(alwaysRules) }, { dispose: false })
}
}
// kilocode_change end
@@ -399,7 +400,7 @@ export const layer = Layer.effect(
existing.saved = true // kilocode_change
if (newRules.length > 0) {
yield* Effect.promise(() => Config.updateGlobal({ permission: toConfig(newRules) }, { dispose: false }))
yield* config.updateGlobal({ permission: toConfig(newRules) }, { dispose: false })
}
yield* drainCovered(
@@ -515,7 +516,7 @@ export function disabled(tools: string[], ruleset: Ruleset): Set<string> {
return result
}
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer))
export const defaultLayer = layer.pipe(Layer.provide(Bus.layer), Layer.provide(Config.defaultLayer)) // kilocode_change
// kilocode_change start — inverse of fromConfig: convert rules back to config format
const SCALAR_ONLY_PERMISSIONS = new Set([
+3 -1
View File
@@ -30,6 +30,7 @@ export const layer = Layer.effect(
const lsp = yield* LSP.Service
const plugin = yield* Plugin.Service
const project = yield* Project.Service
const kilocode = yield* KilocodeBootstrap.Service // kilocode_change - Kilo session bootstrap replaces ShareNext
// const shareNext = yield* ShareNext.Service // kilocode_change - handled by KilocodeBootstrap
const snapshot = yield* Snapshot.Service
const vcs = yield* Vcs.Service
@@ -41,7 +42,7 @@ export const layer = Layer.effect(
yield* config.get()
// Plugin can mutate config so it has to be initialized before anything else.
yield* plugin.init()
yield* Effect.promise(() => KilocodeBootstrap.init()).pipe(Effect.forkDetach) // kilocode_change
yield* kilocode.init().pipe(Effect.catchCause((cause) => Effect.logWarning("kilocode init failed", { cause }))) // kilocode_change
// Each service self-manages its own slow work via Effect.forkScoped against
// its per-instance state scope. We just await materialization here.
yield* Effect.forEach(
@@ -65,6 +66,7 @@ export const defaultLayer: Layer.Layer<Service> = layer.pipe(
LSP.defaultLayer,
Plugin.defaultLayer,
Project.defaultLayer,
KilocodeBootstrap.defaultLayer, // kilocode_change - Kilo session bootstrap replaces ShareNext
// ShareNext.defaultLayer, // kilocode_change - handled by KilocodeBootstrap
Snapshot.defaultLayer,
Vcs.defaultLayer,
+12 -3
View File
@@ -111,11 +111,14 @@ interface State {
export class Service extends Context.Service<Service, Interface>()("@opencode/ProviderAuth") {}
export const layer: Layer.Layer<Service, never, Auth.Service | Plugin.Service> = Layer.effect(
// kilocode_change start
export const layer: Layer.Layer<Service, never, Auth.Service | Plugin.Service | ModelCache.Service> = Layer.effect(
Service,
Effect.gen(function* () {
const auth = yield* Auth.Service
const plugin = yield* Plugin.Service
const cache = yield* ModelCache.Service
// kilocode_change end
const state = yield* InstanceState.make<State>(
Effect.fn("ProviderAuth.state")(function* () {
const plugins = yield* plugin.list()
@@ -231,7 +234,7 @@ export const layer: Layer.Layer<Service, never, Auth.Service | Plugin.Service> =
}
}
Telemetry.trackAuthSuccess(input.providerID)
ModelCache.clear(input.providerID)
yield* cache.clear(input.providerID)
// kilocode_change end
})
@@ -239,8 +242,14 @@ export const layer: Layer.Layer<Service, never, Auth.Service | Plugin.Service> =
}),
)
// kilocode_change start
export const defaultLayer = Layer.suspend(() =>
layer.pipe(Layer.provide(Auth.defaultLayer), Layer.provide(Plugin.defaultLayer)),
layer.pipe(
Layer.provide(Auth.defaultLayer),
Layer.provide(Plugin.defaultLayer),
Layer.provide(ModelCache.defaultLayer),
),
)
// kilocode_change end
export * as ProviderAuth from "./auth"
+259 -335
View File
@@ -1,342 +1,266 @@
// kilocode_change - new file
import { fetchKiloModels, type KiloModelsResult } from "@kilocode/kilo-gateway"
import { Context, Duration, Effect, Layer, Schema } from "effect"
import { FetchHttpClient, HttpClient, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
import { Config } from "../config/config"
import { Auth } from "../auth"
import { makeRuntime } from "../effect/run-service"
import type { Provider } from "./models"
import * as Log from "@opencode-ai/core/util/log"
export namespace ModelCache {
const log = Log.create({ service: "model-cache" })
const auth = makeRuntime(Auth.Service, Auth.defaultLayer)
// Cache structure
const cache = new Map<
string,
{
models: Record<string, any>
timestamp: number
}
>()
const TTL = 5 * 60 * 1000 // 5 minutes
const inFlightRefresh = new Map<string, Promise<Record<string, any>>>()
// Per-provider failure tracking
const failures = new Map<string, KiloModelsResult["error"]>()
/**
* Get the failure state for a provider (undefined = no failure)
*/
export function getFailure(providerID: string): KiloModelsResult["error"] | undefined {
return failures.get(providerID)
}
/**
* Get all provider IDs that have a failure state
*/
export function failedProviders(): string[] {
return [...failures.keys()]
}
/**
* Get cached models if available and not expired
* @param providerID - Provider identifier (e.g., "kilo")
* @returns Cached models or undefined if cache miss or expired
*/
export function get(providerID: string): Record<string, any> | undefined {
const cached = cache.get(providerID)
if (!cached) {
log.debug("cache miss", { providerID })
return undefined
}
const now = Date.now()
const age = now - cached.timestamp
if (age > TTL) {
log.debug("cache expired", { providerID, age })
cache.delete(providerID)
return undefined
}
log.debug("cache hit", { providerID, age })
return cached.models
}
/**
* Fetch models with cache-first approach
* @param providerID - Provider identifier
* @param options - Provider options
* @returns Models from cache or freshly fetched
*/
export async function fetch(providerID: string, options?: any): Promise<Record<string, any>> {
// Check cache first
const cached = get(providerID)
if (cached) {
return cached
}
// Cache miss - fetch models
log.info("fetching models", { providerID })
const authOptions = await getAuthOptions(providerID).catch((err) => {
log.warn("getAuthOptions failed", { providerID, err })
return {}
})
const mergedOptions = { ...authOptions, ...options }
const result = await fetchModels(providerID, mergedOptions)
const { models } = result
if (result.error) {
failures.set(providerID, result.error)
log.warn("model fetch error", { providerID, error: result.error })
} else {
failures.delete(providerID)
}
// Store in cache (even on error, to avoid hammering the API)
cache.set(providerID, {
models,
timestamp: Date.now(),
})
log.info("models fetched and cached", { providerID, count: Object.keys(models).length })
return models
}
/**
* Force refresh models (bypass cache)
* Uses atomic refresh pattern to prevent race conditions
* @param providerID - Provider identifier
* @param options - Provider options
* @returns Freshly fetched models
*/
export async function refresh(providerID: string, options?: any): Promise<Record<string, any>> {
// Check if refresh already in progress
const existing = inFlightRefresh.get(providerID)
if (existing) {
log.debug("refresh already in progress, returning existing promise", { providerID })
return existing
}
// Create new refresh promise
const refreshPromise = (async () => {
log.info("refreshing models", { providerID })
const authOptions = await getAuthOptions(providerID).catch((err) => {
log.warn("getAuthOptions failed during refresh", { providerID, err })
return {}
})
const mergedOptions = { ...authOptions, ...options }
const result = await fetchModels(providerID, mergedOptions)
const { models } = result
if (result.error) {
failures.set(providerID, result.error)
log.warn("model refresh error", { providerID, error: result.error })
} else {
failures.delete(providerID)
}
cache.set(providerID, {
models,
timestamp: Date.now(),
})
log.info("models refreshed", { providerID, count: Object.keys(models).length })
return models
})()
// Track in-flight refresh
inFlightRefresh.set(providerID, refreshPromise)
try {
return await refreshPromise
} finally {
// Clean up in-flight tracking
inFlightRefresh.delete(providerID)
}
}
/**
* Clear cached models for a provider
* @param providerID - Provider identifier
*/
export function clear(providerID: string): void {
const deleted = cache.delete(providerID)
failures.delete(providerID)
if (deleted) {
log.info("cache cleared", { providerID })
} else {
log.debug("no cache to clear", { providerID })
}
}
/**
* Fetch models based on provider type
* @param providerID - Provider identifier
* @param options - Provider options
* @returns Fetched models
*/
async function fetchModels(providerID: string, options: any): Promise<KiloModelsResult> {
if (providerID === "kilo") {
return fetchKiloModels(options)
}
// kilocode_change start
if (providerID === "apertis") {
const models = await fetchApertisModels(options)
return { models }
}
// kilocode_change end
// Other providers not implemented yet
log.debug("provider not implemented", { providerID })
return { models: {} }
}
// kilocode_change start
const APERTIS_BASE_URL = "https://api.apertis.ai/v1"
async function fetchApertisModels(options: any): Promise<Record<string, any>> {
const baseURL = options.baseURL ?? APERTIS_BASE_URL
const apiKey = options.apiKey
if (!apiKey) {
log.debug("no API key for apertis, skipping model fetch")
return {}
}
const url = `${baseURL.replace(/\/+$/, "")}/models`
const response = await fetch(url, {
headers: {
Authorization: `Bearer ${apiKey}`,
},
signal: AbortSignal.timeout(10_000),
})
if (!response.ok) {
log.error("apertis model fetch failed", { status: response.status })
return {}
}
const json = (await response.json()) as { data?: Array<{ id: string; owned_by?: string }> }
const models: Record<string, any> = {}
for (const model of json.data ?? []) {
models[model.id] = {
id: model.id,
name: model.id,
family: model.owned_by ?? "",
release_date: "",
attachment: true,
reasoning: false,
temperature: true,
tool_call: true,
cost: { input: 0, output: 0 },
limit: { context: 128000, output: 4096 },
options: {},
modalities: {
input: ["text", "image"],
output: ["text"],
},
}
}
return models
}
// kilocode_change end
/**
* Get authentication options from multiple sources
* Priority: Config > Auth > Env
* @param providerID - Provider identifier
* @returns Options object with authentication credentials
*/
async function getAuthOptions(providerID: string): Promise<any> {
const options: any = {}
const getAuth = (id: string) => auth.runPromise((svc) => svc.get(id))
if (providerID === "kilo") {
// Get from Config
const config = await Config.get()
const providerConfig = config.provider?.[providerID]
if (providerConfig?.options?.apiKey) {
options.kilocodeToken = providerConfig.options.apiKey
}
// kilocode_change start
if (providerConfig?.options?.kilocodeOrganizationId) {
options.kilocodeOrganizationId = providerConfig.options.kilocodeOrganizationId
}
// kilocode_change end
// Get from Auth
const auth = await getAuth(providerID)
if (auth) {
if (auth.type === "api") {
options.kilocodeToken = auth.key
} else if (auth.type === "oauth") {
options.kilocodeToken = auth.access
// kilocode_change start - read org ID from OAuth accountId for enterprise model filtering
if (auth.accountId) {
options.kilocodeOrganizationId = auth.accountId
}
// kilocode_change end
}
}
// Get from Env (process.env — matches upstream's pattern for sync async helpers)
const env = process.env
if (env.KILO_API_KEY) {
options.kilocodeToken = env.KILO_API_KEY
}
if (env.KILO_ORG_ID) {
options.kilocodeOrganizationId = env.KILO_ORG_ID
}
log.debug("auth options resolved", {
providerID,
hasToken: !!options.kilocodeToken,
hasOrganizationId: !!options.kilocodeOrganizationId,
})
}
// kilocode_change start
if (providerID === "apertis") {
const config = await Config.get()
const providerConfig = config.provider?.[providerID]
if (providerConfig?.options?.apiKey) {
options.apiKey = providerConfig.options.apiKey
}
if (providerConfig?.options?.baseURL) {
options.baseURL = providerConfig.options.baseURL
}
const auth = await getAuth(providerID)
if (auth && auth.type === "api") {
options.apiKey = auth.key
}
const env = process.env
if (env.APERTIS_API_KEY) {
options.apiKey = env.APERTIS_API_KEY
}
if (env.APERTIS_BASE_URL) {
options.baseURL = env.APERTIS_BASE_URL
}
log.debug("apertis auth options resolved", {
providerID,
hasKey: !!options.apiKey,
hasBaseURL: !!options.baseURL,
})
}
// kilocode_change end
return options
}
type Models = Provider["models"]
type KiloOptions = NonNullable<Parameters<typeof fetchKiloModels>[0]>
type Options = { -readonly [K in keyof KiloOptions]?: KiloOptions[K] } & { apiKey?: string }
type Failure = NonNullable<KiloModelsResult["error"]>
type Result = { readonly models: Models; readonly error?: Failure }
type View = { models?: Models; timestamp?: number }
type Cell = {
readonly providerID: string
readonly view: View
readonly cached: Effect.Effect<Result, unknown>
readonly invalidate: Effect.Effect<void>
}
export interface Interface {
readonly getFailure: (providerID: string) => Effect.Effect<Failure | undefined>
readonly failedProviders: () => Effect.Effect<string[]>
readonly get: (providerID: string) => Effect.Effect<Models | undefined>
readonly fetch: (providerID: string, options?: Options) => Effect.Effect<Models, unknown>
readonly refresh: (providerID: string, options?: Options) => Effect.Effect<Models, unknown>
readonly clear: (providerID: string) => Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@kilocode/ModelCache") {}
const log = Log.create({ service: "model-cache" })
const ttl = Duration.minutes(5)
const APERTIS_BASE_URL = "https://api.apertis.ai/v1"
const ApertisItem = Schema.Struct({ id: Schema.String, owned_by: Schema.optional(Schema.String) })
const ApertisResponse = Schema.Struct({ data: Schema.optional(Schema.Array(ApertisItem)) })
type ApertisItem = Schema.Schema.Type<typeof ApertisItem>
export const layer: Layer.Layer<Service, never, Auth.Service | Config.Service | HttpClient.HttpClient> = Layer.effect(
Service,
Effect.gen(function* () {
const auth = yield* Auth.Service
const cfg = yield* Config.Service
const http = yield* HttpClient.HttpClient
const cells = new Map<string, Cell>()
const active = new Map<string, Cell>()
const versions = new Map<string, number>()
const failures = new Map<string, Failure>()
const getFailure = Effect.fn("ModelCache.getFailure")(function* (providerID: string) {
return failures.get(providerID)
})
const failedProviders = Effect.fn("ModelCache.failedProviders")(function* () {
return [...failures.keys()]
})
const aperture = (item: ApertisItem): Models[string] => ({
id: item.id,
name: item.id,
family: item.owned_by ?? "",
release_date: "",
attachment: true,
reasoning: false,
temperature: true,
tool_call: true,
cost: { input: 0, output: 0 },
limit: { context: 128000, output: 4096 },
modalities: { input: ["text", "image"], output: ["text"] },
})
const fetchApertisModels = Effect.fn("ModelCache.fetchApertisModels")(function* (options: Options) {
const baseURL = options.baseURL ?? APERTIS_BASE_URL
if (!options.apiKey) {
log.debug("no API key for apertis, skipping model fetch")
return {}
}
const url = `${baseURL.replace(/\/+$/, "")}/models`
const response = yield* HttpClientRequest.get(url).pipe(
HttpClientRequest.acceptJson,
HttpClientRequest.bearerToken(options.apiKey),
http.execute,
Effect.timeout("10 seconds"),
)
if (response.status < 200 || response.status >= 300) {
log.error("apertis model fetch failed", { status: response.status })
return {}
}
const json = yield* HttpClientResponse.schemaBodyJson(ApertisResponse)(response)
return Object.fromEntries((json.data ?? []).map((item) => [item.id, aperture(item)]))
})
const authOptions = Effect.fn("ModelCache.authOptions")(function* (providerID: string) {
if (providerID !== "kilo" && providerID !== "apertis") return {}
const config = yield* cfg.get()
const options: Options = {}
if (providerID === "kilo") {
const item = config.provider?.[providerID]
if (item?.options?.apiKey) options.kilocodeToken = item.options.apiKey
if (item?.options?.kilocodeOrganizationId) options.kilocodeOrganizationId = item.options.kilocodeOrganizationId
const info = yield* auth.get(providerID)
if (info?.type === "api") options.kilocodeToken = info.key
if (info?.type === "oauth") {
options.kilocodeToken = info.access
if (info.accountId) options.kilocodeOrganizationId = info.accountId
}
if (process.env.KILO_API_KEY) options.kilocodeToken = process.env.KILO_API_KEY
if (process.env.KILO_ORG_ID) options.kilocodeOrganizationId = process.env.KILO_ORG_ID
log.debug("auth options resolved", {
providerID,
hasToken: !!options.kilocodeToken,
hasOrganizationId: !!options.kilocodeOrganizationId,
})
}
if (providerID === "apertis") {
const item = config.provider?.[providerID]
if (item?.options?.apiKey) options.apiKey = item.options.apiKey
if (item?.options?.baseURL) options.baseURL = item.options.baseURL
const info = yield* auth.get(providerID)
if (info?.type === "api") options.apiKey = info.key
if (process.env.APERTIS_API_KEY) options.apiKey = process.env.APERTIS_API_KEY
if (process.env.APERTIS_BASE_URL) options.baseURL = process.env.APERTIS_BASE_URL
log.debug("apertis auth options resolved", {
providerID,
hasKey: !!options.apiKey,
hasBaseURL: !!options.baseURL,
})
}
return options
})
const fetchModels = (providerID: string, options: Options): Effect.Effect<Result, unknown> => {
if (providerID === "kilo") return Effect.tryPromise(() => fetchKiloModels(options))
if (providerID === "apertis") return fetchApertisModels(options).pipe(Effect.map((models) => ({ models })))
log.debug("provider not implemented", { providerID })
return Effect.succeed({ models: {} })
}
const load = Effect.fn("ModelCache.load")(function* (providerID: string, options: Options) {
const resolved = yield* authOptions(providerID).pipe(
Effect.catchCause((cause) =>
Effect.sync(() => {
log.warn("auth options failed", { providerID, cause })
return {}
}),
),
)
return yield* fetchModels(providerID, { ...resolved, ...options })
})
const key = (providerID: string, options?: Options) => {
if (providerID === "kilo") {
return JSON.stringify([providerID, options?.baseURL, options?.kilocodeOrganizationId, options?.kilocodeToken])
}
if (providerID === "apertis") return JSON.stringify([providerID, options?.baseURL, options?.apiKey])
return providerID
}
const cell = Effect.fn("ModelCache.cell")(function* (providerID: string, options: Options = {}) {
const id = key(providerID, options)
const existing = cells.get(id)
if (existing) return existing
const view: View = {}
const [cached, invalidate] = yield* Effect.cachedInvalidateWithTTL(load(providerID, options), ttl)
const next = { providerID, view, cached, invalidate }
cells.set(id, next)
return next
})
// Failed loads are not cached so a temporary outage can recover on the next read.
const evaluate = (entry: Cell) => entry.cached.pipe(Effect.tapCause(() => entry.invalidate))
const commit = (providerID: string, version: number, entry: Cell, result: Result) =>
Effect.sync(() => {
if ((versions.get(providerID) ?? 0) !== version) return result.models
if (result.error) {
failures.set(providerID, result.error)
log.warn("model fetch error", { providerID, error: result.error })
} else {
failures.delete(providerID)
}
entry.view.models = result.models
entry.view.timestamp = Date.now()
active.set(providerID, entry)
log.info("models fetched and cached", { providerID, count: Object.keys(result.models).length })
return result.models
})
const get = Effect.fn("ModelCache.get")(function* (providerID: string) {
const entry = active.get(providerID)
if (!entry?.view.models || entry.view.timestamp === undefined) {
log.debug("cache miss", { providerID })
return
}
const age = Date.now() - entry.view.timestamp
if (age > Duration.toMillis(ttl)) {
log.debug("cache expired", { providerID, age })
entry.view.models = undefined
entry.view.timestamp = undefined
yield* entry.invalidate
return
}
log.debug("cache hit", { providerID, age })
return entry.view.models
})
const fetch = Effect.fn("ModelCache.fetch")(function* (providerID: string, options?: Options) {
const cached = yield* get(providerID)
if (cached) return cached
const version = (versions.get(providerID) ?? 0) + 1
versions.set(providerID, version)
const entry = yield* cell(providerID, options)
log.info("fetching models", { providerID })
const result = yield* evaluate(entry)
return yield* commit(providerID, version, entry, result)
})
const refresh = Effect.fn("ModelCache.refresh")(function* (providerID: string, options?: Options) {
const version = (versions.get(providerID) ?? 0) + 1
versions.set(providerID, version)
const entry = yield* cell(providerID, options)
log.info("refreshing models", { providerID })
yield* entry.invalidate
const result = yield* evaluate(entry)
return yield* commit(providerID, version, entry, result)
})
const clear = Effect.fn("ModelCache.clear")(function* (providerID: string) {
versions.set(providerID, (versions.get(providerID) ?? 0) + 1)
const entries = [...cells.entries()].filter(([, entry]) => entry.providerID === providerID)
yield* Effect.all(
entries.map(([id, entry]) => entry.invalidate.pipe(Effect.tap(() => Effect.sync(() => cells.delete(id))))),
{ discard: true },
)
active.delete(providerID)
failures.delete(providerID)
if (entries.some(([, entry]) => entry.view.models)) {
log.info("cache cleared", { providerID })
return
}
log.debug("no cache to clear", { providerID })
})
return Service.of({ getFailure, failedProviders, get, fetch, refresh, clear })
}),
)
export const defaultLayer = layer.pipe(
Layer.provide(FetchHttpClient.layer),
Layer.provide(Auth.defaultLayer),
Layer.provide(Config.defaultLayer),
)
export * as ModelCache from "./model-cache"
+15 -10
View File
@@ -123,11 +123,15 @@ export interface Interface {
export class Service extends Context.Service<Service, Interface>()("@opencode/ModelsDev") {}
export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClient.HttpClient | Auth.Service> = Layer.effect( // kilocode_change
type Requirements = AppFileSystem.Service | HttpClient.HttpClient | Config.Service | Auth.Service | ModelCache.Service // kilocode_change
export const layer: Layer.Layer<Service, never, Requirements> = Layer.effect(
Service,
Effect.gen(function* () {
const fs = yield* AppFileSystem.Service
const cfg = yield* Config.Service // kilocode_change
const auth = yield* Auth.Service // kilocode_change
const cache = yield* ModelCache.Service // kilocode_change
const http = HttpClient.filterStatusOk(withTransientReadRetry(yield* HttpClient.HttpClient))
const source = Flag.KILO_MODELS_URL || "https://models.dev"
@@ -195,7 +199,7 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
const providers = { ...(yield* cachedGet) }
delete providers["kilo"]
const config = yield* Effect.promise(() => Config.get())
const config = yield* cfg.get()
const disabled = new Set(config.disabled_providers ?? [])
const enabled = config.enabled_providers ? new Set(config.enabled_providers) : undefined
const kiloAllowed = (!enabled || enabled.has("kilo")) && !disabled.has("kilo")
@@ -207,9 +211,8 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
if (kiloAllowed) {
const opts = config.provider?.kilo?.options
const info = yield* auth.get("kilo").pipe(Effect.orDie)
const info = yield* auth.get("kilo").pipe(Effect.catch(() => Effect.succeed(undefined)))
const org = opts?.kilocodeOrganizationId ?? (info?.type === "oauth" ? info.accountId : undefined)
const base = normalizeKiloBaseURL(opts?.baseURL, org)
const fetch = {
...(base ? { baseURL: base } : {}),
@@ -217,10 +220,10 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
}
const [kilo, apertis] = yield* Effect.all(
[
Effect.promise(() => ModelCache.fetch("kilo", fetch).catch(() => ({}))),
cache.fetch("kilo", fetch).pipe(Effect.catch(() => Effect.succeed({}))),
providers["apertis"]
? Effect.succeed(null)
: Effect.promise(() => ModelCache.fetch("apertis", aptFetch).catch(() => ({}))),
: cache.fetch("apertis", aptFetch).pipe(Effect.catch(() => Effect.succeed({}))),
],
{ concurrency: 2 },
)
@@ -234,7 +237,7 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
models: kilo,
}
if (Object.keys(kilo).length === 0) {
yield* Effect.sync(() => void ModelCache.refresh("kilo", fetch).catch(() => {}))
yield* cache.refresh("kilo", fetch).pipe(Effect.ignore, Effect.forkDetach)
}
if (!providers["apertis"] && apertis !== null) {
providers["apertis"] = {
@@ -246,14 +249,14 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
models: apertis,
}
if (Object.keys(apertis).length === 0) {
yield* Effect.sync(() => void ModelCache.refresh("apertis", aptFetch).catch(() => {}))
yield* cache.refresh("apertis", aptFetch).pipe(Effect.ignore, Effect.forkDetach)
}
}
return providers
}
if (!providers["apertis"]) {
const apertis = yield* Effect.promise(() => ModelCache.fetch("apertis", aptFetch).catch(() => ({})))
const apertis = yield* cache.fetch("apertis", aptFetch).pipe(Effect.catch(() => Effect.succeed({})))
providers["apertis"] = {
id: "apertis",
name: "Apertis",
@@ -263,7 +266,7 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
models: apertis,
}
if (Object.keys(apertis).length === 0) {
yield* Effect.sync(() => void ModelCache.refresh("apertis", aptFetch).catch(() => {}))
yield* cache.refresh("apertis", aptFetch).pipe(Effect.ignore, Effect.forkDetach)
}
}
return providers
@@ -299,7 +302,9 @@ export const layer: Layer.Layer<Service, never, AppFileSystem.Service | HttpClie
export const defaultLayer: Layer.Layer<Service> = layer.pipe(
Layer.provide(FetchHttpClient.layer),
Layer.provide(AppFileSystem.defaultLayer),
Layer.provide(Config.defaultLayer), // kilocode_change
Layer.provide(Auth.defaultLayer), // kilocode_change
Layer.provide(ModelCache.defaultLayer), // kilocode_change
)
export * as ModelsDev from "./models"
@@ -102,9 +102,11 @@ export const ConfigRoutes = lazy(() =>
},
},
}),
async (c) => {
return c.json(await Config.warnings())
},
async (c) =>
jsonRequest("ConfigRoutes.warnings", c, function* () {
const cfg = yield* Config.Service
return yield* cfg.warnings()
}),
)
// kilocode_change end
.get(
@@ -16,6 +16,7 @@ export const providerHandlers = HttpApiBuilder.group(InstanceHttpApi, "provider"
const cfg = yield* Config.Service
const provider = yield* Provider.Service
const svc = yield* ProviderAuth.Service
const cache = yield* ModelCache.Service // kilocode_change
const list = Effect.fn("ProviderHttpApi.list")(function* () {
const config = yield* cfg.get()
@@ -32,7 +33,7 @@ export const providerHandlers = HttpApiBuilder.group(InstanceHttpApi, "provider"
connected,
)
// kilocode_change start
const failed = ModelCache.failedProviders()
const failed = yield* cache.failedProviders()
// Note: connected only contains providers with non-empty models after Provider.Service.list(),
// so failed must be checked explicitly for providers whose fetch returned an error.
const failedSet = new Set(failed)
@@ -24,6 +24,7 @@ import { Plugin } from "@/plugin"
import { Project } from "@/project/project"
import { ProviderAuth } from "@/provider/auth"
import { ModelsDev } from "@/provider/models"
import { ModelCache } from "@/provider/model-cache" // kilocode_change
import { Provider } from "@/provider/provider"
import { Pty } from "@/pty"
import { PtyTicket } from "@/pty/ticket"
@@ -163,6 +164,7 @@ export function createRoutes(corsOptions?: CorsOptions) {
LSP.defaultLayer,
Installation.defaultLayer,
MCP.defaultLayer,
ModelCache.defaultLayer, // kilocode_change
ModelsDev.defaultLayer,
Permission.defaultLayer,
Plugin.defaultLayer,
@@ -37,6 +37,7 @@ export const ProviderRoutes = lazy(() =>
jsonRequest("ProviderRoutes.list", c, function* () {
const svc = yield* Provider.Service
const cfg = yield* Config.Service
const cache = yield* ModelCache.Service // kilocode_change
const config = yield* cfg.get()
const all = yield* ModelsDev.Service.use((s) => s.get())
const disabled = new Set(config.disabled_providers ?? [])
@@ -53,7 +54,7 @@ export const ProviderRoutes = lazy(() =>
connected,
)
// kilocode_change start
const failed = ModelCache.failedProviders()
const failed = yield* cache.failedProviders()
// Keep connected or failed providers even when they have 0 models so /connect can re-auth them.
// Note: connected only contains providers whose model list is non-empty after Provider.Service.list(),
// so failed must be checked explicitly for providers whose fetch returned an error.
+20 -8
View File
@@ -78,6 +78,10 @@ const clear = async (wait = false) => {
}
const listDirs = () =>
Effect.runPromise(Config.Service.use((svc) => svc.directories()).pipe(Effect.scoped, Effect.provide(layer)))
// kilocode_change start
const warnings = () =>
Effect.runPromise(Config.Service.use((svc) => svc.warnings()).pipe(Effect.scoped, Effect.provide(layer)))
// kilocode_change end
const ready = () =>
Effect.runPromise(Config.Service.use((svc) => svc.waitForDependencies()).pipe(Effect.scoped, Effect.provide(layer)))
@@ -386,6 +390,7 @@ test("jsonc overrides json in the same directory", async () => {
})
})
// kilocode_change start
test("prefers .kilo directory config over legacy .kilocode", async () => {
await using tmp = await tmpdir({
init: async (dir) => {
@@ -409,11 +414,12 @@ test("prefers .kilo directory config over legacy .kilocode", async () => {
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const config = await Config.get()
const config = await load()
expect(config.model).toBe("new/model")
},
})
})
// kilocode_change end
test("handles environment variable substitution", async () => {
const originalEnv = process.env["TEST_VAR"]
@@ -585,6 +591,7 @@ test("handles file inclusion with replacement tokens", async () => {
})
})
// kilocode_change start
test("validates config schema and reports warning on invalid fields", async () => {
await using tmp = await tmpdir({
init: async (dir) => {
@@ -597,14 +604,16 @@ test("validates config schema and reports warning on invalid fields", async () =
await provideTestInstance({
directory: tmp.path,
fn: async () => {
// kilocode_change - invalid schema surfaces as warnings, not a throw
// invalid schema surfaces as warnings, not a throw
await load()
const warnings = await Config.warnings()
expect(warnings.length).toBeGreaterThan(0)
const issues = await warnings()
expect(issues.length).toBeGreaterThan(0)
},
})
})
// kilocode_change end
// kilocode_change start
test("reports warning for invalid JSON", async () => {
await using tmp = await tmpdir({
init: async (dir) => {
@@ -614,13 +623,14 @@ test("reports warning for invalid JSON", async () => {
await provideTestInstance({
directory: tmp.path,
fn: async () => {
// kilocode_change - invalid JSON surfaces as a warning, not a throw
// invalid JSON surfaces as a warning, not a throw
await load()
const warnings = await Config.warnings()
expect(warnings.length).toBeGreaterThan(0)
const issues = await warnings()
expect(issues.length).toBeGreaterThan(0)
},
})
})
// kilocode_change end
test("handles agent configuration", async () => {
await using tmp = await tmpdir({
@@ -965,6 +975,7 @@ Nested command template`,
})
})
// kilocode_change start
test("prefers .kilo commands over legacy .kilocode commands", async () => {
await using tmp = await tmpdir({
init: async (dir) => {
@@ -988,7 +999,7 @@ Hello from new command`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const config = await Config.get()
const config = await load()
expect(config.command?.["hello"]).toEqual({
description: "New command",
@@ -997,6 +1008,7 @@ Hello from new command`,
},
})
})
// kilocode_change end
test("gets config directories", async () => {
await using tmp = await tmpdir()
@@ -1,13 +1,17 @@
import { afterEach, describe, expect, test } from "bun:test"
import path from "path"
import { Config } from "../../src/config/config"
import { AppRuntime } from "../../src/effect/app-runtime"
import { WithInstance } from "../../src/project/with-instance"
import { Filesystem } from "../../src/util/filesystem"
import { disposeAllInstances, tmpdir } from "../fixture/fixture"
const load = () => AppRuntime.runPromise(Config.Service.use((svc) => svc.get()))
const warnings = () => AppRuntime.runPromise(Config.Service.use((svc) => svc.warnings()))
afterEach(async () => {
await disposeAllInstances()
await Config.invalidate()
await AppRuntime.runPromise(Config.Service.use((svc) => svc.invalidate()))
})
describe("config resilience", () => {
@@ -34,7 +38,7 @@ Valid agent prompt`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const cfg = await Config.get()
const cfg = await load()
expect(cfg.agent?.["skip"]).toBeUndefined()
expect(cfg.agent?.["keep"]).toMatchObject({
@@ -62,8 +66,8 @@ Broken agent prompt`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
await Config.get()
const warns = await Config.warnings()
await load()
const warns = await warnings()
expect(warns.some((w) => w.path.includes("skip.md") && w.message.includes("mode"))).toBe(true)
},
@@ -93,7 +97,7 @@ Valid command template`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const cfg = await Config.get()
const cfg = await load()
expect(cfg.command?.["skip"]).toBeUndefined()
expect(cfg.command?.["keep"]).toEqual({
@@ -120,8 +124,8 @@ Broken command template`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
await Config.get()
const warns = await Config.warnings()
await load()
const warns = await warnings()
expect(warns.some((w) => w.path.includes("skip.md") && w.message.includes("subtask"))).toBe(true)
},
@@ -144,8 +148,8 @@ Broken agent`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
await Config.get()
const warns = await Config.warnings()
await load()
const warns = await warnings()
expect(warns.some((w) => w.path.includes("broken.md") && w.message.includes("invalid"))).toBe(true)
},
@@ -168,8 +172,8 @@ Broken command`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
await Config.get()
const warns = await Config.warnings()
await load()
const warns = await warnings()
expect(warns.some((w) => w.path.includes("broken.md") && w.message.includes("invalid"))).toBe(true)
},
@@ -186,8 +190,8 @@ Broken command`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const cfg = await Config.get()
const warns = await Config.warnings()
const cfg = await load()
const warns = await warnings()
// Config loading should not crash
expect(cfg).toBeDefined()
@@ -207,8 +211,8 @@ Broken command`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
const cfg = await Config.get()
const warns = await Config.warnings()
const cfg = await load()
const warns = await warnings()
expect(cfg).toBeDefined()
expect(warns.some((w) => w.path.includes("kilo.json") && w.message.includes("invalid"))).toBe(true)
@@ -224,8 +228,8 @@ Broken command`,
await WithInstance.provide({
directory: tmp.path,
fn: async () => {
await Config.get()
const warns = await Config.warnings()
await load()
const warns = await warnings()
expect(warns).toEqual([])
},
@@ -4,6 +4,7 @@ import path from "path"
import { ConfigValidation } from "../../src/kilocode/config-validation"
import { WithInstance } from "../../src/project/with-instance"
import { Config } from "../../src/config/config"
import { AppRuntime } from "../../src/effect/app-runtime"
import { Filesystem } from "../../src/util/filesystem"
import { disposeAllInstances, tmpdir } from "../fixture/fixture"
@@ -11,6 +12,8 @@ afterEach(async () => {
await disposeAllInstances()
})
const check = (filepath: string) => ConfigValidation.check(filepath)
describe("ConfigValidation.check", () => {
test("returns empty string for non-config files", async () => {
await using tmp = await tmpdir({ git: true })
@@ -19,7 +22,7 @@ describe("ConfigValidation.check", () => {
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toBe("")
})
@@ -31,7 +34,7 @@ describe("ConfigValidation.check", () => {
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("validated successfully")
@@ -44,7 +47,7 @@ describe("ConfigValidation.check", () => {
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("ERROR")
@@ -59,7 +62,7 @@ describe("ConfigValidation.check", () => {
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("WARNING")
@@ -79,7 +82,7 @@ Do something useful`,
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("validated successfully")
@@ -100,7 +103,7 @@ Do something`,
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("WARNING")
@@ -121,7 +124,7 @@ You are a helpful agent.`,
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toContain("config_validation")
expect(result).toContain("validated successfully")
@@ -134,7 +137,7 @@ You are a helpful agent.`,
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toBe("")
})
@@ -146,7 +149,7 @@ You are a helpful agent.`,
const result = await WithInstance.provide({
directory: tmp.path,
fn: () => ConfigValidation.check(filepath),
fn: () => check(filepath),
})
expect(result).toBe("")
})
@@ -173,8 +176,8 @@ Broken agent`,
directory: tmp.path,
fn: async () => {
// Force config load to populate warnings
await Config.get()
return ConfigValidation.check(filepath)
await AppRuntime.runPromise(Config.Service.use((svc) => svc.get()))
return check(filepath)
},
})
expect(result).toContain("Pre-existing config issues")
@@ -41,6 +41,9 @@ import { Provider } from "../../src/provider/provider"
import { ProviderID } from "../../src/provider/schema"
import { Filesystem } from "../../src/util/filesystem"
import { ModelCache } from "../../src/provider/model-cache"
import { AppRuntime } from "../../src/effect/app-runtime"
const clear = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(id)))
function paid(providers: Awaited<ReturnType<typeof Provider.list>>) {
const item = providers[ProviderID.kilo]
@@ -51,13 +54,12 @@ function paid(providers: Awaited<ReturnType<typeof Provider.list>>) {
const authPath = path.join(Global.Path.data, "auth.json")
test("kilo loader keeps paid models without auth and when config apiKey is present", async () => {
// Reset state that may be stale from other test files sharing this process.
// Persisted auth from other tests and ModelCache's TTL map must not affect this test.
const prev = await Filesystem.readText(authPath).catch(() => undefined)
try {
await Filesystem.write(authPath, JSON.stringify({}))
ModelCache.clear("kilo")
await clear("kilo")
await using base = await tmpdir({
init: async (dir) => {
@@ -115,7 +117,7 @@ test("kilo loader keeps paid models without auth and when auth exists", async ()
try {
await Filesystem.write(authPath, JSON.stringify({}))
ModelCache.clear("kilo")
await clear("kilo")
await using base = await tmpdir({
init: async (dir) => {
@@ -25,6 +25,13 @@ mock.module("@gitlab/opencode-gitlab-auth", () => ({ default: () => ({}) }))
import { tmpdir } from "../fixture/fixture"
import { WithInstance } from "../../src/project/with-instance"
import { ModelCache } from "../../src/provider/model-cache"
import { AppRuntime } from "../../src/effect/app-runtime"
const clear = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(id)))
const fetch = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.fetch(id)))
const failed = () => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.failedProviders()))
const failure = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.getFailure(id)))
const get = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.get(id)))
const CONFIG = JSON.stringify({ $schema: "https://app.kilo.ai/config.json" })
@@ -38,16 +45,16 @@ async function withInstance<T>(fn: () => Promise<T>): Promise<T> {
}
test("401 from gateway sets provider as failed in ModelCache", async () => {
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.failedProviders()).toContain("kilo")
expect(ModelCache.getFailure("kilo")).toMatchObject({ kind: "unauthorized", status: 401 })
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failed()).toContain("kilo")
expect(await failure("kilo")).toMatchObject({ kind: "unauthorized", status: 401 })
})
test("401 from gateway caches empty models (not undefined)", async () => {
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
const cached = ModelCache.get("kilo")
await clear("kilo")
await withInstance(() => fetch("kilo"))
const cached = await get("kilo")
expect(cached).toBeDefined()
expect(Object.keys(cached!)).toHaveLength(0)
})
@@ -0,0 +1,95 @@
// kilocode_change - new file
import { expect, spyOn } from "bun:test"
import { Effect, Layer } from "effect"
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
import { Auth } from "../../src/auth"
import { Bus } from "../../src/bus"
import type { Config } from "../../src/config/config"
import { KiloSessions } from "../../src/kilo-sessions/kilo-sessions"
import { ProjectID } from "../../src/project/schema"
import { Session } from "../../src/session/session"
import { SessionID } from "../../src/session/schema"
import { TestConfig } from "../fixture/config"
import { testEffect } from "../lib/effect"
const it = testEffect(CrossSpawnSpawner.defaultLayer)
function layer(overrides: Partial<Config.Interface> = {}) {
return Layer.merge(
KiloSessions.layer.pipe(Layer.provideMerge(Bus.layer), Layer.provide(TestConfig.layer(overrides))),
Auth.defaultLayer,
)
}
it.instance("initializes once per instance through Config.Service", () => {
let reads = 0
return Effect.gen(function* () {
const sessions = yield* KiloSessions.Service
yield* sessions.init()
yield* sessions.init()
expect(reads).toBe(1)
}).pipe(
Effect.provide(
layer({
getGlobal: () =>
Effect.sync(() => {
reads += 1
return {}
}),
}),
),
)
})
it.instance("does not duplicate created-session subscribers when init is repeated", () => {
const calls: string[] = []
const fetch: typeof globalThis.fetch = Object.assign(
async (input: RequestInfo | URL) => {
const url = String(input)
if (url.endsWith("/api/user")) return new Response("{}", { status: 200 })
if (url.endsWith("/api/session")) {
calls.push(url)
return Response.json({ id: "remote-1", ingestPath: "/api/ingest/session-1" })
}
return new Response("{}", { status: 200 })
},
{ preconnect: globalThis.fetch.preconnect },
)
const request = spyOn(globalThis, "fetch").mockImplementation(fetch)
const id = SessionID.descending("session-created")
return Effect.gen(function* () {
const auth = yield* Auth.Service
const bus = yield* Bus.Service
const sessions = yield* KiloSessions.Service
yield* auth.set("kilo", { type: "api", key: "test-token" })
yield* sessions.init()
yield* sessions.init()
yield* Effect.sleep(50)
yield* bus.publish(Session.Event.Created, {
sessionID: id,
info: {
id,
slug: "test",
projectID: ProjectID.make("project-test"),
directory: "/tmp/test",
title: "test",
version: "test",
time: { created: Date.now(), updated: Date.now() },
},
})
yield* Effect.sleep(50)
expect(calls).toHaveLength(1)
}).pipe(
Effect.ensuring(
Effect.gen(function* () {
const auth = yield* Auth.Service
yield* auth.remove("kilo").pipe(Effect.orDie)
request.mockRestore()
}),
),
Effect.provide(layer()),
)
})
@@ -0,0 +1,192 @@
// kilocode_change - new file
import { expect } from "bun:test"
import { Deferred, Effect, Fiber, Layer, Ref } from "effect"
import { HttpClient, HttpClientResponse } from "effect/unstable/http"
import { Auth } from "../../src/auth"
import { ModelCache } from "../../src/provider/model-cache"
import { TestConfig } from "../fixture/config"
import { testEffect } from "../lib/effect"
type Hit = { readonly url: string }
const auth = Layer.mock(Auth.Service)({
get: () => Effect.succeed(undefined),
})
const it = testEffect(Layer.empty)
function layer(
hits: Ref.Ref<Hit[]>,
cfg = TestConfig.layer(),
access = auth,
gates?: { readonly started: Deferred.Deferred<void>; readonly wait: Deferred.Deferred<void> },
) {
const http = HttpClient.make((request) =>
Effect.gen(function* () {
yield* Ref.update(hits, (list) => [...list, { url: request.url }])
const count = (yield* Ref.get(hits)).length
if (gates && count === 1) {
yield* Deferred.succeed(gates.started, undefined)
yield* Deferred.await(gates.wait)
}
return HttpClientResponse.fromWeb(
request,
Response.json({ data: [{ id: `apertis-${count}`, owned_by: "apertis" }] }),
)
}),
)
return Layer.fresh(ModelCache.layer).pipe(
Layer.provide(Layer.succeed(HttpClient.HttpClient, http)),
Layer.provide(cfg),
Layer.provide(access),
)
}
it.live("fetches Apertis models through the injected HttpClient", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const models = yield* ModelCache.Service.use((cache) =>
cache.fetch("apertis", { apiKey: "test-key", baseURL: "https://apertis.test/v1" }),
).pipe(Effect.provide(layer(hits)))
expect(Object.keys(models)).toEqual(["apertis-1"])
expect((yield* Ref.get(hits)).map((hit) => hit.url)).toEqual(["https://apertis.test/v1/models"])
}),
)
it.live("reuses cached values and refresh invalidates the provider cell", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const run = ModelCache.Service.use((cache) =>
Effect.gen(function* () {
const first = yield* cache.fetch("apertis", { apiKey: "test-key" })
const cached = yield* cache.fetch("apertis", { apiKey: "test-key" })
const refreshed = yield* cache.refresh("apertis", { apiKey: "test-key" })
return { first, cached, refreshed }
}),
).pipe(Effect.provide(layer(hits)))
const out = yield* run
expect(Object.keys(out.first)).toEqual(["apertis-1"])
expect(Object.keys(out.cached)).toEqual(["apertis-1"])
expect(Object.keys(out.refreshed)).toEqual(["apertis-2"])
expect((yield* Ref.get(hits)).length).toBe(2)
}),
)
it.live("keeps concurrent request options isolated", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const started = yield* Deferred.make<void>()
const wait = yield* Deferred.make<void>()
const out = yield* ModelCache.Service.use((cache) =>
Effect.gen(function* () {
const first = yield* cache
.fetch("apertis", { apiKey: "first", baseURL: "https://first.test/v1" })
.pipe(Effect.forkChild)
yield* Deferred.await(started)
const second = yield* cache
.fetch("apertis", { apiKey: "second", baseURL: "https://second.test/v1" })
.pipe(Effect.forkChild)
yield* Effect.sleep("10 millis")
yield* Deferred.succeed(wait, undefined)
const firstModels = yield* Fiber.join(first)
const secondModels = yield* Fiber.join(second)
return { first: firstModels, second: secondModels, current: yield* cache.get("apertis") }
}),
).pipe(Effect.provide(layer(hits, TestConfig.layer(), auth, { started, wait })))
expect(Object.keys(out.first)).toEqual(["apertis-1"])
expect(Object.keys(out.second)).toEqual(["apertis-2"])
expect(out.current).toEqual(out.second)
expect((yield* Ref.get(hits)).map((hit) => hit.url)).toEqual([
"https://first.test/v1/models",
"https://second.test/v1/models",
])
}),
)
it.live("does not let an older fetch override a newer refresh", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const started = yield* Deferred.make<void>()
const wait = yield* Deferred.make<void>()
const models = yield* ModelCache.Service.use((cache) =>
Effect.gen(function* () {
const stale = yield* cache
.fetch("apertis", { apiKey: "first", baseURL: "https://first.test/v1" })
.pipe(Effect.forkChild)
yield* Deferred.await(started)
const fresh = yield* cache.refresh("apertis", { apiKey: "second", baseURL: "https://second.test/v1" })
yield* Deferred.succeed(wait, undefined)
yield* Fiber.join(stale)
return { fresh, current: yield* cache.get("apertis") }
}),
).pipe(Effect.provide(layer(hits, TestConfig.layer(), auth, { started, wait })))
expect(models.current).toEqual(models.fresh)
expect(Object.keys(models.current ?? {})).toEqual(["apertis-2"])
}),
)
it.live("does not restore a fetch that was cleared while pending", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const started = yield* Deferred.make<void>()
const wait = yield* Deferred.make<void>()
const current = yield* ModelCache.Service.use((cache) =>
Effect.gen(function* () {
const pending = yield* cache
.fetch("apertis", { apiKey: "first", baseURL: "https://first.test/v1" })
.pipe(Effect.forkChild)
yield* Deferred.await(started)
yield* cache.clear("apertis")
yield* Deferred.succeed(wait, undefined)
yield* Fiber.join(pending)
return yield* cache.get("apertis")
}),
).pipe(Effect.provide(layer(hits, TestConfig.layer(), auth, { started, wait })))
expect(current).toBeUndefined()
}),
)
it.live("exposes the most recently refreshed provider value", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const models = yield* ModelCache.Service.use((cache) =>
Effect.gen(function* () {
yield* cache.fetch("apertis", { apiKey: "first", baseURL: "https://first.test/v1" })
const refreshed = yield* cache.refresh("apertis", { apiKey: "second", baseURL: "https://second.test/v1" })
const current = yield* cache.get("apertis")
return { refreshed, current }
}),
).pipe(Effect.provide(layer(hits)))
expect(models.current).toEqual(models.refreshed)
expect(Object.keys(models.current ?? {})).toEqual(["apertis-2"])
}),
)
it.live("does not resolve auth or config for unsupported providers", () =>
Effect.gen(function* () {
const hits = yield* Ref.make<Hit[]>([])
const configs = yield* Ref.make(0)
const auths = yield* Ref.make(0)
const cfg = TestConfig.layer({
get: () => Ref.update(configs, (count) => count + 1).pipe(Effect.as({})),
})
const access = Layer.mock(Auth.Service)({
get: () => Ref.update(auths, (count) => count + 1).pipe(Effect.as(undefined)),
})
const models = yield* ModelCache.Service.use((cache) => cache.fetch("openai")).pipe(
Effect.provide(layer(hits, cfg, access)),
)
expect(models).toEqual({})
expect(yield* Ref.get(configs)).toBe(0)
expect(yield* Ref.get(auths)).toBe(0)
expect(yield* Ref.get(hits)).toEqual([])
}),
)
@@ -40,6 +40,11 @@ import { tmpdir } from "../fixture/fixture"
import { WithInstance } from "../../src/project/with-instance"
import { Filesystem } from "../../src/util/filesystem"
import { ModelCache } from "../../src/provider/model-cache"
import { AppRuntime } from "../../src/effect/app-runtime"
const clear = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(id)))
const fetch = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.fetch(id)))
const get = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.get(id)))
const authPath = path.join(Global.Path.data, "auth.json")
@@ -75,10 +80,10 @@ test("model fetch uses accountId from OAuth auth as kilocodeOrganizationId", asy
fn: async () => {
// Reset captured and cache
captured = undefined
ModelCache.clear("kilo")
await clear("kilo")
// Trigger model fetch through the cache
await ModelCache.fetch("kilo")
await fetch("kilo")
// The fetchKiloModels call should have received the organization ID
expect(captured).toBeDefined()
@@ -126,9 +131,9 @@ test("model fetch without OAuth accountId does not set kilocodeOrganizationId",
directory: tmp.path,
fn: async () => {
captured = undefined
ModelCache.clear("kilo")
await clear("kilo")
await ModelCache.fetch("kilo")
await fetch("kilo")
expect(captured).toBeDefined()
expect(captured.kilocodeToken).toBe("test-personal-token")
@@ -176,25 +181,25 @@ test("ModelCache.clear removes cached entry so next fetch hits the network", asy
fn: async () => {
// Populate cache
captured = undefined
ModelCache.clear("kilo")
await ModelCache.fetch("kilo")
await clear("kilo")
await fetch("kilo")
expect(captured).toBeDefined()
// Verify cache is populated — second fetch should NOT call fetchKiloModels
captured = undefined
await ModelCache.fetch("kilo")
await fetch("kilo")
expect(captured).toBeUndefined()
expect(ModelCache.get("kilo")).toBeDefined()
expect(await get("kilo")).toBeDefined()
// Clear the cache
ModelCache.clear("kilo")
await clear("kilo")
// get() should return undefined after clear
expect(ModelCache.get("kilo")).toBeUndefined()
expect(await get("kilo")).toBeUndefined()
// Next fetch should call fetchKiloModels again
captured = undefined
await ModelCache.fetch("kilo")
await fetch("kilo")
expect(captured).toBeDefined()
},
})
@@ -14,7 +14,11 @@ import { provideTmpdirInstance } from "../../fixture/fixture"
import { testEffect } from "../../lib/effect"
const bus = Bus.layer
const env = Layer.mergeAll(Permission.layer.pipe(Layer.provide(bus)), bus, CrossSpawnSpawner.defaultLayer)
const env = Layer.mergeAll(
Permission.layer.pipe(Layer.provide(bus), Layer.provide(Config.defaultLayer)),
bus,
CrossSpawnSpawner.defaultLayer,
)
const it = testEffect(env)
afterAll(async () => {
@@ -22,7 +26,9 @@ afterAll(async () => {
for (const file of ["kilo.jsonc", "kilo.json", "config.json", "opencode.json", "opencode.jsonc"]) {
await fs.rm(path.join(dir, file), { force: true }).catch(() => {})
}
await Config.invalidate()
await Effect.runPromise(
Config.Service.use((svc) => svc.invalidate()).pipe(Effect.scoped, Effect.provide(Config.defaultLayer)),
)
await InstanceRuntime.disposeAllInstances()
})
@@ -14,7 +14,12 @@ import { provideTmpdirInstance } from "../../fixture/fixture"
import { testEffect } from "../../lib/effect"
const bus = Bus.layer
const env = Layer.mergeAll(Permission.layer.pipe(Layer.provide(bus)), bus, CrossSpawnSpawner.defaultLayer)
const env = Layer.mergeAll(
Permission.layer.pipe(Layer.provide(bus), Layer.provide(Config.defaultLayer)),
Config.defaultLayer,
bus,
CrossSpawnSpawner.defaultLayer,
)
const it = testEffect(env)
afterAll(async () => {
@@ -22,7 +27,9 @@ afterAll(async () => {
for (const file of ["kilo.jsonc", "kilo.json", "config.json", "opencode.json", "opencode.jsonc"]) {
await fs.rm(path.join(dir, file), { force: true }).catch(() => {})
}
await Config.invalidate()
await Effect.runPromise(
Config.Service.use((svc) => svc.invalidate()).pipe(Effect.scoped, Effect.provide(Config.defaultLayer)),
)
await InstanceRuntime.disposeAllInstances()
})
@@ -735,7 +742,8 @@ describe("saveAlwaysRules", () => {
yield* reply({ requestID: PermissionID.make("permission_saved_always"), reply: "always" })
yield* Fiber.join(fiber)
const cfg = yield* Effect.promise(() => Config.get())
const config = yield* Config.Service
const cfg = yield* config.get()
expect(cfg.permission?.bash).toMatchObject({ "kilo-permission-8353 test": "allow" })
expect(cfg.permission?.bash).not.toMatchObject({ "kilo-permission-8353 *": "allow" })
@@ -14,7 +14,11 @@ import { provideInstance, provideTmpdirInstance, tmpdirScoped } from "../../fixt
import { testEffect } from "../../lib/effect"
const bus = Bus.layer
const env = Layer.mergeAll(Permission.layer.pipe(Layer.provide(bus)), bus, CrossSpawnSpawner.defaultLayer)
const env = Layer.mergeAll(
Permission.layer.pipe(Layer.provide(bus), Layer.provide(Config.defaultLayer)),
bus,
CrossSpawnSpawner.defaultLayer,
)
const it = testEffect(env)
afterAll(async () => {
@@ -22,7 +26,9 @@ afterAll(async () => {
for (const file of ["kilo.jsonc", "kilo.json", "config.json", "opencode.json", "opencode.jsonc"]) {
await fs.rm(path.join(dir, file), { force: true }).catch(() => {})
}
await Config.invalidate()
await Effect.runPromise(
Config.Service.use((svc) => svc.invalidate()).pipe(Effect.scoped, Effect.provide(Config.defaultLayer)),
)
await InstanceRuntime.disposeAllInstances()
})
@@ -4,7 +4,8 @@
// 2. ModelCache.getFailure() returns the typed error for a failed provider.
// 3. Clear removes failure state.
import { test, expect, mock } from "bun:test"
import { beforeEach, test, expect, mock } from "bun:test"
import { Effect } from "effect"
import path from "path"
import * as Log from "@opencode-ai/core/util/log"
@@ -12,9 +13,13 @@ Log.init({ print: false })
// Stub fetchKiloModels to return controlled typed results.
let stubbedResult: { models: Record<string, any>; error?: { kind: string; status?: number } } = { models: {} }
let stubbedError: Error | undefined
mock.module("@kilocode/kilo-gateway", () => ({
fetchKiloModels: async () => stubbedResult,
fetchKiloModels: async () => {
if (stubbedError) throw stubbedError
return stubbedResult
},
KILO_OPENROUTER_BASE: "https://api.kilo.ai/api/openrouter",
}))
@@ -25,6 +30,12 @@ mock.module("@gitlab/opencode-gitlab-auth", () => ({ default: () => ({}) }))
import { tmpdir } from "../fixture/fixture"
import { WithInstance } from "../../src/project/with-instance"
import { ModelCache } from "../../src/provider/model-cache"
import { AppRuntime } from "../../src/effect/app-runtime"
const clear = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.clear(id)))
const fetch = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.fetch(id)))
const failed = () => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.failedProviders()))
const failure = (id: string) => AppRuntime.runPromise(ModelCache.Service.use((cache) => cache.getFailure(id)))
const CONFIG = JSON.stringify({ $schema: "https://app.kilo.ai/config.json" })
@@ -37,9 +48,13 @@ async function withInstance<T>(fn: () => Promise<T>): Promise<T> {
return WithInstance.provide({ directory: tmp.path, fn })
}
test("failedProviders returns empty array when no fetch has occurred", () => {
ModelCache.clear("kilo")
expect(ModelCache.failedProviders()).not.toContain("kilo")
beforeEach(() => {
stubbedError = undefined
})
test("failedProviders returns empty array when no fetch has occurred", async () => {
await clear("kilo")
expect(await failed()).not.toContain("kilo")
})
test("getFailure returns undefined when fetch succeeds", async () => {
@@ -53,36 +68,47 @@ test("getFailure returns undefined when fetch succeeds", async () => {
},
},
}
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.getFailure("kilo")).toBeUndefined()
expect(ModelCache.failedProviders()).not.toContain("kilo")
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failure("kilo")).toBeUndefined()
expect(await failed()).not.toContain("kilo")
})
test("failedProviders includes provider after auth error", async () => {
stubbedResult = { models: {}, error: { kind: "unauthorized", status: 401 } }
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.failedProviders()).toContain("kilo")
expect(ModelCache.getFailure("kilo")).toMatchObject({ kind: "unauthorized", status: 401 })
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failed()).toContain("kilo")
expect(await failure("kilo")).toMatchObject({ kind: "unauthorized", status: 401 })
})
test("gateway rejection remains recoverable through the Effect error channel", async () => {
stubbedError = new Error("gateway failed")
await clear("kilo")
const models = await withInstance(() =>
AppRuntime.runPromise(
ModelCache.Service.use((cache) => cache.fetch("kilo").pipe(Effect.catch(() => Effect.succeed({})))),
),
)
expect(models).toEqual({})
})
test("clear removes failure state", async () => {
stubbedResult = { models: {}, error: { kind: "network" } }
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.failedProviders()).toContain("kilo")
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failed()).toContain("kilo")
ModelCache.clear("kilo")
expect(ModelCache.failedProviders()).not.toContain("kilo")
expect(ModelCache.getFailure("kilo")).toBeUndefined()
await clear("kilo")
expect(await failed()).not.toContain("kilo")
expect(await failure("kilo")).toBeUndefined()
})
test("failure state is cleared when subsequent fetch succeeds", async () => {
stubbedResult = { models: {}, error: { kind: "unauthorized", status: 401 } }
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.failedProviders()).toContain("kilo")
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failed()).toContain("kilo")
stubbedResult = {
models: {
@@ -94,8 +120,8 @@ test("failure state is cleared when subsequent fetch succeeds", async () => {
},
},
}
ModelCache.clear("kilo")
await withInstance(() => ModelCache.fetch("kilo"))
expect(ModelCache.failedProviders()).not.toContain("kilo")
expect(ModelCache.getFailure("kilo")).toBeUndefined()
await clear("kilo")
await withInstance(() => fetch("kilo"))
expect(await failed()).not.toContain("kilo")
expect(await failure("kilo")).toBeUndefined()
})
@@ -227,19 +227,27 @@ describe("kilocode tool registry indexing", () => {
test("logs indexing bootstrap failures without blocking session bootstrap", async () => {
const logger = Log.create({ service: "kilocode-bootstrap" })
const err = new Error("indexing init failed")
const sessions = spyOn(KiloSessions, "init").mockResolvedValue(undefined)
const calls: string[] = []
const sessions = Layer.succeed(
KiloSessions.Service,
KiloSessions.Service.of({ init: () => Effect.sync(() => calls.push("sessions")) }),
)
const indexing = spyOn(KiloIndexing, "init").mockRejectedValue(err)
const warn = spyOn(logger, "warn").mockImplementation(() => {})
try {
await KilocodeBootstrap.init()
await Effect.runPromise(
KilocodeBootstrap.Service.use((svc) => svc.init()).pipe(
Effect.provide(KilocodeBootstrap.layer.pipe(Layer.provide(sessions))),
Effect.scoped,
),
)
await new Promise((resolve) => setTimeout(resolve, 0))
expect(sessions).toHaveBeenCalledTimes(1)
expect(calls).toEqual(["sessions"])
expect(indexing).toHaveBeenCalledTimes(1)
expect(warn).toHaveBeenCalledWith("indexing bootstrap failed", { err })
} finally {
sessions.mockRestore()
indexing.mockRestore()
warn.mockRestore()
}
+10 -2
View File
@@ -23,7 +23,13 @@ import { testEffect } from "../lib/effect"
import { MessageID, SessionID } from "../../src/session/schema"
const bus = Bus.layer
const env = Layer.mergeAll(Permission.layer.pipe(Layer.provide(bus)), bus, CrossSpawnSpawner.defaultLayer)
// kilocode_change start
const env = Layer.mergeAll(
Permission.layer.pipe(Layer.provide(bus), Layer.provide(Config.defaultLayer)),
bus,
CrossSpawnSpawner.defaultLayer,
)
// kilocode_change end
const it = testEffect(env)
afterEach(async () => {
@@ -36,7 +42,9 @@ afterAll(async () => {
for (const file of ["kilo.jsonc", "kilo.json", "config.json", "opencode.json", "opencode.jsonc"]) {
await fs.rm(path.join(dir, file), { force: true }).catch(() => {})
}
await Config.invalidate()
await Effect.runPromise(
Config.Service.use((svc) => svc.invalidate()).pipe(Effect.scoped, Effect.provide(Config.defaultLayer)),
)
await InstanceRuntime.disposeAllInstances()
})
// kilocode_change end
@@ -8,12 +8,14 @@ import { ProviderAuth } from "@/provider/auth"
import { ProviderID } from "../../src/provider/schema"
import { Plugin } from "@/plugin"
import { Auth } from "@/auth"
import { ModelCache } from "@/provider/model-cache" // kilocode_change
import { Bus } from "@/bus"
import { TestConfig } from "../fixture/config"
function layer(directory: string, plugins: string[]) {
return ProviderAuth.layer.pipe(
Layer.provide(Auth.defaultLayer),
Layer.provide(ModelCache.defaultLayer), // kilocode_change
Layer.provide(
Plugin.layer.pipe(
Layer.provide(Bus.layer),
@@ -5,6 +5,8 @@ import { AppFileSystem } from "@opencode-ai/core/filesystem"
import { Flag } from "@opencode-ai/core/flag/flag"
import { Global } from "@opencode-ai/core/global"
import { ModelsDev } from "../../src/provider/models"
import { ModelCache } from "../../src/provider/model-cache" // kilocode_change
import { Config } from "../../src/config/config" // kilocode_change
import { Auth } from "../../src/auth" // kilocode_change
import { it } from "../lib/effect"
import { rm, writeFile, utimes, mkdir } from "fs/promises"
@@ -90,7 +92,9 @@ const buildLayer = (state: Ref.Ref<MockState>) =>
Layer.fresh(ModelsDev.layer).pipe(
Layer.provide(Layer.succeed(HttpClient.HttpClient, makeMockClient(state))),
Layer.provide(AppFileSystem.defaultLayer),
Layer.provide(Config.defaultLayer), // kilocode_change
Layer.provide(Auth.defaultLayer), // kilocode_change
Layer.provide(ModelCache.defaultLayer), // kilocode_change
)
const writeCache = (data: object, mtimeMs?: number) =>
@@ -120,10 +124,10 @@ const initialState: MockState = {
calls: [],
}
// kilocode_change - skip: upstream tests assert raw-fixture passthrough but Kilo's
// ModelsDev.get() filters/injects providers based on Config.get() (kilo-allowed gating,
// apertis options, kilo provider injection). The test setup doesn't provide an Instance
// context, so Config.get() throws "No context found for instance".
// kilocode_change start - skip: upstream tests assert raw-fixture passthrough but Kilo's
// ModelsDev.get() filters/injects providers based on effect-native config access (kilo-allowed
// gating, apertis options, and Kilo provider injection).
// kilocode_change end
describe.skip("ModelsDev Service", () => {
it.live("get() returns providers from disk when cache file exists", () =>
Effect.gen(function* () {
@@ -2,7 +2,6 @@
import { afterEach, beforeEach, describe, expect, mock, spyOn, test } from "bun:test"
import { $ } from "bun"
import path from "path"
import * as Config from "../../src/config/config"
import { WithInstance } from "../../src/project/with-instance"
import * as Log from "@opencode-ai/core/util/log"
import { resetDatabase } from "../fixture/db"
@@ -29,10 +28,6 @@ describe("experimental.session.list", () => {
try {
await $`git worktree add ${worktree} -b test-branch-${Date.now()}`.cwd(first.path).quiet()
spyOn(Config, "get").mockImplementation(
async () => ({ share: "manual" }) as Awaited<ReturnType<typeof Config.get>>,
)
try {
const { Server } = await import("../../src/server/server")
const { Session } = await import("../../src/session/session")
@@ -99,10 +94,6 @@ describe("experimental.session.list", () => {
try {
await $`git worktree add ${worktree} -b test-branch-sdk-${Date.now()}`.cwd(first.path).quiet()
spyOn(Config, "get").mockImplementation(
async () => ({ share: "manual" }) as Awaited<ReturnType<typeof Config.get>>,
)
try {
const { Server } = await import("../../src/server/server")
const { Session } = await import("../../src/session/session")
@@ -4,7 +4,6 @@ import { $ } from "bun"
import { Effect } from "effect"
import path from "path"
import { WithInstance } from "../../src/project/with-instance"
import * as Config from "../../src/config/config"
import { RecallTool } from "../../src/tool/recall"
import { AppRuntime } from "../../src/effect/app-runtime"
import { resetDatabase } from "../fixture/db"
@@ -43,10 +42,6 @@ describe("tool.recall", () => {
await $`git worktree add ${worktree} -b test-branch-${Date.now()}`.cwd(first.path).quiet()
await Bun.write(path.join(first.path, ".git", "opencode"), "stale-project-id")
spyOn(Config, "get").mockImplementation(
async () => ({ share: "manual" }) as Awaited<ReturnType<typeof Config.get>>,
)
try {
const { Session } = await import("../../src/session/session")
await WithInstance.provide({
@@ -86,8 +81,6 @@ describe("tool.recall", () => {
await using first = await tmpdir({ git: true })
await using second = await tmpdir({ git: true })
spyOn(Config, "get").mockImplementation(async () => ({ share: "manual" }) as Awaited<ReturnType<typeof Config.get>>)
try {
const { Session } = await import("../../src/session/session")
const session = await WithInstance.provide({
@@ -121,10 +114,6 @@ describe("tool.recall", () => {
await $`git worktree add ${worktree} -b test-branch-${Date.now()}`.cwd(first.path).quiet()
await Bun.write(path.join(first.path, ".git", "opencode"), "stale-project-id")
spyOn(Config, "get").mockImplementation(
async () => ({ share: "manual" }) as Awaited<ReturnType<typeof Config.get>>,
)
try {
const { Session } = await import("../../src/session/session")
const session = await WithInstance.provide({