fix(core): await initial plugin readiness (#35755)

This commit is contained in:
Kit Langton 2026-07-08 15:12:45 -04:00 committed by GitHub
commit 5b39972947
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 484 additions and 84 deletions

View file

@ -5,29 +5,22 @@ import { Catalog } from "./catalog"
import { CommandV2 } from "./command"
import { Config } from "./config"
import { LayerNode } from "./effect/layer-node"
import { makeLocationNode, Node } from "./effect/app-node"
import { httpClient } from "./effect/app-node-platform"
import { Node } from "./effect/app-node"
import { EventV2 } from "./event"
import { FileMutation } from "./file-mutation"
import { FileSystem } from "./filesystem"
import { FileSystemSearch } from "./filesystem/search"
import { FSUtil } from "./fs-util"
import { Generate } from "./generate"
import { Form } from "./form"
import { Global } from "./global"
import { LocationWatcher } from "./filesystem/location-watcher"
import { Image } from "./image"
import { LocationWatcher } from "./filesystem/location-watcher"
import { Integration } from "./integration"
import { Location } from "./location"
import { LocationMutation } from "./location-mutation"
import { LocationServiceMap } from "./location-service-map"
import { MCP } from "./mcp/index"
import { ModelsDev } from "./models-dev"
import { Npm } from "./npm"
import { PermissionV2 } from "./permission"
import { PluginV2 } from "./plugin"
import { PluginRuntime } from "./plugin/runtime"
import { SdkPlugins } from "./plugin/sdk"
import { PluginSupervisor } from "./plugin/supervisor"
import { ProjectCopy } from "./project/copy"
import { Pty } from "./pty"
@ -35,7 +28,6 @@ import { QuestionV2 } from "./question"
import { Shell } from "./shell"
import { Reference } from "./reference"
import { ReferenceGuidance } from "./reference/guidance"
import { Ripgrep } from "./ripgrep"
import { SessionRunnerLLM } from "./session/runner/llm"
import { SessionRunnerModel } from "./session/runner/model"
import { SessionCompaction } from "./session/compaction"
@ -51,49 +43,11 @@ import { SessionInstructions } from "./session/instructions"
import { McpTool } from "./tool/mcp"
import { ReadToolFileSystem } from "./tool/read-filesystem"
import { ToolRegistry } from "./tool/registry"
import { WebSearchTool } from "./tool/websearch"
import { ToolOutputStore } from "./tool-output-store"
import { Vcs } from "./vcs"
export { LocationServiceMap } from "./location-service-map"
const pluginSupervisorNode = makeLocationNode({
service: PluginSupervisor.Service,
layer: PluginSupervisor.layer,
deps: [
PluginV2.node,
SdkPlugins.node,
AgentV2.node,
Catalog.node,
CommandV2.node,
Config.node,
EventV2.node,
FileMutation.node,
FileSystem.node,
FSUtil.node,
Global.node,
httpClient,
Image.node,
Integration.node,
Location.node,
LocationMutation.node,
ModelsDev.node,
Npm.node,
PermissionV2.node,
PluginRuntime.node,
Form.node,
ReadToolFileSystem.node,
Reference.node,
Ripgrep.node,
SessionInstructions.node,
SessionTodo.node,
Shell.node,
SkillV2.node,
ToolRegistry.toolsNode,
WebSearchTool.configNode,
],
})
const locationServiceNodes = [
Location.node,
Config.node,
@ -104,7 +58,7 @@ const locationServiceNodes = [
Catalog.node,
AISDK.node,
PluginV2.node,
pluginSupervisorNode,
PluginSupervisor.node,
ProjectCopy.node,
ProjectCopy.refreshNode,
FileSystemSearch.node,

View file

@ -2,18 +2,42 @@ export * as PluginSupervisor from "./supervisor"
import type { Plugin } from "@opencode-ai/plugin/v2/effect/plugin"
import { Event } from "@opencode-ai/schema/config"
import { Context, Effect, Fiber, Layer, Option, Schema, Semaphore, Stream } from "effect"
import { Context, Deferred, Effect, Layer, Option, Schema, Semaphore, Stream } from "effect"
import path from "path"
import { fileURLToPath, pathToFileURL } from "url"
import { AgentV2 } from "../agent"
import { Catalog } from "../catalog"
import { CommandV2 } from "../command"
import { Config } from "../config"
import { ConfigPlugin } from "../config/plugin"
import { makeLocationNode } from "../effect/app-node"
import { httpClient } from "../effect/app-node-platform"
import { EventV2 } from "../event"
import { FileMutation } from "../file-mutation"
import { FileSystem } from "../filesystem"
import { Form } from "../form"
import { FSUtil } from "../fs-util"
import { Global } from "../global"
import { Image } from "../image"
import { Integration } from "../integration"
import { Location } from "../location"
import { LocationMutation } from "../location-mutation"
import { ModelsDev } from "../models-dev"
import { Npm } from "../npm"
import { PermissionV2 } from "../permission"
import { PluginV2 } from "../plugin"
import { PluginPromise } from "../plugin/promise"
import { Reference } from "../reference"
import { Ripgrep } from "../ripgrep"
import { SessionInstructions } from "../session/instructions"
import { SessionTodo } from "../session/todo"
import { Shell } from "../shell"
import { SkillV2 } from "../skill"
import { ReadToolFileSystem } from "../tool/read-filesystem"
import { ToolRegistry } from "../tool/registry"
import { WebSearchTool } from "../tool/websearch"
import { PluginInternal } from "./internal"
import { PluginRuntime } from "./runtime"
import { SdkPlugins } from "./sdk"
const PluginModule = Schema.Struct({
@ -246,7 +270,8 @@ const resolvePackageEntrypoint = Effect.fnUntraced(function* (fs: FSUtil.Interfa
})
export interface Interface {
readonly ready: Effect.Effect<void>
/** Wait for the initial plugin generation and startup updates to settle. */
readonly flush: Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/PluginSupervisor") {}
@ -259,35 +284,96 @@ const layer = Layer.effect(
const config = yield* Config.Service
const events = yield* EventV2.Service
const lock = Semaphore.makeUnsafe(1)
const reload = Effect.fn("PluginSupervisor.reload")(() =>
lock.withPermit(
const ready = yield* Deferred.make<void>()
let observed = 0
let applied = -1
const activate = Effect.fn("PluginSupervisor.activate")(function* (target: number) {
yield* lock.withPermit(
Effect.gen(function* () {
if (applied >= target) return
// Resolve OpenCode's internal plugins with their privileged Location services.
const internal = yield* PluginInternal.list()
// Combine internal plugins with host-contributed SDK plugins in boot order.
const pre = [...internal.pre, ...sdk.all()]
// Read the current layered config before resolving plugin directives and packages.
const entries = yield* config.entries()
const operations = yield* scan(entries)
const operations = yield* scan(yield* config.entries())
// Apply config operations and load enabled package plugins into one ordered generation.
const plugins = yield* resolve(pre, internal.post, operations)
// Replace the active generation in one scoped, batched activation.
yield* registry.activate(plugins)
applied = target
}),
),
)
yield* events.subscribe([Event.Updated, SdkPlugins.Updated]).pipe(
Stream.runForEach(() =>
reload().pipe(Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause }))),
)
})
const updates = yield* events
.subscribe([Event.Updated, SdkPlugins.Updated])
.pipe(Stream.toQueue({ capacity: 1, strategy: "sliding" }))
const signals = yield* Stream.concat(
Stream.succeed(0),
Stream.fromQueue(updates).pipe(Stream.mapEffect(() => Effect.sync(() => ++observed))),
).pipe(Stream.broadcast({ capacity: 1, strategy: "sliding", replay: 1 }))
const attempt = (target: number) =>
activate(target).pipe(
Effect.map(() => observed === target),
Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause }).pipe(Effect.as(false))),
)
yield* signals.pipe(
Stream.runForEach((target) =>
activate(target).pipe(Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause }))),
),
Effect.forkScoped({ startImmediately: true }),
)
const fiber = yield* reload().pipe(
Effect.withSpan("PluginSupervisor.boot"),
yield* signals.pipe(
Stream.debounce("100 millis"),
Stream.mapEffect(attempt),
Stream.filter((settled) => settled),
Stream.take(1),
Stream.runDrain,
Effect.andThen(Deferred.succeed(ready, undefined)),
Effect.forkScoped({ startImmediately: true }),
)
return Service.of({ ready: Fiber.join(fiber) })
return Service.of({ flush: Deferred.await(ready) })
}),
)
const nodeLayer = layer as Layer.Layer<Service, never, PluginInternal.Requirements>
export const node = makeLocationNode({
service: Service,
layer: nodeLayer,
deps: [
PluginV2.node,
SdkPlugins.node,
AgentV2.node,
Catalog.node,
CommandV2.node,
Config.node,
EventV2.node,
FileMutation.node,
FileSystem.node,
FSUtil.node,
Global.node,
httpClient,
Image.node,
Integration.node,
Location.node,
LocationMutation.node,
ModelsDev.node,
Npm.node,
PermissionV2.node,
PluginRuntime.node,
Form.node,
ReadToolFileSystem.node,
Reference.node,
Ripgrep.node,
SessionInstructions.node,
SessionTodo.node,
Shell.node,
SkillV2.node,
ToolRegistry.toolsNode,
WebSearchTool.configNode,
],
})
export { layer }

View file

@ -1,4 +1,5 @@
import { Schema } from "effect"
import { Agent } from "@opencode-ai/schema/agent"
import { SessionMessage } from "./message"
import { SessionSchema } from "./schema"
import { SessionError } from "@opencode-ai/schema/session-error"
@ -12,6 +13,15 @@ export class MessageDecodeError extends Schema.TaggedErrorClass<MessageDecodeErr
}
}
export class AgentNotFoundError extends Schema.TaggedErrorClass<AgentNotFoundError>()("Session.AgentNotFoundError", {
sessionID: SessionSchema.ID,
agent: Agent.ID,
}) {
override get message() {
return `Agent not found: "${this.agent}"`
}
}
export class StepFailedError extends Schema.TaggedErrorClass<StepFailedError>()("Session.StepFailedError", {
error: SessionError.Error,
}) {

View file

@ -3,7 +3,7 @@ export * as SessionRunner from "./index"
import type { LLMError } from "@opencode-ai/llm"
import { Context, Effect } from "effect"
import { SessionSchema } from "../schema"
import type { MessageDecodeError, StepFailedError, UserInterruptedError } from "../error"
import type { AgentNotFoundError, MessageDecodeError, StepFailedError, UserInterruptedError } from "../error"
import { SessionRunnerModel } from "./model"
import type { Instructions } from "../../instructions/index"
import type { ToolOutputStore } from "../../tool-output-store"
@ -12,6 +12,7 @@ export type RunError =
| LLMError
| SessionRunnerModel.Error
| MessageDecodeError
| AgentNotFoundError
| StepFailedError
| UserInterruptedError
| Instructions.InitializationBlocked

View file

@ -48,11 +48,12 @@ import { SessionRunnerSystemPrompt } from "./system-prompt"
import { Snapshot } from "../../snapshot"
import { makeLocationNode } from "../../effect/app-node"
import { llmClient } from "../../effect/app-node-platform"
import { StepFailedError, UserInterruptedError } from "../error"
import { AgentNotFoundError, StepFailedError, UserInterruptedError } from "../error"
import { toSessionError } from "../to-session-error"
import { SessionRunnerRetry } from "./retry"
import type { SessionHooks } from "@opencode-ai/plugin/v2/effect/session"
import { PluginHooks } from "../../plugin/hooks"
import { PluginSupervisor } from "../../plugin/supervisor"
type StepTokens = {
readonly input: number
@ -149,6 +150,7 @@ const layer = Layer.effect(
const db = (yield* Database.Service).db
const compaction = yield* SessionCompaction.Service
const title = yield* SessionTitle.Service
const plugins = yield* PluginSupervisor.Service
// Title generation is a side effect of the first step; it must not delay step continuation.
// Tracked per process so repeated wakes before the second user message arrives don't
// re-fire a redundant LLM call; `SessionTitle` itself is idempotent based on durable history.
@ -209,7 +211,10 @@ const layer = Layer.effect(
const session = yield* getSession(sessionID)
if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID)
return yield* Effect.interrupt
yield* plugins.flush
const agent = yield* agents.select(session.agent)
const agentInfo = agent.info
if (!agentInfo) return yield* new AgentNotFoundError({ sessionID: session.id, agent: session.agent ?? agent.id })
// Establish what the model knows before admitting what the user said, so
// a blocked first step leaves pending inputs untouched.
const checkpoint = yield* InstructionCheckpoint.prepare(
@ -218,9 +223,6 @@ const layer = Layer.effect(
loadInstructions(agent, session.id),
session.id,
)
const toolFibers = yield* FiberSet.make<void, ToolOutputStore.Error | UserInterruptedError>()
const ownedToolFibers: Array<Fiber.Fiber<void, ToolOutputStore.Error | UserInterruptedError>> = []
let needsContinuation = false
let currentStep = step
if (promotion) {
let promoted = 0
@ -236,18 +238,15 @@ const layer = Layer.effect(
const providerMetadataKey = model.route.providerMetadataKey ?? model.provider
const entries = yield* SessionHistory.entriesForRunner(db, session.id, checkpoint.baselineSeq)
const context = entries.map((entry) => entry.message)
const isLastStep = agent.info?.steps !== undefined && currentStep >= agent.info.steps
const isLastStep = agentInfo.steps !== undefined && currentStep >= agentInfo.steps
const toolMaterialization = isLastStep
? undefined
: yield* tools.materialize({ permissions: agent.info?.permissions, model })
: yield* tools.materialize({ permissions: agentInfo.permissions, model })
const promptCacheKey = /^ses_[0-9a-f]{64}$/.test(session.id) ? session.id.slice(4) : session.id
const request = LLM.request({
model,
providerOptions: { openai: { promptCacheKey } },
system: [
agent.info?.system ? agent.info.system : SessionRunnerSystemPrompt.provider(model),
checkpoint.baseline,
]
system: [agentInfo.system ? agentInfo.system : SessionRunnerSystemPrompt.provider(model), checkpoint.baseline]
.filter((part): part is string => part !== undefined && part.length > 0)
.map(SystemPart.make),
messages: [
@ -257,6 +256,9 @@ const layer = Layer.effect(
tools: toolMaterialization?.definitions ?? [],
toolChoice: isLastStep ? "none" : undefined,
})
const toolFibers = yield* FiberSet.make<void, ToolOutputStore.Error | UserInterruptedError>()
const ownedToolFibers: Array<Fiber.Fiber<void, ToolOutputStore.Error | UserInterruptedError>> = []
let needsContinuation = false
const availableTools = new Map(request.tools.map((tool) => [tool.name, tool]))
const requestEvent: SessionHooks["request"] = {
sessionID: session.id,
@ -436,7 +438,7 @@ const layer = Layer.effect(
if (
SessionRunnerRetry.isRetryable(llmFailure) &&
!publisher.hasRetryEvidence() &&
(agent.info?.steps === undefined || currentStep < agent.info.steps)
(agentInfo.steps === undefined || currentStep < agentInfo.steps)
) {
return yield* new SessionRunnerRetry.RetryableFailure({
cause: llmFailure,
@ -700,5 +702,6 @@ export const node = makeLocationNode({
Config.node,
Snapshot.node,
Database.node,
PluginSupervisor.node,
],
})

View file

@ -4,7 +4,7 @@ import { PermissionV2 } from "../permission"
import { QuestionV2 } from "../question"
import { Integration } from "../integration"
import { ToolOutputStore } from "../tool-output-store"
import { StepFailedError, UserInterruptedError } from "./error"
import { AgentNotFoundError, StepFailedError, UserInterruptedError } from "./error"
import { SessionRunnerModel } from "./runner/model"
export function toSessionError(cause: unknown): SessionError.Error {
@ -41,6 +41,7 @@ export function toSessionError(cause: unknown): SessionError.Error {
if (cause instanceof ToolFailure)
return cause.error === undefined ? { type: "tool.execution", message: cause.message } : toSessionError(cause.error)
if (cause instanceof StepFailedError) return cause.error
if (cause instanceof AgentNotFoundError) return { type: "unknown", message: cause.message }
if (cause instanceof UserInterruptedError) return { type: "aborted", message: cause.message }
if (
cause instanceof SessionRunnerModel.ModelNotSelectedError ||