diff --git a/packages/core/src/filesystem.ts b/packages/core/src/filesystem.ts index cdff0337df..a483ccdbd5 100644 --- a/packages/core/src/filesystem.ts +++ b/packages/core/src/filesystem.ts @@ -34,7 +34,24 @@ export type ListInput = typeof ListInput.Type export { FindInput } export const DEFAULT_SEARCH_LIMIT = 100 -export const DEFAULT_SEARCH_TIMEOUT_MS = 30_000 +/** One second of headroom over the hosted provider-side ripgrep kill. */ +export const DEFAULT_SEARCH_TIMEOUT_MS = (Ripgrep.HOSTED_KILL_TIMEOUT_SECONDS + 1) * 1_000 + +/** Applies the shared search budget, failing with a search-scoped message. */ +export const searchTimeout = + (fail: (message: string) => F) => + (effect: Effect.Effect): Effect.Effect => + effect.pipe( + Effect.timeoutOrElse({ + duration: DEFAULT_SEARCH_TIMEOUT_MS, + orElse: () => + Effect.fail( + fail( + `Search timed out after ${DEFAULT_SEARCH_TIMEOUT_MS / 1_000} seconds. Consider using a more specific path or pattern.`, + ), + ), + }), + ) export class GlobInput extends Schema.Class("FileSystem.GlobInput")({ pattern: Schema.String, diff --git a/packages/core/src/location.ts b/packages/core/src/location.ts index b892fc5c11..3d522ff9ae 100644 --- a/packages/core/src/location.ts +++ b/packages/core/src/location.ts @@ -1,3 +1,4 @@ +import path from "path" import { Context, Effect, Layer } from "effect" import { Info, Ref, response } from "@opencode-ai/schema/location" import { Workspace } from "@opencode-ai/schema/workspace" @@ -16,6 +17,12 @@ export interface Interface extends Info { export class Service extends Context.Service()("@opencode/Location") {} +/** + * Path rules for a Location's directory. Hosted directories live in the + * provider filesystem: posix semantics regardless of the host platform. + */ +export const paths = (location: Pick) => (location.workspaceID ? path.posix : path) + export const node = LayerNode.unbound(Service, tags.values.location) const layer = (ref: Ref) => diff --git a/packages/core/src/plugin/supervisor.ts b/packages/core/src/plugin/supervisor.ts index 315e963f72..70ccfe2731 100644 --- a/packages/core/src/plugin/supervisor.ts +++ b/packages/core/src/plugin/supervisor.ts @@ -310,46 +310,44 @@ const layer = Layer.effect( const nodeLayer = layer as Layer.Layer -const nodeDeps = [ - Plugin.node, - SdkPlugins.node, - Agent.node, - Catalog.node, - Command.node, - Config.node, - Credential.node, - Bus.node, - FileMutation.node, - Formatter.node, - FileSystem.node, - FSUtil.node, - Global.node, - httpClient, - Image.node, - Integration.node, - KV.node, - Location.node, - LocationMutation.node, - ModelsDev.node, - Npm.node, - Permission.node, - PluginRuntime.node, - Form.node, - ReadToolFileSystem.node, - Reference.node, - SessionInstructions.node, - Shell.node, - Skill.node, - Tool.node, - Watcher.node, - WebSearch.node, - WellKnown.node, -] as const - export const node = makeLocationNode({ service: Service, layer: nodeLayer, - deps: nodeDeps, + deps: [ + Plugin.node, + SdkPlugins.node, + Agent.node, + Catalog.node, + Command.node, + Config.node, + Credential.node, + Bus.node, + FileMutation.node, + Formatter.node, + FileSystem.node, + FSUtil.node, + Global.node, + httpClient, + Image.node, + Integration.node, + KV.node, + Location.node, + LocationMutation.node, + ModelsDev.node, + Npm.node, + Permission.node, + PluginRuntime.node, + Form.node, + ReadToolFileSystem.node, + Reference.node, + SessionInstructions.node, + Shell.node, + Skill.node, + Tool.node, + Watcher.node, + WebSearch.node, + WellKnown.node, + ], }) export { layer } diff --git a/packages/core/src/ripgrep.ts b/packages/core/src/ripgrep.ts index 73b873bf1c..42cab1b36a 100644 --- a/packages/core/src/ripgrep.ts +++ b/packages/core/src/ripgrep.ts @@ -20,6 +20,13 @@ import { WorkspaceEnvironment } from "./workspace/environment" const ERROR_BYTES = 8 * 1024 const MAX_SUBMATCHES = 100 +/** + * Provider-side kill for hosted searches, since hosted process kill is not + * implemented. Kept below FileSystem.DEFAULT_SEARCH_TIMEOUT_MS so the sandbox + * process dies before the caller's timeout fires. + */ +export const HOSTED_KILL_TIMEOUT_SECONDS = 29 + const RawMatch = Schema.Struct({ type: Schema.Literal("match"), data: Schema.Struct({ @@ -319,8 +326,8 @@ export const hostedNode = makeLocationNode({ [ ...env.shell.args( input.output === "lines" - ? 'set -o pipefail; limit=$1; shift; timeout --signal=KILL 29s rg "$@" | head -n "$limit"' - : `set -o pipefail; limit=$1; shift; timeout --signal=KILL 29s rg "$@" | awk -v limit="$limit" '{ print } /"type":"match"/ { if (++matches >= limit) exit }'`, + ? `set -o pipefail; limit=$1; shift; timeout --signal=KILL ${HOSTED_KILL_TIMEOUT_SECONDS}s rg "$@" | head -n "$limit"` + : `set -o pipefail; limit=$1; shift; timeout --signal=KILL ${HOSTED_KILL_TIMEOUT_SECONDS}s rg "$@" | awk -v limit="$limit" '{ print } /"type":"match"/ { if (++matches >= limit) exit }'`, ), "opencode-ripgrep", String(input.limit + 1), diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 7f86fc3922..2bf93de098 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -352,12 +352,8 @@ const layer = Layer.effect( projectID: project.id, parentID: input.parentID, location, - // Hosted directories live in the provider filesystem: posix - // semantics regardless of the host platform. subpath: RelativePath.make( - location.workspaceID - ? path.posix.relative(project.directory, location.directory) - : path.relative(project.directory, location.directory).replaceAll("\\", "/"), + Location.paths(location).relative(project.directory, location.directory).replaceAll("\\", "/"), ), title: input.title, agent: input.agent, @@ -718,10 +714,10 @@ const layer = Layer.effect( const expanded = value === "~" ? global.home : value.startsWith("~/") ? path.join(global.home, value.slice(2)) : value const directory = AbsolutePath.make(path.resolve(current.location.directory, expanded)) + if (current.location.directory === directory) return const info = yield* fs.stat(directory).pipe(Effect.catch(() => Effect.succeed(undefined))) if (!info) return yield* new DestinationNotFoundError({ directory }) if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) - if (current.location.directory === directory) return const project = yield* projects.resolve(directory) yield* persistProject(project) if ((yield* execution.active).has(input.sessionID)) { diff --git a/packages/core/src/shell.ts b/packages/core/src/shell.ts index 1592e6483c..211352c14f 100644 --- a/packages/core/src/shell.ts +++ b/packages/core/src/shell.ts @@ -91,274 +91,290 @@ interface Backend { readonly validateCwd: (cwd: string) => Effect.Effect } -const layerWith = (backend: Effect.Effect) => Layer.effect( - Service, - Effect.gen(function* () { - const bus = yield* Bus.Service - const location = yield* Location.Service - const global = yield* Global.Service - const spawner = yield* backend - const hooks = yield* PluginHooks.Service - const context = yield* Effect.context() - const runFork = Effect.runForkWith(context) - const sessions = new Map() - const exitOrder: string[] = [] +const layerWith = (backend: Effect.Effect) => + Layer.effect( + Service, + Effect.gen(function* () { + const bus = yield* Bus.Service + const location = yield* Location.Service + const global = yield* Global.Service + const spawner = yield* backend + const hooks = yield* PluginHooks.Service + const context = yield* Effect.context() + const runFork = Effect.runForkWith(context) + const sessions = new Map() + const exitOrder: string[] = [] - const outputDir = path.join(global.data, "shell", location.project.id) - const { mkdir, unlink } = yield* Effect.promise(() => import("fs/promises")) - const { createWriteStream, createReadStream } = yield* Effect.promise(() => import("fs")) - yield* Effect.promise(() => mkdir(outputDir, { recursive: true })) + const outputDir = path.join(global.data, "shell", location.project.id) + const { mkdir, unlink } = yield* Effect.promise(() => import("fs/promises")) + const { createWriteStream, createReadStream } = yield* Effect.promise(() => import("fs")) + yield* Effect.promise(() => mkdir(outputDir, { recursive: true })) - yield* Effect.addFinalizer(() => - Effect.gen(function* () { - for (const session of sessions.values()) { - if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) - // Unblock waiters still pending at teardown; succeed is a no-op once already resolved. - yield* Deferred.fail(session.done, new NotFoundError({ id: Shell.ID.make(session.info.id) })) + yield* Effect.addFinalizer(() => + Effect.gen(function* () { + for (const session of sessions.values()) { + if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) + // Unblock waiters still pending at teardown; succeed is a no-op once already resolved. + yield* Deferred.fail(session.done, new NotFoundError({ id: Shell.ID.make(session.info.id) })) + } + sessions.clear() + exitOrder.length = 0 + }), + ) + + const require = Effect.fn("Shell.require")(function* (id: Shell.ID) { + const session = sessions.get(id) + if (!session) return yield* new NotFoundError({ id }) + return session + }) + + const removeSession = Effect.fnUntraced(function* (id: Shell.ID) { + const session = sessions.get(id) + if (!session) return + sessions.delete(id) + const index = exitOrder.indexOf(id) + if (index !== -1) exitOrder.splice(index, 1) + if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) + // Unblock any wait still pending when the command is removed before it terminated. + yield* Deferred.fail(session.done, new NotFoundError({ id })) + yield* Effect.promise(() => unlink(session.file).catch(() => {})) + yield* bus.publish(Shell.Event.Deleted, { id }) + }) + + const remove = Effect.fn("Shell.remove")(function* (id: Shell.ID) { + yield* require(id) + yield* removeSession(id) + }) + + const list = Effect.fn("Shell.list")(function* () { + return Array.from(sessions.values()) + .filter((session) => session.info.status === "running") + .map((session) => session.info) + }) + + const get = Effect.fn("Shell.get")(function* (id: Shell.ID) { + return (yield* require(id)).info + }) + + const wait = Effect.fn("Shell.wait")(function* (id: Shell.ID) { + return yield* Deferred.await((yield* require(id)).done) + }) + + const timeout = Effect.fn("Shell.timeout")(function* (id: Shell.ID, duration: number) { + const session = yield* require(id) + if (session.info.status !== "running" || !session.timeout) return session.info + yield* session.timeout(duration) + return session.info + }) + + const name = () => spawner.shell.pipe(Effect.map(ShellSelect.name)) + + const output = Effect.fn("Shell.output")(function* (id: Shell.ID, input?: Shell.OutputInput) { + const session = yield* require(id) + const cursor = input?.cursor ?? 0 + const limit = input?.limit ?? 65536 + if (cursor >= session.size) return { output: "", cursor: session.size, size: session.size, truncated: false } + const start = Math.max(0, cursor) + const length = Math.min(limit, session.size - start) + const buffer = Buffer.alloc(length) + const bytesRead = yield* Effect.promise( + () => + new Promise((resolve) => { + const stream = createReadStream(session.file, { start, end: start + length - 1 }) + let offset = 0 + stream.on("data", (chunk: string | Buffer) => { + const bytes = Buffer.from(chunk) + bytes.copy(buffer, offset) + offset += bytes.length + }) + stream.on("end", () => resolve(offset)) + stream.on("error", () => resolve(0)) + }), + ) + return { + output: buffer.subarray(0, bytesRead).toString("utf8"), + cursor: start + bytesRead, + size: session.size, + truncated: false, } - sessions.clear() - exitOrder.length = 0 - }), - ) + }) - const require = Effect.fn("Shell.require")(function* (id: Shell.ID) { - const session = sessions.get(id) - if (!session) return yield* new NotFoundError({ id }) - return session - }) + const create = Effect.fn("Shell.create")(function* ( + input: Shell.CreateInput, + before?: (input: ShellCreateBefore) => Effect.Effect, + ) { + const invocation: ShellCreateBefore = { + command: input.command, + cwd: input.cwd ?? location.directory, + timeout: input.timeout, + shell: yield* spawner.shell, + env: { + ...spawner.env, + TERM: "xterm-256color", + OPENCODE_TERMINAL: "1", + }, + } + yield* hooks.trigger("shell", "create.before", invocation) + if (before) yield* before(invocation) + yield* spawner.validateCwd(invocation.cwd) - const removeSession = Effect.fnUntraced(function* (id: Shell.ID) { - const session = sessions.get(id) - if (!session) return - sessions.delete(id) - const index = exitOrder.indexOf(id) - if (index !== -1) exitOrder.splice(index, 1) - if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) - // Unblock any wait still pending when the command is removed before it terminated. - yield* Deferred.fail(session.done, new NotFoundError({ id })) - yield* Effect.promise(() => unlink(session.file).catch(() => {})) - yield* bus.publish(Shell.Event.Deleted, { id }) - }) + const id = Shell.ID.ascending() + const args = spawner.args(invocation.shell, invocation.command) + const file = path.join(outputDir, `${id}.out`) - const remove = Effect.fn("Shell.remove")(function* (id: Shell.ID) { - yield* require(id) - yield* removeSession(id) - }) + const info: Info = { + id, + status: "running", + command: invocation.command, + cwd: invocation.cwd, + shell: invocation.shell, + file, + metadata: input.metadata ?? {}, + time: { started: Date.now() }, + } - const list = Effect.fn("Shell.list")(function* () { - return Array.from(sessions.values()) - .filter((session) => session.info.status === "running") - .map((session) => session.info) - }) - - const get = Effect.fn("Shell.get")(function* (id: Shell.ID) { - return (yield* require(id)).info - }) - - const wait = Effect.fn("Shell.wait")(function* (id: Shell.ID) { - return yield* Deferred.await((yield* require(id)).done) - }) - - const timeout = Effect.fn("Shell.timeout")(function* (id: Shell.ID, duration: number) { - const session = yield* require(id) - if (session.info.status !== "running" || !session.timeout) return session.info - yield* session.timeout(duration) - return session.info - }) - - const name = () => spawner.shell.pipe(Effect.map(ShellSelect.name)) - - const output = Effect.fn("Shell.output")(function* (id: Shell.ID, input?: Shell.OutputInput) { - const session = yield* require(id) - const cursor = input?.cursor ?? 0 - const limit = input?.limit ?? 65536 - if (cursor >= session.size) return { output: "", cursor: session.size, size: session.size, truncated: false } - const start = Math.max(0, cursor) - const length = Math.min(limit, session.size - start) - const buffer = Buffer.alloc(length) - const bytesRead = yield* Effect.promise( - () => - new Promise((resolve) => { - const stream = createReadStream(session.file, { start, end: start + length - 1 }) - let offset = 0 - stream.on("data", (chunk: string | Buffer) => { - const bytes = Buffer.from(chunk) - bytes.copy(buffer, offset) - offset += bytes.length - }) - stream.on("end", () => resolve(offset)) - stream.on("error", () => resolve(0)) - }), - ) - return { - output: buffer.subarray(0, bytesRead).toString("utf8"), - cursor: start + bytesRead, - size: session.size, - truncated: false, - } - }) - - const create = Effect.fn("Shell.create")(function* ( - input: Shell.CreateInput, - before?: (input: ShellCreateBefore) => Effect.Effect, - ) { - const invocation: ShellCreateBefore = { - command: input.command, - cwd: input.cwd ?? location.directory, - timeout: input.timeout, - shell: yield* spawner.shell, - env: { - ...spawner.env, - TERM: "xterm-256color", - OPENCODE_TERMINAL: "1", - }, - } - yield* hooks.trigger("shell", "create.before", invocation) - if (before) yield* before(invocation) - yield* spawner.validateCwd(invocation.cwd) - - const id = Shell.ID.ascending() - const args = spawner.args(invocation.shell, invocation.command) - const file = path.join(outputDir, `${id}.out`) - - const info: Info = { - id, - status: "running", - command: invocation.command, - cwd: invocation.cwd, - shell: invocation.shell, - file, - metadata: input.metadata ?? {}, - time: { started: Date.now() }, - } - - // Spawn via AppProcess and stream combined output to the file. The handle is scope-bound, so - // the managing fiber keeps its scope open until the command terminates (it awaits `done` at the - // end). `create` returns once `ready` resolves with the registered session. - const ready = Deferred.makeUnsafe() - runFork( - Effect.scoped( - Effect.gen(function* () { - const handle = yield* spawner.spawn( - ChildProcess.make(invocation.shell, args, { - cwd: invocation.cwd, - env: invocation.env, - stdin: "ignore", - detached: spawner.detached, - forceKillAfter: Duration.seconds(3), - }), - ) - const session: Active = { - info: produce(info, (draft) => { - draft.pid = handle.pid - }), - file, - size: 0, - done: Deferred.makeUnsafe(), - } - sessions.set(id, session) - - const stream = createWriteStream(file) - const outputDone = Deferred.makeUnsafe() - const pump = handle.all.pipe( - Stream.runForEach((chunk: Uint8Array) => - Effect.sync(() => { - stream.write(chunk) - session.size += chunk.length + // Spawn via AppProcess and stream combined output to the file. The handle is scope-bound, so + // the managing fiber keeps its scope open until the command terminates (it awaits `done` at the + // end). `create` returns once `ready` resolves with the registered session. + const ready = Deferred.makeUnsafe() + runFork( + Effect.scoped( + Effect.gen(function* () { + const handle = yield* spawner.spawn( + ChildProcess.make(invocation.shell, args, { + cwd: invocation.cwd, + env: invocation.env, + stdin: "ignore", + detached: spawner.detached, + forceKillAfter: Duration.seconds(3), }), - ), - ) - runFork( - Effect.gen(function* () { - yield* pump.pipe(Effect.catch(() => Effect.void)) - yield* Effect.promise( - () => - new Promise((resolve) => { - stream.end(() => resolve()) - }), - ) - yield* Deferred.succeed(outputDone, undefined) - }).pipe(Effect.catch(() => Deferred.succeed(outputDone, undefined))), - ) - yield* Effect.promise( - () => - new Promise((resolve) => { - stream.once("open", () => resolve()) - stream.once("error", () => resolve()) + ) + const session: Active = { + info: produce(info, (draft) => { + draft.pid = handle.pid }), - ) + file, + size: 0, + done: Deferred.makeUnsafe(), + } + sessions.set(id, session) - const finish = (status: Info["status"], exit?: number, beforeWait = Effect.void) => - Effect.gen(function* () { - if (session.info.status !== "running") return - session.info = produce(session.info, (draft) => { - draft.status = status - if (exit !== undefined) draft.exit = exit - draft.time.completed = Date.now() - }) - yield* beforeWait - yield* Deferred.await(outputDone) - // Resolve waiters with the terminal Info before any retention eviction, so an evicted - // session still reports success rather than the removal NotFoundError. This runs before - // the timeout-fiber interrupt below, which on the timeout path would otherwise cancel - // this very fiber (finish is invoked by the timeout fiber) before waiters are resolved. - yield* Deferred.succeed(session.done, session.info) - yield* bus.publish(Shell.Event.Exited, { - id, - ...(exit !== undefined ? { exit } : {}), - status, - }) - exitOrder.push(id) - while (exitOrder.length > EXITED_LIMIT) { - const oldest = exitOrder[0] - if (!oldest) break - yield* removeSession(Shell.ID.make(oldest)) - } - // Cancel a pending timeout once the command exits on its own. Interrupting last avoids - // aborting finish when finish itself runs on the timeout fiber. - if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) - }) + const stream = createWriteStream(file) + const outputDone = Deferred.makeUnsafe() + const pump = handle.all.pipe( + Stream.runForEach((chunk: Uint8Array) => + Effect.sync(() => { + stream.write(chunk) + session.size += chunk.length + }), + ), + ) + runFork( + Effect.gen(function* () { + yield* pump.pipe(Effect.catch(() => Effect.void)) + yield* Effect.promise( + () => + new Promise((resolve) => { + stream.end(() => resolve()) + }), + ) + yield* Deferred.succeed(outputDone, undefined) + }).pipe(Effect.catch(() => Deferred.succeed(outputDone, undefined))), + ) + yield* Effect.promise( + () => + new Promise((resolve) => { + stream.once("open", () => resolve()) + stream.once("error", () => resolve()) + }), + ) - session.timeout = (duration) => - Effect.gen(function* () { - if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) - session.timeoutFiber = undefined - if (duration === 0 || session.info.status !== "running") return - session.timeoutFiber = runFork( - Effect.sleep(Duration.millis(duration)).pipe( - Effect.flatMap(() => - finish("timeout", undefined, handle.kill().pipe(Effect.catch(() => Effect.void))), + const finish = (status: Info["status"], exit?: number, beforeWait = Effect.void) => + Effect.gen(function* () { + if (session.info.status !== "running") return + session.info = produce(session.info, (draft) => { + draft.status = status + if (exit !== undefined) draft.exit = exit + draft.time.completed = Date.now() + }) + yield* beforeWait + yield* Deferred.await(outputDone) + // Resolve waiters with the terminal Info before any retention eviction, so an evicted + // session still reports success rather than the removal NotFoundError. This runs before + // the timeout-fiber interrupt below, which on the timeout path would otherwise cancel + // this very fiber (finish is invoked by the timeout fiber) before waiters are resolved. + yield* Deferred.succeed(session.done, session.info) + yield* bus.publish(Shell.Event.Exited, { + id, + ...(exit !== undefined ? { exit } : {}), + status, + }) + exitOrder.push(id) + while (exitOrder.length > EXITED_LIMIT) { + const oldest = exitOrder[0] + if (!oldest) break + yield* removeSession(Shell.ID.make(oldest)) + } + // Cancel a pending timeout once the command exits on its own. Interrupting last avoids + // aborting finish when finish itself runs on the timeout fiber. + if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) + }) + + session.timeout = (duration) => + Effect.gen(function* () { + if (session.timeoutFiber) yield* Fiber.interrupt(session.timeoutFiber) + session.timeoutFiber = undefined + if (duration === 0 || session.info.status !== "running") return + session.timeoutFiber = runFork( + Effect.sleep(Duration.millis(duration)).pipe( + Effect.flatMap(() => + finish("timeout", undefined, handle.kill().pipe(Effect.catch(() => Effect.void))), + ), + Effect.catch(() => Effect.void), ), - Effect.catch(() => Effect.void), - ), - ) - }) + ) + }) - yield* session.timeout(invocation.timeout) + yield* session.timeout(invocation.timeout) - runFork( - handle.exitCode.pipe( - Effect.flatMap((code) => finish("exited", code)), - Effect.catch(() => Effect.void), - ), - ) + runFork( + handle.exitCode.pipe( + Effect.flatMap((code) => finish("exited", code)), + Effect.catch(() => Effect.void), + ), + ) - yield* bus.publish(Shell.Event.Created, { info }) - yield* Deferred.succeed(ready, session) - // Hold the handle's scope open until the command terminates; closing it earlier would - // release (kill) the process before its exit is observed. - yield* Deferred.await(session.done).pipe(Effect.catch(() => Effect.void)) - }), - ).pipe(Effect.catch(() => Effect.void)), - ) + yield* bus.publish(Shell.Event.Created, { info }) + yield* Deferred.succeed(ready, session) + // Hold the handle's scope open until the command terminates; closing it earlier would + // release (kill) the process before its exit is observed. + yield* Deferred.await(session.done).pipe(Effect.catch(() => Effect.void)) + }), + ).pipe(Effect.catch(() => Effect.void)), + ) - const session = yield* Deferred.await(ready) - return session.info - }) + const session = yield* Deferred.await(ready) + return session.info + }) - return Service.of({ name, create, list, get, wait, timeout, output, remove }) - }), -) + return Service.of({ name, create, list, get, wait, timeout, output, remove }) + }), + ) + +const validateCwdWith = + (stat: (cwd: string) => Effect.Effect<{ readonly type: string }, E>, isNotFound: (error: E) => boolean) => + (cwd: string) => + stat(cwd).pipe( + Effect.mapError( + (error) => + new InvalidCwdError({ path: cwd, reason: isNotFound(error) ? "not_found" : "unavailable", cause: error }), + ), + Effect.flatMap((info) => + info.type === "Directory" + ? Effect.void + : Effect.fail(new InvalidCwdError({ path: cwd, reason: "not_directory" })), + ), + ) export const layer = (options?: ShellSelect.Options) => layerWith( @@ -374,21 +390,10 @@ export const layer = (options?: ShellSelect.Options) => args: ShellSelect.args, env: process.env, detached: process.platform !== "win32", - validateCwd: (cwd) => - fs.stat(cwd).pipe( - Effect.mapError((error) => - new InvalidCwdError({ - path: cwd, - reason: error.reason._tag === "NotFound" ? "not_found" : "unavailable", - cause: error, - }), - ), - Effect.flatMap((info) => - info.type === "Directory" - ? Effect.void - : Effect.fail(new InvalidCwdError({ path: cwd, reason: "not_directory" })), - ), - ), + validateCwd: validateCwdWith( + (cwd) => fs.stat(cwd), + (error) => error.reason._tag === "NotFound", + ), } satisfies Backend }), ) @@ -416,21 +421,10 @@ export const hostedNode = makeLocationNode({ args: (_shell, command) => env.shell.args(command), env: env.shell.environmentOverrides, detached: env.shell.detached, - validateCwd: (cwd) => - env.files.stat(cwd).pipe( - Effect.mapError((error) => - new InvalidCwdError({ - path: cwd, - reason: error._tag === "WorkspaceEnvironment.NotFoundError" ? "not_found" : "unavailable", - cause: error, - }), - ), - Effect.flatMap((info) => - info.type === "Directory" - ? Effect.void - : Effect.fail(new InvalidCwdError({ path: cwd, reason: "not_directory" })), - ), - ), + validateCwd: validateCwdWith( + (cwd) => env.files.stat(cwd), + (error) => error._tag === "WorkspaceEnvironment.NotFoundError", + ), } satisfies Backend }), ), diff --git a/packages/core/src/tool/plugin/glob.ts b/packages/core/src/tool/plugin/glob.ts index d654f55548..b061fe6a48 100644 --- a/packages/core/src/tool/plugin/glob.ts +++ b/packages/core/src/tool/plugin/glob.ts @@ -42,7 +42,7 @@ export const Plugin = { effect: Effect.fn("GlobTool.Plugin")(function* (ctx: PluginContext) { const filesystem = yield* FileSystem.Service const location = yield* Location.Service - const resolve = location.workspaceID ? path.posix.resolve : path.resolve + const resolve = Location.paths(location).resolve const mutation = yield* LocationMutation.Service const permission = yield* Permission.Service @@ -88,15 +88,7 @@ export const Plugin = { limit: limit + 1, }) .pipe( - Effect.timeoutOrElse({ - duration: FileSystem.DEFAULT_SEARCH_TIMEOUT_MS, - orElse: () => - Effect.fail( - new ToolFailure({ - message: `Search timed out after ${FileSystem.DEFAULT_SEARCH_TIMEOUT_MS / 1_000} seconds. Consider using a more specific path or pattern.`, - }), - ), - }), + FileSystem.searchTimeout((message) => new ToolFailure({ message })), Effect.catchTag("FileSystem.SearchPathError", (error) => Effect.fail( new ToolFailure({ diff --git a/packages/core/src/tool/plugin/grep.ts b/packages/core/src/tool/plugin/grep.ts index 49ebe359d8..3db0d2b44d 100644 --- a/packages/core/src/tool/plugin/grep.ts +++ b/packages/core/src/tool/plugin/grep.ts @@ -58,7 +58,7 @@ export const Plugin = { effect: Effect.fn("GrepTool.Plugin")(function* (ctx: PluginContext) { const filesystem = yield* FileSystem.Service const location = yield* Location.Service - const resolve = location.workspaceID ? path.posix.resolve : path.resolve + const resolve = Location.paths(location).resolve const mutation = yield* LocationMutation.Service const permission = yield* Permission.Service @@ -105,15 +105,7 @@ export const Plugin = { limit: limit + 1, }) .pipe( - Effect.timeoutOrElse({ - duration: FileSystem.DEFAULT_SEARCH_TIMEOUT_MS, - orElse: () => - Effect.fail( - new ToolFailure({ - message: `Search timed out after ${FileSystem.DEFAULT_SEARCH_TIMEOUT_MS / 1_000} seconds. Consider using a more specific path or pattern.`, - }), - ), - }), + FileSystem.searchTimeout((message) => new ToolFailure({ message })), Effect.catchTag("FileSystem.SearchPathError", () => Effect.fail(new ToolFailure({ message: `Search path does not exist: ${input.path ?? "."}` })), ), diff --git a/packages/core/src/tool/plugin/patch.ts b/packages/core/src/tool/plugin/patch.ts index abc89c2993..5ad7a75201 100644 --- a/packages/core/src/tool/plugin/patch.ts +++ b/packages/core/src/tool/plugin/patch.ts @@ -116,13 +116,14 @@ export const Plugin = { }) } if (hunk.type === "add") { + const contents = + hunk.contents.endsWith("\n") || hunk.contents === "" ? hunk.contents : `${hunk.contents}\n` prepared.push({ ...hunk, + contents, target, before: "", - after: FileMutation.normalizeText( - hunk.contents.endsWith("\n") || hunk.contents === "" ? hunk.contents : `${hunk.contents}\n`, - ), + after: FileMutation.normalizeText(contents), }) return } @@ -149,7 +150,6 @@ export const Plugin = { }), ), )) - const before = original const update = yield* Effect.try({ try: () => Patch.derive(hunk.path, hunk.chunks, original), catch: (error) => new ToolFailure({ message: `patch verification failed: ${errorMessage(error)}` }), @@ -169,7 +169,7 @@ export const Plugin = { ...hunk, target, content: Patch.joinBom(update.content, update.bom), - before, + before: original, after: content, moveTarget, }) @@ -207,13 +207,7 @@ export const Plugin = { Effect.gen(function* () { if (change.type === "add") { const result = yield* files - .write({ - target: change.target, - content: - change.contents.endsWith("\n") || change.contents === "" - ? change.contents - : `${change.contents}\n`, - }) + .write({ target: change.target, content: change.contents }) .pipe(Effect.mapError((error) => fail(`Failed to write ${change.target.resource}`, error))) formatted.set(change.target.canonical, result.content) applied.push({