import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test" import { $ } from "bun" import fs from "node:fs/promises" import Http from "node:http" import path from "node:path" import { NodeHttpServer } from "@effect/platform-node" import { Effect, Exit, Fiber, Layer, Schema } from "effect" import { HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http" import { eq } from "drizzle-orm" import { GlobalBus, type GlobalEvent } from "@/bus/global" import { Database } from "@opencode-ai/core/database/database" import { ProjectV2 } from "@opencode-ai/core/project" import { ProjectTable } from "@opencode-ai/core/project/sql" import { AbsolutePath } from "@opencode-ai/core/schema" import { Session as SessionNs } from "@/session/session" import { SessionID } from "@/session/schema" import { SessionTable } from "@opencode-ai/core/session/sql" import { SessionProjector } from "@opencode-ai/core/session/projector" import { EventSequenceTable } from "@opencode-ai/core/event/sql" import { resetDatabase } from "../fixture/db" import { disposeAllInstances, provideTmpdirInstance, requireInstance, TestInstance } from "../fixture/fixture" import { testEffect } from "../lib/effect" import { registerAdapter } from "../../src/control-plane/adapters" import { WorkspaceV2 } from "@opencode-ai/core/workspace" import { WorkspaceTable } from "@opencode-ai/core/control-plane/workspace.sql" import type { Target, WorkspaceAdapter, WorkspaceInfo } from "../../src/control-plane/types" import * as Workspace from "../../src/control-plane/workspace" import { InstanceStore } from "@/project/instance-store" import { InstanceBootstrap } from "@/project/bootstrap" import { RuntimeFlags } from "@/effect/runtime-flags" import { Ripgrep } from "@opencode-ai/core/ripgrep" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" import { LayerNode } from "@opencode-ai/core/effect/layer-node" const originalEnv = { OPENCODE_AUTH_CONTENT: process.env.OPENCODE_AUTH_CONTENT, OPENCODE_EXPERIMENTAL_WORKSPACES: process.env.OPENCODE_EXPERIMENTAL_WORKSPACES, OTEL_EXPORTER_OTLP_HEADERS: process.env.OTEL_EXPORTER_OTLP_HEADERS, OTEL_EXPORTER_OTLP_ENDPOINT: process.env.OTEL_EXPORTER_OTLP_ENDPOINT, OTEL_RESOURCE_ATTRIBUTES: process.env.OTEL_RESOURCE_ATTRIBUTES, } const workspaceLayer = (experimentalWorkspaces: boolean) => AppNodeBuilder.build( LayerNode.group([ Workspace.node, SessionNs.node, SessionProjector.node, Database.node, InstanceStore.node, Ripgrep.node, ]), [ [RuntimeFlags.node, RuntimeFlags.layer({ experimentalWorkspaces })], [ InstanceStore.bootstrapNode, Layer.succeed(InstanceBootstrap.Service, InstanceBootstrap.Service.of({ run: Effect.void })), ], ], ) const testServerLayer = Layer.mergeAll( NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }), workspaceLayer(true), ) const it = testEffect(testServerLayer) type RecordedCreate = { info: WorkspaceInfo env: Record from?: WorkspaceInfo } type RecordedAdapter = { adapter: WorkspaceAdapter calls: { configure: WorkspaceInfo[] create: RecordedCreate[] list: number remove: WorkspaceInfo[] target: WorkspaceInfo[] } } type FetchCall = { url: URL method: string headers: Headers bodyText?: string json?: unknown } function unique(prefix: string) { return `${prefix}-${Math.random().toString(36).slice(2)}` } function restoreEnv() { Object.entries(originalEnv).forEach(([key, value]) => { if (value === undefined) { delete process.env[key] return } process.env[key] = value }) } beforeEach(() => { restoreEnv() process.env.OPENCODE_EXPERIMENTAL_WORKSPACES = "true" }) afterEach(async () => { mock.restore() await disposeAllInstances() restoreEnv() await resetDatabase() }) async function initGitRepo(dir: string) { await fs.mkdir(dir, { recursive: true }) await $`git init`.cwd(dir).quiet() await $`git config core.fsmonitor false`.cwd(dir).quiet() await $`git config commit.gpgsign false`.cwd(dir).quiet() await $`git config user.email "test@opencode.test"`.cwd(dir).quiet() await $`git config user.name "Test"`.cwd(dir).quiet() await fs.writeFile(path.join(dir, "tracked.txt"), "base\n") await $`git add tracked.txt`.cwd(dir).quiet() await $`git commit -m "base"`.cwd(dir).quiet() } const startWorkspaceSyncingWithFlag = (projectID: ProjectV2.ID, experimentalWorkspaces: boolean) => Effect.runPromise( Workspace.use.startWorkspaceSyncing(projectID).pipe(Effect.provide(workspaceLayer(experimentalWorkspaces))), ) function captureGlobalEvents() { const events: GlobalEvent[] = [] const handler = (event: GlobalEvent) => events.push(event) GlobalBus.on("event", handler) return { events, dispose() { GlobalBus.off("event", handler) }, } } function expectExitContains(exit: Exit.Exit, ...messages: string[]) { expect(Exit.isFailure(exit)).toBe(true) if (!Exit.isFailure(exit)) return for (const message of messages) expect(String(exit.cause)).toContain(message) } function eventuallyEffect(effect: Effect.Effect, timeout = 1500) { return Effect.gen(function* () { const started = Date.now() let last: unknown while (Date.now() - started < timeout) { const exit = yield* Effect.exit(effect) if (exit._tag === "Success") return last = exit.cause yield* Effect.sleep("10 millis") } throw last ?? new Error("Timed out waiting for condition") }) } function recordedAdapter(input: { target: (info: WorkspaceInfo) => Target | Promise configure?: (info: WorkspaceInfo) => WorkspaceInfo | Promise create?: (info: WorkspaceInfo, env: Record, from?: WorkspaceInfo) => Promise list?: () => Omit[] | Promise[]> remove?: (info: WorkspaceInfo) => Promise }): RecordedAdapter { const calls: RecordedAdapter["calls"] = { configure: [], create: [], list: 0, remove: [], target: [], } return { calls, adapter: { name: "recorded", description: "recorded", configure(info) { calls.configure.push(structuredClone(info)) return input.configure?.(info) ?? info }, async create(info, env, from) { calls.create.push({ info: structuredClone(info), env: { ...env }, from: from ? structuredClone(from) : undefined, }) await input.create?.(info, env, from) }, ...(input.list ? { async list() { calls.list += 1 return input.list?.() ?? [] }, } : {}), async remove(info) { calls.remove.push(structuredClone(info)) await input.remove?.(info) }, target(info) { calls.target.push(structuredClone(info)) return input.target(info) }, }, } } function localAdapter(dir: string, input?: { createDir?: boolean; remove?: (info: WorkspaceInfo) => Promise }) { return recordedAdapter({ configure(info) { return { ...info, directory: dir } }, async create() { if (input?.createDir === false) return await fs.mkdir(dir, { recursive: true }) }, remove: input?.remove, target() { return { type: "local", directory: dir } }, }) } function remoteAdapter(url: string, input?: { directory?: string | null; headers?: HeadersInit }) { return recordedAdapter({ configure(info) { return { ...info, directory: input?.directory ?? info.directory } }, target() { return { type: "remote", url, headers: input?.headers } }, }) } function eventStreamResponse(events: unknown[] = [], keepOpen = true) { const encoder = new TextEncoder() return new Response( new ReadableStream({ start(controller) { if (keepOpen) controller.enqueue(encoder.encode(":\n\n")) events.forEach((event) => controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`))) if (!keepOpen) controller.close() }, }), { status: 200, headers: { "content-type": "text/event-stream" } }, ) } function serverUrl() { return Effect.gen(function* () { return HttpServer.formatAddress((yield* HttpServer.HttpServer).address) }) } function workspaceInfo(projectID: ProjectV2.ID, type: string, input?: Partial): Workspace.Info { return { id: input?.id ?? WorkspaceV2.ID.ascending(), type, name: input?.name ?? unique("workspace"), branch: input?.branch ?? null, directory: input?.directory ?? null, extra: input?.extra ?? null, projectID, timeUsed: input?.timeUsed ?? Date.now(), } } function insertWorkspace(info: Workspace.Info) { return Database.Service.use(({ db }) => db .insert(WorkspaceTable) .values({ id: info.id, type: info.type, branch: info.branch, name: info.name, directory: info.directory, extra: info.extra, project_id: info.projectID, time_used: info.timeUsed, }) .run() .pipe(Effect.orDie), ) } function insertProject(id: ProjectV2.ID, worktree: string) { return Database.Service.use(({ db }) => db .insert(ProjectTable) .values({ id, worktree: AbsolutePath.make(worktree), vcs: null, name: null, time_created: Date.now(), time_updated: Date.now(), sandboxes: [], }) .run() .pipe(Effect.orDie), ) } function attachSessionToWorkspace(sessionID: SessionID, workspaceID: WorkspaceV2.ID) { return Database.Service.use(({ db }) => db .update(SessionTable) .set({ workspace_id: workspaceID }) .where(eq(SessionTable.id, sessionID)) .run() .pipe(Effect.orDie), ) } function sessionSequence(sessionID: SessionID) { return Database.Service.use(({ db }) => db .select({ seq: EventSequenceTable.seq }) .from(EventSequenceTable) .where(eq(EventSequenceTable.aggregate_id, sessionID)) .get() .pipe( Effect.orDie, Effect.map((row) => row?.seq), ), ) } function sessionSequenceOwner(sessionID: SessionID) { return Database.Service.use(({ db }) => db .select({ ownerID: EventSequenceTable.owner_id }) .from(EventSequenceTable) .where(eq(EventSequenceTable.aggregate_id, sessionID)) .get() .pipe( Effect.orDie, Effect.map((row) => row?.ownerID), ), ) } describe("workspace schemas and exports", () => { test("keeps the historical event type names", () => { expect(Workspace.Event.Ready.type).toBe("workspace.ready") expect(Workspace.Event.Failed.type).toBe("workspace.failed") expect(Workspace.Event.Status.type).toBe("workspace.status") }) test("validates create input with workspace id, project id, branch, type, and extra", () => { const input = { id: WorkspaceV2.ID.ascending("wrk_schema_create"), type: "worktree", branch: "feature/schema", projectID: ProjectV2.ID.make("project-schema"), extra: { nested: true }, } const decode = Schema.decodeUnknownSync(Workspace.CreateInput) expect(decode(input)).toEqual(input) expect(() => decode({ ...input, id: 1 })).toThrow() expect(() => decode({ ...input, branch: 1 })).toThrow() }) }) describe("workspace CRUD", () => { it.instance( "get returns undefined for a missing workspace", () => Effect.gen(function* () { const workspace = yield* Workspace.Service expect(yield* workspace.get(WorkspaceV2.ID.ascending("wrk_missing_get"))).toBeUndefined() }), { git: true }, ) it.instance( "list maps database rows, filters by project, and sorts by id", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const otherProjectID = ProjectV2.ID.make("project-other") yield* insertProject(otherProjectID, "/tmp/other") const a = workspaceInfo(instance.project.id, "manual", { id: WorkspaceV2.ID.ascending("wrk_a_list"), branch: "a", directory: "/a", extra: { a: true }, }) const b = workspaceInfo(instance.project.id, "manual", { id: WorkspaceV2.ID.ascending("wrk_b_list"), branch: "b", directory: "/b", extra: ["b"], }) const other = workspaceInfo(otherProjectID, "manual", { id: WorkspaceV2.ID.ascending("wrk_c_list") }) yield* insertWorkspace(b) yield* insertWorkspace(other) yield* insertWorkspace(a) expect(yield* workspace.list(instance.project)).toEqual([a, b]) }), { git: true }, ) it.instance( "create configures, persists, creates, starts local sync, and passes environment", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service process.env.OPENCODE_AUTH_CONTENT = JSON.stringify({ test: { type: "api", key: "secret" } }) process.env.OTEL_EXPORTER_OTLP_HEADERS = "authorization=otel" process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "https://otel.test" process.env.OTEL_RESOURCE_ATTRIBUTES = "service.name=opencode-test" const workspaceID = WorkspaceV2.ID.ascending("wrk_create_local") const type = unique("create-local") const targetDir = path.join(instance.directory, "created-local") const recorded = recordedAdapter({ configure(info) { return { ...info, branch: "configured-branch", name: "Configured Name", directory: targetDir, extra: { configured: true }, } }, async create() { await fs.mkdir(targetDir, { recursive: true }) }, target() { return { type: "local", directory: targetDir } }, }) registerAdapter(instance.project.id, type, recorded.adapter) const info = yield* workspace.create({ id: workspaceID, type, branch: null, projectID: instance.project.id, extra: null, }) expect(info).toEqual({ id: workspaceID, type, branch: "configured-branch", name: "Configured Name", directory: targetDir, extra: { configured: true }, projectID: instance.project.id, timeUsed: info.timeUsed, }) expect(yield* workspace.get(workspaceID)).toEqual(info) expect(yield* workspace.list(instance.project)).toEqual([info]) expect(recorded.calls.configure).toHaveLength(1) expect(recorded.calls.configure[0]).toMatchObject({ id: workspaceID, type, directory: null }) expect(recorded.calls.create).toHaveLength(1) expect(recorded.calls.create[0].info).toEqual({ id: workspaceID, type, branch: "configured-branch", name: "Configured Name", directory: targetDir, extra: { configured: true }, projectID: instance.project.id, }) expect(JSON.parse(recorded.calls.create[0].env.OPENCODE_AUTH_CONTENT ?? "{}")).toEqual({ test: { type: "api", key: "secret" }, }) expect(recorded.calls.create[0].env.OPENCODE_WORKSPACE_ID).toBe(workspaceID) expect(recorded.calls.create[0].env.OPENCODE_EXPERIMENTAL_WORKSPACES).toBe("true") expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_HEADERS).toBe("authorization=otel") expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_ENDPOINT).toBe("https://otel.test") expect(recorded.calls.create[0].env.OTEL_RESOURCE_ATTRIBUTES).toBe("service.name=opencode-test") expect((yield* workspace.status()).find((item) => item.workspaceID === workspaceID)?.status).toBe("connected") yield* workspace.remove(workspaceID) expect((yield* workspace.status()).find((item) => item.workspaceID === workspaceID)?.status).toBeUndefined() }), { git: true }, ) it.instance( "create propagates configure failures and does not insert a workspace", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const type = unique("configure-failure") registerAdapter( instance.project.id, type, recordedAdapter({ configure() { throw new Error("configure exploded") }, target() { return { type: "local", directory: "/unused" } }, }).adapter, ) expectExitContains( yield* Effect.exit(workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })), "configure exploded", ) expect(yield* workspace.list(instance.project)).toEqual([]) }), { git: true }, ) it.instance( "create leaves the inserted row when adapter create fails", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const type = unique("create-failure") const recorded = recordedAdapter({ async create() { throw new Error("create exploded") }, target() { return { type: "local", directory: "/unused" } }, }) registerAdapter(instance.project.id, type, recorded.adapter) expectExitContains( yield* Effect.exit( workspace.create({ type, branch: "branch", projectID: instance.project.id, extra: { x: 1 } }), ), "create exploded", ) const rows = yield* workspace.list(instance.project) expect(rows).toHaveLength(1) expect(rows[0]).toMatchObject({ type, branch: "branch", extra: { x: 1 } }) expect(recorded.calls.target).toHaveLength(0) yield* workspace.remove(rows[0].id) }), { git: true }, ) it.instance( "create returns after a local workspace reports error", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const type = unique("local-error") const missing = path.join(instance.directory, "missing-local-target") const recorded = localAdapter(missing, { createDir: false }) registerAdapter(instance.project.id, type, recorded.adapter) const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null }) expect(info.directory).toBe(missing) expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error") yield* workspace.remove(info.id) }), { git: true }, ) it.instance( "syncList registers adapter-listed workspaces that are missing by name", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const type = unique("list-sync") const existing = workspaceInfo(instance.project.id, type, { id: WorkspaceV2.ID.ascending("wrk_list_sync_existing"), name: "existing", directory: path.join(instance.directory, "existing"), }) yield* insertWorkspace(existing) const discovered = { type, name: "discovered", branch: "feature/discovered", directory: path.join(instance.directory, "discovered"), extra: { source: "adapter" }, projectID: instance.project.id, } const recorded = recordedAdapter({ list() { return [ { type, name: existing.name, branch: "ignored", directory: path.join(instance.directory, "ignored"), extra: null, projectID: instance.project.id, }, discovered, ] }, target(info) { return { type: "local", directory: info.directory ?? instance.directory } }, }) registerAdapter(instance.project.id, type, recorded.adapter) yield* workspace.syncList(instance.project) const synced = (yield* workspace.list(instance.project)).filter((item) => item.name === discovered.name) expect(synced).toHaveLength(1) expect(synced[0]).toMatchObject(discovered) expect(synced[0]?.id).toStartWith("wrk_") expect(yield* workspace.list(instance.project)).toEqual(expect.arrayContaining([existing, synced[0]])) expect(recorded.calls.list).toBe(1) expect(recorded.calls.configure).toHaveLength(0) expect(recorded.calls.create).toHaveLength(0) expect(recorded.calls.target).toHaveLength(1) }), { git: true }, ) it.instance( "syncList calls every registered adapter with a list method", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const typeA = unique("list-sync-a") const typeB = unique("list-sync-b") const adapterA = recordedAdapter({ list() { return [ { type: typeA, name: "adapter-a", branch: null, directory: path.join(instance.directory, "adapter-a"), extra: null, projectID: instance.project.id, }, ] }, target(info) { return { type: "local", directory: info.directory ?? instance.directory } }, }) const adapterB = recordedAdapter({ list() { return [ { type: typeB, name: "adapter-b", branch: null, directory: path.join(instance.directory, "adapter-b"), extra: null, projectID: instance.project.id, }, ] }, target(info) { return { type: "local", directory: info.directory ?? instance.directory } }, }) const noList = recordedAdapter({ target() { return { type: "local", directory: instance.directory } }, }) registerAdapter(instance.project.id, typeA, adapterA.adapter) registerAdapter(instance.project.id, typeB, adapterB.adapter) registerAdapter(instance.project.id, unique("list-sync-none"), noList.adapter) yield* workspace.syncList(instance.project) const synced = yield* workspace.list(instance.project) expect( synced .filter((item) => item.type === typeA || item.type === typeB) .map((item) => item.name) .toSorted(), ).toEqual(["adapter-a", "adapter-b"]) expect(adapterA.calls.list).toBe(1) expect(adapterB.calls.list).toBe(1) expect(noList.calls.list).toBe(0) }), { git: true }, ) it.live("remote create connects to routed event and history endpoints", () => { const calls: FetchCall[] = [] return Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const bodyText = yield* req.text const call = { url: new URL(req.url, "http://localhost"), method: req.method, headers: new Headers(req.headers), bodyText, json: bodyText ? JSON.parse(bodyText) : undefined, } calls.push(call) if (call.url.pathname === "/base/global/event") return HttpServerResponse.fromWeb(eventStreamResponse([], false)) if (call.url.pathname === "/base/sync/history") return yield* HttpServerResponse.json([]) return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( (dir) => Effect.gen(function* () { const workspace = yield* Workspace.Service const instance = yield* requireInstance const type = unique("remote-create") const recorded = remoteAdapter(`${url}/base/?ignored=1#hash`, { directory: dir }) registerAdapter(instance.project.id, type, recorded.adapter) const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null }) expect( calls.map((call) => `${call.method} ${call.url.pathname}${call.url.search}${call.url.hash}`), ).toEqual(["GET /base/global/event", "POST /base/sync/history"]) expect(calls[1].json).toEqual({}) expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("connected") expect(yield* workspace.isSyncing(info.id)).toBe(true) yield* workspace.remove(info.id) expect(yield* workspace.isSyncing(info.id)).toBe(false) expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined() }), { git: true }, ) }) }) it.instance( "remove returns undefined for a missing workspace", () => Effect.gen(function* () { const workspace = yield* Workspace.Service expect(yield* workspace.remove(WorkspaceV2.ID.ascending("wrk_missing_remove"))).toBeUndefined() }), { git: true }, ) it.instance( "remove deletes the workspace, associated sessions, adapter resources, and status", () => { return Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const type = unique("remove-local") const recorded = localAdapter(path.join(dir, "remove-local")) registerAdapter(instance.project.id, type, recorded.adapter) const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null }) const one = yield* sessionSvc.create({}) const two = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(one.id, info.id) yield* attachSessionToWorkspace(two.id, info.id) const removed = yield* workspace.remove(info.id) expect(removed).toEqual(info) expect(yield* workspace.get(info.id)).toBeUndefined() expect(recorded.calls.remove).toEqual([info]) expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined() const { db } = yield* Database.Service expect( yield* db .select({ id: SessionTable.id }) .from(SessionTable) .where(eq(SessionTable.workspace_id, info.id)) .all() .pipe(Effect.orDie), ).toEqual([]) }) }, { git: true }, ) it.instance( "remove still deletes the row when the adapter cannot remove resources", () => Effect.gen(function* () { const instance = yield* requireInstance const workspace = yield* Workspace.Service const type = unique("remove-throws") const info = workspaceInfo(instance.project.id, type, { id: WorkspaceV2.ID.ascending("wrk_remove_throws") }) registerAdapter( instance.project.id, type, recordedAdapter({ async remove() { throw new Error("remove exploded") }, target() { return { type: "local", directory: "/unused" } }, }).adapter, ) yield* insertWorkspace(info) expect(yield* workspace.remove(info.id)).toEqual(info) expect(yield* workspace.get(info.id)).toBeUndefined() }), { git: true }, ) it.instance( "sessionWarp moves a session into a local workspace and claims ownership", () => { return Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const previousType = unique("warp-prev-local") const targetType = unique("warp-target-local") const previous = workspaceInfo(instance.project.id, previousType) const target = workspaceInfo(instance.project.id, targetType) yield* insertWorkspace(previous) yield* insertWorkspace(target) registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-prev-local")).adapter) registerAdapter(instance.project.id, targetType, localAdapter(path.join(dir, "warp-target-local")).adapter) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, previous.id) yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id }) const { db } = yield* Database.Service expect( (yield* db .select({ workspaceID: SessionTable.workspace_id }) .from(SessionTable) .where(eq(SessionTable.id, session.id)) .get() .pipe(Effect.orDie))?.workspaceID, ).toBe(target.id) expect(yield* sessionSequenceOwner(session.id)).toBe(target.id) }) }, { git: true }, ) it.instance( "sessionWarp applies source workspace patch to local target workspace", () => { return Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const previousType = unique("warp-patch-prev-local") const targetType = unique("warp-patch-target-local") const previousDir = path.join(dir, "warp-patch-prev-local") const targetDir = path.join(dir, "warp-patch-target-local") yield* Effect.promise(() => initGitRepo(previousDir)) yield* Effect.promise(() => initGitRepo(targetDir)) yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "tracked.txt"), "changed\n")) yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "new.txt"), "new\n")) const previous = workspaceInfo(instance.project.id, previousType) const target = workspaceInfo(instance.project.id, targetType) yield* insertWorkspace(previous) yield* insertWorkspace(target) registerAdapter(instance.project.id, previousType, localAdapter(previousDir, { createDir: false }).adapter) registerAdapter(instance.project.id, targetType, localAdapter(targetDir, { createDir: false }).adapter) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, previous.id) yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true }) expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "tracked.txt"), "utf8"))).toBe("changed\n") expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "new.txt"), "utf8"))).toBe("new\n") }) }, { git: true }, ) it.instance( "sessionWarp detaches a session to the local project and claims project ownership", () => { return Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const previousType = unique("warp-detach-local") const previous = workspaceInfo(instance.project.id, previousType) yield* insertWorkspace(previous) registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-detach-local")).adapter) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, previous.id) yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id }) const { db } = yield* Database.Service expect( (yield* db .select({ workspaceID: SessionTable.workspace_id }) .from(SessionTable) .where(eq(SessionTable.id, session.id)) .get() .pipe(Effect.orDie))?.workspaceID, ).toBeNull() expect(yield* sessionSequenceOwner(session.id)).toBe(instance.project.id) }) }, { git: true }, ) const itCrossInstance = process.platform === "win32" ? it.instance.skip : it.instance itCrossInstance( "sessionWarp detaches to the source project when invoked from a workspace instance", () => Effect.gen(function* () { const instance = yield* requireInstance const projectID = instance.project.id const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const previousType = unique("warp-detach-workspace-instance") const previous = workspaceInfo(projectID, previousType) yield* insertWorkspace(previous) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, previous.id) const workspaceProjectID = yield* provideTmpdirInstance( (workspaceDir) => Effect.gen(function* () { registerAdapter(projectID, previousType, localAdapter(workspaceDir, { createDir: false }).adapter) const workspaceCtx = yield* requireInstance expect(workspaceCtx.project.id).not.toBe(projectID) yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id }) return workspaceCtx.project.id }), { git: true }, ) const { db } = yield* Database.Service expect( (yield* db .select({ workspaceID: SessionTable.workspace_id }) .from(SessionTable) .where(eq(SessionTable.id, session.id)) .get() .pipe(Effect.orDie))?.workspaceID, ).toBeNull() expect(yield* sessionSequenceOwner(session.id)).toBe(projectID) expect(yield* sessionSequenceOwner(session.id)).not.toBe(workspaceProjectID) }), { git: true }, ) it.live("sessionWarp syncs previous remote history, replays it, steals, and claims the sequence", () => { const calls: FetchCall[] = [] let historySessionID: SessionID | undefined let historySession: SessionNs.Info | undefined let historyNextSeq = 0 return Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const bodyText = yield* req.text const call = { url: new URL(req.url, "http://localhost"), method: req.method, headers: new Headers(req.headers), bodyText, json: bodyText ? JSON.parse(bodyText) : undefined, } calls.push(call) if (call.url.pathname === "/warp-source/sync/history") { return yield* HttpServerResponse.json([ { id: `evt_${unique("warp-source-history")}`, aggregate_id: historySessionID!, seq: historyNextSeq, type: "session.updated.1", data: { sessionID: historySessionID!, info: historySession! }, }, ]) } if (call.url.pathname === "/warp-source/vcs/diff/raw") return HttpServerResponse.text("remote patch") if (call.url.pathname === "/warp-target/sync/replay") return yield* HttpServerResponse.json({ sessionID: "ok" }) if (call.url.pathname === "/warp-target/sync/steal") return yield* HttpServerResponse.json({ sessionID: "ok" }) if (call.url.pathname === "/warp-target/vcs/apply") return yield* HttpServerResponse.json({ applied: true }) return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const previousType = unique("warp-remote-source") const targetType = unique("warp-remote-target") const previous = workspaceInfo(instance.project.id, previousType) const target = workspaceInfo(instance.project.id, targetType, { directory: "remote-target-dir" }) yield* insertWorkspace(previous) yield* insertWorkspace(target) registerAdapter(instance.project.id, previousType, remoteAdapter(`${url}/warp-source`).adapter) registerAdapter(instance.project.id, targetType, remoteAdapter(`${url}/warp-target`).adapter) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, previous.id) historySessionID = session.id historySession = { ...session, workspaceID: previous.id, title: "from source history" } historyNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1 yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true }) expect(calls.map((call) => `${call.method} ${call.url.pathname}`)).toEqual([ "POST /warp-source/sync/history", "GET /warp-source/vcs/diff/raw", "POST /warp-target/vcs/apply", "POST /warp-target/sync/replay", "POST /warp-target/sync/steal", ]) expect(calls[0].json).toEqual({ [session.id]: historyNextSeq - 1 }) expect(calls[2].json).toEqual({ patch: "remote patch" }) expect(calls[3].json).toMatchObject({ directory: "remote-target-dir", events: [ { aggregateID: session.id, seq: 0, type: "session.created.1", }, { aggregateID: session.id, seq: historyNextSeq, type: "session.updated.1", }, ], }) expect(calls[4].json).toEqual({ sessionID: session.id }) expect((yield* sessionSvc.get(session.id)).title).toBe("from source history") expect(yield* sessionSequenceOwner(session.id)).toBe(target.id) }), { git: true }, ) }) }) }) describe("workspace sync state", () => { it.instance( "startWorkspaceSyncing is disabled by the experimental workspace flag", () => Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const type = unique("flag-disabled") const info = workspaceInfo(instance.project.id, type) const session = yield* sessionSvc.create({}) yield* attachSessionToWorkspace(session.id, info.id) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, localAdapter(path.join(dir, "flag-disabled")).adapter) yield* Effect.promise(() => startWorkspaceSyncingWithFlag(instance.project.id, false)) yield* Effect.sleep("25 millis") expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined() }), { git: true }, ) it.instance( "startWorkspaceSyncing starts all workspaces", () => Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const projectID = instance.project.id const firstType = unique("first") const secondType = unique("second") const first = workspaceInfo(projectID, firstType) const second = workspaceInfo(projectID, secondType) yield* Effect.promise(() => fs.mkdir(path.join(dir, "first"), { recursive: true })) yield* Effect.promise(() => fs.mkdir(path.join(dir, "second"), { recursive: true })) yield* insertWorkspace(first) yield* insertWorkspace(second) registerAdapter(projectID, firstType, localAdapter(path.join(dir, "first")).adapter) registerAdapter(projectID, secondType, localAdapter(path.join(dir, "second")).adapter) yield* Effect.addFinalizer(() => Effect.all([workspace.remove(first.id), workspace.remove(second.id)], { discard: true }).pipe(Effect.ignore), ) yield* workspace.startWorkspaceSyncing(projectID) yield* eventuallyEffect( Effect.gen(function* () { const status = yield* workspace.status() expect(status.find((item) => item.workspaceID === first.id)?.status).toBe("connected") expect(status.find((item) => item.workspaceID === second.id)?.status).toBe("connected") }), ) }), { git: true }, ) it.instance( "local start reports error when the target directory is missing", () => Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const type = unique("missing-local") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter( instance.project.id, type, localAdapter(path.join(dir, "missing-target"), { createDir: false }).adapter, ) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { const status = yield* workspace.status() expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("error") }), ) expect(yield* workspace.isSyncing(info.id)).toBe(false) yield* workspace.remove(info.id) }), { git: true }, ) it.instance( "duplicate local status updates are suppressed", () => Effect.gen(function* () { const { directory: dir } = yield* TestInstance const instance = yield* requireInstance const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const captured = captureGlobalEvents() yield* Effect.addFinalizer(() => Effect.sync(() => captured.dispose())) const type = unique("dedupe-local") const info = workspaceInfo(instance.project.id, type) const target = path.join(dir, "dedupe-local") yield* Effect.promise(() => fs.mkdir(target, { recursive: true })) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, localAdapter(target).adapter) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { const status = yield* workspace.status() expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("connected") }), ) expect( captured.events.filter( (event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type, ), ).toHaveLength(1) yield* workspace.remove(info.id) }), { git: true }, ) it.live("remote start emits disconnected, connecting, and connected then refuses duplicate listeners", () => { const calls: FetchCall[] = [] return Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const bodyText = yield* req.text const call = { url: new URL(req.url, "http://localhost"), method: req.method, headers: new Headers(req.headers), bodyText, json: bodyText ? JSON.parse(bodyText) : undefined, } calls.push(call) if (call.url.pathname === "/sync/global/event") return HttpServerResponse.fromWeb(eventStreamResponse()) if (call.url.pathname === "/sync/sync/history") return HttpServerResponse.fromWeb(Response.json([])) return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const captured = captureGlobalEvents() try { const type = unique("remote-start") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sync`).adapter) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe( "connected", ) }), ) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* Effect.sleep("25 millis") expect( captured.events .filter((event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type) .map((event) => event.payload.properties.status), ).toEqual(["disconnected", "connecting", "connected"]) expect(calls.filter((call) => call.url.pathname === "/sync/global/event")).toHaveLength(1) expect(calls.filter((call) => call.url.pathname === "/sync/sync/history")).toHaveLength(1) expect(yield* workspace.isSyncing(info.id)).toBe(true) yield* workspace.remove(info.id) expect(yield* workspace.isSyncing(info.id)).toBe(false) } finally { captured.dispose() } }), { git: true }, ) }) }) it.live("remote connection HTTP failures set error and clear syncing", () => Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest if (new URL(req.url, "http://localhost").pathname === "/failed/global/event") return HttpServerResponse.text("nope", { status: 503 }) return HttpServerResponse.fromWeb(Response.json([])) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const type = unique("remote-connect-fail") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/failed`).adapter) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error") }), ) expect(yield* workspace.isSyncing(info.id)).toBe(false) yield* workspace.remove(info.id) }), { git: true }, ) }), ) it.live("remote history HTTP failures set error", () => Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const url = new URL(req.url, "http://localhost") if (url.pathname === "/history-failed/global/event") return HttpServerResponse.fromWeb(eventStreamResponse([], false)) if (url.pathname === "/history-failed/sync/history") return HttpServerResponse.text("history failed", { status: 500 }) return HttpServerResponse.fromWeb(Response.json([])) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const type = unique("remote-history-fail") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/history-failed`).adapter) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error") }), ) expect(yield* workspace.isSyncing(info.id)).toBe(false) yield* workspace.remove(info.id) }), { git: true }, ) }), ) it.live("sync history sends the local sequence fence and replays returned events in workspace context", () => { const historyBodies: unknown[] = [] let historySessionID: SessionID | undefined let historySession: SessionNs.Info | undefined let historyNextSeq = 0 return Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const bodyText = yield* req.text const url = new URL(req.url, "http://localhost") if (url.pathname === "/history/global/event") return HttpServerResponse.fromWeb(eventStreamResponse()) if (url.pathname === "/history/sync/history") { historyBodies.push(bodyText ? JSON.parse(bodyText) : undefined) return HttpServerResponse.fromWeb( Response.json([ { id: `evt_${unique("history")}`, aggregate_id: historySessionID!, seq: historyNextSeq, type: "session.updated.1", data: { sessionID: historySessionID!, info: historySession! }, }, ]), ) } return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const captured = captureGlobalEvents() try { const type = unique("history-replay") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/history`).adapter) const session = yield* sessionSvc.create({ title: "before history" }) yield* attachSessionToWorkspace(session.id, info.id) historySessionID = session.id historySession = { ...session, workspaceID: info.id, title: "from history" } historyNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1 yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from history") }), ) expect(historyBodies).toEqual([{ [session.id]: historyNextSeq - 1 }]) expect( captured.events.some( (event) => event.workspace === info.id && event.payload.type === "session.updated" && event.payload.properties.sessionID === session.id && event.payload.properties.info.title === "from history", ), ).toBe(true) yield* workspace.remove(info.id) } finally { captured.dispose() } }), { git: true }, ) }) }) it.live("SSE forwards non-heartbeat events and ignores heartbeats", () => Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const url = new URL(req.url, "http://localhost") if (url.pathname === "/sse-forward/global/event") return HttpServerResponse.fromWeb( eventStreamResponse( [ { directory: "remote-dir", project: "remote-project", payload: { type: "server.heartbeat" } }, { directory: "remote-dir", project: "remote-project", payload: { type: "custom.remote", properties: { ok: true } }, }, ], false, ), ) if (url.pathname === "/sse-forward/sync/history") return HttpServerResponse.fromWeb(Response.json([])) return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const captured = captureGlobalEvents() try { const type = unique("sse-forward") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sse-forward`).adapter) yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id) yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.sync(() => expect( captured.events.some( (event) => event.workspace === info.id && event.payload.type === "custom.remote", ), ).toBe(true), ), ) expect( captured.events.some( (event) => event.workspace === info.id && event.payload.type === "server.heartbeat", ), ).toBe(false) expect( captured.events.find((event) => event.workspace === info.id && event.payload.type === "custom.remote"), ).toMatchObject({ directory: "remote-dir", project: "remote-project", payload: { properties: { ok: true } }, }) yield* workspace.remove(info.id) } finally { captured.dispose() } }), { git: true }, ) }), ) it.live("SSE sync events are replayed and forwarded", () => { let sseSessionID: SessionID | undefined let sseSession: SessionNs.Info | undefined let sseNextSeq = 0 return Effect.gen(function* () { yield* HttpServer.serveEffect()( Effect.gen(function* () { const req = yield* HttpServerRequest.HttpServerRequest const url = new URL(req.url, "http://localhost") if (url.pathname === "/sse-sync/global/event") return HttpServerResponse.fromWeb( eventStreamResponse( [ { directory: "remote-dir", project: "remote-project", payload: { type: "sync", syncEvent: { id: `evt_${unique("sse")}`, aggregateID: sseSessionID!, seq: sseNextSeq, type: "session.updated.1", data: { sessionID: sseSessionID!, info: sseSession! }, }, }, }, ], false, ), ) if (url.pathname === "/sse-sync/sync/history") return HttpServerResponse.fromWeb(Response.json([])) return HttpServerResponse.text("unexpected", { status: 500 }) }), ) const url = yield* serverUrl() yield* provideTmpdirInstance( () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionSvc = yield* SessionNs.Service const instance = yield* requireInstance const captured = captureGlobalEvents() try { const type = unique("sse-sync") const info = workspaceInfo(instance.project.id, type) yield* insertWorkspace(info) registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sse-sync`).adapter) const session = yield* sessionSvc.create({ title: "before sse" }) yield* attachSessionToWorkspace(session.id, info.id) sseSessionID = session.id sseSession = { ...session, workspaceID: info.id, title: "from sse" } sseNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1 yield* workspace.startWorkspaceSyncing(instance.project.id) yield* eventuallyEffect( Effect.gen(function* () { expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from sse") }), ) expect( captured.events.some( (event) => event.workspace === info.id && event.payload.type === "sync" && event.payload.syncEvent.seq === sseNextSeq, ), ).toBe(true) yield* workspace.remove(info.id) } finally { captured.dispose() } }), { git: true }, ) }) }) }) describe("workspace waitForSync", () => { it.instance( "returns immediately for an empty fence", () => Effect.gen(function* () { const workspace = yield* Workspace.Service expect(yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_empty"), {})).toBeUndefined() }), { git: true }, ) it.instance( "returns immediately when the stored sequence already satisfies the fence", () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionID = SessionID.descending("ses_wait_done") const { db } = yield* Database.Service yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 4 }).run().pipe(Effect.orDie) expect( yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_done"), { [sessionID]: 4 }), ).toBeUndefined() expect( yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_done_2"), { [sessionID]: 3 }), ).toBeUndefined() }), { git: true }, ) it.instance( "waits until the database reaches the requested sequence and a workspace event arrives", () => Effect.gen(function* () { const workspace = yield* Workspace.Service const workspaceID = WorkspaceV2.ID.ascending("wrk_wait_event") const sessionID = SessionID.descending("ses_wait_event") const { db } = yield* Database.Service yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 1 }).run().pipe(Effect.orDie) yield* Effect.all( [ workspace.waitForSync(workspaceID, { [sessionID]: 2 }), Effect.gen(function* () { yield* Effect.sleep("10 millis") yield* db .update(EventSequenceTable) .set({ seq: 2 }) .where(eq(EventSequenceTable.aggregate_id, sessionID)) .run() .pipe(Effect.orDie) GlobalBus.emit("event", { workspace: workspaceID, payload: { type: "anything" } }) }), ], { concurrency: "unbounded" }, ) }), { git: true }, ) it.instance( "a sync event for a different workspace can also release the fence", () => Effect.gen(function* () { const workspace = yield* Workspace.Service const workspaceID = WorkspaceV2.ID.ascending("wrk_wait_sync_any") const sessionID = SessionID.descending("ses_wait_sync_any") const { db } = yield* Database.Service yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 0 }).run().pipe(Effect.orDie) yield* Effect.all( [ workspace.waitForSync(workspaceID, { [sessionID]: 1 }), Effect.gen(function* () { yield* Effect.sleep("10 millis") yield* db .update(EventSequenceTable) .set({ seq: 1 }) .where(eq(EventSequenceTable.aggregate_id, sessionID)) .run() .pipe(Effect.orDie) GlobalBus.emit("event", { workspace: WorkspaceV2.ID.ascending("wrk_other_workspace"), payload: { type: "sync" }, }) }), ], { concurrency: "unbounded" }, ) }), { git: true }, ) it.instance( "rejects with the abort reason when aborted", () => Effect.gen(function* () { const workspace = yield* Workspace.Service const abort = new AbortController() const reason = new Error("caller aborted") const fiber = yield* Effect.forkChild( workspace.waitForSync( WorkspaceV2.ID.ascending("wrk_wait_abort"), { [SessionID.descending("ses_wait_abort")]: 1 }, abort.signal, ), ) abort.abort(reason) expectExitContains(yield* Fiber.await(fiber), "WorkspaceSyncAbortedError", reason.message) }), { git: true }, ) it.instance( "times out with the requested fence in the error message", () => Effect.gen(function* () { const workspace = yield* Workspace.Service const sessionID = SessionID.descending("ses_wait_timeout") expectExitContains( yield* Effect.exit( workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_timeout"), { [sessionID]: 1 }, undefined, 25), ), `Timed out waiting for sync fence: {"${sessionID}":1}`, ) }), { git: true }, 7000, ) })