| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621 |
- 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 { setTimeout as delay } from "node:timers/promises"
- import { NodeHttpServer } from "@effect/platform-node"
- import { Effect, Layer, Schema } from "effect"
- import { FetchHttpClient, HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
- import { eq } from "drizzle-orm"
- import { AppFileSystem } from "@opencode-ai/core/filesystem"
- import * as Log from "@opencode-ai/core/util/log"
- import { GlobalBus, type GlobalEvent } from "@/bus/global"
- import { Database } from "@/storage/db"
- import { ProjectID } from "@/project/schema"
- import { ProjectTable } from "@/project/project.sql"
- import { Instance } from "@/project/instance"
- import { WithInstance } from "../../src/project/with-instance"
- import { InstanceRef } from "@/effect/instance-ref"
- import { Session as SessionNs } from "@/session/session"
- import { SessionID } from "@/session/schema"
- import { SessionTable } from "@/session/session.sql"
- import { SyncEvent } from "@/sync"
- import { EventSequenceTable } from "@/sync/event.sql"
- import { resetDatabase } from "../fixture/db"
- import { disposeAllInstances, provideTmpdirInstance, TestInstance, tmpdir } from "../fixture/fixture"
- import { testEffect } from "../lib/effect"
- import { registerAdapter } from "../../src/control-plane/adapters"
- import { WorkspaceID } from "../../src/control-plane/schema"
- import { WorkspaceTable } from "../../src/control-plane/workspace.sql"
- import type { Target, WorkspaceAdapter, WorkspaceInfo } from "../../src/control-plane/types"
- import * as Workspace from "../../src/control-plane/workspace"
- import { AppRuntime } from "@/effect/app-runtime"
- import { InstanceStore } from "@/project/instance-store"
- import { InstanceBootstrap } from "@/project/bootstrap"
- import { Auth } from "@/auth"
- import { SessionPrompt } from "@/session/prompt"
- import { Project } from "@/project/project"
- import { Vcs } from "@/project/vcs"
- import { RuntimeFlags } from "@/effect/runtime-flags"
- void Log.init({ print: false })
- 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) =>
- Workspace.layer.pipe(
- Layer.provide(Auth.defaultLayer),
- Layer.provide(SessionNs.defaultLayer),
- Layer.provide(SyncEvent.defaultLayer),
- Layer.provide(SessionPrompt.defaultLayer),
- Layer.provide(Project.defaultLayer),
- Layer.provide(Vcs.defaultLayer),
- Layer.provide(FetchHttpClient.layer),
- Layer.provide(AppFileSystem.defaultLayer),
- Layer.provide(RuntimeFlags.layer({ experimentalWorkspaces })),
- Layer.provide(InstanceStore.defaultLayer.pipe(Layer.provide(InstanceBootstrap.defaultLayer))),
- )
- const testServerLayer = Layer.mergeAll(
- NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }),
- workspaceLayer(true),
- SessionNs.defaultLayer,
- )
- const it = testEffect(testServerLayer)
- type RecordedCreate = {
- info: WorkspaceInfo
- env: Record<string, string | undefined>
- 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(() => {
- Database.close()
- restoreEnv()
- process.env.OPENCODE_EXPERIMENTAL_WORKSPACES = "true"
- })
- afterEach(async () => {
- mock.restore()
- await disposeAllInstances()
- restoreEnv()
- await resetDatabase()
- })
- async function withInstance<T>(fn: (dir: string) => T | Promise<T>) {
- await using tmp = await tmpdir({ git: true })
- return await WithInstance.provide({
- directory: tmp.path,
- fn: () => fn(tmp.path),
- })
- }
- 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 runWorkspace = <A, E>(effect: Effect.Effect<A, E, Workspace.Service>) => AppRuntime.runPromise(effect)
- const createWorkspace = (input: Workspace.CreateInput) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.create(input)))
- const warpWorkspaceSession = (input: Workspace.SessionWarpInput) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.sessionWarp(input)))
- const listWorkspaces = (project: Parameters<Workspace.Interface["list"]>[0]) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.list(project)))
- const syncListWorkspaces = (project: Parameters<Workspace.Interface["syncList"]>[0]) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.syncList(project)))
- const getWorkspace = (id: WorkspaceID) => runWorkspace(Workspace.Service.use((workspace) => workspace.get(id)))
- const removeWorkspace = (id: WorkspaceID) => runWorkspace(Workspace.Service.use((workspace) => workspace.remove(id)))
- const workspaceStatus = () => runWorkspace(Workspace.Service.use((workspace) => workspace.status()))
- const isWorkspaceSyncing = (id: WorkspaceID) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.isSyncing(id)))
- const startWorkspaceSyncing = (projectID: ProjectID) => {
- void runWorkspace(Workspace.Service.use((workspace) => workspace.startWorkspaceSyncing(projectID)))
- }
- const startWorkspaceSyncingWithFlag = (projectID: ProjectID, experimentalWorkspaces: boolean) =>
- Effect.runPromise(
- Workspace.Service.use((workspace) => workspace.startWorkspaceSyncing(projectID)).pipe(
- Effect.provide(workspaceLayer(experimentalWorkspaces)),
- ),
- )
- const waitForWorkspaceSync = (workspaceID: WorkspaceID, state: Record<string, number>, signal?: AbortSignal) =>
- runWorkspace(Workspace.Service.use((workspace) => workspace.waitForSync(workspaceID, state, signal)))
- function captureGlobalEvents() {
- const events: GlobalEvent[] = []
- const handler = (event: GlobalEvent) => events.push(event)
- GlobalBus.on("event", handler)
- return {
- events,
- dispose() {
- GlobalBus.off("event", handler)
- },
- }
- }
- async function eventually<T>(fn: () => T | Promise<T>, timeout = 1500) {
- const started = Date.now()
- let last: unknown
- while (Date.now() - started < timeout) {
- try {
- return await fn()
- } catch (err) {
- last = err
- await delay(10)
- }
- }
- throw last ?? new Error("Timed out waiting for condition")
- }
- function eventuallyEffect(effect: Effect.Effect<void>, 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<Target>
- configure?: (info: WorkspaceInfo) => WorkspaceInfo | Promise<WorkspaceInfo>
- create?: (info: WorkspaceInfo, env: Record<string, string | undefined>, from?: WorkspaceInfo) => Promise<void>
- list?: () => Omit<WorkspaceInfo, "id">[] | Promise<Omit<WorkspaceInfo, "id">[]>
- remove?: (info: WorkspaceInfo) => Promise<void>
- }): 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<void> }) {
- 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<Uint8Array>({
- 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: ProjectID, type: string, input?: Partial<Workspace.Info>): Workspace.Info {
- return {
- id: input?.id ?? WorkspaceID.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) {
- Database.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(),
- )
- }
- function insertProject(id: ProjectID, worktree: string) {
- Database.use((db) =>
- db
- .insert(ProjectTable)
- .values({
- id,
- worktree,
- vcs: null,
- name: null,
- time_created: Date.now(),
- time_updated: Date.now(),
- sandboxes: [],
- })
- .run(),
- )
- }
- function attachSessionToWorkspace(sessionID: SessionID, workspaceID: WorkspaceID) {
- Database.use((db) =>
- db.update(SessionTable).set({ workspace_id: workspaceID }).where(eq(SessionTable.id, sessionID)).run(),
- )
- }
- function sessionSequence(sessionID: SessionID) {
- return Database.use((db) =>
- db
- .select({ seq: EventSequenceTable.seq })
- .from(EventSequenceTable)
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .get(),
- )?.seq
- }
- function sessionSequenceOwner(sessionID: SessionID) {
- return Database.use((db) =>
- db
- .select({ ownerID: EventSequenceTable.owner_id })
- .from(EventSequenceTable)
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .get(),
- )?.ownerID
- }
- function sessionUpdatedType() {
- return SyncEvent.versionedType(SessionNs.Event.Updated.type, SessionNs.Event.Updated.version)
- }
- 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: WorkspaceID.ascending("wrk_schema_create"),
- type: "worktree",
- branch: "feature/schema",
- projectID: ProjectID.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", () => {
- test("get returns undefined for a missing workspace", async () => {
- await withInstance(async () => {
- expect(await getWorkspace(WorkspaceID.ascending("wrk_missing_get"))).toBeUndefined()
- })
- })
- test("list maps database rows, filters by project, and sorts by id", async () => {
- await withInstance(async () => {
- const otherProjectID = ProjectID.make("project-other")
- insertProject(otherProjectID, "/tmp/other")
- const a = workspaceInfo(Instance.project.id, "manual", {
- id: WorkspaceID.ascending("wrk_a_list"),
- branch: "a",
- directory: "/a",
- extra: { a: true },
- })
- const b = workspaceInfo(Instance.project.id, "manual", {
- id: WorkspaceID.ascending("wrk_b_list"),
- branch: "b",
- directory: "/b",
- extra: ["b"],
- })
- const other = workspaceInfo(otherProjectID, "manual", { id: WorkspaceID.ascending("wrk_c_list") })
- insertWorkspace(b)
- insertWorkspace(other)
- insertWorkspace(a)
- expect(await listWorkspaces(Instance.project)).toEqual([a, b])
- })
- })
- test("create configures, persists, creates, starts local sync, and passes environment", async () => {
- await withInstance(async (dir) => {
- 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 = WorkspaceID.ascending("wrk_create_local")
- const type = unique("create-local")
- const targetDir = path.join(dir, "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 = await createWorkspace({
- 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(await getWorkspace(workspaceID)).toEqual(info)
- expect(await listWorkspaces(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((await workspaceStatus()).find((item) => item.workspaceID === workspaceID)?.status).toBe("connected")
- await removeWorkspace(workspaceID)
- expect((await workspaceStatus()).find((item) => item.workspaceID === workspaceID)?.status).toBeUndefined()
- })
- })
- test("create propagates configure failures and does not insert a workspace", async () => {
- await withInstance(async () => {
- const type = unique("configure-failure")
- registerAdapter(
- Instance.project.id,
- type,
- recordedAdapter({
- configure() {
- throw new Error("configure exploded")
- },
- target() {
- return { type: "local", directory: "/unused" }
- },
- }).adapter,
- )
- await expect(
- createWorkspace({ type, branch: null, projectID: Instance.project.id, extra: null }),
- ).rejects.toThrow("configure exploded")
- expect(await listWorkspaces(Instance.project)).toEqual([])
- })
- })
- test("create leaves the inserted row when adapter create fails", async () => {
- await withInstance(async () => {
- 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)
- await expect(
- createWorkspace({ type, branch: "branch", projectID: Instance.project.id, extra: { x: 1 } }),
- ).rejects.toThrow("create exploded")
- const rows = await listWorkspaces(Instance.project)
- expect(rows).toHaveLength(1)
- expect(rows[0]).toMatchObject({ type, branch: "branch", extra: { x: 1 } })
- expect(recorded.calls.target).toHaveLength(0)
- await removeWorkspace(rows[0].id)
- })
- })
- test("create returns after a local workspace reports error", async () => {
- await withInstance(async (dir) => {
- const type = unique("local-error")
- const missing = path.join(dir, "missing-local-target")
- const recorded = localAdapter(missing, { createDir: false })
- registerAdapter(Instance.project.id, type, recorded.adapter)
- const info = await createWorkspace({ type, branch: null, projectID: Instance.project.id, extra: null })
- expect(info.directory).toBe(missing)
- expect((await workspaceStatus()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
- await removeWorkspace(info.id)
- })
- })
- test("syncList registers adapter-listed workspaces that are missing by name", async () => {
- await withInstance(async (dir) => {
- const type = unique("list-sync")
- const existing = workspaceInfo(Instance.project.id, type, {
- id: WorkspaceID.ascending("wrk_list_sync_existing"),
- name: "existing",
- directory: path.join(dir, "existing"),
- })
- insertWorkspace(existing)
- const discovered = {
- type,
- name: "discovered",
- branch: "feature/discovered",
- directory: path.join(dir, "discovered"),
- extra: { source: "adapter" },
- projectID: Instance.project.id,
- }
- const recorded = recordedAdapter({
- list() {
- return [
- {
- type,
- name: existing.name,
- branch: "ignored",
- directory: path.join(dir, "ignored"),
- extra: null,
- projectID: Instance.project.id,
- },
- discovered,
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? dir }
- },
- })
- registerAdapter(Instance.project.id, type, recorded.adapter)
- await syncListWorkspaces(Instance.project)
- const synced = (await listWorkspaces(Instance.project)).filter((item) => item.name === discovered.name)
- expect(synced).toHaveLength(1)
- expect(synced[0]).toMatchObject(discovered)
- expect(synced[0]?.id).toStartWith("wrk_")
- expect(await listWorkspaces(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)
- })
- })
- test("syncList calls every registered adapter with a list method", async () => {
- await withInstance(async (dir) => {
- 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(dir, "adapter-a"),
- extra: null,
- projectID: Instance.project.id,
- },
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? dir }
- },
- })
- const adapterB = recordedAdapter({
- list() {
- return [
- {
- type: typeB,
- name: "adapter-b",
- branch: null,
- directory: path.join(dir, "adapter-b"),
- extra: null,
- projectID: Instance.project.id,
- },
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? dir }
- },
- })
- const noList = recordedAdapter({
- target() {
- return { type: "local", directory: dir }
- },
- })
- registerAdapter(Instance.project.id, typeA, adapterA.adapter)
- registerAdapter(Instance.project.id, typeB, adapterB.adapter)
- registerAdapter(Instance.project.id, unique("list-sync-none"), noList.adapter)
- await syncListWorkspaces(Instance.project)
- const synced = await listWorkspaces(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)
- })
- })
- 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 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 },
- )
- })
- })
- test("remove returns undefined for a missing workspace", async () => {
- await withInstance(async () => {
- expect(await removeWorkspace(WorkspaceID.ascending("wrk_missing_remove"))).toBeUndefined()
- })
- })
- it.instance(
- "remove deletes the workspace, associated sessions, adapter resources, and status",
- () => {
- return Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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({})
- attachSessionToWorkspace(one.id, info.id)
- 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()
- expect(
- Database.use((db) =>
- db.select({ id: SessionTable.id }).from(SessionTable).where(eq(SessionTable.workspace_id, info.id)).all(),
- ),
- ).toEqual([])
- })
- },
- { git: true },
- )
- test("remove still deletes the row when the adapter cannot remove resources", async () => {
- await withInstance(async () => {
- const type = unique("remove-throws")
- const info = workspaceInfo(Instance.project.id, type, { id: WorkspaceID.ascending("wrk_remove_throws") })
- registerAdapter(
- Instance.project.id,
- type,
- recordedAdapter({
- async remove() {
- throw new Error("remove exploded")
- },
- target() {
- return { type: "local", directory: "/unused" }
- },
- }).adapter,
- )
- insertWorkspace(info)
- expect(await removeWorkspace(info.id)).toEqual(info)
- expect(await getWorkspace(info.id)).toBeUndefined()
- })
- })
- 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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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)
- insertWorkspace(previous)
- 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({})
- attachSessionToWorkspace(session.id, previous.id)
- yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id })
- expect(
- Database.use((db) =>
- db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get(),
- )?.workspaceID,
- ).toBe(target.id)
- expect(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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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)
- insertWorkspace(previous)
- 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({})
- 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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const previousType = unique("warp-detach-local")
- const previous = workspaceInfo(instance.project.id, previousType)
- insertWorkspace(previous)
- registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-detach-local")).adapter)
- const session = yield* sessionSvc.create({})
- attachSessionToWorkspace(session.id, previous.id)
- yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id })
- expect(
- Database.use((db) =>
- db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get(),
- )?.workspaceID,
- ).toBeNull()
- expect(sessionSequenceOwner(session.id)).toBe(instance.project.id)
- })
- },
- { git: true },
- )
- test("sessionWarp detaches to the source project when invoked from a workspace instance", async () => {
- await withInstance(async () => {
- const projectID = Instance.project.id
- await using workspaceTmp = await tmpdir({ git: true })
- const previousType = unique("warp-detach-workspace-instance")
- const previous = workspaceInfo(projectID, previousType)
- insertWorkspace(previous)
- registerAdapter(projectID, previousType, localAdapter(workspaceTmp.path, { createDir: false }).adapter)
- const session = await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))
- attachSessionToWorkspace(session.id, previous.id)
- const workspaceProjectID = await WithInstance.provide({
- directory: workspaceTmp.path,
- fn: async () => {
- const id = Instance.project.id
- expect(id).not.toBe(projectID)
- await warpWorkspaceSession({ workspaceID: null, sessionID: session.id })
- return id
- },
- })
- expect(
- Database.use((db) =>
- db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get(),
- )?.workspaceID,
- ).toBeNull()
- expect(sessionSequenceOwner(session.id)).toBe(projectID)
- expect(sessionSequenceOwner(session.id)).not.toBe(workspaceProjectID)
- })
- })
- it.live("sessionWarp syncs previous remote history, replays it, steals, and claims the sequence", () => {
- const calls: FetchCall[] = []
- let historySessionID: SessionID | 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: sessionUpdatedType(),
- data: { sessionID: historySessionID!, info: { title: "from source history" } },
- },
- ])
- }
- 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 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" })
- insertWorkspace(previous)
- 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({})
- attachSessionToWorkspace(session.id, previous.id)
- historySessionID = session.id
- historyNextSeq = (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: SyncEvent.versionedType(SessionNs.Event.Created.type, SessionNs.Event.Created.version),
- },
- {
- aggregateID: session.id,
- seq: historyNextSeq,
- type: sessionUpdatedType(),
- },
- ],
- })
- expect(calls[4].json).toEqual({ sessionID: session.id })
- expect((yield* sessionSvc.get(session.id)).title).toBe("from source history")
- expect(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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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({})
- attachSessionToWorkspace(session.id, info.id)
- 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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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* Effect.sync(() => {
- insertWorkspace(first)
- 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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("missing-local")
- const info = workspaceInfo(instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(
- instance.project.id,
- type,
- localAdapter(path.join(dir, "missing-target"), { createDir: false }).adapter,
- )
- 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* InstanceRef
- if (!instance) return yield* Effect.die(new Error("missing test instance"))
- 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 }))
- insertWorkspace(info)
- registerAdapter(instance.project.id, type, localAdapter(target).adapter)
- 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 captured = captureGlobalEvents()
- try {
- const type = unique("remote-start")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sync`).adapter)
- 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 type = unique("remote-connect-fail")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/failed`).adapter)
- 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 type = unique("remote-history-fail")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/history-failed`).adapter)
- 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 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: sessionUpdatedType(),
- data: { sessionID: historySessionID!, info: { title: "from history" } },
- },
- ]),
- )
- }
- 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 captured = captureGlobalEvents()
- try {
- const type = unique("history-replay")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/history`).adapter)
- const session = yield* sessionSvc.create({ title: "before history" })
- attachSessionToWorkspace(session.id, info.id)
- historySessionID = session.id
- historyNextSeq = (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 === "sync" &&
- event.payload.syncEvent.seq === historyNextSeq,
- ),
- ).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 captured = captureGlobalEvents()
- try {
- const type = unique("sse-forward")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sse-forward`).adapter)
- 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 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: sessionUpdatedType(),
- data: { sessionID: sseSessionID!, info: { title: "from sse" } },
- },
- },
- },
- ],
- 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 captured = captureGlobalEvents()
- try {
- const type = unique("sse-sync")
- const info = workspaceInfo(Instance.project.id, type)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sse-sync`).adapter)
- const session = yield* sessionSvc.create({ title: "before sse" })
- attachSessionToWorkspace(session.id, info.id)
- sseSessionID = session.id
- sseNextSeq = (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", () => {
- test("returns immediately for an empty fence", async () => {
- await withInstance(async () => {
- await expect(waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_empty"), {})).resolves.toBeUndefined()
- })
- })
- test("returns immediately when the stored sequence already satisfies the fence", async () => {
- await withInstance(async () => {
- const sessionID = SessionID.descending("ses_wait_done")
- Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 4 }).run())
- await expect(
- waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_done"), { [sessionID]: 4 }),
- ).resolves.toBeUndefined()
- await expect(
- waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_done_2"), { [sessionID]: 3 }),
- ).resolves.toBeUndefined()
- })
- })
- test("waits until the database reaches the requested sequence and a workspace event arrives", async () => {
- await withInstance(async () => {
- const workspaceID = WorkspaceID.ascending("wrk_wait_event")
- const sessionID = SessionID.descending("ses_wait_event")
- Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 1 }).run())
- const waited = waitForWorkspaceSync(workspaceID, { [sessionID]: 2 })
- await delay(10)
- Database.use((db) =>
- db.update(EventSequenceTable).set({ seq: 2 }).where(eq(EventSequenceTable.aggregate_id, sessionID)).run(),
- )
- GlobalBus.emit("event", { workspace: workspaceID, payload: { type: "anything" } })
- await expect(waited).resolves.toBeUndefined()
- })
- })
- test("a sync event for a different workspace can also release the fence", async () => {
- await withInstance(async () => {
- const workspaceID = WorkspaceID.ascending("wrk_wait_sync_any")
- const sessionID = SessionID.descending("ses_wait_sync_any")
- Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 0 }).run())
- const waited = waitForWorkspaceSync(workspaceID, { [sessionID]: 1 })
- await delay(10)
- Database.use((db) =>
- db.update(EventSequenceTable).set({ seq: 1 }).where(eq(EventSequenceTable.aggregate_id, sessionID)).run(),
- )
- GlobalBus.emit("event", {
- workspace: WorkspaceID.ascending("wrk_other_workspace"),
- payload: { type: "sync" },
- })
- await expect(waited).resolves.toBeUndefined()
- })
- })
- test("rejects with the abort reason when aborted", async () => {
- await withInstance(async () => {
- const abort = new AbortController()
- const reason = new Error("caller aborted")
- const waited = waitForWorkspaceSync(
- WorkspaceID.ascending("wrk_wait_abort"),
- { [SessionID.descending("ses_wait_abort")]: 1 },
- abort.signal,
- )
- abort.abort(reason)
- await expect(waited).rejects.toMatchObject({
- _tag: "WorkspaceSyncAbortedError",
- message: reason.message,
- cause: reason,
- })
- })
- })
- test("times out with the requested fence in the error message", async () => {
- await withInstance(async () => {
- const sessionID = SessionID.descending("ses_wait_timeout")
- await expect(waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_timeout"), { [sessionID]: 1 })).rejects.toThrow(
- `Timed out waiting for sync fence: {"${sessionID}":1}`,
- )
- })
- }, 7000)
- })
|