mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-11 20:19:53 -04:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a44316f9f2 |
@@ -422,7 +422,8 @@ async function loadCatalog(client: OpenCodeClient, cwd: string): Promise<Catalog
|
||||
defaultModel: {
|
||||
providerID: defaultModel.providerID,
|
||||
id: defaultModel.id,
|
||||
variant: defaultModel.variants.find((variant) => variant.id === "default")?.id,
|
||||
variant:
|
||||
defaultModel.variants.find((variant) => variant.id === "default")?.id ?? defaultModel.variants[0]?.id,
|
||||
},
|
||||
modes: agents.map((agent) => ({ id: agent.id, name: agent.name, description: agent.description })),
|
||||
defaultModeID: defaultAgent.id,
|
||||
|
||||
@@ -39,6 +39,8 @@ export const run = Effect.fnUntraced(function* (options: Options) {
|
||||
})
|
||||
|
||||
const processEffect = Effect.fnUntraced(function* (options: Options) {
|
||||
const serviceErrorFormat = process.env.OPENCODE_SERVICE_ERROR_FORMAT
|
||||
delete process.env.OPENCODE_SERVICE_ERROR_FORMAT
|
||||
const global = yield* Global.Service
|
||||
if (options.mode === "service") yield* Effect.sync(() => process.chdir(global.home))
|
||||
return yield* Effect.scoped(
|
||||
@@ -127,15 +129,7 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
|
||||
if (serviceOptions === undefined || port === undefined || !addressInUse(error)) return Effect.fail(error)
|
||||
return recognizeIncumbent(serviceOptions, hostname, port).pipe(
|
||||
Effect.flatMap((found) =>
|
||||
found
|
||||
? Effect.void
|
||||
: Effect.fail(
|
||||
new Error(
|
||||
`Managed service port ${port} on ${hostname} is already in use by another process. ` +
|
||||
"Configure another port with `opencode service set port <port>` and start the service again.",
|
||||
{ cause: error },
|
||||
),
|
||||
),
|
||||
found ? Effect.void : managedPortInUse(hostname, port, error, serviceErrorFormat),
|
||||
),
|
||||
)
|
||||
}),
|
||||
@@ -214,6 +208,17 @@ function serviceURL(hostname: string, port: number) {
|
||||
return `http://${hostname.includes(":") ? `[${hostname}]` : hostname}:${port}`
|
||||
}
|
||||
|
||||
function managedPortInUse(hostname: string, port: number, cause: unknown, format?: string) {
|
||||
const message =
|
||||
`Managed service port ${port} on ${hostname} is already in use by another process. ` +
|
||||
"Configure another port with `opencode service set port <port>` and start the service again."
|
||||
const failure = new Error(message, { cause })
|
||||
if (format !== "plain") return Effect.fail(failure)
|
||||
return Effect.sync(() => process.stderr.write(`OPENCODE_SERVICE_ERROR:${message}\n`)).pipe(
|
||||
Effect.andThen(Effect.fail(failure)),
|
||||
)
|
||||
}
|
||||
|
||||
function truthy(value?: string) {
|
||||
return value === "1" || value?.toLowerCase() === "true"
|
||||
}
|
||||
|
||||
@@ -3,38 +3,6 @@ import type { SessionConfigOption } from "@agentclientprotocol/sdk"
|
||||
import { makeACPFixture, makeSession, secondModel } from "./service-fixture"
|
||||
|
||||
describe("acp service lifecycle", () => {
|
||||
test("does not persist the first catalog variant when no explicit default exists", async () => {
|
||||
const model = { ...secondModel, variants: [{ id: "none" }, { id: "high" }] }
|
||||
await using fixture = makeACPFixture({
|
||||
models: [model],
|
||||
defaultModel: model,
|
||||
fetch(request) {
|
||||
if (request.method === "POST" && request.path === "/api/session") {
|
||||
return Response.json({
|
||||
data: makeSession("ses_default_variant", {
|
||||
model: { providerID: model.providerID, id: model.id },
|
||||
}),
|
||||
})
|
||||
}
|
||||
return undefined
|
||||
},
|
||||
})
|
||||
|
||||
const created = await fixture.service.newSession({ cwd: "/workspace", mcpServers: [] })
|
||||
|
||||
expect(fixture.requests).toContainEqual({
|
||||
method: "POST",
|
||||
path: "/api/session",
|
||||
query: {},
|
||||
body: {
|
||||
location: { directory: "/workspace" },
|
||||
agent: "build",
|
||||
model: { providerID: "test", id: "second-model" },
|
||||
},
|
||||
})
|
||||
expect(currentValue(created, "effort")).toBe("none")
|
||||
})
|
||||
|
||||
test("loads and forks with paginated replay while resume does not replay", async () => {
|
||||
await using fixture = makeACPFixture({
|
||||
fetch(request) {
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
import { ServiceStatus } from "@opencode-ai/protocol/groups/health"
|
||||
import { Effect, FileSystem, Option, Schedule, Schema } from "effect"
|
||||
import { spawn, type ChildProcess } from "node:child_process"
|
||||
import { homedir } from "node:os"
|
||||
import { join } from "node:path"
|
||||
import type { DiscoverOptions, Endpoint, EnsureOptions, StopOptions } from "../service.js"
|
||||
import { ServiceProcess } from "../service-process.js"
|
||||
|
||||
export * from "../service.js"
|
||||
/** Contents of the local service registration file. */
|
||||
@@ -17,11 +17,6 @@ export type Info = import("../service.js").Info
|
||||
// is all a client needs to connect. The daemon's own configuration (port,
|
||||
// persisted password) is CLI-owned and never read here.
|
||||
|
||||
type Contender = {
|
||||
readonly child: ChildProcess
|
||||
readonly error: () => Error | undefined
|
||||
}
|
||||
|
||||
// Read-only lookup: registration file plus health check and version gate.
|
||||
// Never spawns; escalation to ensure() is the caller's policy.
|
||||
/** Discover a healthy, compatible local service without starting one. */
|
||||
@@ -52,11 +47,12 @@ const discoverLocal = Effect.fnUntraced(function* (options: DiscoverOptions) {
|
||||
// becomes discoverable. A contender is never killed merely for slow startup.
|
||||
/** Ensure a healthy, compatible local service is running. */
|
||||
export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOptions = {}) {
|
||||
const contenders = new Set<Contender>()
|
||||
const contenders = new Set<ServiceProcess.Contender>()
|
||||
let timeouts: { readonly info: Info; readonly count: number } | undefined
|
||||
let announced = false
|
||||
let lastSpawn = 0
|
||||
let spawnDelay = 5_000
|
||||
let lastFailure: Error | undefined
|
||||
const announce = (reason: "missing" | "version-mismatch", previousVersion?: string) =>
|
||||
Effect.sync(() => {
|
||||
if (announced) return
|
||||
@@ -67,15 +63,7 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
|
||||
const [command, ...args] = options.command ?? ["opencode", "serve", "--service"]
|
||||
if (command === undefined) return yield* Effect.fail(new Error("Missing service command"))
|
||||
return yield* Effect.try({
|
||||
try: () => {
|
||||
const child = spawn(command, args, { detached: true, stdio: "ignore" })
|
||||
let error: Error | undefined
|
||||
child.once("error", (cause) => {
|
||||
error = new Error("Failed to start server", { cause })
|
||||
})
|
||||
child.unref()
|
||||
return { child, error: () => error }
|
||||
},
|
||||
try: () => ServiceProcess.start(command, args),
|
||||
catch: (cause) => new Error("Failed to start server", { cause }),
|
||||
})
|
||||
})
|
||||
@@ -108,8 +96,9 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
|
||||
return Option.none<LocalService>()
|
||||
} else if (lastSpawn === 0 && info !== undefined) lastSpawn = Date.now()
|
||||
|
||||
const finished = [...contenders].filter(contenderFinished)
|
||||
const failure = finished.map(contenderFailure).find((error): error is Error => error !== undefined)
|
||||
const finished = [...contenders].filter(ServiceProcess.finished)
|
||||
const failure = finished.map(ServiceProcess.failure).find((error): error is Error => error !== undefined)
|
||||
if (failure !== undefined) lastFailure = failure
|
||||
if (finished.some((item) => item.child.exitCode === 0)) {
|
||||
spawnDelay = Math.min(spawnDelay * 2, 30_000)
|
||||
}
|
||||
@@ -129,24 +118,10 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
|
||||
}),
|
||||
)
|
||||
if (Option.isNone(found))
|
||||
return yield* Effect.fail(new Error("Timed out waiting for the background service to start"))
|
||||
return yield* Effect.fail(lastFailure ?? new Error("Timed out waiting for the background service to start"))
|
||||
return found.value.endpoint
|
||||
})
|
||||
|
||||
function contenderFailure(contender: Contender) {
|
||||
const error = contender.error()
|
||||
if (error !== undefined) return error
|
||||
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
|
||||
return new Error(`Server process exited with code ${contender.child.exitCode}`)
|
||||
if (contender.child.signalCode !== null)
|
||||
return new Error(`Server process terminated by ${contender.child.signalCode}`)
|
||||
return undefined
|
||||
}
|
||||
|
||||
function contenderFinished(contender: Contender) {
|
||||
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
|
||||
}
|
||||
|
||||
/** Stop the registered local service. */
|
||||
export const stop = Effect.fn("service.stop")(function* (options: StopOptions = {}) {
|
||||
const existing = yield* find(options)
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import { readFile } from "node:fs/promises"
|
||||
import { spawn, type ChildProcess } from "node:child_process"
|
||||
import { homedir } from "node:os"
|
||||
import { join } from "node:path"
|
||||
import type { DiscoverOptions, Endpoint, Info, EnsureOptions, StopOptions } from "../service.js"
|
||||
import { ServiceProcess } from "../service-process.js"
|
||||
import type { ServiceHealth, ServiceStopResponse } from "./generated/types.js"
|
||||
|
||||
export * from "../service.js"
|
||||
@@ -13,11 +13,6 @@ export * from "../service.js"
|
||||
// intentionally implemented with Node APIs so Promise clients do not need
|
||||
// Effect or @effect/platform-node at runtime.
|
||||
|
||||
type Contender = {
|
||||
readonly child: ChildProcess
|
||||
readonly error: () => Error | undefined
|
||||
}
|
||||
|
||||
/** Discover a healthy, compatible local service without starting one. */
|
||||
export async function discover(options: DiscoverOptions = {}) {
|
||||
return (await discoverLocal(options))?.endpoint
|
||||
@@ -33,11 +28,12 @@ async function discoverLocal(options: DiscoverOptions) {
|
||||
/** Ensure a healthy, compatible local service is running. */
|
||||
export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
|
||||
const deadline = Date.now() + 120_000
|
||||
const contenders = new Set<Contender>()
|
||||
const contenders = new Set<ServiceProcess.Contender>()
|
||||
let timeouts: { readonly info: Info; readonly count: number } | undefined
|
||||
let announced = false
|
||||
let lastSpawn = 0
|
||||
let spawnDelay = 5_000
|
||||
let lastFailure: Error | undefined
|
||||
|
||||
const announce = (reason: "missing" | "version-mismatch", previousVersion?: string) => {
|
||||
if (announced) return
|
||||
@@ -47,21 +43,11 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
|
||||
const spawnContender = () => {
|
||||
const [command, ...args] = options.command ?? ["opencode", "serve", "--service"]
|
||||
if (command === undefined) throw new Error("Missing service command")
|
||||
try {
|
||||
const child = spawn(command, args, { detached: true, stdio: "ignore" })
|
||||
let error: Error | undefined
|
||||
child.once("error", (cause) => {
|
||||
error = new Error("Failed to start server", { cause })
|
||||
})
|
||||
child.unref()
|
||||
return { child, error: () => error }
|
||||
} catch (cause) {
|
||||
throw new Error("Failed to start server", { cause })
|
||||
}
|
||||
return ServiceProcess.start(command, args)
|
||||
}
|
||||
|
||||
while (true) {
|
||||
if (Date.now() >= deadline) throw new Error("Timed out waiting for the background service to start")
|
||||
if (Date.now() >= deadline) throw lastFailure ?? new Error("Timed out waiting for the background service to start")
|
||||
const registration = await registered(options.file, true)
|
||||
if (registration.timedOut && registration.info !== undefined) {
|
||||
timeouts = {
|
||||
@@ -89,8 +75,9 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
|
||||
}
|
||||
} else {
|
||||
if (lastSpawn === 0 && registration.info !== undefined) lastSpawn = Date.now()
|
||||
const finished = [...contenders].filter(contenderFinished)
|
||||
const failure = finished.map(contenderFailure).find((error) => error !== undefined)
|
||||
const finished = [...contenders].filter(ServiceProcess.finished)
|
||||
const failure = finished.map(ServiceProcess.failure).find((error) => error !== undefined)
|
||||
if (failure !== undefined) lastFailure = failure
|
||||
if (finished.some((item) => item.child.exitCode === 0)) {
|
||||
spawnDelay = Math.min(spawnDelay * 2, 30_000)
|
||||
}
|
||||
@@ -107,20 +94,6 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
|
||||
}
|
||||
}
|
||||
|
||||
function contenderFailure(contender: Contender) {
|
||||
const error = contender.error()
|
||||
if (error !== undefined) return error
|
||||
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
|
||||
return new Error(`Server process exited with code ${contender.child.exitCode}`)
|
||||
if (contender.child.signalCode !== null)
|
||||
return new Error(`Server process terminated by ${contender.child.signalCode}`)
|
||||
return undefined
|
||||
}
|
||||
|
||||
function contenderFinished(contender: Contender) {
|
||||
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
|
||||
}
|
||||
|
||||
/** Stop the registered local service. */
|
||||
export async function stop(options: StopOptions = {}) {
|
||||
const existing = await find(options)
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
export * as ServiceProcess from "./service-process"
|
||||
|
||||
import { spawn, type ChildProcess } from "node:child_process"
|
||||
|
||||
const errorPrefix = "OPENCODE_SERVICE_ERROR:"
|
||||
|
||||
export type Contender = {
|
||||
readonly child: ChildProcess
|
||||
readonly error: () => Error | undefined
|
||||
readonly startupError: () => string
|
||||
}
|
||||
|
||||
export function start(command: string, args: ReadonlyArray<string>) {
|
||||
try {
|
||||
const child = spawn(command, args, {
|
||||
detached: true,
|
||||
stdio: ["ignore", "ignore", "pipe"],
|
||||
env: { ...process.env, OPENCODE_SERVICE_ERROR_FORMAT: "plain" },
|
||||
})
|
||||
let error: Error | undefined
|
||||
let pending = ""
|
||||
let startupError = ""
|
||||
child.once("error", (cause) => {
|
||||
error = new Error("Failed to start server", { cause })
|
||||
})
|
||||
child.stderr?.on("data", (chunk) => {
|
||||
const lines = (pending + chunk.toString()).split(/\r?\n/)
|
||||
pending = lines.pop()?.slice(-64 * 1024) ?? ""
|
||||
const message = lines.findLast((line) => line.startsWith(errorPrefix))
|
||||
if (message !== undefined) startupError = message.slice(errorPrefix.length)
|
||||
})
|
||||
unref(child.stderr)
|
||||
child.unref()
|
||||
return { child, error: () => error, startupError: () => startupError } satisfies Contender
|
||||
} catch (cause) {
|
||||
throw new Error("Failed to start server", { cause })
|
||||
}
|
||||
}
|
||||
|
||||
export function failure(contender: Contender) {
|
||||
const error = contender.error()
|
||||
if (error !== undefined) return error
|
||||
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
|
||||
return new Error(contender.startupError() || `Server process exited with code ${contender.child.exitCode}`)
|
||||
if (contender.child.signalCode !== null)
|
||||
return new Error(`Server process terminated by ${contender.child.signalCode}`)
|
||||
return undefined
|
||||
}
|
||||
|
||||
export function finished(contender: Contender) {
|
||||
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
|
||||
}
|
||||
|
||||
function unref(stream: ChildProcess["stderr"]) {
|
||||
if (!stream || !("unref" in stream) || typeof stream.unref !== "function") return
|
||||
stream.unref()
|
||||
}
|
||||
@@ -3,6 +3,11 @@ import { appendFile, rename, writeFile } from "node:fs/promises"
|
||||
const [registration, mode, delay] = process.argv.slice(2)
|
||||
if (registration === undefined || mode === undefined) throw new Error("Missing service fixture arguments")
|
||||
if (mode === "failed") process.exit(1)
|
||||
if (mode === "failed-message") {
|
||||
console.error("sensitive startup detail")
|
||||
console.error("OPENCODE_SERVICE_ERROR:Managed service port is already in use")
|
||||
process.exit(1)
|
||||
}
|
||||
if (mode === "record-start") {
|
||||
await writeFile(registration + ".started", "")
|
||||
process.exit(1)
|
||||
|
||||
@@ -70,6 +70,19 @@ test("reports a failed registered service", async () => {
|
||||
)
|
||||
})
|
||||
|
||||
test("reports the native contender's startup error", async () => {
|
||||
const directory = await temp()
|
||||
const registration = join(directory, "service.json")
|
||||
|
||||
await expect(
|
||||
Service.ensure({
|
||||
file: registration,
|
||||
version: "test",
|
||||
command: [process.execPath, fixture, registration, "failed-message"],
|
||||
}),
|
||||
).rejects.toThrow(/^Managed service port is already in use$/)
|
||||
}, 10_000)
|
||||
|
||||
test("evicts an unresponsive registered service before starting its replacement", async () => {
|
||||
const directory = await temp()
|
||||
const registration = join(directory, "service.json")
|
||||
|
||||
@@ -197,6 +197,20 @@ test("reports a contender that fails to start", async () => {
|
||||
).rejects.toThrow("Server process exited with code 1")
|
||||
}, 10_000)
|
||||
|
||||
test("reports the contender's startup error", async () => {
|
||||
const directory = await temp()
|
||||
const registration = join(directory, "service.json")
|
||||
await expect(
|
||||
run(
|
||||
Service.ensure({
|
||||
file: registration,
|
||||
version: "test",
|
||||
command: [process.execPath, fixture, registration, "failed-message"],
|
||||
}),
|
||||
),
|
||||
).rejects.toThrow(/^Managed service port is already in use$/)
|
||||
}, 10_000)
|
||||
|
||||
test("reports a contender terminated by a signal", async () => {
|
||||
const directory = await temp()
|
||||
const registration = join(directory, "service.json")
|
||||
|
||||
@@ -1,75 +0,0 @@
|
||||
import type { APIEvent } from "@solidjs/start/server"
|
||||
import { and, Database, eq, isNull } from "@opencode-ai/console-core/drizzle/index.js"
|
||||
import { BillingTable, LiteTable } from "@opencode-ai/console-core/schema/billing.sql.js"
|
||||
import { KeyTable } from "@opencode-ai/console-core/schema/key.sql.js"
|
||||
import { LiteData } from "@opencode-ai/console-core/lite.js"
|
||||
import { Subscription } from "@opencode-ai/console-core/subscription.js"
|
||||
|
||||
export async function GET(input: APIEvent) {
|
||||
const token = input.request.headers.get("authorization")?.match(/^Bearer (.+)$/)?.[1]
|
||||
if (!token) return Response.json({ error: "Unauthorized" }, { status: 401 })
|
||||
|
||||
const row = await Database.use((tx) =>
|
||||
tx
|
||||
.select({
|
||||
balance: BillingTable.balance,
|
||||
monthlyLimit: BillingTable.monthlyLimit,
|
||||
monthlyUsage: BillingTable.monthlyUsage,
|
||||
useBalance: BillingTable.lite,
|
||||
rollingUsage: LiteTable.rollingUsage,
|
||||
weeklyUsage: LiteTable.weeklyUsage,
|
||||
goMonthlyUsage: LiteTable.monthlyUsage,
|
||||
timeRollingUpdated: LiteTable.timeRollingUpdated,
|
||||
timeWeeklyUpdated: LiteTable.timeWeeklyUpdated,
|
||||
timeMonthlyUpdated: LiteTable.timeMonthlyUpdated,
|
||||
timeSubscribed: LiteTable.timeCreated,
|
||||
})
|
||||
.from(KeyTable)
|
||||
.innerJoin(BillingTable, eq(BillingTable.workspaceID, KeyTable.workspaceID))
|
||||
.leftJoin(
|
||||
LiteTable,
|
||||
and(
|
||||
eq(LiteTable.workspaceID, KeyTable.workspaceID),
|
||||
eq(LiteTable.userID, KeyTable.userID),
|
||||
isNull(LiteTable.timeDeleted),
|
||||
),
|
||||
)
|
||||
.where(and(eq(KeyTable.key, token), isNull(KeyTable.timeDeleted)))
|
||||
.then((rows) => rows[0]),
|
||||
)
|
||||
if (!row) return Response.json({ error: "Unauthorized" }, { status: 401 })
|
||||
|
||||
const limits = row.timeSubscribed ? LiteData.getLimits() : undefined
|
||||
return Response.json({
|
||||
go:
|
||||
limits && row.timeSubscribed
|
||||
? {
|
||||
useBalance: row.useBalance?.useBalance ?? false,
|
||||
rolling: Subscription.analyzeRollingUsage({
|
||||
limit: limits.rollingLimit,
|
||||
window: limits.rollingWindow,
|
||||
usage: row.rollingUsage ?? 0,
|
||||
timeUpdated: row.timeRollingUpdated ?? new Date(),
|
||||
}),
|
||||
weekly: Subscription.analyzeWeeklyUsage({
|
||||
limit: limits.weeklyLimit,
|
||||
usage: row.weeklyUsage ?? 0,
|
||||
timeUpdated: row.timeWeeklyUpdated ?? new Date(),
|
||||
}),
|
||||
monthly: Subscription.analyzeMonthlyUsage({
|
||||
limit: limits.monthlyLimit,
|
||||
usage: row.goMonthlyUsage ?? 0,
|
||||
timeUpdated: row.timeMonthlyUpdated ?? new Date(),
|
||||
timeSubscribed: row.timeSubscribed,
|
||||
}),
|
||||
}
|
||||
: undefined,
|
||||
zen: {
|
||||
balance: row.balance / 100_000_000,
|
||||
monthly: {
|
||||
usage: (row.monthlyUsage ?? 0) / 100_000_000,
|
||||
limit: row.monthlyLimit ?? undefined,
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -1,8 +1,8 @@
|
||||
{
|
||||
"version": "7",
|
||||
"dialect": "sqlite",
|
||||
"id": "00924d88-1842-4d71-ac74-5682ddc47e1c",
|
||||
"prevIds": ["15060ec5-05f7-4b86-b2a5-9108609432b3"],
|
||||
"id": "15060ec5-05f7-4b86-b2a5-9108609432b3",
|
||||
"prevIds": ["1551a157-8959-4ba9-a52b-4ea3b7b28cae"],
|
||||
"ddl": [
|
||||
{
|
||||
"name": "account_state",
|
||||
@@ -1302,16 +1302,6 @@
|
||||
"entityType": "columns",
|
||||
"table": "session_v2"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": "0",
|
||||
"generated": null,
|
||||
"name": "resume_attempts",
|
||||
"entityType": "columns",
|
||||
"table": "session_v2"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
|
||||
@@ -27,10 +27,8 @@ export const Plugin = define({
|
||||
const changes = yield* PubSub.sliding<string>(1)
|
||||
const lock = Semaphore.makeUnsafe(1)
|
||||
const start = yield* fs.resolve(location.directory)
|
||||
const root = yield* fs.resolve(location.project.directory)
|
||||
const home = yield* fs.resolve(global.home)
|
||||
const project = discovery.project && FSUtil.contains(root, start)
|
||||
const stop = FSUtil.contains(home, start) ? home : root
|
||||
const stop = yield* fs.resolve(location.project.directory)
|
||||
const project = discovery.project && FSUtil.contains(stop, start)
|
||||
const globalFile = yield* fs.resolve(join(global.config, "AGENTS.md"))
|
||||
const loaded: { current: Loaded } = { current: { type: "available", files: [] } }
|
||||
|
||||
|
||||
@@ -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, Predicate, PubSub, Schema, Scope, Stream } from "effect"
|
||||
import { Context, Effect, Layer, Option, PubSub, Scope, Stream } from "effect"
|
||||
import path from "path"
|
||||
import { fileURLToPath } from "url"
|
||||
import { Config } from "../../config"
|
||||
@@ -154,67 +154,19 @@ 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* () {
|
||||
const children = (yield* Effect.forEach(sourceDirectories, (source) =>
|
||||
fs.readDirectoryEntries(path.join(directory, source)).pipe(
|
||||
Effect.orElseSucceed(() => []),
|
||||
Effect.map((entries) =>
|
||||
entries.map((entry) => ({ ...entry, target: path.join(directory, source, entry.name) })),
|
||||
),
|
||||
),
|
||||
))
|
||||
.flat()
|
||||
.sort((a, b) => (a.target < b.target ? -1 : a.target > b.target ? 1 : 0))
|
||||
const targets = yield* Effect.forEach(children, (entry) => discoverChild(fs, entry))
|
||||
return targets.flatMap(Option.toArray).map((target): Operation => ({ type: "add", target, options: {} }))
|
||||
})
|
||||
}
|
||||
|
||||
function discoverChild(fs: FSUtil.Interface, entry: FSUtil.DirEntry & { target: string }) {
|
||||
return Effect.gen(function* () {
|
||||
const source = entry.target.endsWith(".ts") || entry.target.endsWith(".js")
|
||||
if (entry.type === "file" && source) return Option.some(entry.target)
|
||||
if (entry.type === "directory") return yield* discoverPackage(fs, entry.target)
|
||||
if (entry.type !== "symlink") return Option.none<string>()
|
||||
if (source && (yield* fs.isFile(entry.target))) return Option.some(entry.target)
|
||||
if (yield* fs.isDir(entry.target)) return yield* discoverPackage(fs, entry.target)
|
||||
return Option.none<string>()
|
||||
})
|
||||
}
|
||||
|
||||
function discoverPackage(fs: FSUtil.Interface, directory: string) {
|
||||
return Effect.gen(function* () {
|
||||
const root = yield* fs.resolve(directory)
|
||||
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(Predicate.isString)
|
||||
: []
|
||||
return yield* Effect.findFirst(
|
||||
[...configured, "index.ts", "index.js"]
|
||||
.filter((entry) => !path.isAbsolute(entry))
|
||||
.map((entry) => path.resolve(directory, entry))
|
||||
.filter((entry) => FSUtil.contains(directory, entry)),
|
||||
(entry) =>
|
||||
fs
|
||||
.isFile(entry)
|
||||
.pipe(
|
||||
Effect.flatMap((exists) =>
|
||||
exists
|
||||
? fs.resolve(entry).pipe(Effect.map((resolved) => FSUtil.contains(root, resolved)))
|
||||
: Effect.succeed(false),
|
||||
),
|
||||
),
|
||||
)
|
||||
const files = yield* fs
|
||||
.scan(`{${sourceDirectories.join(",")}}/*.{ts,js}`, {
|
||||
cwd: directory,
|
||||
absolute: true,
|
||||
include: "file",
|
||||
dot: true,
|
||||
symlink: true,
|
||||
})
|
||||
.pipe(Effect.orElseSucceed(() => []))
|
||||
return files.sort().map((target): Operation => ({ type: "add", target, options: {} }))
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
-2
@@ -40,7 +40,6 @@ import m37 from "./migration/20260622202450_simplify_session_input"
|
||||
import m38 from "./migration/20260804233008_loose_psylocke"
|
||||
import m39 from "./migration/20260805200742_import_legacy_credentials"
|
||||
import m40 from "./migration/20260808023530_workspace_domain"
|
||||
import m41 from "./migration/20260811161259_execution_claim_attempts"
|
||||
|
||||
export const migrations = [
|
||||
m00,
|
||||
@@ -84,5 +83,4 @@ export const migrations = [
|
||||
m38,
|
||||
m39,
|
||||
m40,
|
||||
m41,
|
||||
] satisfies DatabaseMigration.Migration[]
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
import { Effect } from "effect"
|
||||
import type { DatabaseMigration } from "../migration"
|
||||
|
||||
const migration: DatabaseMigration.Migration = {
|
||||
id: "20260811161259_execution_claim_attempts",
|
||||
up(tx) {
|
||||
return Effect.gen(function* () {
|
||||
yield* tx.run(`ALTER TABLE \`session_v2\` ADD \`resume_attempts\` integer DEFAULT 0 NOT NULL;`)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
export default migration
|
||||
@@ -200,7 +200,6 @@ const schema: Omit<DatabaseMigration.Migration, "id"> = {
|
||||
\`time_compacting\` integer,
|
||||
\`time_archived\` integer,
|
||||
\`time_suspended\` integer,
|
||||
\`resume_attempts\` integer DEFAULT 0 NOT NULL,
|
||||
CONSTRAINT \`fk_session_v2_project_id_project_id_fk\` FOREIGN KEY (\`project_id\`) REFERENCES \`project\`(\`id\`) ON DELETE CASCADE
|
||||
);
|
||||
`)
|
||||
|
||||
@@ -1,32 +1,58 @@
|
||||
import { Database, type SQLQueryBindings } from "bun:sqlite"
|
||||
import { Database } from "bun:sqlite"
|
||||
import { drizzle } from "drizzle-orm/bun-sqlite"
|
||||
import { Context, Effect, Layer } from "effect"
|
||||
import { Context, Effect, Fiber, Layer, Scope, Semaphore, Stream } from "effect"
|
||||
import { identity } from "effect/Function"
|
||||
import { Reactivity } from "effect/unstable/reactivity"
|
||||
import { SqlClient } from "effect/unstable/sql"
|
||||
import { SqlClient, Statement } from "effect/unstable/sql"
|
||||
import type { Connection } from "effect/unstable/sql/SqlConnection"
|
||||
import { classifySqliteError, SqlError } from "effect/unstable/sql/SqlError"
|
||||
import { Sqlite } from "./sqlite"
|
||||
|
||||
const TypeId = "~@opencode-ai/core/database/SqliteBun" as const
|
||||
const ATTR_DB_SYSTEM_NAME = "db.system.name"
|
||||
|
||||
interface Config extends Sqlite.ClientConfig {
|
||||
const TypeId = "~@opencode-ai/core/database/SqliteBun" as const
|
||||
type TypeId = typeof TypeId
|
||||
|
||||
interface SqliteClient extends SqlClient.SqlClient {
|
||||
readonly [TypeId]: TypeId
|
||||
readonly config: Config
|
||||
readonly export: Effect.Effect<Uint8Array, SqlError>
|
||||
readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
|
||||
readonly updateValues: never
|
||||
}
|
||||
|
||||
interface Config {
|
||||
readonly filename: string
|
||||
readonly readonly?: boolean
|
||||
readonly create?: boolean
|
||||
readonly readwrite?: boolean
|
||||
readonly disableWAL?: boolean
|
||||
readonly spanAttributes?: Record<string, unknown>
|
||||
readonly transformResultNames?: (str: string) => string
|
||||
readonly transformQueryNames?: (str: string) => string
|
||||
}
|
||||
|
||||
interface SqliteConnection extends Connection {
|
||||
readonly export: Effect.Effect<Uint8Array, SqlError>
|
||||
readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
|
||||
}
|
||||
|
||||
const make = (options: Config) =>
|
||||
Effect.gen(function* () {
|
||||
const native = (yield* Sqlite.Native) as Database
|
||||
|
||||
const compiler = Statement.makeCompilerSqlite(options.transformQueryNames)
|
||||
const transformRows = options.transformResultNames
|
||||
? Statement.defaultTransforms(options.transformResultNames).array
|
||||
: undefined
|
||||
|
||||
const run = (query: string, params: ReadonlyArray<unknown> = []) =>
|
||||
Effect.withFiber<Array<Record<string, unknown>>, SqlError>((fiber) => {
|
||||
const statement = native.query<Record<string, unknown>, SQLQueryBindings[]>(query)
|
||||
const statement = native.query(query)
|
||||
// @ts-ignore bun-types missing safeIntegers method, fixed in https://github.com/oven-sh/bun/pull/26627
|
||||
statement.safeIntegers(Context.get(fiber.context, SqlClient.SafeIntegers))
|
||||
try {
|
||||
return Effect.succeed(statement.all(...(params as SQLQueryBindings[])) ?? [])
|
||||
return Effect.succeed((statement.all(...(params as any)) ?? []) as Array<Record<string, unknown>>)
|
||||
} catch (cause) {
|
||||
return Effect.fail(
|
||||
new SqlError({
|
||||
@@ -38,11 +64,11 @@ const make = (options: Config) =>
|
||||
|
||||
const runValues = (query: string, params: ReadonlyArray<unknown> = []) =>
|
||||
Effect.withFiber<Array<unknown[]>, SqlError>((fiber) => {
|
||||
const statement = native.query<unknown, SQLQueryBindings[]>(query)
|
||||
const statement = native.query(query)
|
||||
// @ts-ignore bun-types missing safeIntegers method, fixed in https://github.com/oven-sh/bun/pull/26627
|
||||
statement.safeIntegers(Context.get(fiber.context, SqlClient.SafeIntegers))
|
||||
try {
|
||||
return Effect.succeed(statement.values(...(params as SQLQueryBindings[])) ?? [])
|
||||
return Effect.succeed((statement.values(...(params as any)) ?? []) as Array<unknown[]>)
|
||||
} catch (cause) {
|
||||
return Effect.fail(
|
||||
new SqlError({
|
||||
@@ -52,7 +78,25 @@ const make = (options: Config) =>
|
||||
}
|
||||
})
|
||||
|
||||
const connection = Sqlite.makeConnection(run, runValues, {
|
||||
const connection = identity<SqliteConnection>({
|
||||
execute(query, params, transformRows) {
|
||||
return transformRows ? Effect.map(run(query, params), transformRows) : run(query, params)
|
||||
},
|
||||
executeRaw(query, params) {
|
||||
return run(query, params)
|
||||
},
|
||||
executeValues(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeValuesUnprepared(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeUnprepared(query, params, transformRows) {
|
||||
return this.execute(query, params, transformRows)
|
||||
},
|
||||
executeStream() {
|
||||
return Stream.die("executeStream not implemented")
|
||||
},
|
||||
export: Effect.try({
|
||||
try: () => native.serialize(),
|
||||
catch: (cause) =>
|
||||
@@ -60,7 +104,7 @@ const make = (options: Config) =>
|
||||
reason: classifySqliteError(cause, { message: "Failed to export database", operation: "export" }),
|
||||
}),
|
||||
}),
|
||||
loadExtension: (path: string) =>
|
||||
loadExtension: (path) =>
|
||||
Effect.try({
|
||||
try: () => native.loadExtension(path),
|
||||
catch: (cause) =>
|
||||
@@ -70,10 +114,37 @@ const make = (options: Config) =>
|
||||
}),
|
||||
})
|
||||
|
||||
return yield* Sqlite.makeClient(options, connection, TypeId, (acquirer) => ({
|
||||
export: Effect.flatMap(acquirer, (_) => _.export),
|
||||
loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
|
||||
}))
|
||||
const semaphore = yield* Semaphore.make(1)
|
||||
const acquirer = semaphore.withPermits(1)(Effect.succeed(connection))
|
||||
const transactionAcquirer = Effect.uninterruptibleMask((restore) => {
|
||||
const fiber = Fiber.getCurrent()!
|
||||
const scope = Context.getUnsafe(fiber.context, Scope.Scope)
|
||||
return Effect.as(
|
||||
Effect.tap(restore(semaphore.take(1)), () => Scope.addFinalizer(scope, semaphore.release(1))),
|
||||
connection,
|
||||
)
|
||||
})
|
||||
|
||||
const client = Object.assign(
|
||||
(yield* SqlClient.make({
|
||||
acquirer,
|
||||
compiler,
|
||||
transactionAcquirer,
|
||||
spanAttributes: [
|
||||
...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
|
||||
[ATTR_DB_SYSTEM_NAME, "sqlite"],
|
||||
],
|
||||
transformRows,
|
||||
})) as SqliteClient,
|
||||
{
|
||||
[TypeId]: TypeId,
|
||||
config: options,
|
||||
export: Effect.flatMap(acquirer, (_) => _.export),
|
||||
loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
|
||||
},
|
||||
)
|
||||
|
||||
return client
|
||||
})
|
||||
|
||||
const nativeLayer = (config: Config) =>
|
||||
|
||||
@@ -1,14 +1,26 @@
|
||||
import { DatabaseSync, type SQLInputValue } from "node:sqlite"
|
||||
import { drizzle } from "drizzle-orm/node-sqlite"
|
||||
import { Context, Effect, Layer } from "effect"
|
||||
import { Context, Effect, Fiber, Layer, Scope, Semaphore, Stream } from "effect"
|
||||
import { identity } from "effect/Function"
|
||||
import { Reactivity } from "effect/unstable/reactivity"
|
||||
import { SqlClient } from "effect/unstable/sql"
|
||||
import { SqlClient, Statement } from "effect/unstable/sql"
|
||||
import type { Connection } from "effect/unstable/sql/SqlConnection"
|
||||
import { classifySqliteError, SqlError } from "effect/unstable/sql/SqlError"
|
||||
import { Sqlite } from "./sqlite"
|
||||
|
||||
const TypeId = "~@opencode-ai/core/database/SqliteNode" as const
|
||||
const ATTR_DB_SYSTEM_NAME = "db.system.name"
|
||||
|
||||
interface Config extends Sqlite.ClientConfig {
|
||||
const TypeId = "~@opencode-ai/core/database/SqliteNode" as const
|
||||
type TypeId = typeof TypeId
|
||||
|
||||
interface SqliteClient extends SqlClient.SqlClient {
|
||||
readonly [TypeId]: TypeId
|
||||
readonly config: Config
|
||||
readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
|
||||
readonly updateValues: never
|
||||
}
|
||||
|
||||
interface Config {
|
||||
readonly filename: string
|
||||
readonly readonly?: boolean
|
||||
readonly create?: boolean
|
||||
@@ -16,12 +28,24 @@ interface Config extends Sqlite.ClientConfig {
|
||||
readonly disableWAL?: boolean
|
||||
readonly timeout?: number
|
||||
readonly allowExtension?: boolean
|
||||
readonly spanAttributes?: Record<string, unknown>
|
||||
readonly transformResultNames?: (str: string) => string
|
||||
readonly transformQueryNames?: (str: string) => string
|
||||
}
|
||||
|
||||
interface SqliteConnection extends Connection {
|
||||
readonly loadExtension: (path: string) => Effect.Effect<void, SqlError>
|
||||
}
|
||||
|
||||
const make = (options: Config) =>
|
||||
Effect.gen(function* () {
|
||||
const native = (yield* Sqlite.Native) as DatabaseSync
|
||||
|
||||
const compiler = Statement.makeCompilerSqlite(options.transformQueryNames)
|
||||
const transformRows = options.transformResultNames
|
||||
? Statement.defaultTransforms(options.transformResultNames).array
|
||||
: undefined
|
||||
|
||||
const run = (query: string, params: ReadonlyArray<unknown> = []) =>
|
||||
Effect.withFiber<Array<Record<string, unknown>>, SqlError>((fiber) => {
|
||||
const statement = native.prepare(query)
|
||||
@@ -55,8 +79,26 @@ const make = (options: Config) =>
|
||||
}
|
||||
})
|
||||
|
||||
const connection = Sqlite.makeConnection(run, runValues, {
|
||||
loadExtension: (path: string) =>
|
||||
const connection = identity<SqliteConnection>({
|
||||
execute(query, params, transformRows) {
|
||||
return transformRows ? Effect.map(run(query, params), transformRows) : run(query, params)
|
||||
},
|
||||
executeRaw(query, params) {
|
||||
return run(query, params)
|
||||
},
|
||||
executeValues(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeValuesUnprepared(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeUnprepared(query, params, transformRows) {
|
||||
return this.execute(query, params, transformRows)
|
||||
},
|
||||
executeStream() {
|
||||
return Stream.die("executeStream not implemented")
|
||||
},
|
||||
loadExtension: (path) =>
|
||||
Effect.try({
|
||||
try: () => native.loadExtension(path),
|
||||
catch: (cause) =>
|
||||
@@ -66,9 +108,36 @@ const make = (options: Config) =>
|
||||
}),
|
||||
})
|
||||
|
||||
return yield* Sqlite.makeClient(options, connection, TypeId, (acquirer) => ({
|
||||
loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
|
||||
}))
|
||||
const semaphore = yield* Semaphore.make(1)
|
||||
const acquirer = semaphore.withPermits(1)(Effect.succeed(connection))
|
||||
const transactionAcquirer = Effect.uninterruptibleMask((restore) => {
|
||||
const fiber = Fiber.getCurrent()!
|
||||
const scope = Context.getUnsafe(fiber.context, Scope.Scope)
|
||||
return Effect.as(
|
||||
Effect.tap(restore(semaphore.take(1)), () => Scope.addFinalizer(scope, semaphore.release(1))),
|
||||
connection,
|
||||
)
|
||||
})
|
||||
|
||||
const client = Object.assign(
|
||||
(yield* SqlClient.make({
|
||||
acquirer,
|
||||
compiler,
|
||||
transactionAcquirer,
|
||||
spanAttributes: [
|
||||
...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
|
||||
[ATTR_DB_SYSTEM_NAME, "sqlite"],
|
||||
],
|
||||
transformRows,
|
||||
})) as SqliteClient,
|
||||
{
|
||||
[TypeId]: TypeId,
|
||||
config: options,
|
||||
loadExtension: (path: string) => Effect.flatMap(acquirer, (_) => _.loadExtension(path)),
|
||||
},
|
||||
)
|
||||
|
||||
return client
|
||||
})
|
||||
|
||||
const nativeLayer = (config: Config) =>
|
||||
|
||||
@@ -1,100 +1,8 @@
|
||||
export * as Sqlite from "./sqlite"
|
||||
|
||||
import { Context, Effect, Fiber, Scope, Semaphore, Stream } from "effect"
|
||||
import { identity } from "effect/Function"
|
||||
import { SqlClient, Statement } from "effect/unstable/sql"
|
||||
import type { Connection } from "effect/unstable/sql/SqlConnection"
|
||||
import type { SqlError } from "effect/unstable/sql/SqlError"
|
||||
import { Context } from "effect"
|
||||
import type { drizzle } from "drizzle-orm/bun-sqlite"
|
||||
|
||||
export type DrizzleClient = ReturnType<typeof drizzle>
|
||||
export class Native extends Context.Service<Native, unknown>()("@opencode-ai/core/database/SqliteNative") {}
|
||||
export class Drizzle extends Context.Service<Drizzle, DrizzleClient>()("@opencode-ai/core/database/SqliteDrizzle") {}
|
||||
|
||||
export interface ClientConfig {
|
||||
readonly spanAttributes?: Record<string, unknown>
|
||||
readonly transformResultNames?: (str: string) => string
|
||||
readonly transformQueryNames?: (str: string) => string
|
||||
}
|
||||
|
||||
type Run = (
|
||||
query: string,
|
||||
params?: ReadonlyArray<unknown>,
|
||||
) => Effect.Effect<ReadonlyArray<Record<string, unknown>>, SqlError>
|
||||
|
||||
type RunValues = (
|
||||
query: string,
|
||||
params?: ReadonlyArray<unknown>,
|
||||
) => Effect.Effect<ReadonlyArray<ReadonlyArray<unknown>>, SqlError>
|
||||
|
||||
export const makeConnection = <Extensions extends object>(run: Run, runValues: RunValues, extensions: Extensions) =>
|
||||
identity<Connection & Extensions>({
|
||||
execute(query, params, transformRows) {
|
||||
return transformRows ? Effect.map(run(query, params), transformRows) : run(query, params)
|
||||
},
|
||||
executeRaw(query, params) {
|
||||
return run(query, params)
|
||||
},
|
||||
executeValues(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeValuesUnprepared(query, params) {
|
||||
return runValues(query, params)
|
||||
},
|
||||
executeUnprepared(query, params, transformRows) {
|
||||
return this.execute(query, params, transformRows)
|
||||
},
|
||||
executeStream() {
|
||||
return Stream.die("executeStream not implemented")
|
||||
},
|
||||
...extensions,
|
||||
})
|
||||
|
||||
export const makeClient = <
|
||||
Config extends ClientConfig,
|
||||
SqliteConnection extends Connection,
|
||||
const TypeId extends string,
|
||||
Extensions extends object,
|
||||
>(
|
||||
options: Config,
|
||||
connection: SqliteConnection,
|
||||
typeId: TypeId,
|
||||
extensions: (acquirer: Effect.Effect<SqliteConnection, SqlError, Scope.Scope>) => Extensions,
|
||||
) =>
|
||||
Effect.gen(function* () {
|
||||
const semaphore = yield* Semaphore.make(1)
|
||||
const acquirer = semaphore.withPermits(1)(Effect.succeed(connection))
|
||||
const transactionAcquirer = Effect.uninterruptibleMask((restore) => {
|
||||
const fiber = Fiber.getCurrent()!
|
||||
const scope = Context.getUnsafe(fiber.context, Scope.Scope)
|
||||
return Effect.as(
|
||||
Effect.tap(restore(semaphore.take(1)), () => Scope.addFinalizer(scope, semaphore.release(1))),
|
||||
connection,
|
||||
)
|
||||
})
|
||||
const transformRows = options.transformResultNames
|
||||
? Statement.defaultTransforms(options.transformResultNames).array
|
||||
: undefined
|
||||
|
||||
return Object.assign(
|
||||
yield* SqlClient.make({
|
||||
acquirer,
|
||||
compiler: Statement.makeCompilerSqlite(options.transformQueryNames),
|
||||
transactionAcquirer,
|
||||
spanAttributes: [
|
||||
...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
|
||||
["db.system.name", "sqlite"],
|
||||
],
|
||||
transformRows,
|
||||
}),
|
||||
{
|
||||
[typeId]: typeId,
|
||||
config: options,
|
||||
...extensions(acquirer),
|
||||
},
|
||||
) as SqlClient.SqlClient &
|
||||
Record<TypeId, TypeId> & {
|
||||
readonly config: Config
|
||||
readonly updateValues: never
|
||||
} & Extensions
|
||||
})
|
||||
|
||||
@@ -4,7 +4,7 @@ export { Event, ID, Info } from "@opencode-ai/schema/plugin"
|
||||
import { Plugin } from "@opencode-ai/schema/plugin"
|
||||
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
|
||||
import { App } from "./app"
|
||||
import { Context, Effect, Exit, Layer, Logger, References, Scope, Semaphore } from "effect"
|
||||
import { Context, Effect, Exit, Layer, Scope, Semaphore } from "effect"
|
||||
import { Agent } from "./agent"
|
||||
import { AISDK } from "./aisdk"
|
||||
import { Catalog } from "./catalog"
|
||||
@@ -44,12 +44,7 @@ const layer = Layer.effect(
|
||||
const inherit = yield* State.inherit()
|
||||
const loaded = yield* Effect.suspend(() => plugin.effect(host)).pipe(
|
||||
inherit,
|
||||
Effect.updateContext((context: Context.Context<never>) =>
|
||||
Context.make(Scope.Scope, child).pipe(
|
||||
Context.add(Logger.CurrentLoggers, Context.get(context, Logger.CurrentLoggers)),
|
||||
Context.add(References.MinimumLogLevel, Context.get(context, References.MinimumLogLevel)),
|
||||
),
|
||||
),
|
||||
Effect.updateContext((_context: Context.Context<never>) => Context.make(Scope.Scope, child)),
|
||||
Effect.withSpan("Plugin.load", { attributes: { "plugin.id": plugin.id } }),
|
||||
Effect.andThen(bus.publish(Plugin.Event.Added, { id: Plugin.ID.make(plugin.id) })),
|
||||
Effect.onExit((exit) => (Exit.isFailure(exit) ? Scope.close(child, exit) : Effect.void)),
|
||||
|
||||
@@ -101,6 +101,16 @@ export const Plugin = define({
|
||||
item.permissions.push({ action: "question", resource: "*", effect: "allow" })
|
||||
})
|
||||
|
||||
draft.update(Agent.ID.make("plan"), (item) => {
|
||||
item.name = Agent.Name.make("Plan")
|
||||
item.description = "Plan mode. Disallows all edit tools."
|
||||
item.mode = "primary"
|
||||
item.permissions.push(
|
||||
{ action: "question", resource: "*", effect: "allow" },
|
||||
{ action: "edit", resource: "*", effect: "deny" },
|
||||
)
|
||||
})
|
||||
|
||||
draft.update(Agent.ID.make("general"), (item) => {
|
||||
item.name = Agent.Name.make("General")
|
||||
item.description =
|
||||
|
||||
@@ -17,7 +17,6 @@ import { ConfigProviderPlugin } from "../config/plugin/provider"
|
||||
import { ConfigPolicyPlugin } from "../config/plugin/policy"
|
||||
import { ConfigReferencePlugin } from "../config/plugin/reference"
|
||||
import { ConfigSkillPlugin } from "../config/plugin/skill"
|
||||
import { ConfigPluginSource } from "../config/plugin/source"
|
||||
import { ConfigWebSearchPlugin } from "../config/plugin/websearch"
|
||||
import { Bus } from "../bus"
|
||||
import { Environment } from "../environment"
|
||||
@@ -61,7 +60,6 @@ import { WellKnown } from "../wellknown"
|
||||
import { WriteTool } from "../tool/plugin/write"
|
||||
import { AgentPlugin } from "./agent"
|
||||
import { CommandPlugin } from "./command"
|
||||
import { PlanPlugin } from "./plan"
|
||||
import { ModelsDevPlugin } from "./models-dev"
|
||||
import { ProviderPlugins } from "./provider"
|
||||
import { WebSearchPlugins } from "./websearch"
|
||||
@@ -78,7 +76,6 @@ const services = Effect.fn("PluginInternal.services")(function* () {
|
||||
const command = yield* Command.Service
|
||||
const config = yield* Config.Service
|
||||
const credential = yield* Credential.Service
|
||||
const pluginSources = yield* ConfigPluginSource.Service
|
||||
const bus = yield* Bus.Service
|
||||
const environment = yield* Environment.Service
|
||||
const mutation = yield* FileMutation.Service
|
||||
@@ -115,7 +112,6 @@ const services = Effect.fn("PluginInternal.services")(function* () {
|
||||
Context.make(Command.Service, command),
|
||||
Context.make(Config.Service, config),
|
||||
Context.make(Credential.Service, credential),
|
||||
Context.make(ConfigPluginSource.Service, pluginSources),
|
||||
Context.make(Bus.Service, bus),
|
||||
Context.make(Environment.Service, environment),
|
||||
Context.make(FileMutation.Service, mutation),
|
||||
@@ -159,7 +155,6 @@ export const requirements = LayerNode.group([
|
||||
Command.node,
|
||||
Config.node,
|
||||
Credential.node,
|
||||
ConfigPluginSource.node,
|
||||
Bus.node,
|
||||
Environment.node,
|
||||
FileMutation.node,
|
||||
@@ -197,7 +192,6 @@ export type InternalPlugin = Plugin<Requirements | Scope.Scope>
|
||||
const pre = [
|
||||
WellKnownPlugin.Plugin,
|
||||
AgentPlugin.Plugin,
|
||||
PlanPlugin.Plugin,
|
||||
CommandPlugin.Plugin,
|
||||
SkillPlugin.Plugin,
|
||||
...SystemPromptPlugin.Plugins,
|
||||
|
||||
@@ -1,70 +0,0 @@
|
||||
export * as PlanPlugin from "./plan"
|
||||
|
||||
import { ToolFailure } from "@opencode-ai/ai"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
import { Effect, Stream } from "effect"
|
||||
import { Agent } from "../agent"
|
||||
import { SessionEvent } from "../session/event"
|
||||
|
||||
const plan = Agent.ID.make("plan")
|
||||
|
||||
const enter = `<system-reminder>
|
||||
You are in Plan mode. You are not allowed to edit or create files, and you may not ask a subagent to do that either.
|
||||
|
||||
You are in Plan mode until the user switches agents. Plan mode is not changed by user intent, tone, or imperative language. If the user asks you to change files, do not edit. Tell them they need to switch agents.
|
||||
</system-reminder>`
|
||||
|
||||
const leave = `<system-reminder>
|
||||
You are NO LONGER in Plan mode. The previous Plan restrictions no longer apply. Any Plan mode instructions from earlier in this conversation are no longer active.
|
||||
</system-reminder>`
|
||||
|
||||
export const Plugin = define({
|
||||
id: "opencode.plan",
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.agent.transform((draft) => {
|
||||
draft.update(plan, (item) => {
|
||||
item.name = Agent.Name.make("Plan")
|
||||
item.description = "Read-only agent for exploring the codebase and planning work before implementation."
|
||||
item.mode = "primary"
|
||||
item.permissions.push({ action: "question", resource: "*", effect: "allow" })
|
||||
})
|
||||
})
|
||||
|
||||
yield* ctx.tool.hook("execute.before", (event) => {
|
||||
if (event.agent !== plan) return Effect.void
|
||||
if (event.tool !== "edit" && event.tool !== "write" && event.tool !== "patch") return Effect.void
|
||||
return new ToolFailure({
|
||||
message: `Cannot use ${event.tool} in Plan mode. You are in a read-only mode and must not modify files.`,
|
||||
})
|
||||
})
|
||||
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter(
|
||||
(event): event is SessionEvent.Created | SessionEvent.AgentSelected =>
|
||||
event.type === "session.created" || event.type === "session.agent.selected",
|
||||
),
|
||||
Stream.runForEach((event) => {
|
||||
const text = reminder(event)
|
||||
if (!text) return Effect.void
|
||||
return ctx.session
|
||||
.synthetic({
|
||||
sessionID: event.data.sessionID,
|
||||
text,
|
||||
resume: false,
|
||||
})
|
||||
.pipe(Effect.catch(() => Effect.void))
|
||||
}),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
function reminder(event: SessionEvent.Created | SessionEvent.AgentSelected) {
|
||||
if (event.type === "session.created") {
|
||||
if (event.data.agent !== plan) return
|
||||
return enter
|
||||
}
|
||||
if (event.data.agent === event.data.previous) return
|
||||
if (event.data.agent === plan) return enter
|
||||
if (event.data.previous === plan) return leave
|
||||
}
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const AlibabaPlugin = createProviderPlugin({
|
||||
export const AlibabaPlugin = define({
|
||||
id: "opencode.provider.alibaba",
|
||||
package: "@ai-sdk/alibaba",
|
||||
load: async (options) => {
|
||||
const { createAlibaba } = await import("@ai-sdk/alibaba")
|
||||
return createAlibaba(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/alibaba") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/alibaba"))
|
||||
evt.sdk = mod.createAlibaba(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const CoherePlugin = createProviderPlugin({
|
||||
export const CoherePlugin = define({
|
||||
id: "opencode.provider.cohere",
|
||||
package: "@ai-sdk/cohere",
|
||||
load: async (options) => {
|
||||
const { createCohere } = await import("@ai-sdk/cohere")
|
||||
return createCohere(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/cohere") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/cohere"))
|
||||
evt.sdk = mod.createCohere(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const DeepInfraPlugin = createProviderPlugin({
|
||||
export const DeepInfraPlugin = define({
|
||||
id: "opencode.provider.deepinfra",
|
||||
package: "@ai-sdk/deepinfra",
|
||||
load: async (options) => {
|
||||
const { createDeepInfra } = await import("@ai-sdk/deepinfra")
|
||||
return createDeepInfra(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/deepinfra") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/deepinfra"))
|
||||
evt.sdk = mod.createDeepInfra(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,22 +0,0 @@
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
import type { AISDKHooks } from "@opencode-ai/plugin/effect/aisdk"
|
||||
import { Effect } from "effect"
|
||||
|
||||
export function createProviderPlugin(input: {
|
||||
readonly id: string
|
||||
readonly package: string
|
||||
readonly load: (options: AISDKHooks["sdk"]["options"]) => Promise<unknown>
|
||||
}) {
|
||||
return define({
|
||||
id: input.id,
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== input.package) return
|
||||
evt.sdk = yield* Effect.promise(() => input.load(evt.options))
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
}
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const GatewayPlugin = createProviderPlugin({
|
||||
export const GatewayPlugin = define({
|
||||
id: "opencode.provider.gateway",
|
||||
package: "@ai-sdk/gateway",
|
||||
load: async (options) => {
|
||||
const { createGateway } = await import("@ai-sdk/gateway")
|
||||
return createGateway(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/gateway") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/gateway"))
|
||||
evt.sdk = mod.createGateway(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const GroqPlugin = createProviderPlugin({
|
||||
export const GroqPlugin = define({
|
||||
id: "opencode.provider.groq",
|
||||
package: "@ai-sdk/groq",
|
||||
load: async (options) => {
|
||||
const { createGroq } = await import("@ai-sdk/groq")
|
||||
return createGroq(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/groq") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/groq"))
|
||||
evt.sdk = mod.createGroq(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const MistralPlugin = createProviderPlugin({
|
||||
export const MistralPlugin = define({
|
||||
id: "opencode.provider.mistral",
|
||||
package: "@ai-sdk/mistral",
|
||||
load: async (options) => {
|
||||
const { createMistral } = await import("@ai-sdk/mistral")
|
||||
return createMistral(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/mistral") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/mistral"))
|
||||
evt.sdk = mod.createMistral(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const PerplexityPlugin = createProviderPlugin({
|
||||
export const PerplexityPlugin = define({
|
||||
id: "opencode.provider.perplexity",
|
||||
package: "@ai-sdk/perplexity",
|
||||
load: async (options) => {
|
||||
const { createPerplexity } = await import("@ai-sdk/perplexity")
|
||||
return createPerplexity(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/perplexity") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/perplexity"))
|
||||
evt.sdk = mod.createPerplexity(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const TogetherAIPlugin = createProviderPlugin({
|
||||
export const TogetherAIPlugin = define({
|
||||
id: "opencode.provider.togetherai",
|
||||
package: "@ai-sdk/togetherai",
|
||||
load: async (options) => {
|
||||
const { createTogetherAI } = await import("@ai-sdk/togetherai")
|
||||
return createTogetherAI(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "@ai-sdk/togetherai") return
|
||||
const mod = yield* Effect.promise(() => import("@ai-sdk/togetherai"))
|
||||
evt.sdk = mod.createTogetherAI(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import { createProviderPlugin } from "./factory"
|
||||
import { Effect } from "effect"
|
||||
import { define } from "@opencode-ai/plugin/effect/plugin"
|
||||
|
||||
export const VenicePlugin = createProviderPlugin({
|
||||
export const VenicePlugin = define({
|
||||
id: "opencode.provider.venice",
|
||||
package: "venice-ai-sdk-provider",
|
||||
load: async (options) => {
|
||||
const { createVenice } = await import("venice-ai-sdk-provider")
|
||||
return createVenice(options)
|
||||
},
|
||||
effect: Effect.fn(function* (ctx) {
|
||||
yield* ctx.aisdk.hook(
|
||||
"sdk",
|
||||
Effect.fn(function* (evt) {
|
||||
if (evt.package !== "venice-ai-sdk-provider") return
|
||||
const mod = yield* Effect.promise(() => import("venice-ai-sdk-provider"))
|
||||
evt.sdk = mod.createVenice(evt.options)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
@@ -6,8 +6,12 @@ import { define, type Context } from "@opencode-ai/plugin/effect/plugin"
|
||||
import { Effect } from "effect"
|
||||
import { AbsolutePath } from "../schema"
|
||||
import { Skill } from "../skill"
|
||||
import { ConfigPluginSource } from "../config/plugin/source"
|
||||
import { Config } from "../config"
|
||||
import { Location } from "../location"
|
||||
import { FSUtil } from "@opencode-ai/util/fs-util"
|
||||
import os from "os"
|
||||
import path from "path"
|
||||
import { fileURLToPath } from "url"
|
||||
import opencodeContent from "./skill/opencode.md" with { type: "text" }
|
||||
import reportContent from "./skill/report.md" with { type: "text" }
|
||||
|
||||
@@ -68,10 +72,32 @@ const reportContentWithDiagnostics = Effect.fn("SkillPlugin.reportContentWithDia
|
||||
})
|
||||
|
||||
const configuredPlugins = Effect.fn("SkillPlugin.configuredPlugins")(function* () {
|
||||
const sources = yield* ConfigPluginSource.Service
|
||||
return (yield* sources.operations())
|
||||
.map((operation) => (operation.type === "remove" ? `-${operation.target}` : operation.target))
|
||||
.toSorted()
|
||||
const config = yield* Config.Service
|
||||
const fs = yield* FSUtil.Service
|
||||
const location = yield* Location.Service
|
||||
return yield* Effect.forEach(yield* config.entries(), (entry) => {
|
||||
if (entry.type === "document") {
|
||||
const directory = entry.path ? path.dirname(entry.path) : location.directory
|
||||
return Effect.succeed(
|
||||
(entry.info.plugins ?? []).map((item) => {
|
||||
const ref = typeof item === "string" ? { package: item } : item
|
||||
if (ref.package.startsWith("file://")) return fileURLToPath(ref.package)
|
||||
if (ref.package.startsWith("./") || ref.package.startsWith("../")) return path.resolve(directory, ref.package)
|
||||
return ref.package
|
||||
}),
|
||||
)
|
||||
}
|
||||
if (entry.type !== "directory") return Effect.succeed([])
|
||||
return fs
|
||||
.scan("{plugin,plugins}/*.{ts,js}", {
|
||||
cwd: entry.path,
|
||||
absolute: true,
|
||||
include: "file",
|
||||
dot: true,
|
||||
symlink: true,
|
||||
})
|
||||
.pipe(Effect.orElseSucceed(() => []))
|
||||
}).pipe(Effect.map((items) => items.flat().toSorted()))
|
||||
})
|
||||
|
||||
function terminal() {
|
||||
|
||||
@@ -2,7 +2,7 @@ export * as Project from "./project"
|
||||
|
||||
import { Context, Effect, Layer, Schema } from "effect"
|
||||
import { ChildProcess } from "effect/unstable/process"
|
||||
import { asc, desc } from "drizzle-orm"
|
||||
import { asc, desc, isNotNull, isNull, ne, or } from "drizzle-orm"
|
||||
import path from "path"
|
||||
import { AbsolutePath } from "./schema"
|
||||
import { Database } from "./database/database"
|
||||
@@ -13,7 +13,7 @@ import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
import { Hash } from "@opencode-ai/util/hash"
|
||||
import { ProjectDirectories } from "./project/directories"
|
||||
import { ProjectSchema } from "./project/schema"
|
||||
import { ProjectTable, upsertProject } from "./project/sql"
|
||||
import { ProjectTable } from "./project/sql"
|
||||
|
||||
export const ID = ProjectSchema.ID
|
||||
export type ID = ProjectSchema.ID
|
||||
@@ -98,7 +98,19 @@ const layer = Layer.effect(
|
||||
yield* db
|
||||
.transaction((tx) =>
|
||||
Effect.gen(function* () {
|
||||
yield* upsertProject(tx, project)
|
||||
const vcs = project.vcs?.type
|
||||
yield* tx
|
||||
.insert(ProjectTable)
|
||||
.values({ id: project.id, worktree: project.canonical, vcs, sandboxes: [] })
|
||||
.onConflictDoUpdate({
|
||||
target: ProjectTable.id,
|
||||
set: { worktree: project.canonical, vcs: vcs ?? null },
|
||||
setWhere: or(
|
||||
ne(ProjectTable.worktree, project.canonical),
|
||||
vcs ? or(isNull(ProjectTable.vcs), ne(ProjectTable.vcs, vcs)) : isNotNull(ProjectTable.vcs),
|
||||
),
|
||||
})
|
||||
.run()
|
||||
if (!project.vcs) return
|
||||
yield* projectDirectories.create({ projectID: project.id, directory: project.canonical }, tx)
|
||||
if (project.directory === project.canonical) return
|
||||
|
||||
@@ -1,14 +1,8 @@
|
||||
import type { EffectDrizzleSqlite } from "@opencode-ai/effect-drizzle-sqlite"
|
||||
import { isNotNull, isNull, ne, or } from "drizzle-orm"
|
||||
import { sqliteTable, text, integer, primaryKey } from "drizzle-orm/sqlite-core"
|
||||
import { absoluteArrayColumn, absoluteColumn } from "../database/path"
|
||||
import { Timestamps } from "../database/schema.sql"
|
||||
import type { AbsolutePath } from "../schema"
|
||||
import { ProjectSchema } from "./schema"
|
||||
|
||||
type DatabaseClient = EffectDrizzleSqlite.EffectSQLiteDatabase
|
||||
type Transaction = Parameters<Parameters<DatabaseClient["transaction"]>[0]>[0]
|
||||
|
||||
export const ProjectTable = sqliteTable("project", {
|
||||
id: text().$type<ProjectSchema.ID>().primaryKey(),
|
||||
worktree: absoluteColumn().notNull(),
|
||||
@@ -39,22 +33,3 @@ export const ProjectDirectoryTable = sqliteTable(
|
||||
},
|
||||
(table) => [primaryKey({ columns: [table.project_id, table.directory] })],
|
||||
)
|
||||
|
||||
export function upsertProject(
|
||||
db: DatabaseClient | Transaction,
|
||||
project: { readonly id: ProjectSchema.ID; readonly canonical: AbsolutePath; readonly vcs?: ProjectSchema.Vcs },
|
||||
) {
|
||||
const vcs = project.vcs?.type
|
||||
return db
|
||||
.insert(ProjectTable)
|
||||
.values({ id: project.id, worktree: project.canonical, vcs, sandboxes: [] })
|
||||
.onConflictDoUpdate({
|
||||
target: ProjectTable.id,
|
||||
set: { worktree: project.canonical, vcs: vcs ?? null },
|
||||
setWhere: or(
|
||||
ne(ProjectTable.worktree, project.canonical),
|
||||
vcs ? or(isNull(ProjectTable.vcs), ne(ProjectTable.vcs, vcs)) : isNotNull(ProjectTable.vcs),
|
||||
),
|
||||
})
|
||||
.run()
|
||||
}
|
||||
|
||||
@@ -3,7 +3,7 @@ export * from "./session/schema"
|
||||
|
||||
import { Effect, Layer, Schema, Context, Stream, Scope } from "effect"
|
||||
import { ListAnchor } from "@opencode-ai/schema/session"
|
||||
import { and, asc, desc, eq, gt, isNull, like, lt, or, type SQL } from "drizzle-orm"
|
||||
import { and, asc, desc, eq, gt, isNotNull, isNull, like, lt, ne, or, type SQL } from "drizzle-orm"
|
||||
import { Project } from "./project"
|
||||
import { Workspace } from "./workspace"
|
||||
import { Model } from "./model"
|
||||
@@ -21,7 +21,7 @@ import { Agent } from "./agent"
|
||||
import { Money } from "@opencode-ai/schema/money"
|
||||
import { App } from "./app"
|
||||
import { Slug } from "./util/slug"
|
||||
import { upsertProject } from "./project/sql"
|
||||
import { ProjectTable } from "./project/sql"
|
||||
import path from "path"
|
||||
import { fromRow } from "./session/info"
|
||||
import { SessionRunner } from "./session/runner/index"
|
||||
@@ -309,7 +309,22 @@ const layer = Layer.effect(
|
||||
const shellLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
|
||||
const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Info)
|
||||
const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
||||
const persistProject = (project: Project.Resolved) => upsertProject(db, project).pipe(Effect.orDie)
|
||||
const persistProject = (project: Project.Resolved) => {
|
||||
const vcs = project.vcs?.type
|
||||
return db
|
||||
.insert(ProjectTable)
|
||||
.values({ id: project.id, worktree: project.canonical, vcs, sandboxes: [] })
|
||||
.onConflictDoUpdate({
|
||||
target: ProjectTable.id,
|
||||
set: { worktree: project.canonical, vcs: vcs ?? null },
|
||||
setWhere: or(
|
||||
ne(ProjectTable.worktree, project.canonical),
|
||||
vcs ? or(isNull(ProjectTable.vcs), ne(ProjectTable.vcs, vcs)) : isNotNull(ProjectTable.vcs),
|
||||
),
|
||||
})
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
}
|
||||
const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
||||
decodeMessage({ ...row.data, id: row.id, type: row.type }).pipe(
|
||||
Effect.mapError(
|
||||
|
||||
@@ -56,22 +56,16 @@ export const layer = Layer.effect(
|
||||
),
|
||||
Effect.asVoid,
|
||||
)
|
||||
// Write-ahead claim: starting records the durable intent that a turn is in flight, in the same
|
||||
// transaction as the started event. Terminals release it — except shutdown interruption, which
|
||||
// preserves the claim so the next server start resumes the turn. A claim that survives with no
|
||||
// terminal is the signature of a process that died without teardown (crash, SIGKILL, eviction);
|
||||
// recovery is a property of the database, never of a shutdown hook that may not run.
|
||||
const claimOnCommit = (sessionID: SessionSchema.ID) => ({
|
||||
commit: () => store.claim(sessionID),
|
||||
})
|
||||
const releaseOnCommit = (sessionID: SessionSchema.ID) => ({
|
||||
commit: () => store.release(sessionID),
|
||||
// Starting or finishing on its own clears stale suspension; interruption preserves it because
|
||||
// managed-server teardown suspends active Sessions immediately before interrupting their drains.
|
||||
const clearSuspensionOnCommit = (sessionID: SessionSchema.ID) => ({
|
||||
commit: () => Effect.asVoid(store.consumeSuspended(sessionID)),
|
||||
})
|
||||
const coordinator = yield* SessionRunCoordinator.make<SessionSchema.ID, SessionRunner.RunError, InterruptReason>({
|
||||
started: (sessionID) =>
|
||||
reportLifecycle(
|
||||
sessionID,
|
||||
bus.publish(SessionEvent.Execution.Started, { sessionID }, claimOnCommit(sessionID)),
|
||||
bus.publish(SessionEvent.Execution.Started, { sessionID }, clearSuspensionOnCommit(sessionID)),
|
||||
),
|
||||
drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, force) {
|
||||
const session = yield* store.get(sessionID)
|
||||
@@ -92,17 +86,11 @@ export const layer = Layer.effect(
|
||||
Effect.gen(function* () {
|
||||
const outcome = terminal(exit, reason)
|
||||
if (outcome.type === "succeeded") {
|
||||
yield* bus.publish(SessionEvent.Execution.Succeeded, { sessionID }, releaseOnCommit(sessionID))
|
||||
yield* bus.publish(SessionEvent.Execution.Succeeded, { sessionID }, clearSuspensionOnCommit(sessionID))
|
||||
return
|
||||
}
|
||||
if (outcome.type === "interrupted") {
|
||||
// A user cancel (or a superseding execution) releases the claim: the turn must not
|
||||
// resurrect at the next boot. Shutdown interruption keeps it for restart continuity.
|
||||
yield* bus.publish(
|
||||
SessionEvent.Execution.Interrupted,
|
||||
{ sessionID, reason: outcome.reason },
|
||||
outcome.reason === "shutdown" ? undefined : releaseOnCommit(sessionID),
|
||||
)
|
||||
yield* bus.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: outcome.reason })
|
||||
return
|
||||
}
|
||||
yield* bus.publish(
|
||||
@@ -111,7 +99,7 @@ export const layer = Layer.effect(
|
||||
sessionID,
|
||||
error: outcome.error,
|
||||
},
|
||||
releaseOnCommit(sessionID),
|
||||
clearSuspensionOnCommit(sessionID),
|
||||
)
|
||||
}),
|
||||
),
|
||||
|
||||
@@ -5,111 +5,61 @@ import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
import { Bus } from "../../bus"
|
||||
import { SessionEvent } from "../event"
|
||||
import { SessionExecution } from "../execution"
|
||||
import { SessionSchema } from "../schema"
|
||||
import { SessionStore } from "../store"
|
||||
|
||||
const CONTINUE_AFTER_SERVER_RESTART =
|
||||
"The server restarted while you were working. Continue from where you left off without repeating completed work."
|
||||
|
||||
const RESUME_EXHAUSTED = {
|
||||
type: "aborted",
|
||||
message: "Execution was interrupted repeatedly and will not be resumed automatically.",
|
||||
} as const
|
||||
|
||||
export interface Options {
|
||||
/**
|
||||
* Times a single turn may be resumed before it is terminalized instead.
|
||||
* The counter is durable and only a terminal event resets it, so a turn
|
||||
* that keeps dying cannot crash-loop across restarts. Turns that complete
|
||||
* never accumulate: the budget is per-turn, not per-session.
|
||||
*/
|
||||
readonly maxAttempts?: number
|
||||
}
|
||||
|
||||
const DEFAULT_MAX_ATTEMPTS = 10
|
||||
|
||||
export interface Interface {
|
||||
/**
|
||||
* Resumes Sessions whose execution claim was never released — turns orphaned
|
||||
* by a process that died without teardown, or interrupted by a graceful
|
||||
* shutdown (which preserves the claim on purpose). The claim is never
|
||||
* cleared here: only a terminal event releases it, so a death anywhere in
|
||||
* the resume path leaves the same orphaned claim for the next boot.
|
||||
* Marks every execution active in this process for resumption by the next server start.
|
||||
* Call once new work has stopped arriving and before teardown interrupts the drains.
|
||||
*/
|
||||
readonly suspendActiveSessions: Effect.Effect<void>
|
||||
/** Resumes suspended Sessions. Each suspension is consumed atomically, so a Session resumes at most once. */
|
||||
readonly resumeSuspendedSessions: Effect.Effect<void>
|
||||
}
|
||||
|
||||
/**
|
||||
* Recovery for orphaned executions. Claims are written at turn start by
|
||||
* SessionExecution, so this sweep needs no cooperation from the previous
|
||||
* process: crash, SIGKILL, isolate eviction, and graceful restart all leave
|
||||
* the same durable signature.
|
||||
*
|
||||
* The sweep assumes every orphaned claim's owner is dead. The managed-server
|
||||
* protocol guarantees this: a successor is only spawned after the previous
|
||||
* process is confirmed dead (client service `kill`/`evict` poll the PID), the
|
||||
* registration lock admits one managed server at a time, and unregistered
|
||||
* servers sharing the database never sweep. The service is inert until called
|
||||
* — the managed server invokes it at boot; embedders may call it from their
|
||||
* own start-up.
|
||||
* Restart continuity actions for the managed server. The service is inert until called: only the
|
||||
* managed server invokes it, so default, embedded, and stdio servers never suspend or auto-resume.
|
||||
*/
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/SessionRestart") {}
|
||||
|
||||
export const layer = (options?: Options) =>
|
||||
Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const store = yield* SessionStore.Service
|
||||
const execution = yield* SessionExecution.Service
|
||||
const bus = yield* Bus.Service
|
||||
const scope = yield* Effect.scope
|
||||
const maxAttempts = options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS
|
||||
|
||||
const resumeOne = Effect.fnUntraced(function* (sessionID: SessionSchema.ID) {
|
||||
// Durable before the resume runs, so a crash inside the resumed turn is
|
||||
// counted by the next sweep and the budget cannot be dodged.
|
||||
const attempts = yield* store.countResume(sessionID)
|
||||
if (attempts === undefined) return // the Session was deleted since listing
|
||||
if (attempts > maxAttempts) {
|
||||
// Terminalize instead: the release hook clears the claim and resets the
|
||||
// counter atomically with the terminal event.
|
||||
yield* bus.publish(
|
||||
SessionEvent.Execution.Failed,
|
||||
{ sessionID, error: RESUME_EXHAUSTED },
|
||||
{ commit: () => store.release(sessionID) },
|
||||
)
|
||||
return
|
||||
}
|
||||
yield* bus.publish(SessionEvent.Synthetic, {
|
||||
sessionID,
|
||||
text: CONTINUE_AFTER_SERVER_RESTART,
|
||||
description: "Continuing after restart",
|
||||
})
|
||||
// Forked into the service scope so boot never waits on resumed turns;
|
||||
// resuming an already-live Session joins its execution. Drain failures
|
||||
// are logged and durably recorded by the execution layer.
|
||||
yield* execution.resume(sessionID).pipe(Effect.ignore, Effect.forkIn(scope))
|
||||
})
|
||||
|
||||
return Service.of({
|
||||
resumeSuspendedSessions: Effect.gen(function* () {
|
||||
// Child claims never drive recovery (children are not resumed), so a
|
||||
// dead child's claim is noise no terminal will ever release. Clearing
|
||||
// is safe even against a live child: claims are recovery markers, not
|
||||
// locks, and children are excluded from that recovery.
|
||||
yield* store.releaseChildClaims
|
||||
const active = yield* execution.active
|
||||
// Sessions already draining in this process keep their claim; resuming
|
||||
// them would only inject a stray continuation into a live turn.
|
||||
const orphaned = (yield* store.listSuspended()).filter((sessionID) => !active.has(sessionID))
|
||||
yield* Effect.forEach(orphaned, resumeOne, { concurrency: "unbounded", discard: true })
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const store = yield* SessionStore.Service
|
||||
const execution = yield* SessionExecution.Service
|
||||
const bus = yield* Bus.Service
|
||||
return Service.of({
|
||||
suspendActiveSessions: Effect.gen(function* () {
|
||||
yield* store.suspend(yield* execution.active)
|
||||
}),
|
||||
resumeSuspendedSessions: Effect.gen(function* () {
|
||||
const sessions = yield* store.listSuspended()
|
||||
yield* Effect.forEach(
|
||||
sessions,
|
||||
(sessionID) =>
|
||||
Effect.gen(function* () {
|
||||
if (!(yield* store.consumeSuspended(sessionID))) return
|
||||
yield* bus.publish(SessionEvent.Synthetic, {
|
||||
sessionID,
|
||||
text: CONTINUE_AFTER_SERVER_RESTART,
|
||||
description: "Continuing after restart",
|
||||
})
|
||||
// Drain failures are already logged and durably recorded by the execution layer.
|
||||
yield* Effect.ignore(execution.resume(sessionID))
|
||||
}),
|
||||
{ concurrency: "unbounded", discard: true },
|
||||
)
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const node = makeGlobalNode({
|
||||
service: Service,
|
||||
layer: layer(),
|
||||
layer,
|
||||
deps: [SessionStore.node, SessionExecution.node, Bus.node],
|
||||
})
|
||||
|
||||
@@ -58,9 +58,7 @@ export const SessionTable = sqliteTable(
|
||||
...Timestamps,
|
||||
time_compacting: integer(),
|
||||
time_archived: integer(),
|
||||
/** The execution claim timestamp (historical column name; see SessionStore.claim). */
|
||||
time_suspended: integer(),
|
||||
resume_attempts: integer().notNull().default(0),
|
||||
},
|
||||
(table) => [
|
||||
index("session_v2_project_idx").on(table.project_id),
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
export * as SessionStore from "./store"
|
||||
|
||||
import { and, eq, isNotNull, isNull, sql } from "drizzle-orm"
|
||||
import { and, eq, inArray, isNotNull, isNull } from "drizzle-orm"
|
||||
import { Context, Effect, Layer, Schema } from "effect"
|
||||
import { Database } from "../database/database"
|
||||
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
@@ -17,32 +17,10 @@ export interface Interface {
|
||||
readonly message: (
|
||||
messageID: SessionMessage.ID,
|
||||
) => Effect.Effect<{ readonly sessionID: Session.ID; readonly message: SessionMessage.Info } | undefined>
|
||||
/**
|
||||
* Top-level Sessions holding an execution claim. Child (subagent) Sessions
|
||||
* are excluded: a resumed parent re-runs its tool call and spawns fresh
|
||||
* children, so resuming orphaned children would duplicate their work.
|
||||
*/
|
||||
readonly listSuspended: () => Effect.Effect<ReadonlyArray<Session.ID>>
|
||||
/**
|
||||
* Records the execution claim: the durable write-ahead intent that a turn is
|
||||
* (or was) in flight. Set when execution starts; a claim that survives to the
|
||||
* next boot marks a turn that never completed — its process crashed or shut
|
||||
* down mid-turn.
|
||||
*/
|
||||
readonly claim: (sessionID: Session.ID) => Effect.Effect<void>
|
||||
/** Releases the claim and resets resume accounting. Terminal events call this on commit. */
|
||||
readonly release: (sessionID: Session.ID) => Effect.Effect<void>
|
||||
/**
|
||||
* Clears orphaned child (subagent) claims. Children are never resumed
|
||||
* independently, so a dead child's claim is noise no terminal will ever
|
||||
* release.
|
||||
*/
|
||||
readonly releaseChildClaims: Effect.Effect<void>
|
||||
/**
|
||||
* Durably counts one more resume of an orphaned claim, returning the new
|
||||
* total — or undefined when the Session no longer exists.
|
||||
*/
|
||||
readonly countResume: (sessionID: Session.ID) => Effect.Effect<number | undefined>
|
||||
/** Clears suspension, reporting whether this caller consumed it. At most one concurrent caller receives true. */
|
||||
readonly consumeSuspended: (sessionID: Session.ID) => Effect.Effect<boolean>
|
||||
readonly suspend: (sessionIDs: Iterable<Session.ID>) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/SessionStore") {}
|
||||
@@ -79,52 +57,35 @@ const layer = Layer.effect(
|
||||
return yield* db
|
||||
.select({ sessionID: SessionTable.id })
|
||||
.from(SessionTable)
|
||||
.where(and(isNotNull(SessionTable.time_suspended), isNull(SessionTable.parent_id)))
|
||||
.where(isNotNull(SessionTable.time_suspended))
|
||||
.all()
|
||||
.pipe(
|
||||
Effect.orDie,
|
||||
Effect.map((rows) => rows.map((row) => row.sessionID)),
|
||||
)
|
||||
}),
|
||||
claim: Effect.fn("SessionStore.claim")(function* (sessionID) {
|
||||
// The null guard makes re-claiming a still-claimed Session a zero-row
|
||||
// no-op (a resumed turn re-claims through the same started hook).
|
||||
// Claim bookkeeping never counts as user activity: time_updated is
|
||||
// pinned so session ordering only moves on real changes.
|
||||
consumeSuspended: Effect.fn("SessionStore.consumeSuspended")(function* (sessionID) {
|
||||
return (
|
||||
(yield* db
|
||||
.update(SessionTable)
|
||||
.set({ time_suspended: null })
|
||||
.where(and(eq(SessionTable.id, sessionID), isNotNull(SessionTable.time_suspended)))
|
||||
.returning({ sessionID: SessionTable.id })
|
||||
.get()
|
||||
.pipe(Effect.orDie)) !== undefined
|
||||
)
|
||||
}),
|
||||
suspend: Effect.fn("SessionStore.suspend")(function* (sessionIDs) {
|
||||
const ids = Array.from(sessionIDs)
|
||||
if (ids.length === 0) return
|
||||
// The null guard preserves the original suspension time if a Session is somehow suspended twice.
|
||||
yield* db
|
||||
.update(SessionTable)
|
||||
.set({ time_suspended: Date.now(), time_updated: sql`${SessionTable.time_updated}` })
|
||||
.where(and(eq(SessionTable.id, sessionID), isNull(SessionTable.time_suspended)))
|
||||
.set({ time_suspended: Date.now() })
|
||||
.where(and(inArray(SessionTable.id, ids), isNull(SessionTable.time_suspended)))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
}),
|
||||
release: Effect.fn("SessionStore.release")(function* (sessionID) {
|
||||
yield* db
|
||||
.update(SessionTable)
|
||||
.set({ time_suspended: null, resume_attempts: 0, time_updated: sql`${SessionTable.time_updated}` })
|
||||
.where(eq(SessionTable.id, sessionID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
}),
|
||||
releaseChildClaims: db
|
||||
.update(SessionTable)
|
||||
.set({ time_suspended: null, resume_attempts: 0, time_updated: sql`${SessionTable.time_updated}` })
|
||||
.where(and(isNotNull(SessionTable.time_suspended), isNotNull(SessionTable.parent_id)))
|
||||
.run()
|
||||
.pipe(Effect.orDie, Effect.asVoid, Effect.withSpan("SessionStore.releaseChildClaims")),
|
||||
countResume: Effect.fn("SessionStore.countResume")(function* (sessionID) {
|
||||
const row = yield* db
|
||||
.update(SessionTable)
|
||||
.set({
|
||||
resume_attempts: sql`${SessionTable.resume_attempts} + 1`,
|
||||
time_updated: sql`${SessionTable.time_updated}`,
|
||||
})
|
||||
.where(eq(SessionTable.id, sessionID))
|
||||
.returning({ attempts: SessionTable.resume_attempts })
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
return row?.attempts
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -3,7 +3,7 @@ export * as SessionTransfer from "./transfer"
|
||||
import { SessionTransfer } from "@opencode-ai/schema/session-transfer"
|
||||
import { Tool } from "@opencode-ai/schema/tool"
|
||||
import { Skill } from "@opencode-ai/schema/skill"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { eq, isNotNull, isNull, ne, or } from "drizzle-orm"
|
||||
import { Context, DateTime, Effect, Layer, Schema } from "effect"
|
||||
import path from "path"
|
||||
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
@@ -12,7 +12,7 @@ import { Bus } from "../bus"
|
||||
import { Database } from "../database/database"
|
||||
import { Location } from "../location"
|
||||
import { Project } from "../project"
|
||||
import { upsertProject } from "../project/sql"
|
||||
import { ProjectTable } from "../project/sql"
|
||||
import { AbsolutePath, RelativePath } from "../schema"
|
||||
import { Session } from "../session"
|
||||
import { Slug } from "../util/slug"
|
||||
@@ -49,7 +49,22 @@ const layer = Layer.effect(
|
||||
const sessions = yield* Session.Service
|
||||
const encodeMessage = Schema.encodeSync(SessionMessage.Info)
|
||||
|
||||
const persistProject = (project: Project.Resolved) => upsertProject(db, project).pipe(Effect.orDie)
|
||||
const persistProject = (project: Project.Resolved) => {
|
||||
const vcs = project.vcs?.type
|
||||
return db
|
||||
.insert(ProjectTable)
|
||||
.values({ id: project.id, worktree: project.canonical, vcs, sandboxes: [] })
|
||||
.onConflictDoUpdate({
|
||||
target: ProjectTable.id,
|
||||
set: { worktree: project.canonical, vcs: vcs ?? null },
|
||||
setWhere: or(
|
||||
ne(ProjectTable.worktree, project.canonical),
|
||||
vcs ? or(isNull(ProjectTable.vcs), ne(ProjectTable.vcs, vcs)) : isNotNull(ProjectTable.vcs),
|
||||
),
|
||||
})
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
}
|
||||
|
||||
return Service.of({
|
||||
export: Effect.fn("SessionTransfer.export")(function* (input) {
|
||||
|
||||
@@ -165,6 +165,7 @@ describe("Agent", () => {
|
||||
"compaction",
|
||||
"explore",
|
||||
"general",
|
||||
"plan",
|
||||
"summary",
|
||||
"title",
|
||||
])
|
||||
|
||||
@@ -168,75 +168,6 @@ describe("PluginSupervisor config", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.live("loads auto-discovered plugin package entrypoints in order", () =>
|
||||
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("package-exports")
|
||||
expect(ids).toContain("package-module")
|
||||
expect(ids).toContain("package-main")
|
||||
expect(ids).toContain("package-index")
|
||||
}),
|
||||
false,
|
||||
async (directory) => {
|
||||
await Promise.all([
|
||||
writeDiscoveredPackage(directory, "exports", { exports: "./entry.ts" }, { "entry.ts": "package-exports" }),
|
||||
writeDiscoveredPackage(
|
||||
directory,
|
||||
"module",
|
||||
{ exports: "./missing.js", module: "./entry.js" },
|
||||
{ "entry.js": "package-module" },
|
||||
),
|
||||
writeDiscoveredPackage(
|
||||
directory,
|
||||
"main",
|
||||
{ exports: { import: "./missing.js" }, module: "./missing.js", main: "./entry.js" },
|
||||
{ "entry.js": "package-main" },
|
||||
),
|
||||
writeDiscoveredPackage(directory, "index", undefined, { "index.js": "package-index" }),
|
||||
])
|
||||
},
|
||||
),
|
||||
)
|
||||
|
||||
it.live("keeps auto-discovered package entrypoints inside the package directory", () =>
|
||||
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("contained-fallback")
|
||||
expect(ids).toContain("symlink-fallback")
|
||||
expect(ids).not.toContain("escaped-entrypoint")
|
||||
}),
|
||||
false,
|
||||
async (directory) => {
|
||||
await fs.mkdir(path.join(directory, ".opencode"), { recursive: true })
|
||||
await fs.writeFile(path.join(directory, ".opencode", "escape.js"), discoveredPlugin("escaped-entrypoint"))
|
||||
await writeDiscoveredPackage(
|
||||
directory,
|
||||
"contained",
|
||||
{ exports: "../../escape.js" },
|
||||
{ "index.js": "contained-fallback" },
|
||||
)
|
||||
await writeDiscoveredPackage(
|
||||
directory,
|
||||
"symlink",
|
||||
{ exports: "./entry.js" },
|
||||
{ "index.js": "symlink-fallback" },
|
||||
)
|
||||
await fs.symlink(
|
||||
path.join(directory, ".opencode", "escape.js"),
|
||||
path.join(directory, ".opencode", "plugins", "symlink", "entry.js"),
|
||||
)
|
||||
},
|
||||
),
|
||||
)
|
||||
|
||||
staticIt.live("uses only internal and SDK plugins when the static source is wired", () =>
|
||||
Effect.gen(function* () {
|
||||
const sdk = yield* SdkPlugins.Service
|
||||
@@ -458,21 +389,3 @@ export default Plugin.define({
|
||||
})
|
||||
`
|
||||
}
|
||||
|
||||
function discoveredPlugin(id: string) {
|
||||
return `export default { id: ${JSON.stringify(id)}, setup() {} }`
|
||||
}
|
||||
|
||||
async function writeDiscoveredPackage(
|
||||
directory: string,
|
||||
name: string,
|
||||
manifest: Record<string, unknown> | undefined,
|
||||
files: Record<string, string>,
|
||||
) {
|
||||
const plugin = path.join(directory, ".opencode", "plugins", name)
|
||||
await fs.mkdir(plugin, { recursive: true })
|
||||
await Promise.all([
|
||||
...(manifest ? [fs.writeFile(path.join(plugin, "package.json"), JSON.stringify(manifest))] : []),
|
||||
...Object.entries(files).map(([file, id]) => fs.writeFile(path.join(plugin, file), discoveredPlugin(id))),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -24,7 +24,6 @@ const it = testEffect(Layer.empty)
|
||||
|
||||
const instructionLayer = (input: {
|
||||
config?: string
|
||||
home?: string
|
||||
locationServiceLayer: Layer.Layer<Location.Service>
|
||||
filesystemLayer?: Layer.Layer<FSUtil.Service>
|
||||
project?: boolean
|
||||
@@ -35,15 +34,7 @@ const instructionLayer = (input: {
|
||||
LayerNode.group([InstructionDiscovery.node, Bus.node, FSUtil.node, Global.node, Location.node, Watcher.node]),
|
||||
[
|
||||
[InstructionDiscovery.node, InstructionDiscovery.configured({ project: input.project })],
|
||||
[
|
||||
Global.node,
|
||||
input.config || input.home
|
||||
? Global.layerWith({
|
||||
...(input.config ? { config: input.config } : {}),
|
||||
...(input.home ? { home: input.home } : {}),
|
||||
})
|
||||
: tempGlobalLayer,
|
||||
],
|
||||
[Global.node, input.config ? Global.layerWith({ config: input.config }) : tempGlobalLayer],
|
||||
[Location.node, input.locationServiceLayer],
|
||||
[Watcher.node, watcher],
|
||||
...(input.filesystemLayer ? [[FSUtil.node, input.filesystemLayer] as const] : []),
|
||||
@@ -121,13 +112,10 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) => {
|
||||
const global = path.join(tmp.path, "global")
|
||||
const home = path.join(tmp.path, "home")
|
||||
const shared = path.join(home, "code")
|
||||
const project = path.join(shared, "repo")
|
||||
const project = path.join(tmp.path, "project")
|
||||
const directory = path.join(project, "packages", "core")
|
||||
const outside = path.join(tmp.path, "AGENTS.md")
|
||||
const globalFile = path.join(global, "AGENTS.md")
|
||||
const sharedFile = path.join(shared, "AGENTS.md")
|
||||
const projectFile = path.join(project, "AGENTS.md")
|
||||
const packageFile = path.join(directory, "AGENTS.md")
|
||||
return Effect.gen(function* () {
|
||||
@@ -136,7 +124,6 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
await fs.mkdir(directory, { recursive: true })
|
||||
await fs.writeFile(outside, "outside")
|
||||
await fs.writeFile(globalFile, "global")
|
||||
await fs.writeFile(sharedFile, "shared")
|
||||
await fs.writeFile(projectFile, "project")
|
||||
await fs.writeFile(packageFile, "package")
|
||||
})
|
||||
@@ -148,20 +135,13 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
{ path: packageFile, type: "file" },
|
||||
{ path: path.join(project, "packages", "AGENTS.md"), type: "file" },
|
||||
{ path: projectFile, type: "file" },
|
||||
{ path: sharedFile, type: "file" },
|
||||
{ path: path.join(home, "AGENTS.md"), type: "file" },
|
||||
])
|
||||
expect(yield* watcher.subscriptions()).not.toContainEqual({
|
||||
path: path.join(tmp.path, "AGENTS.md"),
|
||||
type: "file",
|
||||
})
|
||||
const initialized = yield* readInitial(yield* discovery.load())
|
||||
expect(initialized.text).toBe(
|
||||
[
|
||||
`Instructions from: ${globalFile}\nglobal`,
|
||||
`Instructions from: ${packageFile}\npackage`,
|
||||
`Instructions from: ${projectFile}\nproject`,
|
||||
`Instructions from: ${sharedFile}\nshared`,
|
||||
].join("\n\n"),
|
||||
)
|
||||
expect(initialized.text).not.toContain("outside")
|
||||
@@ -179,7 +159,6 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
"These instructions replace all previously loaded ambient instructions.",
|
||||
`Instructions from: ${globalFile}\nglobal`,
|
||||
`Instructions from: ${projectFile}\nproject`,
|
||||
`Instructions from: ${sharedFile}\nshared`,
|
||||
].join("\n\n"),
|
||||
)
|
||||
|
||||
@@ -187,8 +166,6 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
yield* emitAndWait({ type: "delete", path: globalFile })
|
||||
yield* Effect.promise(() => fs.rm(projectFile))
|
||||
yield* emitAndWait({ type: "delete", path: projectFile })
|
||||
yield* Effect.promise(() => fs.rm(sharedFile))
|
||||
yield* emitAndWait({ type: "delete", path: sharedFile })
|
||||
expect((yield* readUpdate(yield* discovery.load(), initialized)).text).toBe(
|
||||
"Previously loaded instructions no longer apply.",
|
||||
)
|
||||
@@ -196,7 +173,6 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
Effect.provide(
|
||||
instructionLayer({
|
||||
config: global,
|
||||
home,
|
||||
locationServiceLayer: Layer.succeed(
|
||||
Location.Service,
|
||||
Location.Service.of(
|
||||
@@ -239,17 +215,15 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.live("discovers a newly created instruction file above the project root", () =>
|
||||
it.live("discovers a newly created instruction file in an intermediate directory", () =>
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) => {
|
||||
const home = path.join(tmp.path, "home")
|
||||
const shared = path.join(home, "code")
|
||||
const project = path.join(shared, "repo")
|
||||
const intermediate = path.join(shared, "AGENTS.md")
|
||||
const directory = path.join(project, "core")
|
||||
const project = path.join(tmp.path, "project")
|
||||
const intermediate = path.join(project, "packages", "AGENTS.md")
|
||||
const directory = path.join(project, "packages", "core")
|
||||
const projectFile = path.join(project, "AGENTS.md")
|
||||
return Effect.gen(function* () {
|
||||
yield* Effect.promise(() => fs.mkdir(directory, { recursive: true }))
|
||||
@@ -261,7 +235,7 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
yield* emitAndWait({ type: "create", path: intermediate })
|
||||
|
||||
expect((yield* readInitial(yield* discovery.load())).text).toBe(
|
||||
[`Instructions from: ${projectFile}\nproject`, `Instructions from: ${intermediate}\nintermediate`].join(
|
||||
[`Instructions from: ${intermediate}\nintermediate`, `Instructions from: ${projectFile}\nproject`].join(
|
||||
"\n\n",
|
||||
),
|
||||
)
|
||||
@@ -269,48 +243,6 @@ describe("ConfigInstructionPlugin.Plugin", () => {
|
||||
Effect.provide(
|
||||
instructionLayer({
|
||||
config: path.join(tmp.path, "global"),
|
||||
home,
|
||||
locationServiceLayer: Layer.succeed(
|
||||
Location.Service,
|
||||
Location.Service.of(
|
||||
location(
|
||||
{ directory: AbsolutePath.make(directory) },
|
||||
{ projectDirectory: AbsolutePath.make(project) },
|
||||
),
|
||||
),
|
||||
),
|
||||
}),
|
||||
),
|
||||
)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("stops instruction candidates at the project root outside home", () =>
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) => {
|
||||
const global = path.join(tmp.path, "global")
|
||||
const home = path.join(tmp.path, "home")
|
||||
const project = path.join(tmp.path, "scratch", "repo")
|
||||
const directory = path.join(project, "packages", "core")
|
||||
return Effect.gen(function* () {
|
||||
yield* Effect.promise(() => fs.mkdir(directory, { recursive: true }))
|
||||
yield* start()
|
||||
const watcher = yield* Watcher.Test
|
||||
expect(yield* watcher.subscriptions()).toEqual([
|
||||
{ path: path.join(global, "AGENTS.md"), type: "file" },
|
||||
{ path: path.join(directory, "AGENTS.md"), type: "file" },
|
||||
{ path: path.join(project, "packages", "AGENTS.md"), type: "file" },
|
||||
{ path: path.join(project, "AGENTS.md"), type: "file" },
|
||||
])
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
instructionLayer({
|
||||
config: global,
|
||||
home,
|
||||
locationServiceLayer: Layer.succeed(
|
||||
Location.Service,
|
||||
Location.Service.of(
|
||||
|
||||
@@ -1,18 +1,18 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { NodeFileSystem } from "@effect/platform-node"
|
||||
import { Config } from "@opencode-ai/core/config"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { ConfigPluginSource } from "@opencode-ai/core/config/plugin/source"
|
||||
import { Effect, Layer, Stream } from "effect"
|
||||
import { FSUtil } from "@opencode-ai/util/fs-util"
|
||||
import { Location } from "@opencode-ai/core/location"
|
||||
import { Effect, Stream } from "effect"
|
||||
import { SkillPlugin } from "@opencode-ai/core/plugin/skill"
|
||||
import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { Skill } from "@opencode-ai/core/skill"
|
||||
import { location } from "../fixture/location"
|
||||
import { testEffect } from "../lib/effect"
|
||||
import { host } from "./host"
|
||||
|
||||
const it = testEffect(AppNodeBuilder.build(Skill.node))
|
||||
const sources = (operations: readonly ConfigPluginSource.Operation[] = []) =>
|
||||
Layer.succeed(
|
||||
ConfigPluginSource.Service,
|
||||
ConfigPluginSource.Service.of({ operations: () => Effect.succeed(operations), changes: () => Stream.never }),
|
||||
)
|
||||
|
||||
describe("SkillPlugin.Plugin", () => {
|
||||
it.effect("registers built-in skills", () =>
|
||||
@@ -27,7 +27,15 @@ describe("SkillPlugin.Plugin", () => {
|
||||
reload: skill.reload,
|
||||
},
|
||||
}),
|
||||
).pipe(Effect.provide(sources()))
|
||||
).pipe(
|
||||
Effect.provide(Config.testLayer()),
|
||||
Effect.provideService(
|
||||
Location.Service,
|
||||
Location.Service.of(location({ directory: AbsolutePath.make(import.meta.dir) })),
|
||||
),
|
||||
Effect.provide(AppNodeBuilder.build(FSUtil.node)),
|
||||
Effect.provide(NodeFileSystem.layer),
|
||||
)
|
||||
const skills = yield* skill.list()
|
||||
const report = skills.find((item) => item.id === "report")
|
||||
|
||||
@@ -50,30 +58,4 @@ describe("SkillPlugin.Plugin", () => {
|
||||
expect(report?.content).toContain("- install/channel: beta")
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("reports canonical configured plugin sources with existing labels and ordering", () =>
|
||||
Effect.gen(function* () {
|
||||
const skill = yield* Skill.Service
|
||||
yield* SkillPlugin.Plugin.effect(
|
||||
host({
|
||||
skill: {
|
||||
list: () => Effect.die("unused skill.list"),
|
||||
transform: skill.transform,
|
||||
reload: skill.reload,
|
||||
},
|
||||
}),
|
||||
)
|
||||
const report = (yield* skill.list()).find((item) => item.id === "report")
|
||||
expect(report?.content).toContain("- Active plugins: -disabled, local.ts, package-plugin, package-plugin")
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
sources([
|
||||
{ type: "add", target: "package-plugin", options: {} },
|
||||
{ type: "remove", target: "disabled" },
|
||||
{ type: "add", target: "local.ts", options: {}, mtime: 1 },
|
||||
{ type: "add", target: "package-plugin", options: { enabled: true } },
|
||||
]),
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
@@ -18,7 +18,6 @@ import { SessionRunner } from "@opencode-ai/core/session/runner"
|
||||
import { SessionTable } from "@opencode-ai/core/session/sql"
|
||||
import { SessionStore } from "@opencode-ai/core/session/store"
|
||||
import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionStore.node])))
|
||||
@@ -50,90 +49,58 @@ describe("SessionExecution lifecycle", () => {
|
||||
})
|
||||
})
|
||||
|
||||
it.effect("the sweep only lists claimed top-level Sessions", () =>
|
||||
it.effect("atomically consumes each suspension at most once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const store = yield* SessionStore.Service
|
||||
const parent = Session.ID.make("ses_recover_parent")
|
||||
const child = Session.ID.make("ses_recover_child")
|
||||
const idle = Session.ID.make("ses_recover_idle")
|
||||
yield* seedSessions(database, [parent], { time_suspended: Date.now() })
|
||||
yield* seedSessions(database, [idle])
|
||||
// An orphaned child is never resumed: the resumed parent re-runs its
|
||||
// tool call and spawns a fresh child instead.
|
||||
yield* seedSessions(database, [child], { time_suspended: Date.now(), parent_id: parent })
|
||||
const first = Session.ID.make("ses_recover_first")
|
||||
const second = Session.ID.make("ses_recover_second")
|
||||
yield* seedSessions(database, [first, second], { time_suspended: Date.now() })
|
||||
|
||||
expect(yield* store.listSuspended()).toEqual([parent])
|
||||
|
||||
// The sweep clears orphaned child claims outright; parents keep theirs.
|
||||
yield* store.releaseChildClaims
|
||||
expect(yield* claims(database)).toEqual({ [parent]: true, [child]: false, [idle]: false })
|
||||
expect(yield* store.consumeSuspended(first)).toBe(true)
|
||||
expect(yield* store.consumeSuspended(first)).toBe(false)
|
||||
expect(yield* store.consumeSuspended(second)).toBe(true)
|
||||
expect(yield* suspensions(database)).toEqual({ [first]: false, [second]: false })
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("claims at execution start, releases on completion, and preserves through teardown", () =>
|
||||
it.effect("suspension survives teardown interruption and clears when a drain finishes on its own", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const interrupted = Session.ID.make("ses_claim_interrupted")
|
||||
const completed = Session.ID.make("ses_claim_completed")
|
||||
const interrupted = Session.ID.make("ses_suspend_interrupted")
|
||||
const completed = Session.ID.make("ses_suspend_completed")
|
||||
yield* seedSessions(database, [interrupted, completed])
|
||||
|
||||
// Each drain signals once it runs; the claim commits before the drain starts.
|
||||
const interruptedRunning = yield* Deferred.make<void>()
|
||||
const completedRunning = yield* Deferred.make<void>()
|
||||
const draining = yield* Deferred.make<void>()
|
||||
const release = yield* Deferred.make<void>()
|
||||
const scope = yield* Scope.make()
|
||||
const context = yield* buildExecution(scope, ({ sessionID }) =>
|
||||
sessionID === completed
|
||||
? Deferred.succeed(completedRunning, undefined).pipe(Effect.andThen(Deferred.await(release)))
|
||||
: Deferred.succeed(interruptedRunning, undefined).pipe(Effect.andThen(Effect.never)),
|
||||
? Deferred.await(release)
|
||||
: Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never)),
|
||||
)
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
yield* execution.resume(interrupted).pipe(Effect.forkScoped)
|
||||
const completing = yield* execution.resume(completed).pipe(Effect.forkIn(scope))
|
||||
yield* Deferred.await(interruptedRunning)
|
||||
yield* Deferred.await(completedRunning)
|
||||
yield* Deferred.await(draining)
|
||||
|
||||
// The write-ahead claim exists WHILE the turns run — no shutdown hook involved.
|
||||
expect(yield* claims(database)).toEqual({ [interrupted]: true, [completed]: true })
|
||||
yield* restart.suspendActiveSessions
|
||||
expect(yield* suspensions(database)).toEqual({ [interrupted]: true, [completed]: true })
|
||||
|
||||
// A drain that finishes on its own releases its claim.
|
||||
// A drain that finishes on its own after suspension clears its stale suspension.
|
||||
yield* Deferred.succeed(release, undefined)
|
||||
yield* Fiber.join(completing)
|
||||
yield* execution.awaitIdle(completed)
|
||||
expect((yield* claims(database))[completed]).toBe(false)
|
||||
expect((yield* suspensions(database))[completed]).toBe(false)
|
||||
|
||||
// Teardown interruption (graceful twin of an unclean death) preserves the claim
|
||||
// for the next server start.
|
||||
// Teardown interruption preserves suspension for the next server start.
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
expect((yield* claims(database))[interrupted]).toBe(true)
|
||||
expect((yield* suspensions(database))[interrupted]).toBe(true)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("a user interrupt releases the claim so the turn never resurrects", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const sessionID = Session.ID.make("ses_claim_user_cancel")
|
||||
yield* seedSessions(database, [sessionID])
|
||||
|
||||
const draining = yield* Deferred.make<void>()
|
||||
const scope = yield* Scope.make()
|
||||
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
|
||||
const context = yield* buildExecution(scope, () =>
|
||||
Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never)),
|
||||
)
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
yield* execution.resume(sessionID).pipe(Effect.forkScoped)
|
||||
yield* Deferred.await(draining)
|
||||
expect((yield* claims(database))[sessionID]).toBe(true)
|
||||
|
||||
yield* execution.interrupt(sessionID)
|
||||
yield* execution.awaitIdle(sessionID)
|
||||
expect((yield* claims(database))[sessionID]).toBe(false)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("starts every claimed execution without waiting for earlier drains to finish", () =>
|
||||
it.effect("starts every suspended execution without waiting for earlier drains to finish", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const sessionIDs = Array.from({ length: 5 }, (_, index) => Session.ID.make(`ses_resume_concurrent_${index}`))
|
||||
@@ -158,7 +125,7 @@ describe("SessionExecution lifecycle", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("resumes each claimed Session at most once", () =>
|
||||
it.effect("resumes each suspended Session at most once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const bus = yield* Bus.Service
|
||||
@@ -167,22 +134,14 @@ describe("SessionExecution lifecycle", () => {
|
||||
yield* seedSessions(database, [first, second], { time_suspended: Date.now() })
|
||||
|
||||
const drained: string[] = []
|
||||
const bothDraining = yield* Deferred.make<void>()
|
||||
const continued: SessionEvent.Synthetic[] = []
|
||||
const scope = yield* Scope.make()
|
||||
const context = yield* buildExecution(scope, ({ sessionID }) =>
|
||||
Effect.sync(() => {
|
||||
drained.push(sessionID)
|
||||
if (drained.length === 2) Deferred.doneUnsafe(bothDraining, Effect.void)
|
||||
}),
|
||||
)
|
||||
const context = yield* buildExecution(scope, ({ sessionID }) => Effect.sync(() => void drained.push(sessionID)))
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
yield* bus.project(SessionEvent.Synthetic, (event) => Effect.sync(() => void continued.push(event)))
|
||||
|
||||
// The sweep forks resumed drains, so completion is observed through the executions.
|
||||
yield* restart.resumeSuspendedSessions
|
||||
yield* Deferred.await(bothDraining)
|
||||
yield* Effect.forEach([first, second], execution.awaitIdle, { discard: true })
|
||||
expect(drained.toSorted()).toEqual([first, second])
|
||||
expect(continued.map((event) => event.data).toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual(
|
||||
@@ -192,9 +151,7 @@ describe("SessionExecution lifecycle", () => {
|
||||
description: "Continuing after restart",
|
||||
})),
|
||||
)
|
||||
// Drains completed naturally, so claims are released and counters reset.
|
||||
expect(yield* claims(database)).toEqual({ [first]: false, [second]: false })
|
||||
expect(yield* attempts(database, first)).toBe(0)
|
||||
expect(yield* suspensions(database)).toEqual({ [first]: false, [second]: false })
|
||||
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained.length).toBe(2)
|
||||
@@ -202,104 +159,17 @@ describe("SessionExecution lifecycle", () => {
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("terminalizes a turn that exhausts its resume budget instead of crash-looping", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const bus = yield* Bus.Service
|
||||
const sessionID = Session.ID.make("ses_resume_exhausted")
|
||||
// A claim from a dead process, already resumed twice without completing.
|
||||
yield* seedSessions(database, [sessionID], { time_suspended: Date.now(), resume_attempts: 2 })
|
||||
|
||||
const drained: string[] = []
|
||||
const failures: SessionEvent.Execution.Failed[] = []
|
||||
const scope = yield* Scope.make()
|
||||
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
|
||||
const context = yield* buildExecution(scope, ({ sessionID: id }) => Effect.sync(() => void drained.push(id)), {
|
||||
maxAttempts: 2,
|
||||
})
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
yield* bus.project(SessionEvent.Execution.Failed, (event) => Effect.sync(() => void failures.push(event)))
|
||||
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained).toEqual([])
|
||||
expect(failures.map((event) => event.data.error.type)).toEqual(["aborted"])
|
||||
// The terminal released the claim and reset the counter atomically.
|
||||
expect(yield* claims(database)).toEqual({ [sessionID]: false })
|
||||
expect(yield* attempts(database, sessionID)).toBe(0)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("counts every resume durably and never consumes the claim it recovers", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const sessionID = Session.ID.make("ses_resume_counted")
|
||||
yield* seedSessions(database, [sessionID], { time_suspended: Date.now() })
|
||||
|
||||
const draining = yield* Deferred.make<void>()
|
||||
const scope = yield* Scope.make()
|
||||
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
|
||||
// The drain never terminalizes (mirrors a process that will die mid-turn).
|
||||
const context = yield* buildExecution(scope, () =>
|
||||
Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never)),
|
||||
)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
yield* restart.resumeSuspendedSessions.pipe(Effect.forkIn(scope))
|
||||
yield* Deferred.await(draining)
|
||||
|
||||
// The attempt is durable before the drain runs, and the claim is held
|
||||
// throughout: a crash anywhere in the resume path leaves both intact.
|
||||
expect(yield* attempts(database, sessionID)).toBe(1)
|
||||
expect((yield* claims(database))[sessionID]).toBe(true)
|
||||
|
||||
// Teardown (a graceful shutdown's interrupt) preserves both, so the next
|
||||
// boot counts attempt 2 against the same turn.
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
expect((yield* claims(database))[sessionID]).toBe(true)
|
||||
expect(yield* attempts(database, sessionID)).toBe(1)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("the sweep leaves Sessions already draining in this process untouched", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const bus = yield* Bus.Service
|
||||
const sessionID = Session.ID.make("ses_resume_local_active")
|
||||
yield* seedSessions(database, [sessionID])
|
||||
|
||||
const draining = yield* Deferred.make<void>()
|
||||
const continued: SessionEvent.Synthetic[] = []
|
||||
const scope = yield* Scope.make()
|
||||
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
|
||||
const context = yield* buildExecution(scope, () =>
|
||||
Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never)),
|
||||
)
|
||||
const execution = Context.get(context, SessionExecution.Service)
|
||||
const restart = Context.get(context, SessionRestart.Service)
|
||||
yield* bus.project(SessionEvent.Synthetic, (event) => Effect.sync(() => void continued.push(event)))
|
||||
|
||||
// A live local turn holds a claim; the sweep must not count, continue, or terminalize it.
|
||||
yield* execution.resume(sessionID).pipe(Effect.forkScoped)
|
||||
yield* Deferred.await(draining)
|
||||
yield* restart.resumeSuspendedSessions
|
||||
|
||||
expect(continued).toEqual([])
|
||||
expect(yield* attempts(database, sessionID)).toBe(0)
|
||||
expect((yield* claims(database))[sessionID]).toBe(true)
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
function seedSessions(
|
||||
database: Database.Service["Service"],
|
||||
sessionIDs: ReadonlyArray<Session.ID>,
|
||||
values: Partial<Pick<typeof SessionTable.$inferInsert, "time_suspended" | "resume_attempts" | "parent_id">> = {},
|
||||
values: { time_suspended?: number } = {},
|
||||
) {
|
||||
return Effect.gen(function* () {
|
||||
yield* database.db
|
||||
.insert(ProjectTable)
|
||||
.values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
|
||||
.onConflictDoNothing()
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
yield* database.db
|
||||
@@ -320,35 +190,19 @@ function seedSessions(
|
||||
})
|
||||
}
|
||||
|
||||
function claims(database: Database.Service["Service"]) {
|
||||
function suspensions(database: Database.Service["Service"]) {
|
||||
return database.db
|
||||
.select({ id: SessionTable.id, claimed: SessionTable.time_suspended })
|
||||
.select({ id: SessionTable.id, suspended: SessionTable.time_suspended })
|
||||
.from(SessionTable)
|
||||
.all()
|
||||
.pipe(
|
||||
Effect.orDie,
|
||||
Effect.map((rows) => Object.fromEntries(rows.map((row) => [row.id, row.claimed !== null]))),
|
||||
)
|
||||
}
|
||||
|
||||
function attempts(database: Database.Service["Service"], sessionID: Session.ID) {
|
||||
return database.db
|
||||
.select({ attempts: SessionTable.resume_attempts })
|
||||
.from(SessionTable)
|
||||
.where(eq(SessionTable.id, sessionID))
|
||||
.get()
|
||||
.pipe(
|
||||
Effect.orDie,
|
||||
Effect.map((row) => row?.attempts),
|
||||
Effect.map((rows) => Object.fromEntries(rows.map((row) => [row.id, row.suspended !== null]))),
|
||||
)
|
||||
}
|
||||
|
||||
/** Builds the local execution layer plus the restart actions against the test harness services. */
|
||||
function buildExecution(
|
||||
scope: Scope.Closeable,
|
||||
drain: SessionRunner.Interface["drain"],
|
||||
options?: SessionRestart.Options,
|
||||
) {
|
||||
function buildExecution(scope: Scope.Closeable, drain: SessionRunner.Interface["drain"]) {
|
||||
return Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const bus = yield* Bus.Service
|
||||
@@ -364,7 +218,7 @@ function buildExecution(
|
||||
),
|
||||
)
|
||||
return yield* Layer.buildWithScope(
|
||||
SessionRestart.layer(options).pipe(
|
||||
SessionRestart.layer.pipe(
|
||||
Layer.provideMerge(SessionExecution.layer),
|
||||
Layer.provide(Layer.succeed(Database.Service, database)),
|
||||
Layer.provide(Layer.succeed(Bus.Service, bus)),
|
||||
|
||||
@@ -63,7 +63,6 @@ const session = (
|
||||
time_compacting: 3,
|
||||
time_archived: null,
|
||||
time_suspended: null,
|
||||
resume_attempts: 0,
|
||||
...overrides,
|
||||
})
|
||||
|
||||
|
||||
@@ -13,19 +13,6 @@ const session = yield * opencode.sessions.get({ sessionID })
|
||||
|
||||
It also exports `Tool` for plugins that add tools with `ctx.tool.transform(...)`. Embedded plugins run through the ordinary discovery flow and register tools into each Location's `ToolRegistry` through the normal `Tools.Service.register(...)` path. Closing the owning Effect Scope releases router resources, location services, fibers, and scoped tool registrations.
|
||||
|
||||
Embedded hosts are silent by default. Set `log` to receive structured log entries at the selected minimum level:
|
||||
|
||||
```ts
|
||||
const opencode =
|
||||
yield *
|
||||
OpenCode.create({
|
||||
log: {
|
||||
level: "warn",
|
||||
emit: (entry) => console.error(entry.message, entry.attributes, entry.cause),
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
`sessions.events({ sessionID, after })` replays durable events after the optional aggregate sequence, then emits newly committed durable events. `sessions.interrupt(...)` targets execution owned by this host, and `sessions.message(...)` retrieves one projected Session message.
|
||||
|
||||
The same constructor is available as a service Layer:
|
||||
@@ -39,4 +26,4 @@ const program = Effect.gen(function* () {
|
||||
yield * program.pipe(Effect.provide(OpenCode.layer))
|
||||
```
|
||||
|
||||
`OpenCode.layer` adapts the silent default `OpenCode.create()` for dependency injection; use `OpenCode.layerWith(options)` to configure the host.
|
||||
`OpenCode.layer` adapts `OpenCode.create()` for dependency injection; it does not define another host implementation.
|
||||
|
||||
@@ -1,71 +0,0 @@
|
||||
import { Context, Formatter, Layer, Logger, References } from "effect"
|
||||
|
||||
export type LogLevel = "trace" | "debug" | "info" | "warn" | "error" | "fatal"
|
||||
|
||||
export type LogEntry = {
|
||||
readonly level: LogLevel
|
||||
readonly message: string
|
||||
readonly attributes?: Readonly<Record<string, unknown>>
|
||||
readonly cause?: unknown
|
||||
}
|
||||
|
||||
export type LogWriter = (entry: LogEntry) => void
|
||||
|
||||
export type LogOptions = {
|
||||
readonly level?: LogLevel
|
||||
readonly emit: LogWriter
|
||||
}
|
||||
|
||||
const levels: Record<LogLevel, Logger.Options<unknown>["logLevel"]> = {
|
||||
trace: "Trace",
|
||||
debug: "Debug",
|
||||
info: "Info",
|
||||
warn: "Warn",
|
||||
error: "Error",
|
||||
fatal: "Fatal",
|
||||
}
|
||||
|
||||
function normalizeLevel(level: Logger.Options<unknown>["logLevel"]): LogLevel | undefined {
|
||||
const output = Object.fromEntries(Object.entries(levels).map(([name, effect]) => [effect, name]))
|
||||
return output[level] as LogLevel
|
||||
}
|
||||
|
||||
export function layer(log?: LogOptions) {
|
||||
const logger = Logger.make((options) => {
|
||||
if (!log) return
|
||||
const level = normalizeLevel(options.logLevel)
|
||||
if (!level) return
|
||||
const entry = Logger.formatStructured.log(options)
|
||||
const values = Array.isArray(entry.message) ? entry.message : [entry.message]
|
||||
const [message, ...data] = values
|
||||
const details =
|
||||
data.length === 1 && !Array.isArray(data[0]) ? (data[0] as Readonly<Record<string, unknown>>) : undefined
|
||||
const { cause: detailCause, ...detailAttributes } = details ?? {}
|
||||
const attributes = {
|
||||
...entry.annotations,
|
||||
...detailAttributes,
|
||||
...(Object.keys(entry.spans).length > 0 ? { spans: entry.spans } : {}),
|
||||
...(!details && data.length > 0 ? { data: data.length === 1 ? data[0] : data } : {}),
|
||||
}
|
||||
try {
|
||||
log.emit({
|
||||
level,
|
||||
message: typeof message === "string" ? message : Formatter.format(message),
|
||||
...(Object.keys(attributes).length > 0 ? { attributes } : {}),
|
||||
...(entry.cause === undefined && detailCause === undefined ? {} : { cause: entry.cause ?? detailCause }),
|
||||
})
|
||||
} catch {
|
||||
// A host logger must not break OpenCode operations.
|
||||
}
|
||||
})
|
||||
return Layer.merge(
|
||||
Logger.layer([logger], { mergeWithExisting: false }),
|
||||
Layer.succeed(References.MinimumLogLevel, levels[log?.level ?? "info"]),
|
||||
)
|
||||
}
|
||||
|
||||
export function context(source: Context.Context<never>) {
|
||||
return Context.make(Logger.CurrentLoggers, Context.get(source, Logger.CurrentLoggers)).pipe(
|
||||
Context.add(References.MinimumLogLevel, Context.get(source, References.MinimumLogLevel)),
|
||||
)
|
||||
}
|
||||
@@ -2,27 +2,18 @@ import { OpenCode } from "@opencode-ai/client/effect"
|
||||
import { SdkPlugins } from "@opencode-ai/core/plugin/sdk"
|
||||
import { createEmbeddedRoutes } from "@opencode-ai/server/routes"
|
||||
import type { ServerOptions } from "@opencode-ai/server/options"
|
||||
import { Context, Effect, Layer, ManagedRuntime, Scope } from "effect"
|
||||
import { FetchHttpClient, HttpEffect, HttpRouter, HttpServer, HttpServerRequest } from "effect/unstable/http"
|
||||
import * as Logging from "./logging"
|
||||
import { Context, Effect, Layer, ManagedRuntime } from "effect"
|
||||
import { FetchHttpClient, HttpEffect, HttpRouter, HttpServer } from "effect/unstable/http"
|
||||
|
||||
export type { LogEntry, LogLevel, LogOptions, LogWriter } from "./logging"
|
||||
import type { LogOptions } from "./logging"
|
||||
|
||||
export type CreateOptions = ServerOptions & {
|
||||
readonly log?: LogOptions
|
||||
}
|
||||
|
||||
export const create = Effect.fn("OpenCode.create")(function* (options: CreateOptions = {}) {
|
||||
const { log, ...server } = options
|
||||
export const create = Effect.fn("OpenCode.create")(function* (options: ServerOptions = {}) {
|
||||
const runtime = yield* Effect.acquireRelease(
|
||||
Effect.sync(() =>
|
||||
ManagedRuntime.make(
|
||||
createEmbeddedRoutes({
|
||||
...server,
|
||||
app: { ...server.app, name: server.app?.name ?? "sdk" },
|
||||
database: { path: ":memory:", ...server.database },
|
||||
}).pipe(Layer.provide(HttpServer.layerServices), Layer.provideMerge(Logging.layer(log))),
|
||||
...options,
|
||||
app: { ...options.app, name: options.app?.name ?? "sdk" },
|
||||
database: { path: ":memory:", ...options.database },
|
||||
}).pipe(Layer.provide(HttpServer.layerServices)),
|
||||
),
|
||||
),
|
||||
(runtime) => runtime.disposeEffect,
|
||||
@@ -30,9 +21,7 @@ export const create = Effect.fn("OpenCode.create")(function* (options: CreateOpt
|
||||
const context = yield* runtime.contextEffect
|
||||
const plugins = Context.get(context, SdkPlugins.Service)
|
||||
const router = Context.get(context, HttpRouter.HttpRouter)
|
||||
const handler = HttpEffect.toWebHandlerWith<never, HttpServerRequest.HttpServerRequest | Scope.Scope>(
|
||||
Logging.context(context),
|
||||
)(router.asHttpEffect())
|
||||
const handler = HttpEffect.toWebHandler(router.asHttpEffect())
|
||||
const fetch = Object.assign((input: RequestInfo | URL, init?: RequestInit) => handler(new Request(input, init)), {
|
||||
preconnect: () => undefined,
|
||||
}) satisfies typeof globalThis.fetch
|
||||
|
||||
@@ -247,10 +247,12 @@ function unavailable(status: Status.State) {
|
||||
}
|
||||
|
||||
/**
|
||||
* The managed server owns restart continuity: at boot it resumes Sessions whose execution claim was
|
||||
* never released. Claims are written when execution starts (see SessionExecution), so recovery covers
|
||||
* graceful restarts and unclean deaths alike — no shutdown hook participates.
|
||||
* The managed server owns restart continuity: it resumes Sessions the previous server suspended and
|
||||
* suspends its own active Sessions on graceful shutdown. Suspension runs while the drains are still
|
||||
* alive: connections close first, this finalizer runs next, and Session execution teardown follows.
|
||||
*/
|
||||
const installRestartContinuity = Effect.fnUntraced(function* (restart: SessionRestart.Interface) {
|
||||
yield* Effect.forkScoped(restart.resumeSuspendedSessions)
|
||||
// Registered after the fork so suspension observes still-running resumed drains during teardown.
|
||||
yield* Effect.addFinalizer(() => restart.suspendActiveSessions)
|
||||
})
|
||||
|
||||
@@ -7,13 +7,13 @@ import Config from "@npmcli/config"
|
||||
import { definitions, flatten, nerfDarts, shorthands } from "@npmcli/config/lib/definitions/index.js"
|
||||
import { Effect } from "effect"
|
||||
|
||||
const npmPath = fileURLToPath(new URL("..", import.meta.url))
|
||||
|
||||
export const load = (dir: string) =>
|
||||
Effect.tryPromise({
|
||||
try: async () => {
|
||||
const config = new Config({
|
||||
// Resolved per call: on workerd import.meta.url is undefined and building
|
||||
// this URL at module scope fails startup validation; npm config never runs there.
|
||||
npmPath: fileURLToPath(new URL("..", import.meta.url)),
|
||||
npmPath,
|
||||
cwd: dir,
|
||||
env: { ...process.env },
|
||||
argv: [process.execPath, process.execPath, "--prefix", dir],
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
export * as Observability from "./observability.js"
|
||||
|
||||
import * as NodeFileSystem from "@effect/platform-node/NodeFileSystem"
|
||||
import { NodeFileSystem } from "@effect/platform-node"
|
||||
import { LayerNode } from "./effect/layer-node.js"
|
||||
import { Effect, Layer, Logger, References, Schema } from "effect"
|
||||
import { FetchHttpClient } from "effect/unstable/http"
|
||||
@@ -50,10 +50,4 @@ export function layer(
|
||||
).pipe(Layer.catchCause(() => local))
|
||||
}
|
||||
|
||||
// Layer.suspend: constructing the loggers eagerly at module scope performs
|
||||
// I/O (file logger, run id) that workerd forbids in global scope.
|
||||
export const node = LayerNode.make({
|
||||
name: "observability",
|
||||
layer: Layer.suspend(() => layer()),
|
||||
deps: [],
|
||||
})
|
||||
export const node = LayerNode.make({ name: "observability", layer: layer(), deps: [] })
|
||||
|
||||
@@ -3,7 +3,7 @@ import path from "path"
|
||||
import { Global } from "../global.js"
|
||||
import { runID } from "./shared.js"
|
||||
|
||||
function formatter(id: string = runID()) {
|
||||
function formatter(id: string = runID) {
|
||||
return Logger.map(Logger.formatStructured, (output) => {
|
||||
const messages = Array.isArray(output.message) ? output.message : [output.message]
|
||||
return [
|
||||
@@ -51,7 +51,7 @@ export function file(local = true, channel = "local") {
|
||||
return path.join(Global.Path.log, `opencode-${channel.replace(/[^a-zA-Z0-9._-]/g, "-")}.log`)
|
||||
}
|
||||
|
||||
export function fileLogger(target = file(), id: string = runID()) {
|
||||
export function fileLogger(target = file(), id: string = runID) {
|
||||
// Do not set batchWindow to 0; it causes high idle CPU usage.
|
||||
return Effect.gen(function* () {
|
||||
const fs = yield* FileSystem.FileSystem
|
||||
|
||||
@@ -54,8 +54,8 @@ export function resource(app: App = { client: "opencode", version: "unknown", ch
|
||||
...resourceAttributes(),
|
||||
"deployment.environment.name": app.channel,
|
||||
"opencode.client": app.client,
|
||||
"opencode.run": runID(),
|
||||
"service.instance.id": runID(),
|
||||
"opencode.run": runID,
|
||||
"service.instance.id": runID,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,8 +1 @@
|
||||
// Lazy: workerd forbids generating random values in global scope, so the id
|
||||
// materializes on first call (inside a handler) and stays stable afterwards.
|
||||
let generated: string | undefined
|
||||
|
||||
export function runID(): string {
|
||||
generated ??= crypto.randomUUID().slice(0, 8)
|
||||
return generated
|
||||
}
|
||||
export const runID = crypto.randomUUID().slice(0, 8)
|
||||
|
||||
@@ -19,9 +19,8 @@ V2 loads:
|
||||
|
||||
1. The global file at `$XDG_CONFIG_HOME/opencode/AGENTS.md`, normally
|
||||
`~/.config/opencode/AGENTS.md`.
|
||||
2. Every `AGENTS.md` from the current Location up to and including the home
|
||||
directory when the Location is inside it. For Locations outside home, the
|
||||
scan stops at the project root.
|
||||
2. Every `AGENTS.md` from the current Location up to and including the project
|
||||
root.
|
||||
|
||||
For example, when the Location is `packages/web`, OpenCode can load all three
|
||||
project files below:
|
||||
@@ -36,9 +35,9 @@ my-project/
|
||||
```
|
||||
|
||||
The files are combined rather than selecting a single winner. They are rendered
|
||||
in this order: global, then files from the Location toward home or the project
|
||||
in this order: global, then project files from the Location toward the project
|
||||
root. OpenCode does not resolve conflicts between their contents, so keep broad
|
||||
guidance global and put scoped guidance in the relevant directory.
|
||||
guidance global and put scoped guidance in the relevant project directory.
|
||||
|
||||
If the Location is outside the project root, only the global file is loaded.
|
||||
Setting `OPENCODE_DISABLE_PROJECT_CONFIG=1` also skips project `AGENTS.md`
|
||||
|
||||
@@ -504,8 +504,8 @@ require a V2 rewrite. See [Skills](/skills).
|
||||
|
||||
### Instruction files
|
||||
|
||||
Existing `AGENTS.md` files stay in place. V2 discovers the global `~/.config/opencode/AGENTS.md` and ambient `AGENTS.md`
|
||||
files from the current directory up to home. For projects outside home, discovery stops at the project root.
|
||||
Existing `AGENTS.md` files stay in place. V2 discovers the global `~/.config/opencode/AGENTS.md` and project `AGENTS.md`
|
||||
files from the current directory up to the project root.
|
||||
|
||||
If a V1 setup relied on a `CLAUDE.md` fallback, move that guidance into the applicable `AGENTS.md`. V2 currently only
|
||||
discovers `AGENTS.md`; because non-API V1 behavior is intended to remain compatible, also run `/report` with the affected
|
||||
|
||||
Reference in New Issue
Block a user