Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton 93bfa34b57 fix(core): discover local plugin packages 2026-08-11 12:08:04 -04:00
6 changed files with 136 additions and 48 deletions
+15 -28
View File
@@ -297,7 +297,7 @@ export const classifyHttpFailure = (input: {
})
}
export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray<string | RegExp>) => {
const toHttpError = (redactedNames: ReadonlyArray<string | RegExp>) => (error: unknown) => {
const transportError = (input: {
readonly message: string
readonly kind?: string | undefined
@@ -314,35 +314,23 @@ export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray<string
}),
})
const cause =
HttpClientError.isHttpClientError(error) && "cause" in error.reason
? error.reason.cause
: error instanceof Error
? error.cause
: undefined
const code = [cause, error]
.map((value) => (typeof value === "object" && value !== null ? Reflect.get(value, "code") : undefined))
.find((value): value is string => typeof value === "string")
const request = HttpClientError.isHttpClientError(error) && "request" in error ? error.request : undefined
const raw = cause instanceof Error ? cause.message : error instanceof Error ? error.message : undefined
const detail = raw && request ? redactBody(raw, secretValues(request)) : raw
const message = code && detail && !detail.includes(code) ? `${code}: ${detail}` : detail
if (Cause.isTimeoutError(error) || Cause.isTimeoutError(cause))
return transportError({ message: message ?? "HTTP transport timed out", kind: code ?? "Timeout", request })
if (!HttpClientError.isHttpClientError(error)) {
return transportError({ message: message ?? "HTTP transport failed", kind: code, request })
if (Cause.isTimeoutError(error)) {
return transportError({ message: error.message, kind: "Timeout" })
}
if (!HttpClientError.isHttpClientError(error)) {
return transportError({ message: error instanceof Error ? error.message : "HTTP transport failed" })
}
const request = "request" in error ? error.request : undefined
if (error.reason._tag === "TransportError") {
return transportError({
message: message ?? error.reason.description ?? "HTTP transport failed",
kind: code ?? error.reason._tag,
message: error.reason.description ?? "HTTP transport failed",
kind: error.reason._tag,
request,
})
}
return transportError({
message: message ?? `HTTP transport failed: ${error.reason._tag}`,
kind: code ?? error.reason._tag,
message: `HTTP transport failed: ${error.reason._tag}`,
kind: error.reason._tag,
request,
})
}
@@ -355,16 +343,15 @@ export const layer: Layer.Layer<Service, never, HttpClient.HttpClient> = Layer.e
Effect.gen(function* () {
const redactedNames = yield* Headers.CurrentRedactedNames
if (!middleware)
return yield* http.execute(request).pipe(
Effect.mapError((error) => mapHttpError(error, redactedNames)),
Effect.flatMap(statusError(request, redactedNames)),
)
return yield* http
.execute(request)
.pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames)))
const response = yield* middleware(request, (input) =>
http
.execute(input)
.pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))),
).pipe(Effect.mapError((error) => mapHttpError(error, redactedNames)))
).pipe(Effect.mapError(toHttpError(redactedNames)))
return yield* statusError(response.request, redactedNames)(response)
})
return Service.of({
+9 -6
View File
@@ -2,7 +2,6 @@ import { Effect, Stream } from "effect"
import { Headers, HttpClientRequest } from "effect/unstable/http"
import { Auth } from "../auth"
import { render as renderEndpoint } from "../endpoint"
import { mapHttpError } from "../executor"
import { Framing } from "../framing"
import type { HttpMiddleware, Transport, TransportPrepareInput } from "./index"
import * as ProviderShared from "../../protocols/shared"
@@ -87,16 +86,20 @@ export const httpJson = <Body, Frame>(input: HttpJsonInput<Body, Frame>): HttpJs
middleware: prepareInput.middleware,
}
}),
frames: (prepared, _request, runtime) =>
frames: (prepared, request, runtime) =>
Stream.unwrap(
runtime.http
.execute(prepared.request, prepared.middleware)
.pipe(
Effect.map((response) =>
Stream.unwrap(
Effect.map(Headers.CurrentRedactedNames, (redactedNames) =>
prepared.framing.frame(
response.stream.pipe(Stream.mapError((error) => mapHttpError(error, redactedNames))),
prepared.framing.frame(
response.stream.pipe(
Stream.mapError((error) =>
ProviderShared.eventError(
`${request.model.provider}/${request.model.route.id}`,
`Failed to read ${request.model.provider}/${request.model.route.id} stream`,
ProviderShared.errorText(error),
),
),
),
),
+2 -2
View File
@@ -63,14 +63,14 @@ export const dynamicResponse = (handler: Handler) => runtimeLayer(handlerLayer(h
* Layer that emits the supplied SSE chunks and then aborts mid-stream. Used to
* exercise transport errors that surface during parsing.
*/
export const truncatedStream = (chunks: ReadonlyArray<string>, error: Error = new Error("connection reset")) =>
export const truncatedStream = (chunks: ReadonlyArray<string>) =>
dynamicResponse((input) =>
Effect.sync(() => {
const encoder = new TextEncoder()
const stream = new ReadableStream({
start(controller) {
for (const chunk of chunks) controller.enqueue(encoder.encode(chunk))
controller.error(error)
controller.error(new Error("connection reset"))
},
})
return input.respond(stream, { headers: SSE_HEADERS })
+4 -10
View File
@@ -1221,18 +1221,12 @@ describe("OpenAI Chat route", () => {
it.effect("surfaces transport errors that occur mid-stream", () =>
Effect.gen(function* () {
const layer = truncatedStream(
[`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`],
Object.assign(new Error("socket closed unexpectedly"), { code: "ECONNRESET" }),
)
const layer = truncatedStream([
`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`,
])
const error = yield* LLMClient.generate(request).pipe(Effect.provide(layer), Effect.flip)
expect(error.reason).toMatchObject({
_tag: "Transport",
message: "ECONNRESET: socket closed unexpectedly",
kind: "ECONNRESET",
url: "https://api.openai.test/v1/chat/completions",
})
expect(error.message).toContain("Failed to read openai/openai-chat stream")
}),
)
+39 -2
View File
@@ -4,7 +4,7 @@ import { Directory, Document, type Entry } from "@opencode-ai/schema/config"
import { ConfigPlugin } from "@opencode-ai/schema/config/plugin"
import { FSUtil } from "@opencode-ai/util/fs-util"
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
import { Context, Effect, Layer, Option, PubSub, Scope, Stream } from "effect"
import { Context, Effect, Layer, Option, PubSub, Schema, Scope, Stream } from "effect"
import path from "path"
import { fileURLToPath } from "url"
import { Config } from "../../config"
@@ -154,6 +154,12 @@ const scan = Effect.fn("ConfigPluginSource.scan")(function* (
})
const sourceDirectories = ["plugin", "plugins"] as const
const Package = Schema.Struct({
exports: Schema.optional(Schema.Unknown),
module: Schema.optional(Schema.Unknown),
main: Schema.optional(Schema.Unknown),
})
const decodePackage = Schema.decodeUnknownOption(Package)
function discoverDirectory(fs: FSUtil.Interface, directory: string) {
return Effect.gen(function* () {
@@ -166,7 +172,38 @@ function discoverDirectory(fs: FSUtil.Interface, directory: string) {
symlink: true,
})
.pipe(Effect.orElseSucceed(() => []))
return files.sort().map((target): Operation => ({ type: "add", target, options: {} }))
const children = yield* fs
.scan(`{${sourceDirectories.join(",")}}/*`, {
cwd: directory,
absolute: true,
include: "all",
dot: true,
symlink: true,
})
.pipe(Effect.orElseSucceed(() => []))
const directories = yield* Effect.filter(children.sort(), fs.isDir)
const packages = yield* Effect.forEach(directories, (child) => discoverPackage(fs, child))
return [...files.sort(), ...packages.filter((target): target is string => typeof target === "string")].map(
(target): Operation => ({ type: "add", target, options: {} }),
)
})
}
function discoverPackage(fs: FSUtil.Interface, directory: string) {
return Effect.gen(function* () {
const manifest = yield* fs
.readJson(path.join(directory, "package.json"))
.pipe(Effect.map(decodePackage), Effect.orElseSucceed(Option.none))
const configured = Option.isSome(manifest)
? [manifest.value.exports, manifest.value.module, manifest.value.main].filter(
(entry): entry is string => typeof entry === "string",
)
: []
const target = yield* Effect.findFirst(
[...configured, "index.ts", "index.js"].map((entry) => path.resolve(directory, entry)),
fs.isFile,
)
return Option.getOrUndefined(target)
})
}
+67
View File
@@ -168,6 +168,69 @@ describe("PluginSupervisor config", () => {
),
)
it.live("loads auto-discovered plugin packages from package metadata", () =>
withLocation(
undefined,
Effect.gen(function* () {
yield* ready()
const plugins = yield* Plugin.Service
expect((yield* plugins.list()).map((plugin) => String(plugin.id))).toContain("package-metadata")
}),
false,
async (directory) => {
const plugin = path.join(directory, ".opencode", "plugins", "package-metadata")
await fs.mkdir(plugin, { recursive: true })
await fs.writeFile(path.join(plugin, "package.json"), JSON.stringify({ exports: "./entry.ts" }))
await fs.writeFile(path.join(plugin, "entry.ts"), discoveredPlugin("package-metadata"))
},
),
)
it.live("loads auto-discovered plugin packages from index fallback", () =>
withLocation(
undefined,
Effect.gen(function* () {
yield* ready()
const plugins = yield* Plugin.Service
expect((yield* plugins.list()).map((plugin) => String(plugin.id))).toContain("index-fallback")
}),
false,
async (directory) => {
const plugin = path.join(directory, ".opencode", "plugins", "index-fallback")
await fs.mkdir(plugin, { recursive: true })
await fs.writeFile(path.join(plugin, "index.js"), discoveredPlugin("index-fallback"))
},
),
)
it.live("prefers package metadata over index fallback", () =>
withLocation(
undefined,
Effect.gen(function* () {
yield* ready()
const plugins = yield* Plugin.Service
const ids = (yield* plugins.list()).map((plugin) => String(plugin.id))
expect(ids).toContain("metadata-precedence")
expect(ids).not.toContain("module-collision")
expect(ids).not.toContain("main-collision")
expect(ids).not.toContain("index-collision")
}),
false,
async (directory) => {
const plugin = path.join(directory, ".opencode", "plugins", "collision")
await fs.mkdir(plugin, { recursive: true })
await fs.writeFile(
path.join(plugin, "package.json"),
JSON.stringify({ exports: "./entry.js", module: "./module.js", main: "./main.js" }),
)
await fs.writeFile(path.join(plugin, "entry.js"), discoveredPlugin("metadata-precedence"))
await fs.writeFile(path.join(plugin, "module.js"), discoveredPlugin("module-collision"))
await fs.writeFile(path.join(plugin, "main.js"), discoveredPlugin("main-collision"))
await fs.writeFile(path.join(plugin, "index.js"), discoveredPlugin("index-collision"))
},
),
)
staticIt.live("uses only internal and SDK plugins when the static source is wired", () =>
Effect.gen(function* () {
const sdk = yield* SdkPlugins.Service
@@ -389,3 +452,7 @@ export default Plugin.define({
})
`
}
function discoveredPlugin(id: string) {
return `export default { id: ${JSON.stringify(id)}, setup() {} }`
}