| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526 |
- import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"
- 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 } from "effect"
- import { HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
- import { asc, eq } from "drizzle-orm"
- import * as Log from "@opencode-ai/core/util/log"
- import { Flag } from "@opencode-ai/core/flag/flag"
- 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 { Session as SessionNs } from "@/session/session"
- import { SessionID, MessageID, PartID } from "@/session/schema"
- import { SessionTable } from "@/session/session.sql"
- import { ModelID, ProviderID } from "@/provider/schema"
- import { SyncEvent } from "@/sync"
- import { EventSequenceTable, EventTable } from "@/sync/event.sql"
- import { resetDatabase } from "../fixture/db"
- import { disposeAllInstances, provideTmpdirInstance, 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 WorkspaceOld from "../../src/control-plane/workspace"
- import { AppRuntime } from "@/effect/app-runtime"
- void Log.init({ print: false })
- const testServerLayer = Layer.mergeAll(
- NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }),
- WorkspaceOld.defaultLayer,
- SessionNs.defaultLayer,
- )
- const it = testEffect(testServerLayer)
- const originalWorkspacesFlag = Flag.OPENCODE_EXPERIMENTAL_WORKSPACES
- const originalEnv = {
- OPENCODE_AUTH_CONTENT: process.env.OPENCODE_AUTH_CONTENT,
- 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,
- }
- type RecordedCreate = {
- info: WorkspaceInfo
- env: Record<string, string | undefined>
- from?: WorkspaceInfo
- }
- type RecordedAdapter = {
- adapter: WorkspaceAdapter
- calls: {
- configure: WorkspaceInfo[]
- create: RecordedCreate[]
- 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()
- Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = true
- restoreEnv()
- })
- afterEach(async () => {
- mock.restore()
- await disposeAllInstances()
- Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = originalWorkspacesFlag
- restoreEnv()
- await resetDatabase()
- })
- async function withInstance<T>(fn: (dir: string) => T | Promise<T>) {
- await using tmp = await tmpdir({ git: true })
- return Instance.provide({
- directory: tmp.path,
- fn: () => fn(tmp.path),
- })
- }
- const runWorkspace = <A, E>(effect: Effect.Effect<A, E, WorkspaceOld.Service>) => AppRuntime.runPromise(effect)
- const createWorkspace = (input: WorkspaceOld.CreateInput) =>
- runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.create(input)))
- const restoreWorkspaceSession = (input: WorkspaceOld.SessionRestoreInput) =>
- runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.sessionRestore(input)))
- const listWorkspaces = (project: Parameters<WorkspaceOld.Interface["list"]>[0]) =>
- runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.list(project)))
- const getWorkspace = (id: WorkspaceID) => runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.get(id)))
- const removeWorkspace = (id: WorkspaceID) => runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.remove(id)))
- const workspaceStatus = () => runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.status()))
- const isWorkspaceSyncing = (id: WorkspaceID) =>
- runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.isSyncing(id)))
- const startWorkspaceSyncing = (projectID: ProjectID) => {
- void runWorkspace(WorkspaceOld.Service.use((workspace) => workspace.startWorkspaceSyncing(projectID)))
- }
- const waitForWorkspaceSync = (workspaceID: WorkspaceID, state: Record<string, number>, signal?: AbortSignal) =>
- runWorkspace(WorkspaceOld.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>
- remove?: (info: WorkspaceInfo) => Promise<void>
- }): RecordedAdapter {
- const calls: RecordedAdapter["calls"] = {
- configure: [],
- create: [],
- 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)
- },
- 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<WorkspaceInfo>): WorkspaceInfo {
- 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,
- }
- }
- function insertWorkspace(info: WorkspaceInfo) {
- 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,
- })
- .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 eventRows(sessionID: SessionID) {
- return Database.use((db) =>
- db
- .select({ seq: EventTable.seq, type: EventTable.type, data: EventTable.data })
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, sessionID))
- .orderBy(asc(EventTable.seq))
- .all(),
- )
- }
- function sessionUpdatedType() {
- return SyncEvent.versionedType(SessionNs.Event.Updated.type, SessionNs.Event.Updated.version)
- }
- function replaceSessionEvents(sessionID: SessionID, count: number) {
- Database.use((db) => {
- db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, sessionID)).run()
- if (count === 0) return
- db.insert(EventSequenceTable)
- .values({ aggregate_id: sessionID, seq: count - 1 })
- .run()
- db.insert(EventTable)
- .values(
- Array.from({ length: count }, (_, i) => ({
- id: `evt_${unique(`manual-${i}`)}`,
- aggregate_id: sessionID,
- seq: i,
- type: sessionUpdatedType(),
- data: { sessionID, info: { title: `manual ${i}` } },
- })),
- )
- .run()
- })
- }
- describe("workspace-old schemas and exports", () => {
- test("keeps the historical event type names", () => {
- expect(WorkspaceOld.Event.Ready.type).toBe("workspace.ready")
- expect(WorkspaceOld.Event.Failed.type).toBe("workspace.failed")
- expect(WorkspaceOld.Event.Restore.type).toBe("workspace.restore")
- expect(WorkspaceOld.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 },
- }
- expect(WorkspaceOld.CreateInput.zod.parse(input)).toEqual(input)
- expect(() => WorkspaceOld.CreateInput.zod.parse({ ...input, id: "bad" })).toThrow()
- expect(() => WorkspaceOld.CreateInput.zod.parse({ ...input, branch: 1 })).toThrow()
- })
- test("validates session restore input", () => {
- const input = {
- workspaceID: WorkspaceID.ascending("wrk_schema_restore"),
- sessionID: SessionID.descending("ses_schema_restore"),
- }
- expect(WorkspaceOld.SessionRestoreInput.zod.parse(input)).toEqual(input)
- expect(() => WorkspaceOld.SessionRestoreInput.zod.parse({ ...input, workspaceID: "bad" })).toThrow()
- expect(() => WorkspaceOld.SessionRestoreInput.zod.parse({ ...input, sessionID: "bad" })).toThrow()
- })
- })
- describe("workspace-old 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,
- })
- 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(info)
- 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)
- })
- })
- 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* WorkspaceOld.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()
- })
- })
- test("remove deletes the workspace, associated sessions, adapter resources, and status", async () => {
- await withInstance(async (dir) => {
- const type = unique("remove-local")
- const recorded = localAdapter(path.join(dir, "remove-local"))
- registerAdapter(Instance.project.id, type, recorded.adapter)
- const info = await createWorkspace({ type, branch: null, projectID: Instance.project.id, extra: null })
- const one = await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))
- const two = await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))
- attachSessionToWorkspace(one.id, info.id)
- attachSessionToWorkspace(two.id, info.id)
- const removed = await removeWorkspace(info.id)
- expect(removed).toEqual(info)
- expect(await getWorkspace(info.id)).toBeUndefined()
- expect(recorded.calls.remove).toEqual([info])
- expect((await workspaceStatus()).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([])
- })
- })
- 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()
- })
- })
- })
- describe("workspace-old sync state", () => {
- test("startWorkspaceSyncing is disabled by the experimental workspace flag", async () => {
- await withInstance(async (dir) => {
- Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = false
- const type = unique("flag-disabled")
- const info = workspaceInfo(Instance.project.id, type)
- const session = await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))
- attachSessionToWorkspace(session.id, info.id)
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, localAdapter(path.join(dir, "flag-disabled")).adapter)
- startWorkspaceSyncing(Instance.project.id)
- await delay(25)
- expect((await workspaceStatus()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
- })
- })
- test("startWorkspaceSyncing starts only workspaces with sessions", async () => {
- await withInstance(async (dir) => {
- const withSessionType = unique("with-session")
- const withoutSessionType = unique("without-session")
- const withSession = workspaceInfo(Instance.project.id, withSessionType)
- const withoutSession = workspaceInfo(Instance.project.id, withoutSessionType)
- const withSessionDir = path.join(dir, "with-session")
- const withoutSessionDir = path.join(dir, "without-session")
- await fs.mkdir(withSessionDir, { recursive: true })
- await fs.mkdir(withoutSessionDir, { recursive: true })
- insertWorkspace(withSession)
- insertWorkspace(withoutSession)
- registerAdapter(Instance.project.id, withSessionType, localAdapter(withSessionDir).adapter)
- registerAdapter(Instance.project.id, withoutSessionType, localAdapter(withoutSessionDir).adapter)
- attachSessionToWorkspace(
- (await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))).id,
- withSession.id,
- )
- startWorkspaceSyncing(Instance.project.id)
- await eventually(() =>
- workspaceStatus().then((status) =>
- expect(status.find((item) => item.workspaceID === withSession.id)?.status).toBe("connected"),
- ),
- )
- expect((await workspaceStatus()).find((item) => item.workspaceID === withoutSession.id)?.status).toBeUndefined()
- await removeWorkspace(withSession.id)
- await removeWorkspace(withoutSession.id)
- })
- })
- test("local start reports error when the target directory is missing", async () => {
- await withInstance(async (dir) => {
- 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(
- (await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))).id,
- info.id,
- )
- startWorkspaceSyncing(Instance.project.id)
- await eventually(() =>
- workspaceStatus().then((status) =>
- expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("error"),
- ),
- )
- expect(await isWorkspaceSyncing(info.id)).toBe(false)
- await removeWorkspace(info.id)
- })
- })
- test("duplicate local status updates are suppressed", async () => {
- await withInstance(async (dir) => {
- const captured = captureGlobalEvents()
- try {
- const type = unique("dedupe-local")
- const info = workspaceInfo(Instance.project.id, type)
- const target = path.join(dir, "dedupe-local")
- await fs.mkdir(target, { recursive: true })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, localAdapter(target).adapter)
- attachSessionToWorkspace(
- (await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))).id,
- info.id,
- )
- startWorkspaceSyncing(Instance.project.id)
- startWorkspaceSyncing(Instance.project.id)
- await eventually(() =>
- workspaceStatus().then((status) =>
- expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("connected"),
- ),
- )
- expect(
- captured.events.filter(
- (event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Status.type,
- ),
- ).toHaveLength(1)
- await removeWorkspace(info.id)
- } finally {
- captured.dispose()
- }
- })
- })
- 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* WorkspaceOld.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 === WorkspaceOld.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* WorkspaceOld.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* WorkspaceOld.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* WorkspaceOld.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)).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* WorkspaceOld.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* WorkspaceOld.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)).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-old 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)
- })
- describe("workspace-old sessionRestore", () => {
- test("throws when the workspace is missing", async () => {
- await withInstance(async () => {
- await expect(
- restoreWorkspaceSession({
- workspaceID: WorkspaceID.ascending("wrk_restore_missing"),
- sessionID: SessionID.descending("ses_restore_missing_workspace"),
- }),
- ).rejects.toThrow("Workspace not found: wrk_restore_missing")
- })
- })
- test("throws when switching a missing session fails", async () => {
- await withInstance(async (dir) => {
- const type = unique("restore-missing-session")
- const info = workspaceInfo(Instance.project.id, type, { directory: dir })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, localAdapter(dir).adapter)
- await expect(
- restoreWorkspaceSession({ workspaceID: info.id, sessionID: SessionID.descending("ses_missing_restore") }),
- ).rejects.toThrow("NotFoundError")
- await removeWorkspace(info.id)
- })
- })
- it.live("posts remote replay batches of 10, emits progress, and includes the workspace update event", () => {
- const replay: 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,
- }
- if (call.url.pathname === "/restore/sync/replay") {
- replay.push(call)
- return HttpServerResponse.fromWeb(Response.json({ ok: true }))
- }
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const workspace = yield* WorkspaceOld.Service
- const sessionSvc = yield* SessionNs.Service
- const captured = captureGlobalEvents()
- try {
- const type = unique("restore-remote")
- const info = workspaceInfo(Instance.project.id, type, { directory: dir })
- insertWorkspace(info)
- registerAdapter(
- Instance.project.id,
- type,
- remoteAdapter(`${url}/restore/?ignored=1#hash`, {
- directory: dir,
- headers: { authorization: "Bearer restore" },
- }).adapter,
- )
- const session = yield* sessionSvc.create({ title: "restore remote" })
- replaceSessionEvents(session.id, 24)
- const result = yield* workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id })
- expect(result).toEqual({ total: 3 })
- expect(replay).toHaveLength(3)
- expect(replay.map((call) => call.url.pathname + call.url.search + call.url.hash)).toEqual([
- "/restore/sync/replay",
- "/restore/sync/replay",
- "/restore/sync/replay",
- ])
- expect(replay.every((call) => call.headers.get("authorization") === "Bearer restore")).toBe(true)
- expect(replay.every((call) => call.headers.get("content-type") === "application/json")).toBe(true)
- expect(replay.map((call) => (call.json as { events: unknown[] }).events.length)).toEqual([10, 10, 5])
- expect(replay.map((call) => (call.json as { directory: string }).directory)).toEqual([dir, dir, dir])
- expect(
- replay.flatMap((call) =>
- (call.json as { events: Array<{ seq: number }> }).events.map((event) => event.seq),
- ),
- ).toEqual(Array.from({ length: 25 }, (_, i) => i))
- expect(
- (replay[2].json as { events: Array<{ seq: number; type: string; data: unknown }> }).events.at(-1),
- ).toMatchObject({
- seq: 24,
- type: sessionUpdatedType(),
- data: { sessionID: session.id, info: { workspaceID: info.id } },
- })
- expect((yield* sessionSvc.get(session.id)).workspaceID).toBe(info.id)
- expect(
- captured.events
- .filter(
- (event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Restore.type,
- )
- .map((event) => event.payload.properties.step),
- ).toEqual([0, 1, 2, 3])
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- })
- })
- it.live("remote restore sends an empty directory string when the workspace directory is null", () => {
- const replay: FetchCall[] = []
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- replay.push({
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- })
- return HttpServerResponse.fromWeb(Response.json({ ok: true }))
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* WorkspaceOld.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("restore-null-dir")
- const info = workspaceInfo(Instance.project.id, type, { directory: null })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/null-dir`, { directory: null }).adapter)
- const session = yield* sessionSvc.create({ title: "null dir" })
- replaceSessionEvents(session.id, 0)
- expect(yield* workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id })).toEqual({
- total: 1,
- })
- expect((replay[0].json as { directory: string }).directory).toBe("")
- expect((replay[0].json as { events: unknown[] }).events).toHaveLength(1)
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- })
- })
- it.live("remote restore failures include status and body and do not emit completed batch progress", () => {
- const replay: FetchCall[] = []
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- replay.push({
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- })
- return HttpServerResponse.text("replay failed", { status: 503 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const workspace = yield* WorkspaceOld.Service
- const sessionSvc = yield* SessionNs.Service
- const captured = captureGlobalEvents()
- try {
- const type = unique("restore-remote-fail")
- const info = workspaceInfo(Instance.project.id, type, { directory: dir })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/fail`, { directory: dir }).adapter)
- const session = yield* sessionSvc.create({ title: "restore fail" })
- replaceSessionEvents(session.id, 11)
- const error = yield* Effect.flip(
- workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id }),
- )
- expect((error as Error).message).toContain(
- `Failed to replay session ${session.id} into workspace ${info.id}: HTTP 503 replay failed`,
- )
- expect(replay).toHaveLength(1)
- expect(
- captured.events
- .filter(
- (event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Restore.type,
- )
- .map((event) => event.payload.properties.step),
- ).toEqual([0])
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- })
- })
- it.live("local restore replays batches and emits progress", () =>
- provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const workspace = yield* WorkspaceOld.Service
- const sessionSvc = yield* SessionNs.Service
- const captured = captureGlobalEvents()
- try {
- const type = unique("restore-local")
- const info = workspaceInfo(Instance.project.id, type, { directory: dir })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, localAdapter(dir).adapter)
- const session = yield* sessionSvc.create({ title: "restore local" })
- replaceSessionEvents(session.id, 20)
- expect(yield* workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id })).toEqual({
- total: 3,
- })
- expect((yield* sessionSvc.get(session.id)).workspaceID).toBe(info.id)
- expect(eventRows(session.id).map((row) => row.seq)).toEqual(Array.from({ length: 21 }, (_, i) => i))
- expect(
- captured.events
- .filter(
- (event) => event.workspace === info.id && event.payload.type === WorkspaceOld.Event.Restore.type,
- )
- .map((event) => event.payload.properties.step),
- ).toEqual([0, 1, 2, 3])
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- ),
- )
- it.live("session restore includes real message and part events in sequence order", () => {
- const replay: FetchCall[] = []
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- replay.push({
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- })
- return HttpServerResponse.fromWeb(Response.json({ ok: true }))
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const workspace = yield* WorkspaceOld.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("restore-real-events")
- const info = workspaceInfo(Instance.project.id, type, { directory: dir })
- insertWorkspace(info)
- registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/real`, { directory: dir }).adapter)
- const session = yield* sessionSvc.create({ title: "real events" })
- for (let i = 0; i < 3; i++) {
- const msg = yield* sessionSvc.updateMessage({
- id: MessageID.ascending(),
- role: "user",
- sessionID: session.id,
- agent: "build",
- model: { providerID: ProviderID.make("test"), modelID: ModelID.make("test") },
- time: { created: Date.now() },
- })
- yield* sessionSvc.updatePart({
- id: PartID.ascending(),
- sessionID: session.id,
- messageID: msg.id,
- type: "text",
- text: `message ${i}`,
- })
- }
- const before = eventRows(session.id)
- expect(yield* workspace.sessionRestore({ workspaceID: info.id, sessionID: session.id })).toEqual({
- total: 1,
- })
- const posted = (replay[0].json as { events: Array<{ seq: number; type: string }> }).events
- expect(posted.map((event) => event.seq)).toEqual([...before.map((row) => row.seq), before.at(-1)!.seq + 1])
- expect(posted.map((event) => event.type).slice(0, -1)).toEqual(before.map((row) => row.type))
- expect(posted.at(-1)?.type).toBe(sessionUpdatedType())
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- })
- })
- })
|