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
167 changes: 137 additions & 30 deletions packages/core/src/plugin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { Plugin } from "@opencode-ai/schema/plugin"
import { Node } from "@opencode-ai/util/effect/app-node"
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
import type { PersistentPty } from "./persistent-pty.js"
import { Cause, Context, Effect, Exit, Latch, Layer, Logger, References, Scope, Semaphore } from "effect"
import { Cause, Context, Effect, Exit, Latch, Layer, Logger, Queue, References, Scope, Semaphore } from "effect"
import { Bus } from "./bus.js"
import { KV } from "./kv.js"
import { PluginHost } from "./plugin/host.js"
Expand All @@ -26,43 +26,64 @@ const layer = Layer.effect(
const lock = Semaphore.makeUnsafe(1)
const ready = yield* Latch.make(true)
const pending = new Set<object>()
const hold = () =>
Effect.sync(() => {
const token = {}
pending.add(token)
ready.closeUnsafe()
return Effect.sync(() => {
if (pending.delete(token) && pending.size === 0) ready.openUnsafe()
})
let closed = false
const holdUnsafe = () => {
if (closed) return Effect.void
const token = {}
pending.add(token)
ready.closeUnsafe()
return Effect.sync(() => {
if (pending.delete(token) && pending.size === 0) ready.openUnsafe()
})
}
const hold = () => Effect.sync(holdUnsafe)
const pendingFailures = yield* Queue.unbounded<PendingFailure>()
let discovered: readonly Failure[] = []
let inventory: Plugin.Info[] = []
const list = Effect.fn("Plugin.list")(function* () {
return inventory
})
const host = yield* PluginHost.make({ list })
const load = Effect.fnUntraced(function* (plugin: Generation) {
const child = yield* Scope.fork(scope)
const activation: Activation = { plugin, scope: yield* Scope.fork(scope) }
const inherit = yield* State.inherit()
const loaded = yield* Effect.suspend(() =>
const grouped = State.group((failure, refresh) => {
activation.failure = {
error: `Plugin disabled after ${failure.state}.transform failed. Check server logs for details.`,
ref: `err_${crypto.randomUUID().slice(0, 8)}`,
}
Queue.offerUnsafe(pendingFailures, {
plugin,
scope: activation.scope,
failure,
refresh,
ref: activation.failure.ref,
release: holdUnsafe(),
})
})
const exit = yield* Effect.suspend(() =>
plugin.effect({ ...host, storage: PluginHost.storage(kv, plugin.id) }),
).pipe(
grouped,
inherit,
Effect.updateContext((context: Context.Context<never>) =>
Context.make(Scope.Scope, child).pipe(
Context.make(Scope.Scope, activation.scope).pipe(
Context.add(Logger.CurrentLoggers, Context.get(context, Logger.CurrentLoggers)),
Context.add(References.MinimumLogLevel, Context.get(context, References.MinimumLogLevel)),
),
),
Effect.withSpan("Plugin.load", { attributes: { "plugin.id": plugin.id } }),
Effect.onExit((exit) => (Exit.isFailure(exit) ? Scope.close(child, exit) : Effect.void)),
Effect.onExit((exit) =>
Exit.isFailure(exit) && !activation.failure ? Scope.close(activation.scope, exit) : Effect.void,
),
Effect.exit,
)
if (Exit.isSuccess(loaded)) return { scope: child } as const
if (activation.failure || Exit.isSuccess(exit)) return { activation } as const
yield* Effect.logWarning("failed to load plugin", {
"plugin.id": plugin.id,
cause: loaded.cause,
cause: exit.cause,
})
return { error: Cause.pretty(loaded.cause) } as const
return { error: Cause.pretty(exit.cause) } as const
})

const activate = Effect.fn("Plugin.activate")(function* (
Expand All @@ -81,6 +102,8 @@ const layer = Layer.effect(
() =>
lock.withPermit(
Effect.gen(function* () {
if (closed) return
discovered = failures
const current = Array.from(active.values())
const changed = definitions.findIndex((definition, index) => {
const entry = current[index]
Expand Down Expand Up @@ -108,29 +131,36 @@ const layer = Layer.effect(
([id, slot]) =>
Effect.gen(function* () {
active.delete(id)
if (slot.loaded) yield* Scope.close(slot.loaded.scope, Exit.void)
if (slot.activation && !slot.activation.failure)
yield* Scope.close(slot.activation.scope, Exit.void)
}),
{ discard: true },
)
for (const definition of definitions.slice(prefix)) {
const loaded = yield* load(definition)
if (loaded.scope !== undefined) {
const slot = previous.get(definition.id)
// Reordering healthy registrations does not authorize retrying a failed revision.
if (slot?.activation?.failure && slot.plugin.revision === definition.revision) {
active.set(definition.id, { ...slot, plugin: definition })
continue
}
const result = yield* load(definition)
if (result.activation !== undefined) {
active.set(definition.id, {
plugin: definition,
loaded: { plugin: definition, scope: loaded.scope },
activation: result.activation,
})
continue
}
active.set(definition.id, { plugin: definition, error: loaded.error })
active.set(definition.id, { plugin: definition, error: result.error })

const fallback = previous.get(definition.id)?.loaded
if (!fallback) continue
const fallback = slot?.activation
if (!fallback || fallback.failure) continue
const restored = yield* load(fallback.plugin)
if (restored.scope !== undefined) {
if (restored.activation !== undefined) {
active.set(definition.id, {
plugin: definition,
loaded: { plugin: fallback.plugin, scope: restored.scope },
error: loaded.error,
activation: restored.activation,
error: result.error,
})
continue
}
Expand All @@ -149,9 +179,68 @@ const layer = Layer.effect(
)
})

yield* Queue.take(pendingFailures).pipe(
Effect.flatMap((item) =>
Effect.gen(function* () {
yield* Effect.logWarning("disabled plugin after transform failure", {
"plugin.id": item.plugin.id,
state: item.failure.state,
ref: item.ref,
cause: Cause.die(item.failure.cause),
})
yield* lock.withPermit(
Effect.gen(function* () {
if (closed) return
// Failure is already recorded on its exact activation, so an old queued item
// cannot disable a replacement and teardown need not wait for this worker.
inventory = [...Array.from(active.values()).map(slotInfo), ...discovered]
const refreshed = yield* State.batch(item.refresh).pipe(Effect.exit)
yield* bus.publish(Plugin.Event.Updated, {})
if (Exit.isFailure(refreshed))
yield* Effect.logWarning("failed to refresh state after disabling plugin", {
"plugin.id": item.plugin.id,
ref: item.ref,
cause: refreshed.cause,
})
}),
)
}).pipe(
// Cleanup must also be scheduled if an inventory observer fails. User finalizers
// may await readiness, so never join them under the activation lock or readiness hold.
Effect.ensuring(
Scope.close(item.scope, Exit.void).pipe(
Effect.catchCause((cause) =>
Effect.logWarning("failed to clean up disabled plugin", {
"plugin.id": item.plugin.id,
ref: item.ref,
cause,
}),
),
Effect.forkScoped({ startImmediately: true }),
),
),
Effect.catchCauseIf(
(cause) => !Cause.hasInterrupts(cause),
(cause) =>
Effect.logError("failed to report disabled plugin", {
"plugin.id": item.plugin.id,
ref: item.ref,
cause,
}),
),
Effect.ensuring(item.release),
),
),
Effect.forever,
Effect.forkScoped,
)

const close = (exit: Exit.Exit<unknown, unknown>) =>
lock.withPermit(
Effect.gen(function* () {
closed = true
pending.clear()
ready.openUnsafe()
active.clear()
yield* State.shutdown(Scope.close(scope, exit))
}),
Expand All @@ -168,19 +257,37 @@ const layer = Layer.effect(
}),
)

// `plugin` is the definition the slot was last asked to run; `loaded` is the generation actually
// running, which stays an older fallback while the requested revision keeps failing setup.
// `plugin` is the requested definition; `activation` is its last activation, which may have
// failed or be an older fallback while the requested revision keeps failing setup.
type Slot = {
readonly plugin: Generation
readonly loaded?: { readonly plugin: Generation; readonly scope: Scope.Closeable }
readonly activation?: Activation
readonly error?: string
}

// Share the activation across slot snapshots so teardown sees failures synchronously,
// including failures discovered after activate() has captured its previous slots.
type Activation = {
readonly plugin: Generation
readonly scope: Scope.Closeable
failure?: { readonly error: string; readonly ref: string }
}

type PendingFailure = {
readonly plugin: Generation
readonly scope: Scope.Closeable
readonly failure: State.Failure
readonly refresh: Effect.Effect<void>
readonly ref: string
readonly release: Effect.Effect<void>
}

function slotInfo(slot: Slot): Plugin.Info {
const failure = slot.activation?.failure ?? (slot.error === undefined ? undefined : { error: slot.error })
return {
id: Plugin.ID.make(slot.plugin.id),
source: slot.plugin.source ?? { type: "builtin" },
state: slot.error === undefined ? { status: "active" } : { status: "failed", error: slot.error },
state: failure === undefined ? { status: "active" } : { status: "failed", ...failure },
features: { server: true, ...slot.plugin.features },
}
}
Expand Down
105 changes: 86 additions & 19 deletions packages/core/src/state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,49 @@ export interface Transformable<Editor> {
readonly reload: Reload
}

export interface Failure {
readonly state: string
readonly cause: unknown
}

type GroupedRegistration = {
readonly remove: () => boolean
readonly notify: Effect.Effect<void>
}

type RegistrationGroup = {
failed: boolean
readonly registrations: Set<GroupedRegistration>
readonly report: (failure: Failure, refresh: Effect.Effect<void>) => void
}

const CurrentGroup = Context.Reference<RegistrationGroup | undefined>("@opencode/State/CurrentGroup", {
defaultValue: () => undefined,
})

/**
* Groups registrations without coupling State to plugin identity or asynchronous cleanup.
* A failed group is detached synchronously; its supervisor must run refresh and close its scope.
*/
export function group(report: RegistrationGroup["report"]) {
const group: RegistrationGroup = { failed: false, registrations: new Set(), report }
return <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.provideService(effect, CurrentGroup, group)
}

function disable(group: RegistrationGroup, failure: Failure) {
if (group.failed) return
group.failed = true
const notifications = new Set<Effect.Effect<void>>()
for (const registration of group.registrations) {
registration.remove()
notifications.add(registration.notify)
}
group.report(
failure,
Effect.forEach(notifications, (notify) => notify, { discard: true }),
)
}

type Batch = {
active: boolean
readonly shutdown: boolean
Expand Down Expand Up @@ -112,19 +155,38 @@ export interface Interface<State, Editor> extends Transformable<Editor> {

export function create<State, Editor>(options: Options<State, Editor>): Interface<State, Editor> {
let state = options.initial()
const transforms: { run: TransformCallback<Editor> }[] = []
const transforms = new Set<{ run: TransformCallback<Editor>; group: RegistrationGroup | undefined }>()
let dirty = false
let closed = false
let version = 0

const invalidate = () => {
dirty = true
version++
}

const get = () => {
if (closed || !dirty) return state
const next = options.initial()
const editor = options.editor(next)
for (const transform of transforms) transform.run(editor)
// Only a complete fold becomes visible; a throwing callback leaves the previous value and stays dirty.
state = next
dirty = false
return state
while (true) {
const started = version
const next = options.initial()
const editor = options.editor(next)
for (const transform of transforms) {
try {
transform.run(editor)
} catch (cause) {
if (!transform.group) throw cause
disable(transform.group, { state: options.name ?? "anonymous", cause })
}
// A nested read can disable a group that already contributed to this candidate.
if (version !== started) break
}
if (version !== started) continue
// Ungrouped failures still propagate; grouped failures restart from a fresh candidate.
state = next
dirty = false
return state
}
}

// One stable value per State, so a batch's notification Set holds it at most once.
Expand All @@ -137,7 +199,7 @@ export function create<State, Editor>(options: Options<State, Editor>): Interfac
const changed = Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
if (closed) return
dirty = true
invalidate()
const batch = yield* CurrentBatch
if (batch?.active) {
if (batch.shutdown) {
Expand All @@ -156,18 +218,23 @@ export function create<State, Editor>(options: Options<State, Editor>): Interfac
transform: Effect.fn("State.transform")(function* (update) {
yield* Effect.annotateCurrentSpan("state", options.name ?? "anonymous")
const scope = yield* Scope.Scope
const group = yield* CurrentGroup
if (group?.failed) return { dispose: Effect.void }
return yield* Effect.uninterruptible(
Effect.gen(function* () {
const transform = { run: update }
const dispose = Effect.uninterruptible(
Effect.suspend(() => {
const index = transforms.indexOf(transform)
if (index < 0) return Effect.void
transforms.splice(index, 1)
return changed
}),
)
transforms.push(transform)
const transform = { run: update, group }
const registration: GroupedRegistration = {
remove: () => {
if (!transforms.delete(transform)) return false
group?.registrations.delete(registration)
invalidate()
return true
},
notify: changed,
}
const dispose = Effect.uninterruptible(Effect.suspend(() => (registration.remove() ? changed : Effect.void)))
transforms.add(transform)
group?.registrations.add(registration)
yield* Scope.addFinalizer(scope, dispose)
yield* changed
return { dispose }
Expand Down
Loading
Loading