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
118 changes: 90 additions & 28 deletions packages/opencode/src/kilocode/agent-manager/service.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,11 @@
import { Bus } from "@/bus"
import { InstanceState } from "@/effect/instance-state"
import { InstanceRef } from "@/effect/instance-ref"
import { registerDisposer } from "@/effect/instance-registry"
import { Identifier } from "@/id/id"
import type { InstanceContext } from "@/project/instance-context"
import * as Log from "@opencode-ai/core/util/log"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { Context, Deferred, Duration, Effect, Layer, Schema } from "effect"
import { Context, Deferred, Duration, Effect, Layer, LayerMap, Schema } from "effect"
import { ErrorCode, Event, type Failure, type Request, RequestID, type Result } from "./protocol"

const log = Log.create({ service: "agent-manager-host" })
Expand Down Expand Up @@ -33,11 +35,17 @@ interface Entry {
}

interface State {
context: InstanceContext
closed: boolean
pending: Map<RequestID, Entry>
close: () => Effect.Effect<void>
}

class StateService extends Context.Service<StateService, State>()("@kilocode/AgentManager.State") {}

function matches(request: Request, result: Result) {
if (request.operation === "overview") return result.operation === "overview"
if (!("sessionID" in result)) return false
return result.operation === request.operation && result.sessionID === request.targetSessionID
}

Expand All @@ -58,31 +66,84 @@ export function layer(timeout: Duration.Input = "60 seconds") {
Service,
Effect.gen(function* () {
const bus = yield* Bus.Service
const state = yield* InstanceState.make<State>(
Effect.fn("AgentManager.state")(function* () {
const state = { pending: new Map<RequestID, Entry>() }
yield* Effect.addFinalizer(() =>
const active = new Map<string, Set<State>>()
const states = yield* LayerMap.make(
(context: InstanceContext) =>
Layer.effect(
StateService,
Effect.gen(function* () {
for (const entry of state.pending.values()) {
yield* bus.publish(Event.Cancelled, {
requestID: entry.info.id,
sessionID: entry.info.sessionID,
reason: "disposed",
})
yield* Deferred.fail(
entry.deferred,
new HostError({ code: "disconnected", detail: "The Agent Manager host disconnected" }),
)
const data = {
context,
closed: false,
pending: new Map<RequestID, Entry>(),
}
state.pending.clear()
const done = Effect.gen(function* () {
if (data.closed) return
data.closed = true
for (const entry of data.pending.values()) {
yield* bus
.publish(Event.Cancelled, {
requestID: entry.info.id,
sessionID: entry.info.sessionID,
reason: "disposed",
})
.pipe(Effect.ignore)
yield* Deferred.fail(
entry.deferred,
new HostError({ code: "disconnected", detail: "The Agent Manager host disconnected" }),
)
}
data.pending.clear()
}).pipe(Effect.provideService(InstanceRef, context))
const state: State = { ...data, close: () => done }
const current = active.get(context.directory) ?? new Set<State>()
current.add(state)
active.set(context.directory, current)
yield* Effect.addFinalizer(() =>
done.pipe(
Effect.ensuring(
Effect.sync(() => {
const current = active.get(context.directory)
if (!current) return
current.delete(state)
if (current.size === 0) active.delete(context.directory)
}),
),
),
)
return StateService.of(state)
}),
)
return state
}),
).pipe(Layer.provide(Layer.succeed(InstanceRef, context))),
{ idleTimeToLive: Duration.infinity },
)

const cancel = Effect.fn("AgentManager.cancel")(function* (id: RequestID, reason: "cancelled" | "timeout") {
const pending = (yield* InstanceState.get(state)).pending
const off = registerDisposer(async (directory) => {
const current = active.get(directory)
if (!current) return
await Effect.runPromise(
Effect.forEach(
[...current],
(state) => state.close().pipe(Effect.ensuring(states.invalidate(state.context))),
{ concurrency: "unbounded", discard: true },
),
)
})
yield* Effect.addFinalizer(() => Effect.sync(off))

const state = Effect.fn("AgentManager.state")(function* () {
const context = yield* InstanceRef
if (!context) return yield* Effect.die(new Error("Agent Manager instance context not available"))
return yield* Effect.gen(function* () {
return yield* StateService
}).pipe(Effect.provide(states.get(context)))
})

const cancel = Effect.fn("AgentManager.cancel")(function* (
current: State,
id: RequestID,
reason: "cancelled" | "timeout",
) {
const pending = current.pending
const entry = pending.get(id)
if (!entry) return
pending.delete(id)
Expand All @@ -100,7 +161,8 @@ export function layer(timeout: Duration.Input = "60 seconds") {
})

const request: Interface["request"] = Effect.fn("AgentManager.request")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const current = yield* state()
const pending = current.pending
const id = RequestID.make(Identifier.create("amr", "ascending"))
const deferred = yield* Deferred.make<Result, HostError>()
const info = { ...input, id } as Request
Expand All @@ -110,18 +172,18 @@ export function layer(timeout: Duration.Input = "60 seconds") {
return yield* Deferred.await(deferred).pipe(
Effect.timeoutOrElse({
duration: timeout,
orElse: () => cancel(id, "timeout").pipe(Effect.andThen(Deferred.await(deferred))),
orElse: () => cancel(current, id, "timeout").pipe(Effect.andThen(Deferred.await(deferred))),
}),
)
}).pipe(Effect.ensuring(cancel(id, "cancelled")))
}).pipe(Effect.ensuring(cancel(current, id, "cancelled")))
})

const list: Interface["list"] = Effect.fn("AgentManager.list")(function* () {
return Array.from((yield* InstanceState.get(state)).pending.values(), (entry) => entry.info)
return Array.from((yield* state()).pending.values(), (entry) => entry.info)
})

const reply: Interface["reply"] = Effect.fn("AgentManager.reply")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const pending = (yield* state()).pending
const entry = pending.get(input.requestID)
if (!entry) {
log.warn("reply for unknown request", { requestID: input.requestID })
Expand All @@ -133,7 +195,7 @@ export function layer(timeout: Duration.Input = "60 seconds") {
})

const reject: Interface["reject"] = Effect.fn("AgentManager.reject")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const pending = (yield* state()).pending
const entry = pending.get(input.requestID)
if (!entry) {
log.warn("rejection for unknown request", { requestID: input.requestID })
Expand Down
59 changes: 59 additions & 0 deletions packages/opencode/test/kilocode/agent-manager-service.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { expect } from "bun:test"
import { Bus } from "@/bus"
import { GlobalBus, type GlobalEvent } from "@/bus/global"
import { InstanceRef } from "@/effect/instance-ref"
import { disposeInstance } from "@/effect/instance-registry"
import { Event, type Request } from "@/kilocode/agent-manager/protocol"
import { AgentManager, HostError } from "@/kilocode/agent-manager/service"
Expand All @@ -10,12 +11,25 @@ import { TestInstance } from "../fixture/fixture"
import { testEffect } from "../lib/effect"

const it = testEffect(AgentManager.layer("20 millis").pipe(Layer.provideMerge(Bus.layer)))
const isolation = testEffect(AgentManager.layer("20 millis").pipe(Layer.provideMerge(Bus.layer)))
const sessionID = SessionID.make("ses_agent_manager_test")

function request(manager: AgentManager.Interface) {
return manager.request({ operation: "overview", sessionID })
}

function context(directory: string) {
return {
directory,
worktree: directory,
project: { id: directory, worktree: directory, vcs: "git", sandboxes: [] },
} as never
}

function provide(directory: string, value = context(directory)) {
return <A, E, R>(effect: Effect.Effect<A, E, R>) => effect.pipe(Effect.provideService(InstanceRef, value))
}

it.instance(
"publishes, lists, and completes a correlated request",
() =>
Expand Down Expand Up @@ -51,6 +65,50 @@ it.instance(
{ git: true },
)

isolation.live("keeps pending requests isolated by instance context", () =>
Effect.gen(function* () {
const manager = yield* AgentManager.Service
const one = "/agent-manager-test-one"
const two = "/agent-manager-test-two"
const oneContext = context(one)
const twoContext = context(two)
const first = yield* request(manager).pipe(provide(one, oneContext), Effect.forkScoped)
const second = yield* request(manager).pipe(provide(two, twoContext), Effect.forkScoped)
const pendingOne = yield* manager
.list()
.pipe(provide(one, oneContext), Effect.repeat({ until: (items) => items.length === 1 }))
const pendingTwo = yield* manager
.list()
.pipe(provide(two, twoContext), Effect.repeat({ until: (items) => items.length === 1 }))

expect(pendingOne[0].id).not.toBe(pendingTwo[0].id)
expect(
yield* manager
.reply({
requestID: pendingOne[0].id,
result: { operation: "overview", overview: { sections: [], ungrouped: [] } },
})
.pipe(provide(two, twoContext), Effect.flip),
).toMatchObject({ _tag: "AgentManager.NotFoundError" })

yield* manager
.reply({
requestID: pendingOne[0].id,
result: { operation: "overview", overview: { sections: [], ungrouped: [] } },
})
.pipe(provide(one, oneContext))
yield* manager
.reject({
requestID: pendingTwo[0].id,
error: { code: "stale_session", message: "The target session is stale" },
})
.pipe(provide(two, twoContext))

expect(yield* Fiber.join(first)).toMatchObject({ operation: "overview" })
expect((yield* Fiber.join(second).pipe(Effect.flip)).code).toBe("stale_session")
}),
)

it.instance(
"propagates host rejection and rejects mismatched prompt acknowledgements",
() =>
Expand Down Expand Up @@ -118,6 +176,7 @@ it.instance(
yield* Effect.promise(() => disposeInstance(instance.directory))
const error = yield* Fiber.join(fiber).pipe(Effect.flip)
expect(error.code).toBe("disconnected")
expect(yield* manager.list()).toEqual([])
}),
{ git: true },
)
1 change: 0 additions & 1 deletion script/architecture-allowlist.json
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
"description": "Legacy InstanceState.make usages in Kilo-owned code (packages/opencode/src/kilocode/**, packages/opencode/src/kilo-*/**). Target: encapsulate in scoped Effect Services in packages/core.",
"allowed": {
"packages/opencode/src/kilo-sessions/kilo-sessions.ts": { "count": 1, "owner": "session-runtime", "reason": "Kilo session coordination state" },
"packages/opencode/src/kilocode/agent-manager/service.ts": { "count": 1, "owner": "agent-manager", "reason": "Agent Manager multi-project worktree state" },
"packages/opencode/src/kilocode/background-process/index.ts": { "count": 1, "owner": "process-runtime", "reason": "Directory-keyed background process registry" },
"packages/opencode/src/kilocode/interactive-terminal/index.ts": { "count": 1, "owner": "terminal-runtime", "reason": "Interactive terminal manager state" },
"packages/opencode/src/kilocode/notebook/service.ts": { "count": 1, "owner": "notebook-runtime", "reason": "Notebook cell execution service state" },
Expand Down
Loading