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
115 changes: 85 additions & 30 deletions packages/opencode/src/kilocode/notebook/service.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
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 { Deferred, Duration, Effect, Layer, Schema, Context } from "effect"
import { capture } from "@/kilocode/instance"
import { Context, Deferred, Duration, Effect, Layer, LayerMap, Schema } from "effect"
import * as Log from "@opencode-ai/core/util/log"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { ErrorCode, Event, type Failure, type Request, RequestID, type Result } from "./protocol"
Expand Down Expand Up @@ -36,8 +38,17 @@ interface Entry {
}
interface State {
pending: Map<RequestID, Entry>
dispose: () => Effect.Effect<void>
}

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

const context = Effect.gen(function* () {
const ctx = (yield* InstanceRef) ?? capture()
if (!ctx) return yield* Effect.die(new Error("Instance context not provided"))
return ctx
})

function matches(request: Request, result: Result) {
if (request.path !== result.requestPath) return false
if (request.operation === "read") return result.operation === "read"
Expand All @@ -63,31 +74,67 @@ export function layer(timeout: Duration.Input = "10 minutes") {
Service,
Effect.gen(function* () {
const bus = yield* Bus.Service
const state = yield* InstanceState.make<State>(
Effect.fn("Notebook.state")(function* () {
const state = { pending: new Map<RequestID, Entry>() }
yield* Effect.addFinalizer(() =>
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",
})
const states = new Map<string, State>()
const stateLayer = (directory: string) =>
Layer.effect(
StateService,
Effect.gen(function* () {
const instance = (yield* InstanceRef) ?? capture()
if (!instance) return yield* Effect.die(new Error("Instance context not provided"))
const pending = new Map<RequestID, Entry>()
const dispose = Effect.fn("Notebook.dispose")(function* () {
const entries = Array.from(pending.values())
pending.clear()
for (const entry of entries) {
yield* bus
.publish(Event.Cancelled, {
requestID: entry.info.id,
sessionID: entry.info.sessionID,
reason: "disposed" as const,
})
.pipe(Effect.provideService(InstanceRef, instance))
yield* Deferred.fail(
entry.deferred,
new HostError({ code: "disconnected", detail: "The notebook host disconnected" }),
)
}
state.pending.clear()
}),
)
return state
}),
})
const state: State = { pending, dispose }
states.set(directory, state)
yield* Effect.addFinalizer(() =>
Effect.gen(function* () {
if (states.get(directory) === state) states.delete(directory)
yield* state.dispose()
}),
)
return StateService.of(state)
}),
)
const map = yield* LayerMap.make((directory: string) => stateLayer(directory), { idleTimeToLive: "10 minutes" })
const off = registerDisposer((directory) =>
Effect.runPromise(
Effect.gen(function* () {
const state = states.get(directory)
yield* map.invalidate(directory)
if (state) {
yield* state.dispose()
if (states.get(directory) === state) states.delete(directory)
}
}),
),
)
yield* Effect.addFinalizer(() => Effect.sync(off))

const use = <A, E>(effect: Effect.Effect<A, E, StateService>): Effect.Effect<A, E> => {
const run = Effect.gen(function* () {
const ctx = yield* context
return yield* effect.pipe(Effect.provide(map.get(ctx.directory)))
})
return run.pipe(Effect.scoped)
}

const cancel = Effect.fn("Notebook.cancel")(function* (id: RequestID, reason: "cancelled" | "timeout") {
const pending = (yield* InstanceState.get(state)).pending
const pending = (yield* StateService).pending
const entry = pending.get(id)
if (!entry) return
pending.delete(id)
Expand All @@ -102,8 +149,8 @@ export function layer(timeout: Duration.Input = "10 minutes") {
)
})

const request: Interface["request"] = Effect.fn("Notebook.request")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const request = Effect.fn("Notebook.request")(function* (input: Input) {
const pending = (yield* StateService).pending
const id = RequestID.make(Identifier.create("nbr", "ascending"))
const deferred = yield* Deferred.make<Result, HostError>()
const info = { ...input, id } as Request
Expand All @@ -119,20 +166,20 @@ export function layer(timeout: Duration.Input = "10 minutes") {
}).pipe(Effect.ensuring(cancel(id, "cancelled")))
})

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

const cancelSession: Interface["cancelSession"] = Effect.fn("Notebook.cancelSession")(function* (sessionID) {
const pending = (yield* InstanceState.get(state)).pending
const cancelSession = Effect.fn("Notebook.cancelSession")(function* (sessionID: Request["sessionID"]) {
const pending = (yield* StateService).pending
const ids = Array.from(pending.values())
.filter((entry) => entry.info.sessionID === sessionID)
.map((entry) => entry.info.id)
yield* Effect.forEach(ids, (id) => cancel(id, "cancelled"), { discard: true })
})

const reply: Interface["reply"] = Effect.fn("Notebook.reply")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const reply = Effect.fn("Notebook.reply")(function* (input: Parameters<Interface["reply"]>[0]) {
const pending = (yield* StateService).pending
const entry = pending.get(input.requestID)
if (!entry) {
log.warn("reply for unknown request", { requestID: input.requestID })
Expand All @@ -141,10 +188,11 @@ export function layer(timeout: Duration.Input = "10 minutes") {
if (!matches(entry.info, input.result)) return yield* new InvalidReplyError({ requestID: input.requestID })
pending.delete(input.requestID)
yield* Deferred.succeed(entry.deferred, input.result)
return yield* Effect.void
})

const reject: Interface["reject"] = Effect.fn("Notebook.reject")(function* (input) {
const pending = (yield* InstanceState.get(state)).pending
const reject = Effect.fn("Notebook.reject")(function* (input: Parameters<Interface["reject"]>[0]) {
const pending = (yield* StateService).pending
const entry = pending.get(input.requestID)
if (!entry) {
log.warn("rejection for unknown request", { requestID: input.requestID })
Expand All @@ -161,9 +209,16 @@ export function layer(timeout: Duration.Input = "10 minutes") {
currentRevision: input.error.currentRevision,
}),
)
return yield* Effect.void
})

return Service.of({ request, list, cancelSession, reply, reject })
return Service.of({
request: (input) => use(request(input)),
list: () => use(list()),
cancelSession: (sessionID) => use(cancelSession(sessionID)),
reply: (input) => use(reply(input)),
reject: (input) => use(reject(input)),
})
}),
)
}
Expand Down
48 changes: 48 additions & 0 deletions packages/opencode/test/kilocode/notebook-service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { Effect, Fiber, Layer, Queue, Schema } from "effect"
import { TestInstance } from "../fixture/fixture"
import { disposeInstance } from "@/effect/instance-registry"
import { testEffect } from "../lib/effect"
import { InstanceRef } from "@/effect/instance-ref"

const it = testEffect(Notebook.layer("20 millis").pipe(Layer.provideMerge(Bus.layer)))
const sessionID = SessionID.make("ses_notebook_test")
Expand Down Expand Up @@ -199,11 +200,58 @@ it.instance(
Effect.gen(function* () {
const notebook = yield* Notebook.Service
const instance = yield* TestInstance
const cancelled = yield* Queue.unbounded<GlobalEvent>()
const handler = (event: GlobalEvent) => {
if (event.directory === instance.directory && event.payload?.type === Event.Cancelled.type)
Queue.offerUnsafe(cancelled, event)
}
GlobalBus.on("event", handler)
yield* Effect.addFinalizer(() => Effect.sync(() => GlobalBus.off("event", handler)))
const fiber = yield* request(notebook).pipe(Effect.forkChild)
yield* notebook.list().pipe(Effect.repeat({ until: (items) => items.length === 1 }))
yield* Effect.promise(() => disposeInstance(instance.directory))
const event = yield* Queue.take(cancelled).pipe(Effect.timeout("2 seconds"))
expect(event.payload?.properties).toMatchObject({ reason: "disposed" })
const err = yield* Fiber.join(fiber).pipe(Effect.flip)
expect(err.code).toBe("disconnected")
expect(yield* notebook.list()).toEqual([])
}),
{ git: true },
)

it.instance(
"isolates pending requests by directory and recreates disposed state",
() =>
Effect.gen(function* () {
const notebook = yield* Notebook.Service
const first = yield* TestInstance
const current = yield* InstanceRef
const second = current
? { ...current, directory: `${current.directory}-second` }
: yield* Effect.die(new Error("missing instance context"))
const firstFiber = yield* request(notebook).pipe(Effect.forkChild)
const firstPending = yield* notebook.list().pipe(Effect.repeat({ until: (items) => items.length === 1 }))

const secondResult = yield* Effect.gen(function* () {
expect(yield* notebook.list()).toEqual([])
const fiber = yield* request(notebook).pipe(Effect.forkChild)
const pending = yield* notebook.list().pipe(Effect.repeat({ until: (items) => items.length === 1 }))
yield* notebook.reply({
requestID: pending[0].id,
result: { operation: "read", path: "analysis.ipynb", requestPath: "analysis.ipynb", revision: "content:second", cells: [] },
})
return yield* Fiber.join(fiber)
}).pipe(Effect.provideService(InstanceRef, second))
expect(secondResult).toMatchObject({ revision: "content:second" })
expect(yield* notebook.list()).toEqual([firstPending[0]])

yield* Effect.promise(() => disposeInstance(first.directory))
expect((yield* Fiber.join(firstFiber).pipe(Effect.flip)).code).toBe("disconnected")

const replacement = yield* request(notebook).pipe(Effect.forkChild)
const replacementPending = yield* notebook.list().pipe(Effect.repeat({ until: (items) => items.length === 1 }))
expect(replacementPending[0].id).not.toBe(firstPending[0].id)
yield* Fiber.interrupt(replacement)
}),
{ git: true },
)
1 change: 0 additions & 1 deletion script/architecture-allowlist.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
"allowed": {
"packages/opencode/src/kilo-sessions/kilo-sessions.ts": { "count": 1, "owner": "session-runtime", "reason": "Kilo session coordination state" },
"packages/opencode/src/kilocode/background-process/index.ts": { "count": 1, "owner": "process-runtime", "reason": "Directory-keyed background process registry" },
"packages/opencode/src/kilocode/notebook/service.ts": { "count": 1, "owner": "notebook-runtime", "reason": "Notebook cell execution service state" },
"packages/opencode/src/kilocode/watcher.ts": { "count": 1, "owner": "watcher-runtime", "reason": "Eager location watcher subscription" }
}
},
Expand Down
Loading