mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-08 10:09:52 -04:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 82370fdd63 | |||
| e707ffa84f | |||
| 77faeea584 | |||
| b31cec37e3 | |||
| 2c80a77906 | |||
| ec0e5373e1 | |||
| 334b278547 |
+60
-105
@@ -1,15 +1,11 @@
|
||||
{
|
||||
"version": "7",
|
||||
"dialect": "sqlite",
|
||||
"id": "2d214a71-3b0a-48c1-a667-741952c4e188",
|
||||
"id": "2aefa7b6-6847-4490-96e8-6680159f757c",
|
||||
"prevIds": [
|
||||
"f14a9b18-8207-487e-a3d3-227e629ba9ad"
|
||||
"6ff49c08-7759-48fd-beca-6086853fce79"
|
||||
],
|
||||
"ddl": [
|
||||
{
|
||||
"name": "workspace",
|
||||
"entityType": "tables"
|
||||
},
|
||||
{
|
||||
"name": "account_state",
|
||||
"entityType": "tables"
|
||||
@@ -75,84 +71,8 @@
|
||||
"entityType": "tables"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "id",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "type",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": "''",
|
||||
"generated": null,
|
||||
"name": "name",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "branch",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "directory",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "extra",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "project_id",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "time_used",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
"name": "workspace",
|
||||
"entityType": "tables"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
@@ -1385,18 +1305,53 @@
|
||||
"table": "session_v2"
|
||||
},
|
||||
{
|
||||
"columns": [
|
||||
"project_id"
|
||||
],
|
||||
"tableTo": "project",
|
||||
"columnsTo": [
|
||||
"id"
|
||||
],
|
||||
"onUpdate": "NO ACTION",
|
||||
"onDelete": "CASCADE",
|
||||
"nameExplicit": false,
|
||||
"name": "fk_workspace_project_id_project_id_fk",
|
||||
"entityType": "fks",
|
||||
"type": "text",
|
||||
"notNull": false,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "id",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "provider",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "binding",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "created_at",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"notNull": true,
|
||||
"autoincrement": false,
|
||||
"default": null,
|
||||
"generated": null,
|
||||
"name": "last_used_at",
|
||||
"entityType": "columns",
|
||||
"table": "workspace"
|
||||
},
|
||||
{
|
||||
@@ -1564,15 +1519,6 @@
|
||||
"entityType": "pks",
|
||||
"table": "instruction_entry"
|
||||
},
|
||||
{
|
||||
"columns": [
|
||||
"id"
|
||||
],
|
||||
"nameExplicit": false,
|
||||
"name": "workspace_pk",
|
||||
"table": "workspace",
|
||||
"entityType": "pks"
|
||||
},
|
||||
{
|
||||
"columns": [
|
||||
"id"
|
||||
@@ -1690,6 +1636,15 @@
|
||||
"table": "session_v2",
|
||||
"entityType": "pks"
|
||||
},
|
||||
{
|
||||
"columns": [
|
||||
"id"
|
||||
],
|
||||
"nameExplicit": false,
|
||||
"name": "workspace_pk",
|
||||
"table": "workspace",
|
||||
"entityType": "pks"
|
||||
},
|
||||
{
|
||||
"columns": [
|
||||
{
|
||||
|
||||
@@ -1,20 +0,0 @@
|
||||
import { sqliteTable, text, integer } from "drizzle-orm/sqlite-core"
|
||||
import { ProjectTable } from "../project/sql"
|
||||
import { Project } from "../project"
|
||||
import { Workspace } from "../workspace"
|
||||
|
||||
export const WorkspaceTable = sqliteTable("workspace", {
|
||||
id: text().$type<Workspace.ID>().primaryKey(),
|
||||
type: text().notNull(),
|
||||
name: text().notNull().default(""),
|
||||
branch: text(),
|
||||
directory: text(),
|
||||
extra: text({ mode: "json" }),
|
||||
project_id: text()
|
||||
.$type<Project.ID>()
|
||||
.notNull()
|
||||
.references(() => ProjectTable.id, { onDelete: "cascade" }),
|
||||
time_used: integer()
|
||||
.notNull()
|
||||
.$default(() => Date.now()),
|
||||
})
|
||||
+1
@@ -42,5 +42,6 @@ export const migrations: DatabaseMigration.Migration[] = (
|
||||
import("./migration/20260622202450_simplify_session_input"),
|
||||
import("./migration/20260804233008_loose_psylocke"),
|
||||
import("./migration/20260805200742_import_legacy_credentials"),
|
||||
import("./migration/20260808023530_workspace_domain"),
|
||||
])
|
||||
).map((module) => module.default)
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
import { Effect } from "effect"
|
||||
import type { DatabaseMigration } from "../migration"
|
||||
|
||||
const migration: DatabaseMigration.Migration = {
|
||||
id: "20260808023530_workspace_domain",
|
||||
up(tx) {
|
||||
return Effect.gen(function* () {
|
||||
yield* tx.run(`DROP TABLE \`workspace\`;`)
|
||||
yield* tx.run(`
|
||||
CREATE TABLE \`workspace\` (
|
||||
\`id\` text PRIMARY KEY,
|
||||
\`provider\` text NOT NULL,
|
||||
\`binding\` text NOT NULL,
|
||||
\`created_at\` integer NOT NULL,
|
||||
\`last_used_at\` integer NOT NULL
|
||||
);
|
||||
`)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
export default migration
|
||||
@@ -4,19 +4,6 @@ import type { DatabaseMigration } from "./migration"
|
||||
const schema: Omit<DatabaseMigration.Migration, "id"> = {
|
||||
up(tx) {
|
||||
return Effect.gen(function* () {
|
||||
yield* tx.run(`
|
||||
CREATE TABLE \`workspace\` (
|
||||
\`id\` text PRIMARY KEY,
|
||||
\`type\` text NOT NULL,
|
||||
\`name\` text DEFAULT '' NOT NULL,
|
||||
\`branch\` text,
|
||||
\`directory\` text,
|
||||
\`extra\` text,
|
||||
\`project_id\` text NOT NULL,
|
||||
\`time_used\` integer NOT NULL,
|
||||
CONSTRAINT \`fk_workspace_project_id_project_id_fk\` FOREIGN KEY (\`project_id\`) REFERENCES \`project\`(\`id\`) ON DELETE CASCADE
|
||||
);
|
||||
`)
|
||||
yield* tx.run(`
|
||||
CREATE TABLE \`account_state\` (
|
||||
\`id\` integer PRIMARY KEY,
|
||||
@@ -216,6 +203,15 @@ const schema: Omit<DatabaseMigration.Migration, "id"> = {
|
||||
CONSTRAINT \`fk_session_v2_project_id_project_id_fk\` FOREIGN KEY (\`project_id\`) REFERENCES \`project\`(\`id\`) ON DELETE CASCADE
|
||||
);
|
||||
`)
|
||||
yield* tx.run(`
|
||||
CREATE TABLE \`workspace\` (
|
||||
\`id\` text PRIMARY KEY,
|
||||
\`provider\` text NOT NULL,
|
||||
\`binding\` text NOT NULL,
|
||||
\`created_at\` integer NOT NULL,
|
||||
\`last_used_at\` integer NOT NULL
|
||||
);
|
||||
`)
|
||||
yield* tx.run(`CREATE UNIQUE INDEX \`event_aggregate_seq_idx\` ON \`event\` (\`aggregate_id\`,\`seq\`);`)
|
||||
yield* tx.run(`CREATE INDEX \`event_aggregate_type_seq_idx\` ON \`event\` (\`aggregate_id\`,\`type\`,\`seq\`);`)
|
||||
yield* tx.run(
|
||||
|
||||
@@ -5,6 +5,8 @@ import { ChildProcessSpawner } from "effect/unstable/process/ChildProcessSpawner
|
||||
import type { Files } from "./files"
|
||||
import { makeFiles } from "./index"
|
||||
import { makeLocalDriver } from "./local"
|
||||
import { Location } from "../location"
|
||||
import { Workspace } from "../workspace"
|
||||
|
||||
export interface Interface {
|
||||
readonly files: Files
|
||||
@@ -17,10 +19,25 @@ const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const spawner = yield* ChildProcessSpawner
|
||||
return Service.of({ files: makeFiles(makeLocalDriver(spawner)), spawner })
|
||||
const location = yield* Location.Service
|
||||
const workspace = yield* Workspace.Service
|
||||
const driver = location.workspaceID
|
||||
? yield* workspace.connect(location.workspaceID).pipe(
|
||||
// Environment has no error channel; an unknown or destroyed placement is a configuration defect by design.
|
||||
Effect.mapError(
|
||||
(cause) => new Error(`Failed to bind Environment to workspace ${location.workspaceID}`, { cause }),
|
||||
),
|
||||
Effect.orDie,
|
||||
)
|
||||
: makeLocalDriver(spawner)
|
||||
return Service.of({ files: makeFiles(driver), spawner: driver.spawner })
|
||||
}),
|
||||
)
|
||||
|
||||
export const node = makeLocationNode({ service: Service, layer, deps: [CrossSpawnSpawner.node] })
|
||||
export const node = makeLocationNode({
|
||||
service: Service,
|
||||
layer,
|
||||
deps: [CrossSpawnSpawner.node, Location.node, Workspace.node],
|
||||
})
|
||||
|
||||
export * as EnvironmentService from "./environment"
|
||||
|
||||
@@ -48,13 +48,15 @@ export const readText = Effect.fn("FileMutation.readText")(function* (files: Fil
|
||||
return Bom.decodeBytes((yield* files.read(target)).bytes)
|
||||
})
|
||||
|
||||
export const syncTextBom = Effect.fn("FileMutation.syncTextBom")((files: Files, target: string, bom: boolean) =>
|
||||
Effect.gen(function* () {
|
||||
const synced = Bom.syncBytes((yield* files.read(target)).bytes, bom)
|
||||
if (synced.bytes) yield* files.write(target, synced.bytes)
|
||||
return synced.text
|
||||
}).pipe(Effect.uninterruptible),
|
||||
)
|
||||
export const syncTextBom = Effect.fn("FileMutation.syncTextBom")(function* (
|
||||
files: Files,
|
||||
target: string,
|
||||
bom: boolean,
|
||||
) {
|
||||
const synced = Bom.syncBytes((yield* files.read(target)).bytes, bom)
|
||||
if (synced.bytes) yield* files.write(target, synced.bytes)
|
||||
return synced.text
|
||||
})
|
||||
|
||||
/** Share transaction locks across Location graphs that address the same file. */
|
||||
const transactionLocks = KeyedMutex.makeUnsafe<string>()
|
||||
@@ -68,10 +70,15 @@ const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const environment = yield* Environment.Service
|
||||
const locks = KeyedMutex.makeUnsafe<string>()
|
||||
const withLock: Interface["withLock"] = (targets) => (effect) =>
|
||||
[...new Set(targets.map(FSUtil.resolve))]
|
||||
.sort()
|
||||
.reduceRight((result, target) => transactionLocks.withLock(target)(result), effect)
|
||||
const withTargetLock =
|
||||
(target: Target) =>
|
||||
<A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
||||
locks.withLock(target.absolute)(Effect.uninterruptible(effect))
|
||||
|
||||
const writeResult = (target: Target, existed: boolean): WriteResult => ({
|
||||
operation: "write",
|
||||
@@ -81,32 +88,36 @@ const layer = Layer.effect(
|
||||
})
|
||||
|
||||
const write = Effect.fn("FileMutation.write")((input: WriteInput) =>
|
||||
Effect.gen(function* () {
|
||||
const existed = yield* environment.files.stat(input.target.absolute).pipe(
|
||||
Effect.as(true),
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(false)),
|
||||
)
|
||||
yield* environment.files.write(
|
||||
input.target.absolute,
|
||||
typeof input.content === "string" ? new TextEncoder().encode(input.content) : input.content,
|
||||
)
|
||||
return writeResult(input.target, existed)
|
||||
}).pipe(Effect.uninterruptible),
|
||||
withTargetLock(input.target)(
|
||||
Effect.gen(function* () {
|
||||
const existed = yield* environment.files.stat(input.target.absolute).pipe(
|
||||
Effect.as(true),
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(false)),
|
||||
)
|
||||
yield* environment.files.write(
|
||||
input.target.absolute,
|
||||
typeof input.content === "string" ? new TextEncoder().encode(input.content) : input.content,
|
||||
)
|
||||
return writeResult(input.target, existed)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
const writeTextPreservingBom = Effect.fn("FileMutation.writeTextPreservingBom")((input: TextWriteInput) =>
|
||||
Effect.gen(function* () {
|
||||
const next = Bom.split(input.content)
|
||||
const current = yield* environment.files.read(input.target.absolute, { offset: 0, length: 3 }).pipe(
|
||||
Effect.map((result) => result.bytes),
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(undefined)),
|
||||
)
|
||||
yield* environment.files.write(
|
||||
input.target.absolute,
|
||||
new TextEncoder().encode(Bom.join(next.text, Boolean(current && Bom.has(current)) || next.bom)),
|
||||
)
|
||||
return writeResult(input.target, current !== undefined)
|
||||
}).pipe(Effect.uninterruptible),
|
||||
withTargetLock(input.target)(
|
||||
Effect.gen(function* () {
|
||||
const next = Bom.split(input.content)
|
||||
const current = yield* environment.files.read(input.target.absolute, { offset: 0, length: 3 }).pipe(
|
||||
Effect.map((result) => result.bytes),
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(undefined)),
|
||||
)
|
||||
yield* environment.files.write(
|
||||
input.target.absolute,
|
||||
new TextEncoder().encode(Bom.join(next.text, Boolean(current && Bom.has(current)) || next.bom)),
|
||||
)
|
||||
return writeResult(input.target, current !== undefined)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
return Service.of({ withLock, write, writeTextPreservingBom })
|
||||
|
||||
@@ -16,7 +16,6 @@ import { SessionPendingTable, SessionMessageTable, SessionTable } from "./sql"
|
||||
import { Slug } from "../util/slug"
|
||||
import { Money } from "@opencode-ai/schema/money"
|
||||
import type { SessionSchema } from "./schema"
|
||||
import { WorkspaceTable } from "../control-plane/workspace.sql"
|
||||
|
||||
type DatabaseService = Database.Interface["db"]
|
||||
type CurrentDurableEvent = Extract<SessionEvent.Event, { readonly durable: object }>
|
||||
@@ -376,13 +375,6 @@ const layer = Layer.effectDiscard(
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!stored) return yield* Effect.die(new SessionAlreadyProjected())
|
||||
if (!event.data.location.workspaceID) return
|
||||
yield* db
|
||||
.update(WorkspaceTable)
|
||||
.set({ time_used: Date.now() })
|
||||
.where(eq(WorkspaceTable.id, event.data.location.workspaceID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
}),
|
||||
)
|
||||
yield* bus.project(SessionEvent.Moved, (event) =>
|
||||
|
||||
@@ -69,43 +69,34 @@ export const Plugin = {
|
||||
id: context.id,
|
||||
}
|
||||
const target = yield* mutation.resolve({ path: input.path, kind: "file" })
|
||||
return yield* fileMutation.withLock([target.absolute])(
|
||||
Effect.gen(function* () {
|
||||
const external = target.externalDirectory
|
||||
if (external)
|
||||
yield* permission.assert({
|
||||
...LocationMutation.externalDirectoryPermission(external),
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source,
|
||||
})
|
||||
const current = yield* FileMutation.readText(environment.files, target.absolute).pipe(
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(undefined)),
|
||||
)
|
||||
const next = Bom.split(input.content)
|
||||
const preview = fileDiff(
|
||||
target.resource,
|
||||
current?.text ?? "",
|
||||
next.text,
|
||||
current ? "modified" : "added",
|
||||
)
|
||||
yield* permission.assert({
|
||||
action: "edit",
|
||||
resources: [target.resource],
|
||||
save: ["*"],
|
||||
metadata: { files: [preview] },
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source,
|
||||
})
|
||||
const result = yield* fileMutation.writeTextPreservingBom({ target, content: input.content })
|
||||
const bom = (yield* FileMutation.readText(environment.files, target.absolute)).bom
|
||||
if (yield* formatter.file(target.absolute)) {
|
||||
yield* FileMutation.syncTextBom(environment.files, target.absolute, bom)
|
||||
}
|
||||
return result
|
||||
}),
|
||||
const external = target.externalDirectory
|
||||
if (external)
|
||||
yield* permission.assert({
|
||||
...LocationMutation.externalDirectoryPermission(external),
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source,
|
||||
})
|
||||
const current = yield* FileMutation.readText(environment.files, target.absolute).pipe(
|
||||
Effect.catchTag("Environment.NotFound", () => Effect.succeed(undefined)),
|
||||
)
|
||||
const next = Bom.split(input.content)
|
||||
const preview = fileDiff(target.resource, current?.text ?? "", next.text, current ? "modified" : "added")
|
||||
yield* permission.assert({
|
||||
action: "edit",
|
||||
resources: [target.resource],
|
||||
save: ["*"],
|
||||
metadata: { files: [preview] },
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source,
|
||||
})
|
||||
const result = yield* fileMutation.writeTextPreservingBom({ target, content: input.content })
|
||||
const bom = (yield* FileMutation.readText(environment.files, target.absolute)).bom
|
||||
if (yield* formatter.file(target.absolute)) {
|
||||
yield* FileMutation.syncTextBom(environment.files, target.absolute, bom)
|
||||
}
|
||||
return result
|
||||
}).pipe(
|
||||
Effect.map((output) => ({ output, content: toModelOutput(output) })),
|
||||
Effect.mapError((error) => new ToolFailure({ message: `Unable to write ${input.path}`, error })),
|
||||
|
||||
@@ -1,6 +1,219 @@
|
||||
export * as Workspace from "./workspace"
|
||||
|
||||
import { Workspace } from "@opencode-ai/schema/workspace"
|
||||
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { Clock, Context, Duration, Effect, Exit, Layer, Ref, Schedule, Schema, Scope } from "effect"
|
||||
import { systemError } from "effect/PlatformError"
|
||||
import { make } from "effect/unstable/process/ChildProcessSpawner"
|
||||
import type { Driver as EnvironmentDriver } from "./environment/driver"
|
||||
import { Database } from "./database/database"
|
||||
import { KeyedMutex } from "./effect/keyed-mutex"
|
||||
import { WorkspaceDriver } from "./workspace/driver"
|
||||
import { WorkspaceTable } from "./workspace/sql"
|
||||
|
||||
export const ID = Workspace.ID
|
||||
export type ID = typeof ID.Type
|
||||
export type ID = Workspace.ID
|
||||
|
||||
export class Info extends Schema.Class<Info>("Workspace.Info")({
|
||||
id: ID,
|
||||
provider: Schema.String,
|
||||
binding: WorkspaceDriver.Binding,
|
||||
createdAt: Schema.Number,
|
||||
lastUsedAt: Schema.Number,
|
||||
}) {}
|
||||
|
||||
export class NotFound extends Schema.TaggedErrorClass<NotFound>()("Workspace.NotFound", { workspaceID: ID }) {}
|
||||
|
||||
export interface Interface {
|
||||
readonly create: (provider: string) => Effect.Effect<Info, WorkspaceDriver.Error | WorkspaceDriver.ProviderNotFound>
|
||||
readonly connect: (
|
||||
workspaceID: ID,
|
||||
) => Effect.Effect<EnvironmentDriver, NotFound | WorkspaceDriver.Error | WorkspaceDriver.ProviderNotFound>
|
||||
readonly destroy: (
|
||||
workspaceID: ID,
|
||||
) => Effect.Effect<void, NotFound | WorkspaceDriver.Error | WorkspaceDriver.ProviderNotFound>
|
||||
}
|
||||
|
||||
export interface Options {
|
||||
readonly idleThreshold?: Duration.Input
|
||||
readonly pollInterval?: Duration.Input
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/Workspace") {}
|
||||
|
||||
interface Connection {
|
||||
readonly driver: WorkspaceDriver.Interface
|
||||
readonly environment: EnvironmentDriver
|
||||
readonly saveBinding: (binding: WorkspaceDriver.Binding) => Effect.Effect<void>
|
||||
readonly lastActivity: Ref.Ref<number>
|
||||
readonly active: Ref.Ref<number>
|
||||
readonly scope: Scope.Closeable
|
||||
}
|
||||
|
||||
export const configured = (options: Options = {}) =>
|
||||
makeGlobalNode({
|
||||
service: Service,
|
||||
layer: layer(options),
|
||||
deps: [Database.node, WorkspaceDriver.node],
|
||||
})
|
||||
|
||||
const layer = (options: Options) =>
|
||||
Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const db = (yield* Database.Service).db
|
||||
const registry = yield* WorkspaceDriver.RegistryService
|
||||
const lifetime = yield* Scope.Scope
|
||||
const connections = new Map<ID, Connection>()
|
||||
const locks = KeyedMutex.makeUnsafe<ID>()
|
||||
const idleThreshold = Duration.toMillis(options.idleThreshold ?? Duration.minutes(20))
|
||||
|
||||
const load = Effect.fn("Workspace.load")(function* (workspaceID: ID) {
|
||||
const row = yield* db
|
||||
.select()
|
||||
.from(WorkspaceTable)
|
||||
.where(eq(WorkspaceTable.id, workspaceID))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!row) return yield* new NotFound({ workspaceID })
|
||||
return row
|
||||
})
|
||||
|
||||
const open = Effect.fn("Workspace.open")(function* (workspaceID: ID) {
|
||||
const existing = connections.get(workspaceID)
|
||||
if (existing) return existing
|
||||
|
||||
const row = yield* load(workspaceID)
|
||||
const driver = yield* registry.get(row.provider)
|
||||
const saveBinding = (value: WorkspaceDriver.Binding) =>
|
||||
db
|
||||
.update(WorkspaceTable)
|
||||
.set({ binding: value })
|
||||
.where(eq(WorkspaceTable.id, workspaceID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
const scope = yield* Scope.fork(lifetime)
|
||||
const environment = yield* driver.connect({ workspaceID, binding: row.binding, saveBinding }).pipe(
|
||||
Effect.provideService(Scope.Scope, scope),
|
||||
Effect.onError((cause) => Scope.close(scope, Exit.failCause(cause))),
|
||||
)
|
||||
const now = yield* Clock.currentTimeMillis
|
||||
const connection: Connection = {
|
||||
driver,
|
||||
environment,
|
||||
saveBinding,
|
||||
lastActivity: yield* Ref.make(now),
|
||||
active: yield* Ref.make(0),
|
||||
scope,
|
||||
}
|
||||
connections.set(workspaceID, connection)
|
||||
yield* db
|
||||
.update(WorkspaceTable)
|
||||
.set({ last_used_at: now })
|
||||
.where(eq(WorkspaceTable.id, workspaceID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
return connection
|
||||
})
|
||||
|
||||
yield* Effect.gen(function* () {
|
||||
const now = yield* Clock.currentTimeMillis
|
||||
yield* Effect.forEach(
|
||||
[...connections.entries()],
|
||||
([workspaceID, expected]) =>
|
||||
locks.withLock(workspaceID)(
|
||||
Effect.gen(function* () {
|
||||
const connection = connections.get(workspaceID)
|
||||
if (connection !== expected || (yield* Ref.get(connection.active)) > 0) return
|
||||
const lastActivity = yield* Ref.get(connection.lastActivity)
|
||||
if (now - lastActivity < idleThreshold) return
|
||||
const row = yield* load(workspaceID)
|
||||
// Deliberate: a racing spawn blocks, then wakes cleanly. Unlocking mid-suspend could reattach a sandbox being terminated.
|
||||
yield* connection.driver.suspendForIdle({
|
||||
workspaceID,
|
||||
binding: row.binding,
|
||||
saveBinding: connection.saveBinding,
|
||||
})
|
||||
yield* db
|
||||
.update(WorkspaceTable)
|
||||
.set({ last_used_at: lastActivity })
|
||||
.where(eq(WorkspaceTable.id, workspaceID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
connections.delete(workspaceID)
|
||||
yield* Scope.close(connection.scope, Exit.void)
|
||||
}).pipe(Effect.catchCause((cause) => Effect.logError("workspace idle suspension failed", cause))),
|
||||
),
|
||||
{ concurrency: "unbounded", discard: true },
|
||||
)
|
||||
}).pipe(Effect.repeat(Schedule.spaced(options.pollInterval ?? Duration.minutes(1))), Effect.forkScoped)
|
||||
|
||||
return Service.of({
|
||||
create: Effect.fn("Workspace.create")(function* (provider) {
|
||||
const driver = yield* registry.get(provider)
|
||||
const workspaceID = ID.create()
|
||||
const result = yield* driver.create({ workspaceID })
|
||||
const now = yield* Clock.currentTimeMillis
|
||||
yield* db
|
||||
.insert(WorkspaceTable)
|
||||
.values({ id: workspaceID, provider, binding: result.binding, created_at: now, last_used_at: now })
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
return new Info({ id: workspaceID, provider, binding: result.binding, createdAt: now, lastUsedAt: now })
|
||||
}),
|
||||
connect: Effect.fn("Workspace.connect")(function* (workspaceID) {
|
||||
const spawner = make((command) =>
|
||||
Effect.acquireRelease(
|
||||
locks.withLock(workspaceID)(
|
||||
Effect.gen(function* () {
|
||||
const connection = yield* open(workspaceID).pipe(
|
||||
Effect.mapError((cause) =>
|
||||
systemError({
|
||||
_tag: "Unknown",
|
||||
module: "Workspace",
|
||||
method: "spawn",
|
||||
description: `Failed to wake workspace ${workspaceID}`,
|
||||
cause,
|
||||
}),
|
||||
),
|
||||
)
|
||||
yield* Ref.set(connection.lastActivity, yield* Clock.currentTimeMillis)
|
||||
yield* Ref.update(connection.active, (active) => active + 1)
|
||||
return connection
|
||||
}),
|
||||
),
|
||||
(connection) =>
|
||||
locks.withLock(workspaceID)(
|
||||
Effect.gen(function* () {
|
||||
yield* Ref.update(connection.active, (active) => active - 1)
|
||||
yield* Ref.set(connection.lastActivity, yield* Clock.currentTimeMillis)
|
||||
}),
|
||||
),
|
||||
).pipe(Effect.flatMap((connection) => connection.environment.spawner.spawn(command))),
|
||||
)
|
||||
// Overrides are connection-bound; per-spawn routing is required before any driver ships them, so they are deliberately omitted.
|
||||
return { spawner }
|
||||
}),
|
||||
destroy: Effect.fn("Workspace.destroy")(function* (workspaceID) {
|
||||
yield* locks.withLock(workspaceID)(
|
||||
Effect.gen(function* () {
|
||||
const row = yield* load(workspaceID)
|
||||
const connection = connections.get(workspaceID)
|
||||
connections.delete(workspaceID)
|
||||
if (connection) yield* Scope.close(connection.scope, Exit.void)
|
||||
const driver = yield* registry.get(row.provider)
|
||||
yield* driver.destroy({ workspaceID, binding: row.binding })
|
||||
yield* db.delete(WorkspaceTable).where(eq(WorkspaceTable.id, workspaceID)).run().pipe(Effect.orDie)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const node = configured()
|
||||
|
||||
// TODO(workspace-plan): add the boot janitor and ~23h safety snapshot rotation in a later PR.
|
||||
// TODO(workspace-plan): make cold wake interruptible with a re-pin loop against janitor races.
|
||||
// TODO(workspace-plan): consider RcMap at end-of-series consolidation; idle suspend and destroy need distinct finalizers.
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
export * as WorkspaceDriver from "./driver"
|
||||
|
||||
import { Workspace } from "@opencode-ai/schema/workspace"
|
||||
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
|
||||
import { Context, Effect, Layer, Schema } from "effect"
|
||||
import type { Scope } from "effect"
|
||||
import type { Driver as EnvironmentDriver } from "../environment/driver"
|
||||
|
||||
/**
|
||||
* Smallest provider-owned JSON value required to reconnect to the same
|
||||
* provider resource. Core stores it opaquely and hands it back; only the
|
||||
* owning driver reads inside.
|
||||
*/
|
||||
export const Binding = Schema.Record(Schema.String, Schema.Json)
|
||||
export type Binding = typeof Binding.Type
|
||||
|
||||
export class Error extends Schema.TaggedErrorClass<Error>()("WorkspaceDriver.Error", {
|
||||
message: Schema.optional(Schema.String),
|
||||
cause: Schema.optional(Schema.Defect()),
|
||||
}) {}
|
||||
|
||||
export class ProviderNotFound extends Schema.TaggedErrorClass<ProviderNotFound>()("WorkspaceDriver.ProviderNotFound", {
|
||||
provider: Schema.String,
|
||||
}) {}
|
||||
|
||||
export interface Interface {
|
||||
readonly create: (input: {
|
||||
readonly workspaceID: Workspace.ID
|
||||
}) => Effect.Effect<{ readonly binding: Binding }, Error>
|
||||
readonly connect: (input: {
|
||||
readonly workspaceID: Workspace.ID
|
||||
readonly binding: Binding
|
||||
readonly saveBinding: (binding: Binding) => Effect.Effect<void>
|
||||
}) => Effect.Effect<EnvironmentDriver, Error, Scope.Scope>
|
||||
readonly suspendForIdle: (input: {
|
||||
readonly workspaceID: Workspace.ID
|
||||
readonly binding: Binding
|
||||
readonly saveBinding: (binding: Binding) => Effect.Effect<void>
|
||||
}) => Effect.Effect<void, Error>
|
||||
readonly destroy: (input: {
|
||||
readonly workspaceID: Workspace.ID
|
||||
readonly binding: Binding
|
||||
}) => Effect.Effect<void, Error>
|
||||
}
|
||||
|
||||
export const make = (driver: Interface) => driver
|
||||
|
||||
export interface Registry {
|
||||
readonly get: (provider: string) => Effect.Effect<Interface, ProviderNotFound>
|
||||
}
|
||||
|
||||
export class RegistryService extends Context.Service<RegistryService, Registry>()(
|
||||
"@opencode/WorkspaceDriverRegistry",
|
||||
) {}
|
||||
|
||||
export const registry = (drivers: Readonly<Record<string, Interface>>): Registry => ({
|
||||
get: (provider) => {
|
||||
const driver = drivers[provider]
|
||||
return driver ? Effect.succeed(driver) : Effect.fail(new ProviderNotFound({ provider }))
|
||||
},
|
||||
})
|
||||
|
||||
export const registryNode = (drivers: Readonly<Record<string, Interface>>) =>
|
||||
makeGlobalNode({
|
||||
service: RegistryService,
|
||||
layer: Layer.succeed(RegistryService, RegistryService.of(registry(drivers))),
|
||||
deps: [],
|
||||
})
|
||||
|
||||
export const node = registryNode({})
|
||||
@@ -0,0 +1,11 @@
|
||||
import { Workspace } from "@opencode-ai/schema/workspace"
|
||||
import { integer, sqliteTable, text } from "drizzle-orm/sqlite-core"
|
||||
import type { WorkspaceDriver } from "./driver"
|
||||
|
||||
export const WorkspaceTable = sqliteTable("workspace", {
|
||||
id: text().$type<Workspace.ID>().primaryKey(),
|
||||
provider: text().notNull(),
|
||||
binding: text({ mode: "json" }).$type<WorkspaceDriver.Binding>().notNull(),
|
||||
created_at: integer().notNull(),
|
||||
last_used_at: integer().notNull(),
|
||||
})
|
||||
@@ -9,11 +9,12 @@ import { Environment } from "@opencode-ai/core/environment"
|
||||
import { Location } from "@opencode-ai/core/location"
|
||||
import { LocationMutation } from "@opencode-ai/core/location-mutation"
|
||||
import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { type EnvironmentFilesTransform, transformEnvironmentFiles } from "./fixture/environment"
|
||||
import { location } from "./fixture/location"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { it } from "./lib/effect"
|
||||
|
||||
function provide(directory: string, environmentLayer = LayerNode.compile(Environment.node)) {
|
||||
function provide(directory: string, transformFiles: EnvironmentFilesTransform = () => ({})) {
|
||||
const activeLocation = Layer.succeed(
|
||||
Location.Service,
|
||||
Location.Service.of(location({ directory: AbsolutePath.make(directory) })),
|
||||
@@ -21,7 +22,7 @@ function provide(directory: string, environmentLayer = LayerNode.compile(Environ
|
||||
return Effect.provide(
|
||||
AppNodeBuilder.build(LayerNode.group([LocationMutation.node, FileMutation.node]), [
|
||||
[Location.node, activeLocation],
|
||||
[Environment.node, environmentLayer],
|
||||
[Environment.node, transformEnvironmentFiles(activeLocation, transformFiles)],
|
||||
]),
|
||||
)
|
||||
}
|
||||
@@ -109,6 +110,49 @@ describe("FileMutation", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.live("serializes concurrent writes to the same absolute target", () =>
|
||||
withTmp((directory) =>
|
||||
Effect.gen(function* () {
|
||||
const targetPath = path.join(directory, "shared.txt")
|
||||
yield* Effect.promise(() => fs.writeFile(targetPath, "initial"))
|
||||
const firstStarted = yield* Deferred.make<void>()
|
||||
const releaseFirst = yield* Deferred.make<void>()
|
||||
const secondStarted = yield* Deferred.make<void>()
|
||||
let writes = 0
|
||||
const filesystem = instrumentWrites((write) =>
|
||||
Effect.gen(function* () {
|
||||
writes++
|
||||
if (writes === 1) {
|
||||
yield* Deferred.succeed(firstStarted, undefined)
|
||||
yield* Deferred.await(releaseFirst)
|
||||
} else {
|
||||
yield* Deferred.succeed(secondStarted, undefined)
|
||||
}
|
||||
yield* write
|
||||
}),
|
||||
)
|
||||
|
||||
yield* Effect.gen(function* () {
|
||||
const mutation = yield* LocationMutation.Service
|
||||
const files = yield* FileMutation.Service
|
||||
const firstPlan = yield* mutation.resolve({ path: "shared.txt" })
|
||||
const secondPlan = yield* mutation.resolve({ path: "shared.txt" })
|
||||
const first = yield* files.write({ target: firstPlan, content: "first" }).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(firstStarted)
|
||||
const second = yield* files.write({ target: secondPlan, content: "second" }).pipe(Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
expect(yield* Deferred.isDone(secondStarted)).toBe(false)
|
||||
|
||||
yield* Deferred.succeed(releaseFirst, undefined)
|
||||
yield* Deferred.await(secondStarted)
|
||||
yield* Fiber.join(first)
|
||||
yield* Fiber.join(second)
|
||||
expect(yield* Effect.promise(() => fs.readFile(targetPath, "utf8"))).toBe("second")
|
||||
}).pipe(provide(directory, filesystem))
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("shares transaction locks across Location service instances", () =>
|
||||
withTmp((directory) =>
|
||||
Effect.gen(function* () {
|
||||
@@ -159,4 +203,46 @@ describe("FileMutation", () => {
|
||||
}).pipe(provide(directory)),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("allows distinct absolute targets to proceed independently", () =>
|
||||
withTmp((directory) =>
|
||||
Effect.gen(function* () {
|
||||
const firstStarted = yield* Deferred.make<void>()
|
||||
const releaseFirst = yield* Deferred.make<void>()
|
||||
const secondFinished = yield* Deferred.make<void>()
|
||||
const secondPath = path.join(directory, "second.txt")
|
||||
let writes = 0
|
||||
const filesystem = instrumentWrites((write) =>
|
||||
++writes === 1
|
||||
? Deferred.succeed(firstStarted, undefined).pipe(
|
||||
Effect.andThen(Deferred.await(releaseFirst)),
|
||||
Effect.andThen(write),
|
||||
)
|
||||
: write.pipe(Effect.andThen(Deferred.succeed(secondFinished, undefined))),
|
||||
)
|
||||
|
||||
yield* Effect.gen(function* () {
|
||||
const mutation = yield* LocationMutation.Service
|
||||
const files = yield* FileMutation.Service
|
||||
const firstPlan = yield* mutation.resolve({ path: "first.txt" })
|
||||
const secondPlan = yield* mutation.resolve({ path: "second.txt" })
|
||||
const first = yield* files.write({ target: firstPlan, content: "first" }).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(firstStarted)
|
||||
const second = yield* files.write({ target: secondPlan, content: "second" }).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(secondFinished)
|
||||
expect(yield* Effect.promise(() => fs.readFile(secondPath, "utf8"))).toBe("second")
|
||||
|
||||
yield* Deferred.succeed(releaseFirst, undefined)
|
||||
yield* Fiber.join(first)
|
||||
yield* Fiber.join(second)
|
||||
}).pipe(provide(directory, filesystem))
|
||||
}),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
function instrumentWrites(
|
||||
run: <E>(write: Effect.Effect<void, E>, target: string) => Effect.Effect<void, E>,
|
||||
): EnvironmentFilesTransform {
|
||||
return (files) => ({ write: (target, content) => run(files.write(target, content), target) })
|
||||
}
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { Environment } from "@opencode-ai/core/environment"
|
||||
import { Location } from "@opencode-ai/core/location"
|
||||
import { Effect, Layer } from "effect"
|
||||
|
||||
export type EnvironmentFilesTransform = (files: Environment.Files) => Partial<Environment.Files>
|
||||
|
||||
export function transformEnvironmentFiles(
|
||||
location: Layer.Layer<Location.Service>,
|
||||
transform: EnvironmentFilesTransform = () => ({}),
|
||||
) {
|
||||
return Layer.effect(
|
||||
Environment.Service,
|
||||
Effect.gen(function* () {
|
||||
const current = yield* Environment.Service
|
||||
return Environment.Service.of({
|
||||
...current,
|
||||
files: { ...current.files, ...transform(current.files) },
|
||||
})
|
||||
}),
|
||||
).pipe(Layer.provide(AppNodeBuilder.build(Environment.node, [[Location.node, location]])))
|
||||
}
|
||||
@@ -3,12 +3,15 @@ import fs from "fs/promises"
|
||||
import path from "path"
|
||||
import { Effect } from "effect"
|
||||
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { Location } from "@opencode-ai/core/location"
|
||||
import { Ripgrep } from "@opencode-ai/core/ripgrep"
|
||||
import { RelativePath } from "@opencode-ai/core/schema"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { testEffect } from "./lib/effect"
|
||||
import { tempLocationLayer } from "./fixture/location"
|
||||
|
||||
const it = testEffect(LayerNode.compile(Ripgrep.node))
|
||||
const it = testEffect(AppNodeBuilder.build(Ripgrep.node, [[Location.node, tempLocationLayer]]))
|
||||
|
||||
describe("Ripgrep", () => {
|
||||
it.live("globs files as an array", () =>
|
||||
@@ -129,7 +132,9 @@ describe("Ripgrep", () => {
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) =>
|
||||
Effect.gen(function* () {
|
||||
yield* Effect.promise(() => fs.writeFile(path.join(tmp.path, "generated.ts"), `Cloudflare${"x".repeat(70 * 1024)}\n`))
|
||||
yield* Effect.promise(() =>
|
||||
fs.writeFile(path.join(tmp.path, "generated.ts"), `Cloudflare${"x".repeat(70 * 1024)}\n`),
|
||||
)
|
||||
|
||||
const matches = yield* (yield* Ripgrep.Service).grep({
|
||||
cwd: tmp.path,
|
||||
|
||||
@@ -14,6 +14,7 @@ import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { Session } from "@opencode-ai/core/session"
|
||||
import { Tool } from "@opencode-ai/core/tool"
|
||||
import { EditTool } from "@opencode-ai/core/tool/plugin/edit"
|
||||
import { transformEnvironmentFiles } from "./fixture/environment"
|
||||
import { location } from "./fixture/location"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
|
||||
@@ -80,29 +81,6 @@ const reset = () => {
|
||||
formatFile = () => Effect.succeed(false)
|
||||
}
|
||||
|
||||
const environment = Layer.effect(
|
||||
Environment.Service,
|
||||
Effect.gen(function* () {
|
||||
const current = yield* Environment.Service
|
||||
return Environment.Service.of({
|
||||
...current,
|
||||
files: {
|
||||
...current.files,
|
||||
read: (target, range) =>
|
||||
current.files
|
||||
.read(target, range)
|
||||
.pipe(
|
||||
Effect.tap((result) =>
|
||||
Effect.sync(() => reads++).pipe(Effect.andThen(Effect.suspend(() => afterRead(target, result.bytes)))),
|
||||
),
|
||||
),
|
||||
write: (target, content) =>
|
||||
Effect.sync(() => writes.push(target)).pipe(Effect.andThen(current.files.write(target, content))),
|
||||
},
|
||||
})
|
||||
}),
|
||||
).pipe(Layer.provide(LayerNode.compile(Environment.node)))
|
||||
|
||||
const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) => Effect.Effect<A, E, R>) => {
|
||||
const activeLocation = Layer.succeed(
|
||||
Location.Service,
|
||||
@@ -115,7 +93,23 @@ const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) =
|
||||
AppNodeBuilder.build(
|
||||
LayerNode.group([Tool.node, Tool.node, LocationMutation.node, FileMutation.node, editToolNode]),
|
||||
[
|
||||
[Environment.node, environment],
|
||||
[
|
||||
Environment.node,
|
||||
transformEnvironmentFiles(activeLocation, (files) => ({
|
||||
read: (target, range) =>
|
||||
files
|
||||
.read(target, range)
|
||||
.pipe(
|
||||
Effect.tap((result) =>
|
||||
Effect.sync(() => reads++).pipe(
|
||||
Effect.andThen(Effect.suspend(() => afterRead(target, result.bytes))),
|
||||
),
|
||||
),
|
||||
),
|
||||
write: (target, content) =>
|
||||
Effect.sync(() => writes.push(target)).pipe(Effect.andThen(files.write(target, content))),
|
||||
})),
|
||||
],
|
||||
[Location.node, activeLocation],
|
||||
[Formatter.node, formatter],
|
||||
[Permission.node, permission],
|
||||
|
||||
@@ -14,6 +14,7 @@ import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { Session } from "@opencode-ai/core/session"
|
||||
import { Tool } from "@opencode-ai/core/tool"
|
||||
import { PatchTool } from "@opencode-ai/core/tool/plugin/patch"
|
||||
import { transformEnvironmentFiles } from "./fixture/environment"
|
||||
import { location } from "./fixture/location"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
|
||||
@@ -82,34 +83,6 @@ const reset = () => {
|
||||
formatFile = () => Effect.succeed(false)
|
||||
}
|
||||
|
||||
const environment = Layer.effect(
|
||||
Environment.Service,
|
||||
Effect.gen(function* () {
|
||||
const current = yield* Environment.Service
|
||||
return Environment.Service.of({
|
||||
...current,
|
||||
files: {
|
||||
...current.files,
|
||||
read: (target, range) =>
|
||||
Effect.sync(() => {
|
||||
if (!editApproved) readsBeforeEditApproval++
|
||||
}).pipe(Effect.andThen(current.files.read(target, range))),
|
||||
remove: (target) => {
|
||||
if (failRemoveTarget && path.basename(target) === failRemoveTarget) return Effect.die("forced remove failure")
|
||||
if (failRemoveErrorTarget && path.basename(target) === failRemoveErrorTarget)
|
||||
return Effect.fail(new Environment.Failed({ path: target, cause: new Error("forced remove failure") }))
|
||||
return current.files.remove(target)
|
||||
},
|
||||
write: (target, content) => {
|
||||
if (failWriteTarget && path.basename(target) === failWriteTarget)
|
||||
return Effect.fail(new Environment.Failed({ path: target, cause: new Error("forced write failure") }))
|
||||
return current.files.write(target, content)
|
||||
},
|
||||
},
|
||||
})
|
||||
}),
|
||||
).pipe(Layer.provide(LayerNode.compile(Environment.node)))
|
||||
|
||||
const withTool = <A, E, R>(
|
||||
directory: string,
|
||||
body: (registry: Tool.Interface) => Effect.Effect<A, E, R>,
|
||||
@@ -126,7 +99,27 @@ const withTool = <A, E, R>(
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
AppNodeBuilder.build(LayerNode.group([Tool.node, FileMutation.node, patchToolNode]), [
|
||||
[Environment.node, environment],
|
||||
[
|
||||
Environment.node,
|
||||
transformEnvironmentFiles(activeLocation, (files) => ({
|
||||
read: (target, range) =>
|
||||
Effect.sync(() => {
|
||||
if (!editApproved) readsBeforeEditApproval++
|
||||
}).pipe(Effect.andThen(files.read(target, range))),
|
||||
remove: (target) => {
|
||||
if (failRemoveTarget && path.basename(target) === failRemoveTarget)
|
||||
return Effect.die("forced remove failure")
|
||||
if (failRemoveErrorTarget && path.basename(target) === failRemoveErrorTarget)
|
||||
return Effect.fail(new Environment.Failed({ path: target, cause: new Error("forced remove failure") }))
|
||||
return files.remove(target)
|
||||
},
|
||||
write: (target, content) => {
|
||||
if (failWriteTarget && path.basename(target) === failWriteTarget)
|
||||
return Effect.fail(new Environment.Failed({ path: target, cause: new Error("forced write failure") }))
|
||||
return files.write(target, content)
|
||||
},
|
||||
})),
|
||||
],
|
||||
[Location.node, activeLocation],
|
||||
[Formatter.node, formatter],
|
||||
[Permission.node, permission],
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import fs from "fs/promises"
|
||||
import path from "path"
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Deferred, Effect, Fiber, Layer } from "effect"
|
||||
import { Effect, Layer } from "effect"
|
||||
import { FileMutation } from "@opencode-ai/core/file-mutation"
|
||||
import { Formatter } from "@opencode-ai/core/formatter"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
@@ -13,8 +13,8 @@ import { Permission } from "@opencode-ai/core/permission"
|
||||
import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
import { Session } from "@opencode-ai/core/session"
|
||||
import { Tool } from "@opencode-ai/core/tool"
|
||||
import { EditTool } from "@opencode-ai/core/tool/plugin/edit"
|
||||
import { WriteTool } from "@opencode-ai/core/tool/plugin/write"
|
||||
import { transformEnvironmentFiles } from "./fixture/environment"
|
||||
import { location } from "./fixture/location"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
|
||||
@@ -27,25 +27,10 @@ const writeToolNode = makeLocationNode({
|
||||
deps: [Tool.node, LocationMutation.node, FileMutation.node, Environment.node, Formatter.node, Permission.node],
|
||||
})
|
||||
|
||||
const editToolNode = makeLocationNode({
|
||||
name: "test/edit-tool-plugin",
|
||||
layer: Layer.effectDiscard(registerToolPlugin(EditTool.Plugin)),
|
||||
deps: [
|
||||
Tool.node,
|
||||
LocationMutation.node,
|
||||
FileMutation.node,
|
||||
Environment.node,
|
||||
Formatter.node,
|
||||
Location.node,
|
||||
Permission.node,
|
||||
],
|
||||
})
|
||||
|
||||
const sessionID = Session.ID.make("ses_write_tool_test")
|
||||
const assertions: Permission.AssertInput[] = []
|
||||
const writes: string[] = []
|
||||
let formatFile = (_target: string): Effect.Effect<boolean> => Effect.succeed(false)
|
||||
let afterPermission = (_input: Permission.AssertInput): Effect.Effect<void> => Effect.void
|
||||
let denyAction: string | undefined
|
||||
|
||||
const permission = Layer.succeed(
|
||||
@@ -53,7 +38,6 @@ const permission = Layer.succeed(
|
||||
Permission.Service.of({
|
||||
assert: (input) =>
|
||||
Effect.sync(() => assertions.push(input)).pipe(
|
||||
Effect.andThen(Effect.suspend(() => afterPermission(input))),
|
||||
Effect.andThen(
|
||||
input.action === denyAction
|
||||
? Effect.fail(
|
||||
@@ -82,30 +66,10 @@ const reset = () => {
|
||||
assertions.length = 0
|
||||
writes.length = 0
|
||||
formatFile = () => Effect.succeed(false)
|
||||
afterPermission = () => Effect.void
|
||||
denyAction = undefined
|
||||
}
|
||||
|
||||
const environment = Layer.effect(
|
||||
Environment.Service,
|
||||
Effect.gen(function* () {
|
||||
const current = yield* Environment.Service
|
||||
return Environment.Service.of({
|
||||
...current,
|
||||
files: {
|
||||
...current.files,
|
||||
write: (target, content) =>
|
||||
Effect.sync(() => writes.push(target)).pipe(Effect.andThen(current.files.write(target, content))),
|
||||
},
|
||||
})
|
||||
}),
|
||||
).pipe(Layer.provide(LayerNode.compile(Environment.node)))
|
||||
|
||||
const withTool = <A, E, R>(
|
||||
directory: string,
|
||||
body: (registry: Tool.Interface) => Effect.Effect<A, E, R>,
|
||||
options?: { edit?: boolean },
|
||||
) => {
|
||||
const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) => Effect.Effect<A, E, R>) => {
|
||||
const activeLocation = Layer.succeed(
|
||||
Location.Service,
|
||||
Location.Service.of(location({ directory: AbsolutePath.make(directory) })),
|
||||
@@ -115,15 +79,15 @@ const withTool = <A, E, R>(
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
AppNodeBuilder.build(
|
||||
LayerNode.group([
|
||||
Tool.node,
|
||||
LocationMutation.node,
|
||||
FileMutation.node,
|
||||
writeToolNode,
|
||||
...(options?.edit ? [editToolNode] : []),
|
||||
]),
|
||||
LayerNode.group([Tool.node, Tool.node, LocationMutation.node, FileMutation.node, writeToolNode]),
|
||||
[
|
||||
[Environment.node, environment],
|
||||
[
|
||||
Environment.node,
|
||||
transformEnvironmentFiles(activeLocation, (files) => ({
|
||||
write: (target, content) =>
|
||||
Effect.sync(() => writes.push(target)).pipe(Effect.andThen(files.write(target, content))),
|
||||
})),
|
||||
],
|
||||
[Location.node, activeLocation],
|
||||
[Formatter.node, formatter],
|
||||
[Permission.node, permission],
|
||||
@@ -139,12 +103,6 @@ const call = (input: typeof WriteTool.Input.Type, id = "call-write") => ({
|
||||
call: { type: "tool-call" as const, id, name: "write", input },
|
||||
})
|
||||
|
||||
const editCall = (input: typeof EditTool.Input.Type, id = "call-edit") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "edit", input },
|
||||
})
|
||||
|
||||
const it = testEffect(Layer.empty)
|
||||
|
||||
describe("WriteTool", () => {
|
||||
@@ -454,59 +412,4 @@ describe("WriteTool", () => {
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("serializes write and edit transactions across Location service instances", () =>
|
||||
Effect.acquireUseRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => {
|
||||
reset()
|
||||
const target = path.join(tmp.path, "shared.txt")
|
||||
return Effect.gen(function* () {
|
||||
yield* Effect.promise(() => fs.writeFile(target, "initial"))
|
||||
const formatting = yield* Deferred.make<void>()
|
||||
const releaseFormatting = yield* Deferred.make<void>()
|
||||
const editApproved = yield* Deferred.make<void>()
|
||||
let formats = 0
|
||||
formatFile = () =>
|
||||
++formats === 1
|
||||
? Deferred.succeed(formatting, undefined).pipe(
|
||||
Effect.andThen(Deferred.await(releaseFormatting)),
|
||||
Effect.as(false),
|
||||
)
|
||||
: Effect.succeed(false)
|
||||
afterPermission = (input) =>
|
||||
input.source?.id === "call-serialized-edit" && input.action === "edit"
|
||||
? Deferred.succeed(editApproved, undefined).pipe(Effect.asVoid)
|
||||
: Effect.void
|
||||
|
||||
const write = yield* withTool(
|
||||
tmp.path,
|
||||
(registry) =>
|
||||
executeTool(registry, call({ path: "shared.txt", content: "before" }, "call-serialized-write")),
|
||||
{ edit: true },
|
||||
).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(formatting)
|
||||
const edit = yield* withTool(
|
||||
tmp.path,
|
||||
(registry) =>
|
||||
executeTool(
|
||||
registry,
|
||||
editCall({ path: "shared.txt", oldString: "before", newString: "after" }, "call-serialized-edit"),
|
||||
),
|
||||
{ edit: true },
|
||||
).pipe(Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
expect(yield* Deferred.isDone(editApproved)).toBe(false)
|
||||
|
||||
yield* Deferred.succeed(releaseFormatting, undefined)
|
||||
expect((yield* Fiber.join(write)).status).toBe("completed")
|
||||
expect((yield* Fiber.join(edit)).status).toBe("completed")
|
||||
expect(yield* Deferred.isDone(editApproved)).toBe(true)
|
||||
expect(yield* Effect.promise(() => fs.readFile(target, "utf8"))).toBe("after")
|
||||
})
|
||||
},
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
),
|
||||
)
|
||||
|
||||
})
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
import { beforeEach, expect } from "bun:test"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { Database } from "@opencode-ai/core/database/database"
|
||||
import { makeMemoryDriver } from "@opencode-ai/core/environment"
|
||||
import { Workspace } from "@opencode-ai/core/workspace"
|
||||
import { WorkspaceDriver } from "@opencode-ai/core/workspace/driver"
|
||||
import { WorkspaceTable } from "@opencode-ai/core/workspace/sql"
|
||||
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
|
||||
import { eq } from "drizzle-orm"
|
||||
import { Effect } from "effect"
|
||||
import { TestClock } from "effect/testing"
|
||||
import { ChildProcess } from "effect/unstable/process"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const calls: Array<{ readonly operation: string; readonly binding?: WorkspaceDriver.Binding }> = []
|
||||
const memory = makeMemoryDriver()
|
||||
let failConnect = false
|
||||
|
||||
const driver = WorkspaceDriver.make({
|
||||
create: ({ workspaceID }) => {
|
||||
calls.push({ operation: "create" })
|
||||
return Effect.succeed({ binding: { workspaceID, generation: 0 } })
|
||||
},
|
||||
connect: ({ binding }) => {
|
||||
calls.push({ operation: "connect", binding })
|
||||
if (failConnect) return Effect.fail(new WorkspaceDriver.Error({ message: "wake failed" }))
|
||||
return Effect.succeed(memory)
|
||||
},
|
||||
suspendForIdle: ({ binding, saveBinding }) => {
|
||||
calls.push({ operation: "suspendForIdle", binding })
|
||||
return saveBinding({ ...binding, generation: Number(binding.generation) + 1, suspended: true })
|
||||
},
|
||||
destroy: ({ binding }) => {
|
||||
calls.push({ operation: "destroy", binding })
|
||||
return Effect.void
|
||||
},
|
||||
})
|
||||
|
||||
const it = testEffect(
|
||||
AppNodeBuilder.build(
|
||||
LayerNode.group([Database.node, Workspace.configured({ idleThreshold: "5 minutes", pollInterval: "1 minute" })]),
|
||||
[[WorkspaceDriver.node, WorkspaceDriver.registryNode({ fake: driver })]],
|
||||
),
|
||||
)
|
||||
|
||||
beforeEach(() => {
|
||||
calls.splice(0)
|
||||
failConnect = false
|
||||
})
|
||||
|
||||
it.effect("persists the workspace lifecycle and reconnects after idle suspension", () =>
|
||||
Effect.gen(function* () {
|
||||
const workspace = yield* Workspace.Service
|
||||
const created = yield* workspace.create("fake")
|
||||
|
||||
expect(created.id.startsWith("wrk_")).toBe(true)
|
||||
expect(created.binding).toEqual({ workspaceID: created.id, generation: 0 })
|
||||
|
||||
const environment = yield* workspace.connect(created.id)
|
||||
expect(calls.map((call) => call.operation)).toEqual(["create"])
|
||||
|
||||
yield* TestClock.adjust("4 minutes")
|
||||
yield* Effect.scoped(environment.spawner.spawn(ChildProcess.make("activity"))).pipe(Effect.exit)
|
||||
yield* TestClock.adjust("4 minutes")
|
||||
expect(calls.map((call) => call.operation)).toEqual(["create", "connect"])
|
||||
|
||||
yield* TestClock.adjust("2 minutes")
|
||||
expect(calls.map((call) => call.operation)).toEqual(["create", "connect", "suspendForIdle"])
|
||||
|
||||
const stored = yield* Database.Service.use(({ db }) =>
|
||||
db.select().from(WorkspaceTable).where(eq(WorkspaceTable.id, created.id)).get(),
|
||||
).pipe(Effect.orDie)
|
||||
expect(stored?.binding).toEqual({ workspaceID: created.id, generation: 1, suspended: true })
|
||||
expect(stored?.last_used_at).toBe(4 * 60 * 1000)
|
||||
|
||||
yield* Effect.scoped(environment.spawner.spawn(ChildProcess.make("wake"))).pipe(Effect.exit)
|
||||
expect(calls.map((call) => call.operation)).toEqual(["create", "connect", "suspendForIdle", "connect"])
|
||||
expect(calls.at(-1)?.binding).toEqual({ workspaceID: created.id, generation: 1, suspended: true })
|
||||
|
||||
yield* workspace.destroy(created.id)
|
||||
expect(calls.at(-1)?.operation).toBe("destroy")
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("surfaces wake failures through the spawn error channel", () =>
|
||||
Effect.gen(function* () {
|
||||
const workspace = yield* Workspace.Service
|
||||
const created = yield* workspace.create("fake")
|
||||
const environment = yield* workspace.connect(created.id)
|
||||
yield* Effect.scoped(environment.spawner.spawn(ChildProcess.make("connect"))).pipe(Effect.exit)
|
||||
|
||||
yield* TestClock.adjust("6 minutes")
|
||||
failConnect = true
|
||||
|
||||
const error = yield* Effect.scoped(environment.spawner.spawn(ChildProcess.make("wake"))).pipe(Effect.flip)
|
||||
expect(error).toMatchObject({
|
||||
_tag: "PlatformError",
|
||||
reason: {
|
||||
_tag: "Unknown",
|
||||
module: "Workspace",
|
||||
method: "spawn",
|
||||
description: `Failed to wake workspace ${created.id}`,
|
||||
},
|
||||
})
|
||||
}),
|
||||
)
|
||||
@@ -28,6 +28,7 @@ import { SessionRestart } from "@opencode-ai/core/session/execution/restart"
|
||||
import { PluginRuntime } from "@opencode-ai/core/plugin/runtime"
|
||||
import { SdkPlugins } from "@opencode-ai/core/plugin/sdk"
|
||||
import { WellKnown } from "@opencode-ai/core/wellknown"
|
||||
import { WorkspaceDriver } from "@opencode-ai/core/workspace/driver"
|
||||
import { Watcher } from "@opencode-ai/core/filesystem/watcher"
|
||||
import { HttpRouter } from "effect/unstable/http"
|
||||
import { HttpApiBuilder } from "effect/unstable/httpapi"
|
||||
@@ -43,6 +44,7 @@ import { formLocationLayer } from "./middleware/form-location"
|
||||
import { sessionLocationLayer } from "./middleware/session-location"
|
||||
import { ServerInfo } from "./server-info"
|
||||
import type { ServerOptions } from "./options"
|
||||
import { modalWorkspaceDriver, provider as modalProvider } from "./workspace/modal-workspace"
|
||||
|
||||
const applicationServices = LayerNode.group([
|
||||
Database.node,
|
||||
@@ -115,6 +117,10 @@ function makeRoutes<AuthError, AuthServices>(
|
||||
],
|
||||
[PluginRuntime.node, PluginRuntime.layerWithCell(pluginRuntimeCell)],
|
||||
[PluginRuntime.providerNode, PluginRuntime.providerNodeWithCell(pluginRuntimeCell)],
|
||||
[
|
||||
WorkspaceDriver.node,
|
||||
WorkspaceDriver.registryNode({ [modalProvider]: modalWorkspaceDriver({ app: "opencode-workspaces" }) }),
|
||||
],
|
||||
]
|
||||
const serviceLayer = options.simulation
|
||||
? Layer.unwrap(
|
||||
|
||||
@@ -0,0 +1,129 @@
|
||||
import { WorkspaceDriver } from "@opencode-ai/core/workspace/driver"
|
||||
import { Effect, Option, Schema } from "effect"
|
||||
import type { App, Image, ModalClient, ModalClientParams, Sandbox } from "modal"
|
||||
import { createModalSandboxWithClient, makeModalDriver, type ModalImageSpec, openModalClient } from "./modal"
|
||||
|
||||
export const provider = "modal"
|
||||
|
||||
export const ModalBinding = Schema.Struct({
|
||||
sandboxId: Schema.optional(Schema.String),
|
||||
snapshotImageId: Schema.optional(Schema.String),
|
||||
})
|
||||
export type ModalBinding = typeof ModalBinding.Type
|
||||
|
||||
export interface ModalWorkspaceOptions {
|
||||
readonly app: string
|
||||
readonly client?: ModalClientParams
|
||||
readonly image?: ModalImageSpec
|
||||
}
|
||||
|
||||
export const modalWorkspaceDriver = (options: ModalWorkspaceOptions): WorkspaceDriver.Interface => {
|
||||
const name = (workspaceID: string) => `ws-${workspaceID}`
|
||||
const decodeBinding = Schema.decodeUnknownOption(ModalBinding)
|
||||
let clientPromise: Promise<ModalClient> | undefined
|
||||
let appPromise: Promise<App> | undefined
|
||||
// The SDK client and app handle are shared for the process lifetime of this driver.
|
||||
const client = () => (clientPromise ??= openModalClient(options.client))
|
||||
const app = () =>
|
||||
(appPromise ??= client().then((value) => value.apps.fromName(options.app, { createIfMissing: true })))
|
||||
|
||||
const attempt = <A>(run: () => Promise<A>) =>
|
||||
Effect.tryPromise({ try: run, catch: (cause) => new WorkspaceDriver.Error({ cause }) })
|
||||
|
||||
const binding = (value: WorkspaceDriver.Binding): ModalBinding => Option.getOrElse(decodeBinding(value), () => ({}))
|
||||
|
||||
const live = async (lookup: () => Promise<Sandbox>) => {
|
||||
const { NotFoundError } = await import("modal")
|
||||
const sandbox = await lookup().catch((error) => {
|
||||
if (error instanceof NotFoundError) return undefined
|
||||
throw error
|
||||
})
|
||||
if (sandbox && (await sandbox.poll()) === null) return sandbox
|
||||
}
|
||||
|
||||
const findLive = async (modalClient: ModalClient, value: ModalBinding, workspaceID: string) => {
|
||||
if (value.sandboxId) {
|
||||
const sandboxID = value.sandboxId
|
||||
const sandbox = await live(() => modalClient.sandboxes.fromId(sandboxID))
|
||||
if (sandbox) return sandbox
|
||||
}
|
||||
// Name fallback is valid only before the first snapshot; afterward a live named sandbox is stale by design.
|
||||
if (value.snapshotImageId) return
|
||||
return live(() => modalClient.sandboxes.fromName(options.app, name(workspaceID)))
|
||||
}
|
||||
|
||||
const createSandbox = async (workspaceID: string, image?: Image) => {
|
||||
const { AlreadyExistsError } = await import("modal")
|
||||
const modalClient = await client()
|
||||
return createModalSandboxWithClient(
|
||||
modalClient,
|
||||
await app(),
|
||||
{
|
||||
image: options.image,
|
||||
sandbox: {
|
||||
name: name(workspaceID),
|
||||
tags: { workspace: workspaceID },
|
||||
timeoutMs: 24 * 60 * 60 * 1000,
|
||||
},
|
||||
},
|
||||
image,
|
||||
).catch((error) => {
|
||||
if (error instanceof AlreadyExistsError) return modalClient.sandboxes.fromName(options.app, name(workspaceID))
|
||||
throw error
|
||||
})
|
||||
}
|
||||
|
||||
const deleteImage = (modalClient: ModalClient, imageID?: string) =>
|
||||
imageID ? attempt(() => modalClient.images.delete(imageID)).pipe(Effect.ignore) : Effect.void
|
||||
|
||||
const terminate = (sandbox?: Sandbox) =>
|
||||
sandbox ? attempt(() => sandbox.terminate({ wait: true })).pipe(Effect.ignore) : Effect.void
|
||||
|
||||
return WorkspaceDriver.make({
|
||||
create: ({ workspaceID }) =>
|
||||
attempt(async () => {
|
||||
const sandbox = await createSandbox(workspaceID)
|
||||
return { binding: { sandboxId: sandbox.sandboxId } }
|
||||
}),
|
||||
connect: ({ workspaceID, binding: value, saveBinding }) =>
|
||||
Effect.gen(function* () {
|
||||
const modalBinding = binding(value)
|
||||
const modalClient = yield* attempt(client)
|
||||
const sandbox = yield* attempt(async () => {
|
||||
const existing = await findLive(modalClient, modalBinding, workspaceID)
|
||||
const image =
|
||||
existing || !modalBinding.snapshotImageId
|
||||
? undefined
|
||||
: await modalClient.images.fromId(modalBinding.snapshotImageId)
|
||||
return existing ?? createSandbox(workspaceID, image)
|
||||
})
|
||||
if (modalBinding.sandboxId !== sandbox.sandboxId) {
|
||||
yield* saveBinding({ ...modalBinding, sandboxId: sandbox.sandboxId })
|
||||
}
|
||||
return makeModalDriver(sandbox)
|
||||
}),
|
||||
suspendForIdle: ({ workspaceID, binding: value, saveBinding }) =>
|
||||
Effect.gen(function* () {
|
||||
const modalBinding = binding(value)
|
||||
const modalClient = yield* attempt(client)
|
||||
const sandbox = yield* attempt(() => findLive(modalClient, modalBinding, workspaceID))
|
||||
if (!sandbox) return
|
||||
const snapshot = yield* attempt(() => sandbox.snapshotFilesystem({ ttlMs: null }))
|
||||
yield* saveBinding({ snapshotImageId: snapshot.imageId })
|
||||
yield* Effect.all([deleteImage(modalClient, modalBinding.snapshotImageId), terminate(sandbox)], {
|
||||
concurrency: "unbounded",
|
||||
discard: true,
|
||||
})
|
||||
}),
|
||||
destroy: ({ workspaceID, binding: value }) =>
|
||||
Effect.gen(function* () {
|
||||
const modalBinding = binding(value)
|
||||
const modalClient = yield* attempt(client)
|
||||
const sandbox = yield* attempt(() => findLive(modalClient, modalBinding, workspaceID))
|
||||
yield* Effect.all([terminate(sandbox), deleteImage(modalClient, modalBinding.snapshotImageId)], {
|
||||
concurrency: "unbounded",
|
||||
discard: true,
|
||||
})
|
||||
}),
|
||||
})
|
||||
}
|
||||
@@ -3,7 +3,7 @@ import { systemError } from "effect/PlatformError"
|
||||
import type { Command, KillOptions } from "effect/unstable/process/ChildProcess"
|
||||
import { ExitCode, make, makeHandle, ProcessId } from "effect/unstable/process/ChildProcessSpawner"
|
||||
import type { Driver } from "@opencode-ai/core/environment"
|
||||
import type { ModalClientParams, Sandbox, SandboxCreateParams } from "modal"
|
||||
import type { App, Image, ModalClient, ModalClientParams, Sandbox, SandboxCreateParams } from "modal"
|
||||
|
||||
const INNER_WRAPPER = `
|
||||
pidfile=$1
|
||||
@@ -44,13 +44,16 @@ export interface ModalImageSpec {
|
||||
readonly dockerfileCommands: ReadonlyArray<string>
|
||||
}
|
||||
|
||||
export interface ModalSandboxOptions {
|
||||
readonly app: string
|
||||
readonly client?: ModalClientParams
|
||||
export interface ModalSandboxCreateOptions {
|
||||
readonly image?: ModalImageSpec
|
||||
readonly sandbox?: SandboxCreateParams
|
||||
}
|
||||
|
||||
export interface ModalSandboxOptions extends ModalSandboxCreateOptions {
|
||||
readonly app: string
|
||||
readonly client?: ModalClientParams
|
||||
}
|
||||
|
||||
/**
|
||||
* Ubuntu supplies the GNU coreutils and findutils required by the derived Files
|
||||
* scripts. Busybox images do not satisfy the Environment contract.
|
||||
@@ -64,19 +67,11 @@ export const ubuntuImage: ModalImageSpec = {
|
||||
|
||||
/** Creates a Modal sandbox lazily, keeping the SDK off the server startup path when Modal is unused. */
|
||||
export const createModalSandbox = async (options: ModalSandboxOptions) => {
|
||||
const { ModalClient } = await import("modal")
|
||||
const client = new ModalClient(options.client)
|
||||
const client = await openModalClient(options.client)
|
||||
const app = await client.apps.fromName(options.app, { createIfMissing: true })
|
||||
const imageSpec = options.image ?? ubuntuImage
|
||||
const image = client.images.fromRegistry(imageSpec.registry).dockerfileCommands([...imageSpec.dockerfileCommands])
|
||||
// Always Modal's Full-VM runtime (beta, enabled per account): a real kernel
|
||||
// with real device nodes, so workspaces can run Docker and other
|
||||
// kernel-dependent workloads. Costs versus gVisor, measured Aug 2026:
|
||||
// per-exec floor ~285-535ms versus ~90-165ms, and filesystem snapshots only
|
||||
// (no memory snapshots — acceptable; fs-snapshot is the persistence design).
|
||||
const sandbox = await client.sandboxes.create(app, image, {
|
||||
...options.sandbox,
|
||||
experimentalOptions: { ...options.sandbox?.experimentalOptions, vm_runtime: true },
|
||||
const sandbox = await createModalSandboxWithClient(client, app, {
|
||||
image: options.image,
|
||||
sandbox: options.sandbox,
|
||||
})
|
||||
return {
|
||||
driver: makeModalDriver(sandbox),
|
||||
@@ -85,6 +80,32 @@ export const createModalSandbox = async (options: ModalSandboxOptions) => {
|
||||
}
|
||||
}
|
||||
|
||||
export const openModalClient = async (params?: ModalClientParams) => {
|
||||
const { ModalClient } = await import("modal")
|
||||
return new ModalClient(params)
|
||||
}
|
||||
|
||||
export const createModalSandboxWithClient = async (
|
||||
client: ModalClient,
|
||||
app: App,
|
||||
options: ModalSandboxCreateOptions,
|
||||
existingImage?: Image,
|
||||
) => {
|
||||
const imageSpec = options.image ?? ubuntuImage
|
||||
const image =
|
||||
existingImage ??
|
||||
client.images.fromRegistry(imageSpec.registry).dockerfileCommands([...imageSpec.dockerfileCommands])
|
||||
// Always Modal's Full-VM runtime (beta, enabled per account): a real kernel
|
||||
// with real device nodes, so workspaces can run Docker and other
|
||||
// kernel-dependent workloads. Costs versus gVisor, measured Aug 2026:
|
||||
// per-exec floor ~285-535ms versus ~90-165ms, and filesystem snapshots only
|
||||
// (no memory snapshots — acceptable; fs-snapshot is the persistence design).
|
||||
return client.sandboxes.create(app, image, {
|
||||
...options.sandbox,
|
||||
experimentalOptions: { ...options.sandbox?.experimentalOptions, vm_runtime: true },
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Adapts Modal exec to the Environment driver. Files intentionally has no native
|
||||
* overrides: exec latency dominates payload work (VM runtime floor measured
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
import fs from "node:fs"
|
||||
import os from "node:os"
|
||||
import path from "node:path"
|
||||
import { expect, test } from "bun:test"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { makeFiles } from "@opencode-ai/core/environment"
|
||||
import { Workspace } from "@opencode-ai/core/workspace"
|
||||
import { WorkspaceDriver } from "@opencode-ai/core/workspace/driver"
|
||||
import { Effect, Layer } from "effect"
|
||||
import { TestClock } from "effect/testing"
|
||||
import { modalWorkspaceDriver, provider } from "../src/workspace/modal-workspace"
|
||||
|
||||
const enabled =
|
||||
!!process.env.OPENCODE_TEST_MODAL &&
|
||||
((!!process.env.MODAL_TOKEN_ID && !!process.env.MODAL_TOKEN_SECRET) ||
|
||||
fs.existsSync(path.join(os.homedir(), ".modal.toml")))
|
||||
|
||||
const testLayer = Layer.provideMerge(
|
||||
AppNodeBuilder.build(Workspace.configured({ idleThreshold: "1 minute", pollInterval: "1 minute" }), [
|
||||
[
|
||||
WorkspaceDriver.node,
|
||||
WorkspaceDriver.registryNode({ [provider]: modalWorkspaceDriver({ app: "opencode-workspace-tests" }) }),
|
||||
],
|
||||
]),
|
||||
TestClock.layer(),
|
||||
)
|
||||
const modalTest = enabled ? test : test.skip
|
||||
|
||||
modalTest(
|
||||
"wakes a workspace from its filesystem snapshot",
|
||||
() =>
|
||||
Effect.runPromise(
|
||||
Effect.gen(function* () {
|
||||
const workspace = yield* Workspace.Service
|
||||
yield* Effect.acquireUseRelease(
|
||||
workspace.create(provider),
|
||||
(created) =>
|
||||
Effect.gen(function* () {
|
||||
const environment = yield* workspace.connect(created.id)
|
||||
const files = makeFiles(environment)
|
||||
const file = `/tmp/opencode-workspace-${crypto.randomUUID()}.txt`
|
||||
yield* files.write(file, new TextEncoder().encode("survived snapshot"))
|
||||
|
||||
yield* TestClock.adjust("2 minutes")
|
||||
|
||||
const restored = yield* files.read(file)
|
||||
expect(new TextDecoder().decode(restored.bytes)).toBe("survived snapshot")
|
||||
}),
|
||||
(created) => workspace.destroy(created.id).pipe(Effect.ignore),
|
||||
)
|
||||
}).pipe(Effect.scoped, Effect.provide(testLayer)),
|
||||
),
|
||||
180_000,
|
||||
)
|
||||
Reference in New Issue
Block a user