Compare commits

...
20 changed files with 815 additions and 236 deletions
+60 -105
View File
@@ -1,15 +1,11 @@
{
"version": "7",
"dialect": "sqlite",
"id": "2d214a71-3b0a-48c1-a667-741952c4e188",
"id": "6ff49c08-7759-48fd-beca-6086853fce79",
"prevIds": [
"f14a9b18-8207-487e-a3d3-227e629ba9ad"
"2d214a71-3b0a-48c1-a667-741952c4e188"
],
"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": false,
"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
View File
@@ -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,
\`created_at\` integer NOT NULL,
\`last_used_at\` integer NOT NULL
);
`)
})
},
}
export default migration
+9 -13
View File
@@ -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,
\`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(
+13 -2
View File
@@ -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,19 @@ 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(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"
-8
View File
@@ -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) =>
+204 -1
View File
@@ -1,6 +1,209 @@
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 { make as makeSpawner } 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: Schema.NullOr(WorkspaceDriver.Binding),
createdAt: Schema.Number,
lastUsedAt: Schema.Number,
}) {}
export class NotFound extends Schema.TaggedErrorClass<NotFound>()("Workspace.NotFound", { workspaceID: ID }) {}
export class BindingNotFound extends Schema.TaggedErrorClass<BindingNotFound>()("Workspace.BindingNotFound", {
workspaceID: ID,
}) {}
export interface Interface {
readonly create: (provider: string) => Effect.Effect<Info, WorkspaceDriver.Error | WorkspaceDriver.ProviderNotFound>
readonly connect: (
workspaceID: ID,
) => Effect.Effect<
EnvironmentDriver,
NotFound | BindingNotFound | WorkspaceDriver.Error | WorkspaceDriver.ProviderNotFound,
Scope.Scope
>
readonly destroy: (
workspaceID: ID,
) => Effect.Effect<void, NotFound | BindingNotFound | 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 binding: Ref.Ref<WorkspaceDriver.Binding>
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)
if (!row.binding) return yield* new BindingNotFound({ workspaceID })
const driver = yield* registry.get(row.provider)
const binding = yield* Ref.make(row.binding)
const saveBinding = (value: WorkspaceDriver.Binding) =>
db
.update(WorkspaceTable)
.set({ binding: value })
.where(eq(WorkspaceTable.id, workspaceID))
.run()
.pipe(Effect.orDie, Effect.andThen(Ref.set(binding, value)))
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 connection: Connection = {
driver,
environment,
binding,
saveBinding,
lastActivity: yield* Ref.make(yield* Clock.currentTimeMillis),
active: yield* Ref.make(0),
scope,
}
connections.set(workspaceID, connection)
yield* db
.update(WorkspaceTable)
.set({ last_used_at: yield* Clock.currentTimeMillis })
.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
if (now - (yield* Ref.get(connection.lastActivity)) < idleThreshold) return
yield* connection.driver.suspendForIdle({
binding: yield* Ref.get(connection.binding),
saveBinding: connection.saveBinding,
})
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) {
yield* Scope.Scope
const initial = yield* locks.withLock(workspaceID)(open(workspaceID))
const spawner = makeSpawner((command) =>
Effect.acquireRelease(
locks.withLock(workspaceID)(
Effect.gen(function* () {
const connection = yield* open(workspaceID).pipe(Effect.orDie)
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))),
)
return { spawner, overrides: initial.environment.overrides }
}),
destroy: Effect.fn("Workspace.destroy")(function* (workspaceID) {
yield* locks.withLock(workspaceID)(
Effect.gen(function* () {
const row = yield* load(workspaceID)
if (!row.binding) return yield* new BindingNotFound({ 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({ binding: connection ? yield* Ref.get(connection.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.
+63
View File
@@ -0,0 +1,63 @@
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 binding: Binding
readonly saveBinding: (binding: Binding) => Effect.Effect<void>
}) => Effect.Effect<void, Error>
readonly destroy: (input: { 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 node = makeGlobalNode({
service: RegistryService,
layer: Layer.succeed(RegistryService, RegistryService.of(registry({}))),
deps: [],
})
+11
View File
@@ -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>(),
created_at: integer().notNull(),
last_used_at: integer().notNull(),
})
+18 -6
View File
@@ -13,7 +13,7 @@ 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, environmentLayer?: Layer.Layer<Environment.Service>) {
const activeLocation = Layer.succeed(
Location.Service,
Location.Service.of(location({ directory: AbsolutePath.make(directory) })),
@@ -21,7 +21,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, environmentLayer ?? AppNodeBuilder.build(Environment.node, [[Location.node, activeLocation]])],
]),
)
}
@@ -118,7 +118,7 @@ describe("FileMutation", () => {
const releaseFirst = yield* Deferred.make<void>()
const secondStarted = yield* Deferred.make<void>()
let writes = 0
const filesystem = instrumentWrites((write) =>
const filesystem = instrumentWrites(directory, (write) =>
Effect.gen(function* () {
writes++
if (writes === 1) {
@@ -211,7 +211,7 @@ describe("FileMutation", () => {
const secondFinished = yield* Deferred.make<void>()
const secondPath = path.join(directory, "second.txt")
let writes = 0
const filesystem = instrumentWrites((write) =>
const filesystem = instrumentWrites(directory, (write) =>
++writes === 1
? Deferred.succeed(firstStarted, undefined).pipe(
Effect.andThen(Deferred.await(releaseFirst)),
@@ -240,7 +240,10 @@ describe("FileMutation", () => {
)
})
function instrumentWrites(run: <E>(write: Effect.Effect<void, E>, target: string) => Effect.Effect<void, E>) {
function instrumentWrites(
directory: string,
run: <E>(write: Effect.Effect<void, E>, target: string) => Effect.Effect<void, E>,
) {
return Layer.effect(
Environment.Service,
Effect.gen(function* () {
@@ -253,5 +256,14 @@ function instrumentWrites(run: <E>(write: Effect.Effect<void, E>, target: string
},
})
}),
).pipe(Layer.provide(LayerNode.compile(Environment.node)))
).pipe(
Layer.provide(
AppNodeBuilder.build(Environment.node, [
[
Location.node,
Layer.succeed(Location.Service, Location.Service.of(location({ directory: AbsolutePath.make(directory) }))),
],
]),
),
)
}
+7 -2
View File
@@ -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,
+25 -22
View File
@@ -80,28 +80,31 @@ 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)))),
const environment = (activeLocation: Layer.Layer<Location.Service>) =>
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)))
write: (target, content) =>
Effect.sync(() => writes.push(target)).pipe(Effect.andThen(current.files.write(target, content))),
},
})
}),
).pipe(Layer.provide(AppNodeBuilder.build(Environment.node, [[Location.node, activeLocation]])))
const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) => Effect.Effect<A, E, R>) => {
const activeLocation = Layer.succeed(
@@ -115,7 +118,7 @@ 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, environment(activeLocation)],
[Location.node, activeLocation],
[Formatter.node, formatter],
[Permission.node, permission],
+29 -27
View File
@@ -82,33 +82,35 @@ 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)
const environment = (activeLocation: Layer.Layer<Location.Service>) =>
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)
},
},
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)))
})
}),
).pipe(Layer.provide(AppNodeBuilder.build(Environment.node, [[Location.node, activeLocation]])))
const withTool = <A, E, R>(
directory: string,
@@ -126,7 +128,7 @@ const withTool = <A, E, R>(
}).pipe(
Effect.provide(
AppNodeBuilder.build(LayerNode.group([Tool.node, FileMutation.node, patchToolNode]), [
[Environment.node, environment],
[Environment.node, environment(activeLocation)],
[Location.node, activeLocation],
[Formatter.node, formatter],
[Permission.node, permission],
+16 -15
View File
@@ -68,20 +68,21 @@ const reset = () => {
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 environment = (activeLocation: Layer.Layer<Location.Service>) =>
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(AppNodeBuilder.build(Environment.node, [[Location.node, activeLocation]])))
const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) => Effect.Effect<A, E, R>) => {
const activeLocation = Layer.succeed(
@@ -95,7 +96,7 @@ const withTool = <A, E, R>(directory: string, body: (registry: Tool.Interface) =
AppNodeBuilder.build(
LayerNode.group([Tool.node, Tool.node, LocationMutation.node, FileMutation.node, writeToolNode]),
[
[Environment.node, environment],
[Environment.node, environment(activeLocation)],
[Location.node, activeLocation],
[Formatter.node, formatter],
[Permission.node, permission],
+99
View File
@@ -0,0 +1,99 @@
import { beforeEach, expect } from "bun:test"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
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 { Database } from "@opencode-ai/core/database/database"
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
import { eq } from "drizzle-orm"
import { Effect, Layer } 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()
const driver = WorkspaceDriver.make({
create: ({ workspaceID }) => {
calls.push({ operation: "create" })
return Effect.succeed({ binding: { workspaceID, generation: 0 } })
},
connect: ({ binding }) => {
calls.push({ operation: "connect", binding })
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 registryNode = makeGlobalNode({
service: WorkspaceDriver.RegistryService,
layer: Layer.succeed(
WorkspaceDriver.RegistryService,
WorkspaceDriver.RegistryService.of(WorkspaceDriver.registry({ fake: driver })),
),
deps: [],
})
const it = testEffect(
AppNodeBuilder.build(
LayerNode.group([Database.node, Workspace.configured({ idleThreshold: "5 minutes", pollInterval: "1 minute" })]),
[[WorkspaceDriver.node, registryNode]],
),
)
beforeEach(() => calls.splice(0))
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", "connect"])
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 })
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("roundtrips nullable bindings through the workspace table", () =>
Effect.gen(function* () {
const id = Workspace.ID.create()
yield* Database.Service.use(({ db }) =>
db.insert(WorkspaceTable).values({ id, provider: "fake", binding: null, created_at: 10, last_used_at: 20 }).run(),
).pipe(Effect.orDie)
const row = yield* Database.Service.use(({ db }) => db.select().from(WorkspaceTable).get()).pipe(Effect.orDie)
expect(row).toEqual({ id, provider: "fake", binding: null, created_at: 10, last_used_at: 20 })
}),
)
+3
View File
@@ -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 { modalWorkspaceRegistryNode } from "./workspace/modal-workspace"
const applicationServices = LayerNode.group([
Database.node,
@@ -115,6 +117,7 @@ function makeRoutes<AuthError, AuthServices>(
],
[PluginRuntime.node, PluginRuntime.layerWithCell(pluginRuntimeCell)],
[PluginRuntime.providerNode, PluginRuntime.providerNodeWithCell(pluginRuntimeCell)],
[WorkspaceDriver.node, modalWorkspaceRegistryNode({ app: "opencode-workspaces" })],
]
const serviceLayer = options.simulation
? Layer.unwrap(
@@ -0,0 +1,134 @@
import { WorkspaceDriver } from "@opencode-ai/core/workspace/driver"
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
import { Effect, Layer } from "effect"
import type { Image, ModalClient, ModalClientParams, Sandbox } from "modal"
import { createModalSandboxWithClient, makeModalDriver, type ModalImageSpec } from "./modal"
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 openClient = Effect.tryPromise({
try: async () => {
const { ModalClient } = await import("modal")
return new ModalClient(options.client)
},
catch: (cause) => new WorkspaceDriver.Error({ cause }),
})
const attempt = <A>(run: () => Promise<A>) =>
Effect.tryPromise({ try: run, catch: (cause) => new WorkspaceDriver.Error({ cause }) })
const useClient = <A, E>(run: (client: ModalClient) => Effect.Effect<A, E>) =>
Effect.acquireUseRelease(openClient, run, (client) => Effect.sync(() => client.close()))
const findLive = async (client: ModalClient, binding: WorkspaceDriver.Binding, workspaceID?: string) => {
const { NotFoundError } = await import("modal")
if (typeof binding.sandboxId === "string") {
const sandbox = await client.sandboxes.fromId(binding.sandboxId).catch((error) => {
if (error instanceof NotFoundError) return undefined
throw error
})
if (sandbox && (await sandbox.poll()) === null) return sandbox
}
if (!workspaceID || typeof binding.snapshotImageId === "string") return
const sandbox = await client.sandboxes.fromName(options.app, name(workspaceID)).catch((error) => {
if (error instanceof NotFoundError) return undefined
throw error
})
if (sandbox && (await sandbox.poll()) === null) return sandbox
}
const createSandbox = async (client: ModalClient, workspaceID: string, image?: Image) => {
const { AlreadyExistsError } = await import("modal")
return createModalSandboxWithClient(
client,
{
app: options.app,
client: options.client,
image: options.image,
sandbox: {
name: name(workspaceID),
tags: { workspace: workspaceID },
timeoutMs: 24 * 60 * 60 * 1000,
},
},
image,
).catch((error) => {
if (error instanceof AlreadyExistsError) return client.sandboxes.fromName(options.app, name(workspaceID))
throw error
})
}
const deleteImage = (client: ModalClient, imageID: unknown) =>
typeof imageID === "string" ? attempt(() => client.images.delete(imageID)).pipe(Effect.ignore) : Effect.void
const terminate = (sandbox: Sandbox | undefined) =>
sandbox ? attempt(() => sandbox.terminate({ wait: true })).pipe(Effect.ignore) : Effect.void
return WorkspaceDriver.make({
create: ({ workspaceID }) =>
useClient((client) =>
attempt(async () => {
const sandbox = await createSandbox(client, workspaceID)
return { binding: { sandboxId: sandbox.sandboxId } }
}),
),
connect: ({ workspaceID, binding, saveBinding }) =>
Effect.acquireRelease(openClient, (client) => Effect.sync(() => client.close())).pipe(
Effect.flatMap((client) =>
Effect.gen(function* () {
const sandbox = yield* attempt(async () => {
const existing = await findLive(client, binding, workspaceID)
const image =
existing || typeof binding.snapshotImageId !== "string"
? undefined
: await client.images.fromId(binding.snapshotImageId)
return existing ?? createSandbox(client, workspaceID, image)
})
if (binding.sandboxId !== sandbox.sandboxId) {
yield* saveBinding({
sandboxId: sandbox.sandboxId,
...(typeof binding.snapshotImageId === "string" ? { snapshotImageId: binding.snapshotImageId } : {}),
})
}
return makeModalDriver(sandbox)
}),
),
),
suspendForIdle: ({ binding, saveBinding }) =>
useClient((client) =>
Effect.gen(function* () {
const sandbox = yield* attempt(() => findLive(client, binding))
if (!sandbox) return
const snapshot = yield* attempt(() => sandbox.snapshotFilesystem({ ttlMs: null }))
yield* saveBinding({ snapshotImageId: snapshot.imageId })
yield* deleteImage(client, binding.snapshotImageId)
yield* terminate(sandbox)
}),
),
destroy: ({ binding }) =>
useClient((client) =>
Effect.gen(function* () {
const sandbox = yield* attempt(() => findLive(client, binding))
yield* terminate(sandbox)
yield* deleteImage(client, binding.snapshotImageId)
}),
),
})
}
export const modalWorkspaceRegistryNode = (options: ModalWorkspaceOptions) =>
makeGlobalNode({
service: WorkspaceDriver.RegistryService,
layer: Layer.succeed(
WorkspaceDriver.RegistryService,
WorkspaceDriver.RegistryService.of(WorkspaceDriver.registry({ modal: modalWorkspaceDriver(options) })),
),
deps: [],
})
+50 -15
View File
@@ -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 { Image, ModalClient, ModalClientParams, Sandbox, SandboxCreateParams } from "modal"
const INNER_WRAPPER = `
pidfile=$1
@@ -13,15 +13,30 @@ trap 'rm -f -- "$pidfile"' EXIT
"$@"
`
// Modal's VM runtime accepts process-group signals without delivering them
// (kill(-pgid) returns 0 and nothing dies; direct-pid signals work), so the
// group is enumerated from /proc and each member is signalled directly. The
// second pass catches children forked between scan and signal.
const KILL = `
pidfile=$1
sig=$2
i=0
while [ ! -s "$1" ] && [ "$i" -lt 250 ]; do sleep 0.02; i=$((i + 1)); done
if [ -s "$1" ]; then
pid=$(cat "$1")
/bin/kill "-$2" "-$pid" 2>/dev/null || true
else
exit 47
fi
while [ ! -s "$pidfile" ] && [ "$i" -lt 250 ]; do sleep 0.02; i=$((i + 1)); done
[ -s "$pidfile" ] || exit 47
target=$(cat "$pidfile")
pass=0
while [ "$pass" -lt 2 ]; do
for stat in /proc/[0-9]*/stat; do
[ -e "$stat" ] || continue
pid=\${stat#/proc/}
pid=\${pid%/stat}
set -- $(sed "s/.*) //" "$stat" 2>/dev/null)
if [ "\${3:-}" = "$target" ]; then
/bin/kill "-$sig" "$pid" 2>/dev/null || true
fi
done
pass=$((pass + 1))
done
`
export interface ModalImageSpec {
@@ -51,10 +66,7 @@ export const ubuntuImage: ModalImageSpec = {
export const createModalSandbox = async (options: ModalSandboxOptions) => {
const { ModalClient } = await import("modal")
const client = new ModalClient(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])
const sandbox = await client.sandboxes.create(app, image, options.sandbox)
const sandbox = await createModalSandboxWithClient(client, options)
return {
driver: makeModalDriver(sandbox),
sandbox,
@@ -62,14 +74,37 @@ export const createModalSandbox = async (options: ModalSandboxOptions) => {
}
}
export const createModalSandboxWithClient = async (
client: ModalClient,
options: ModalSandboxOptions,
existingImage?: Image,
) => {
const app = await client.apps.fromName(options.app, { createIfMissing: true })
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: Modal exec and filesystem tools share the same roughly 175ms floor,
* so the derived exec defaults are the simplest implementation with no measured loss.
* overrides: exec latency dominates payload work (VM runtime floor measured
* ~285-535ms per exec, Aug 2026), so the derived exec defaults are the simplest
* implementation with no measured loss.
*
* Modal cannot signal a ContainerProcess. Each command therefore starts a new
* process group and records its leader in a unique pid file; kill runs a second
* sandbox command that signals that group. Pid files are removed best-effort.
* sandbox command that enumerates that group from /proc and signals each member
* directly (see KILL). Pid files are removed best-effort.
*/
export const makeModalDriver = (sandbox: Sandbox): Driver => {
const spawn = Effect.fnUntraced(function* (command: Command) {
@@ -0,0 +1,51 @@
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 { modalWorkspaceRegistryNode } 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, modalWorkspaceRegistryNode({ 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("modal"),
(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,
)