Skip to content
Open
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
27 changes: 19 additions & 8 deletions packages/opencode/src/cli/tui/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import { writeHeapSnapshot } from "node:v8"
import { Heap } from "@/cli/heap"
import { AppRuntime } from "@/effect/app-runtime"
import { Effect } from "effect"
import { disposeAllInstancesAndEmitGlobalDisposed } from "@/server/global-lifecycle"
import { awaitSessionsIdle, disposeAllInstancesAndEmitGlobalDisposed } from "@/server/global-lifecycle"

Heap.start()

Expand All @@ -26,6 +26,7 @@ GlobalBus.on("event", (event) => {
})

let server: Awaited<ReturnType<typeof Server.listen>> | undefined
let reloading: Promise<void> | undefined

export const rpc = {
async fetch(input: { url: string; method: string; headers: Record<string, string>; body?: string }) {
Expand Down Expand Up @@ -61,13 +62,23 @@ export const rpc = {
await upgrade().catch(() => {})
},
async reload() {
await AppRuntime.runPromise(
Effect.gen(function* () {
const cfg = yield* Config.Service
yield* cfg.invalidate()
yield* disposeAllInstancesAndEmitGlobalDisposed({ swallowErrors: true })
}),
)
// SIGUSR2 arrives from desktop environments on theme changes, so a reload
// can land mid-run. Swapping config in disposes every instance, which
// cancels the session that is currently working — wait for it to finish
// instead. Signals that arrive while waiting join the pending reload.
if (!reloading) {
reloading = AppRuntime.runPromise(
Effect.gen(function* () {
yield* awaitSessionsIdle()
const cfg = yield* Config.Service
yield* cfg.invalidate()
yield* disposeAllInstancesAndEmitGlobalDisposed({ swallowErrors: true })
}),
).finally(() => {
reloading = undefined
})
}
await reloading
},
async shutdown() {
await InstanceRuntime.disposeAllInstances()
Expand Down
9 changes: 9 additions & 0 deletions packages/opencode/src/project/instance-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ export interface LoadInput {

export interface Interface {
readonly load: (input: LoadInput) => Effect.Effect<InstanceContext>
readonly list: () => Effect.Effect<InstanceContext[]>
readonly reload: (input: LoadInput) => Effect.Effect<InstanceContext>
readonly dispose: (ctx: InstanceContext) => Effect.Effect<void>
readonly disposeDirectory: (directory: string) => Effect.Effect<void>
Expand Down Expand Up @@ -123,6 +124,13 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
).pipe(Effect.withSpan("InstanceStore.load"))
}

const list = Effect.fn("InstanceStore.list")(function* () {
const exits = yield* Effect.forEach([...cache.values()], (entry) =>
Deferred.await(entry.deferred).pipe(Effect.exit),
)
return exits.filter(Exit.isSuccess).map((exit) => exit.value)
})

const reload = (input: LoadInput): Effect.Effect<InstanceContext> => {
const directory = FSUtil.resolve(input.directory)
return Effect.uninterruptibleMask((restore) =>
Expand Down Expand Up @@ -193,6 +201,7 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser

return Service.of({
load,
list,
reload,
dispose,
disposeDirectory,
Expand Down
22 changes: 22 additions & 0 deletions packages/opencode/src/server/global-lifecycle.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { GlobalBus } from "@/bus/global"
import { InstanceRef } from "@/effect/instance-ref"
import { InstanceStore } from "@/project/instance-store"
import { SessionStatus } from "@/session/status"
import { Effect } from "effect"
import { Event } from "./event"

Expand All @@ -25,4 +27,24 @@ export const disposeAllInstancesAndEmitGlobalDisposed = Effect.fn("Server.dispos
},
)

// Disposing an instance cancels every session runner it owns, so a config
// reload that lands while the model is streaming aborts the run. Callers that
// reload on an external trigger (SIGUSR2, config writes) wait here first so the
// reload is deferred rather than dropped.
export const awaitSessionsIdle = Effect.fn("Server.awaitSessionsIdle")(function* () {
while (yield* sessionsBusy) {
yield* Effect.sleep(IDLE_POLL_INTERVAL)
}
})

const IDLE_POLL_INTERVAL = "250 millis"

const sessionsBusy = Effect.gen(function* () {
const store = yield* InstanceStore.Service
const status = yield* SessionStatus.Service
const instances = yield* store.list()
const active = yield* Effect.forEach(instances, (ctx) => status.list().pipe(Effect.provideService(InstanceRef, ctx)))
return active.some((sessions) => sessions.size > 0)
})

export * as GlobalLifecycle from "./global-lifecycle"
13 changes: 13 additions & 0 deletions packages/opencode/test/project/instance.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,19 @@ describe("InstanceStore", () => {
}),
)

it.live("lists loaded instance contexts", () =>
Effect.gen(function* () {
const first = yield* tmpdirScoped({ git: true })
const second = yield* tmpdirScoped({ git: true })
const store = yield* InstanceStore.Service

yield* store.load({ directory: first })
yield* store.load({ directory: second })

expect((yield* store.list()).map((ctx) => ctx.directory).toSorted()).toEqual([first, second].toSorted())
}),
)

it.live("runs bootstrap with InstanceRef provided", () =>
Effect.gen(function* () {
const dir = yield* tmpdirScoped({ git: true })
Expand Down
54 changes: 54 additions & 0 deletions packages/opencode/test/server/global-lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
import { describe, expect } from "bun:test"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
import { Effect, Layer } from "effect"
import { InstanceRef } from "../../src/effect/instance-ref"
import { InstanceBootstrap } from "../../src/project/bootstrap-service"
import { InstanceStore } from "../../src/project/instance-store"
import { SessionID } from "../../src/session/schema"
import { SessionStatus } from "../../src/session/status"
import { awaitSessionsIdle } from "../../src/server/global-lifecycle"
import { tmpdirScoped } from "../fixture/fixture"
import { awaitWithTimeout, testEffect } from "../lib/effect"

const noopBootstrap = Layer.succeed(InstanceBootstrap.Service, InstanceBootstrap.Service.of({ run: Effect.void }))

const it = testEffect(
LayerNode.compile(LayerNode.group([InstanceStore.node, SessionStatus.node, CrossSpawnSpawner.node]), [
[InstanceStore.bootstrapNode, noopBootstrap],
]),
)

const sessionID = SessionID.make("ses_global_lifecycle")

describe("awaitSessionsIdle", () => {
it.live("resolves when no instance has a busy session", () =>
Effect.gen(function* () {
const dir = yield* tmpdirScoped({ git: true })
const store = yield* InstanceStore.Service
yield* store.load({ directory: dir })

yield* awaitWithTimeout(awaitSessionsIdle(), "awaitSessionsIdle blocked while idle")
}),
)

it.live("waits for a busy session to go idle", () =>
Effect.gen(function* () {
const dir = yield* tmpdirScoped({ git: true })
const store = yield* InstanceStore.Service
const ctx = yield* store.load({ directory: dir })
const status = yield* SessionStatus.Service

yield* status.set(sessionID, { type: "busy" }).pipe(Effect.provideService(InstanceRef, ctx))

const blocked = yield* awaitSessionsIdle().pipe(
Effect.as(false),
Effect.timeoutOrElse({ duration: "500 millis", orElse: () => Effect.succeed(true) }),
)
expect(blocked).toBe(true)

yield* status.set(sessionID, { type: "idle" }).pipe(Effect.provideService(InstanceRef, ctx))
yield* awaitWithTimeout(awaitSessionsIdle(), "awaitSessionsIdle did not resolve after the session went idle")
}),
)
})
Loading