Skip to content
Closed
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
30 changes: 22 additions & 8 deletions packages/core/src/plugin/supervisor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ export * as PluginSupervisor from "./supervisor.js"
export { Service, type Interface } from "./supervisor-service.js"

import { Event } from "@opencode-ai/schema/config"
import { Cause, Effect, Latch, Layer, Stream } from "effect"
import { Cause, Effect, Exit, Latch, Layer, Stream } from "effect"
import path from "path"
import { ConfigPluginSource } from "../config/plugin/source.js"
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
Expand Down Expand Up @@ -93,9 +93,8 @@ const resolve = Effect.fn("PluginSupervisor.resolve")(function* (
}
})

export const layer = Layer.effect(
Service,
Effect.gen(function* () {
function make(failOnError = false) {
return Effect.gen(function* () {
const registry = yield* Plugin.Service
const sdk = yield* SdkPlugins.Service
const instance = yield* InstancePlugins.Service
Expand All @@ -107,6 +106,7 @@ export const layer = Layer.effect(
let outdated = new Set<string>()
let generation = 0
let observed = 0
let activation = Exit.void

const activate = Effect.fn("PluginSupervisor.activate")(function* () {
const current = ++generation
Expand Down Expand Up @@ -184,16 +184,22 @@ export const layer = Layer.effect(
Stream.debounce("100 millis"),
Stream.runForEach((target) =>
Effect.gen(function* () {
yield* activate().pipe(Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause })))
activation = yield* Effect.exit(activate())
if (Exit.isFailure(activation))
yield* Effect.logError("failed to reload plugins", { cause: activation.cause })
if (observed === target) yield* ready.open
}),
),
Effect.forkScoped({ startImmediately: true }),
)
yield* Effect.sleep("24 hours").pipe(Effect.andThen(activate()), Effect.forever, Effect.forkScoped)
return Service.of({ awaitActivation: ready.await })
}),
)
return Service.of({
awaitActivation: failOnError ? ready.await.pipe(Effect.andThen(() => activation)) : ready.await,
})
})
}

export const layer = Layer.effect(Service, make())

const nodeDeps = [
Plugin.node,
Expand All @@ -212,3 +218,11 @@ function pluginSource(target: string): Plugin.Source {
}

export const node = makeLocationNode({ service: Service, layer, deps: nodeDeps })

/**
* Opt into propagating failed plugin generations through awaitActivation instead of only logging them.
* Individual plugin setup failures remain in the registry inventory.
*/
export function configured(options: { readonly failOnError?: boolean } = {}) {
return makeLocationNode({ service: Service, layer: Layer.effect(Service, make(options.failOnError)), deps: nodeDeps })
}
171 changes: 134 additions & 37 deletions packages/core/test/config/plugin.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ import { PluginSupervisor } from "@opencode-ai/core/plugin/supervisor"
import { Model } from "@opencode-ai/core/model"
import { Provider } from "@opencode-ai/core/provider"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { Cause, Effect, Fiber, Layer, Logger, Option, Schedule, Stream } from "effect"
import { Cause, Effect, Exit, Fiber, Layer, Logger, Option, Schedule, Stream } from "effect"
import { Database } from "../../src/database/database"
import { tmpdir } from "../fixture/tmpdir"
import { tempGlobalLayer } from "../fixture/global"
Expand Down Expand Up @@ -104,12 +104,6 @@ const coldNpm = makeGlobalNode({
),
deps: [Global.node],
})
const coldIt = testEffect(
AppNodeBuilder.build(
LayerNode.group([Database.node, Bus.node, SdkPlugins.node, LocationServiceMap.node, Global.node]),
[Global.node.replace(tempGlobalLayer), Npm.node.replace(coldNpm)],
),
)
describe("PluginSupervisor config", () => {
it.live("applies selectors in order", () =>
withLocation(
Expand Down Expand Up @@ -508,20 +502,6 @@ describe("PluginSupervisor config", () => {
),
)

it.live("unblocks awaitActivation when plugin activation fails", () =>
Effect.gen(function* () {
const sdk = yield* SdkPlugins.Service
yield* sdk.register(define({ id: "duplicate-id", effect: () => Effect.void }))
yield* sdk.register(define({ id: "duplicate-id", effect: () => Effect.void }))
yield* withLocation(
undefined,
Effect.gen(function* () {
yield* ready().pipe(Effect.timeout("2 seconds"))
}),
)
}),
)

updateIt.live("marks active package plugins as outdated after a background check", () =>
withLocation(
{ plugins: ["outdated-plugin"] },
Expand All @@ -538,23 +518,140 @@ describe("PluginSupervisor config", () => {
}),
),
)
})

coldIt.live("activates available plugins before a missing package finishes installing", () =>
withLocation(
{ plugins: ["cold-plugin"] },
Effect.gen(function* () {
const global = yield* Global.Service
yield* waitForFile(path.join(global.tmp, "cold-plugin", "started"))
const plugins = yield* Plugin.Service
expect((yield* plugins.list()).map((plugin) => String(plugin.id))).toContain("opencode.provider.openai")
const supervisor = yield* PluginSupervisor.Service
expect(Option.isNone(yield* supervisor.awaitActivation.pipe(Effect.timeoutOption("20 millis")))).toBeTrue()
yield* Effect.promise(() => Bun.write(path.join(global.tmp, "cold-plugin", "release"), ""))
yield* supervisor.awaitActivation.pipe(Effect.timeout("2 seconds"))
expect((yield* plugins.list()).map((plugin) => String(plugin.id))).toContain("cold-plugin")
}),
),
)
describe("PluginSupervisor awaitActivation", () => {
for (const [mode, node] of [
["default", PluginSupervisor.node],
["options omitted", PluginSupervisor.configured()],
["failOnError false", PluginSupervisor.configured({ failOnError: false })],
["strict", PluginSupervisor.configured({ failOnError: true })],
] as const) {
const activationIt = testEffect(
AppNodeBuilder.build(
LayerNode.group([Database.node, Bus.node, SdkPlugins.node, LocationServiceMap.node, Global.node]),
[Global.node.replace(tempGlobalLayer), Npm.node.replace(coldNpm), PluginSupervisor.node.replace(node)],
),
)

activationIt.live(`settles a host/builtin generation collision (${mode})`, () => {
const output: string[] = []
return Effect.gen(function* () {
const sdk = yield* SdkPlugins.Service
// Two registrations in the host store would overwrite, not collide.
yield* sdk.register(define({ id: "opencode.agent", effect: () => Effect.void }))
yield* withLocation(
undefined,
Effect.gen(function* () {
const exit = yield* ready().pipe(Effect.exit, Effect.timeout("2 seconds"))
expect(Exit.isFailure(exit)).toBe(mode === "strict")
if (Exit.isFailure(exit)) {
expect(Cause.hasDies(exit.cause)).toBeTrue()
expect(Cause.pretty(exit.cause)).toContain("Duplicate plugin ID: opencode.agent")
}
expect(yield* ready().pipe(Effect.exit, Effect.timeout("2 seconds"))).toEqual(exit)
expect(output.filter((line) => line.includes("failed to reload plugins"))).toEqual([
expect.stringContaining("Duplicate plugin ID: opencode.agent"),
])
const plugins = yield* Plugin.Service
expect(yield* plugins.list()).toEqual([])
}),
)
}).pipe(Effect.provide(Logger.layer([Logger.map(Logger.formatSimple, (line) => output.push(line))])))
})

if (mode !== "default" && mode !== "strict") continue

activationIt.live(`waits for package activation without cancelling it when a waiter is interrupted (${mode})`, () =>
withLocation(
{ plugins: ["cold-plugin"] },
Effect.gen(function* () {
const global = yield* Global.Service
yield* waitForFile(path.join(global.tmp, "cold-plugin", "started"))
const plugins = yield* Plugin.Service
expect((yield* plugins.list()).map((plugin) => String(plugin.id))).toContain("opencode.provider.openai")
const supervisor = yield* PluginSupervisor.Service
const waiter = yield* supervisor.awaitActivation.pipe(Effect.forkScoped({ startImmediately: true }))
expect(Option.isNone(yield* Fiber.await(waiter).pipe(Effect.timeoutOption("20 millis")))).toBeTrue()
yield* Fiber.interrupt(waiter)
expect(Exit.hasInterrupts(yield* Fiber.await(waiter))).toBeTrue()
expect(Option.isNone(yield* supervisor.awaitActivation.pipe(Effect.timeoutOption("20 millis")))).toBeTrue()
yield* Effect.promise(() => Bun.write(path.join(global.tmp, "cold-plugin", "release"), ""))
yield* supervisor.awaitActivation.pipe(Effect.timeout("2 seconds"))
yield* supervisor.awaitActivation.pipe(Effect.timeout("2 seconds"))
expect((yield* plugins.list()).find((plugin) => plugin.id === "cold-plugin")?.state).toEqual({
status: "active",
})
}),
),
)

if (mode !== "strict") continue

activationIt.live("keeps individual setup failures in inventory even with strict waits", () =>
withLocation(
{ plugins: ["-*", path.join(import.meta.dir, "../plugin/fixtures/failing")] },
Effect.gen(function* () {
yield* ready().pipe(Effect.timeout("2 seconds"))
const plugins = yield* Plugin.Service
expect(yield* plugins.list()).toEqual([
expect.objectContaining({
id: "failing-plugin",
state: { status: "failed", error: expect.stringContaining("plugin failed") },
}),
])
}),
),
)

activationIt.live("clears a failed reload only after the latest coalesced activation settles", () => {
const output: string[] = []
return withLocation(
{ plugins: ["cold-plugin"] },
Effect.gen(function* () {
const supervisor = yield* PluginSupervisor.Service
const sdk = yield* SdkPlugins.Service
const bus = yield* Bus.Service
const location = yield* Location.Service
const global = yield* Global.Service
yield* waitForFile(path.join(global.tmp, "cold-plugin", "started"))

// Every requested generation would collide with the builtin, even though the host store overwrites.
yield* Effect.forEach([1, 2, 3], () =>
sdk.register(define({ id: "opencode.agent", effect: () => Effect.void })),
)
const waiter = yield* supervisor.awaitActivation.pipe(
Effect.exit,
Effect.forkScoped({ startImmediately: true }),
)
expect(Option.isNone(yield* Fiber.await(waiter).pipe(Effect.timeoutOption("20 millis")))).toBeTrue()
yield* Effect.promise(() => Bun.write(path.join(global.tmp, "cold-plugin", "release"), ""))
const failed = yield* Fiber.join(waiter).pipe(Effect.timeout("2 seconds"))
expect(Exit.isFailure(failed) ? Cause.pretty(failed.cause) : "").toContain(
"Duplicate plugin ID: opencode.agent",
)
expect(yield* supervisor.awaitActivation.pipe(Effect.exit)).toEqual(failed)

// Wait for the repaired generation to reach the registry before testing readiness again.
const changed = yield* bus
.subscribe(Plugin.Event.Updated)
.pipe(Stream.take(1), Stream.runDrain, Effect.forkScoped({ startImmediately: true }))
yield* Effect.promise(() =>
Bun.write(
path.join(location.directory, "opencode.json"),
JSON.stringify({ plugins: ["-opencode.agent", "cold-plugin"] }),
),
)
yield* Fiber.join(changed).pipe(Effect.timeout("2 seconds"))
yield* supervisor.awaitActivation.pipe(Effect.timeout("2 seconds"))
yield* supervisor.awaitActivation.pipe(Effect.timeout("2 seconds"))
expect(output.filter((line) => line.includes("failed to reload plugins"))).toEqual([
expect.stringContaining("Duplicate plugin ID: opencode.agent"),
])
}),
).pipe(Effect.provide(Logger.layer([Logger.map(Logger.formatSimple, (line) => output.push(line))])))
})
}
})

const ready = Effect.fnUntraced(function* () {
Expand Down
Loading