workspace.test.ts 62 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621
  1. import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"
  2. import { $ } from "bun"
  3. import fs from "node:fs/promises"
  4. import Http from "node:http"
  5. import path from "node:path"
  6. import { setTimeout as delay } from "node:timers/promises"
  7. import { NodeHttpServer } from "@effect/platform-node"
  8. import { Effect, Layer, Schema } from "effect"
  9. import { FetchHttpClient, HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
  10. import { eq } from "drizzle-orm"
  11. import { AppFileSystem } from "@opencode-ai/core/filesystem"
  12. import * as Log from "@opencode-ai/core/util/log"
  13. import { GlobalBus, type GlobalEvent } from "@/bus/global"
  14. import { Database } from "@/storage/db"
  15. import { ProjectID } from "@/project/schema"
  16. import { ProjectTable } from "@/project/project.sql"
  17. import { Instance } from "@/project/instance"
  18. import { WithInstance } from "../../src/project/with-instance"
  19. import { InstanceRef } from "@/effect/instance-ref"
  20. import { Session as SessionNs } from "@/session/session"
  21. import { SessionID } from "@/session/schema"
  22. import { SessionTable } from "@/session/session.sql"
  23. import { SyncEvent } from "@/sync"
  24. import { EventSequenceTable } from "@/sync/event.sql"
  25. import { resetDatabase } from "../fixture/db"
  26. import { disposeAllInstances, provideTmpdirInstance, TestInstance, tmpdir } from "../fixture/fixture"
  27. import { testEffect } from "../lib/effect"
  28. import { registerAdapter } from "../../src/control-plane/adapters"
  29. import { WorkspaceID } from "../../src/control-plane/schema"
  30. import { WorkspaceTable } from "../../src/control-plane/workspace.sql"
  31. import type { Target, WorkspaceAdapter, WorkspaceInfo } from "../../src/control-plane/types"
  32. import * as Workspace from "../../src/control-plane/workspace"
  33. import { AppRuntime } from "@/effect/app-runtime"
  34. import { InstanceStore } from "@/project/instance-store"
  35. import { InstanceBootstrap } from "@/project/bootstrap"
  36. import { Auth } from "@/auth"
  37. import { SessionPrompt } from "@/session/prompt"
  38. import { Project } from "@/project/project"
  39. import { Vcs } from "@/project/vcs"
  40. import { RuntimeFlags } from "@/effect/runtime-flags"
  41. void Log.init({ print: false })
  42. const originalEnv = {
  43. OPENCODE_AUTH_CONTENT: process.env.OPENCODE_AUTH_CONTENT,
  44. OPENCODE_EXPERIMENTAL_WORKSPACES: process.env.OPENCODE_EXPERIMENTAL_WORKSPACES,
  45. OTEL_EXPORTER_OTLP_HEADERS: process.env.OTEL_EXPORTER_OTLP_HEADERS,
  46. OTEL_EXPORTER_OTLP_ENDPOINT: process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
  47. OTEL_RESOURCE_ATTRIBUTES: process.env.OTEL_RESOURCE_ATTRIBUTES,
  48. }
  49. const workspaceLayer = (experimentalWorkspaces: boolean) =>
  50. Workspace.layer.pipe(
  51. Layer.provide(Auth.defaultLayer),
  52. Layer.provide(SessionNs.defaultLayer),
  53. Layer.provide(SyncEvent.defaultLayer),
  54. Layer.provide(SessionPrompt.defaultLayer),
  55. Layer.provide(Project.defaultLayer),
  56. Layer.provide(Vcs.defaultLayer),
  57. Layer.provide(FetchHttpClient.layer),
  58. Layer.provide(AppFileSystem.defaultLayer),
  59. Layer.provide(RuntimeFlags.layer({ experimentalWorkspaces })),
  60. Layer.provide(InstanceStore.defaultLayer.pipe(Layer.provide(InstanceBootstrap.defaultLayer))),
  61. )
  62. const testServerLayer = Layer.mergeAll(
  63. NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }),
  64. workspaceLayer(true),
  65. SessionNs.defaultLayer,
  66. )
  67. const it = testEffect(testServerLayer)
  68. type RecordedCreate = {
  69. info: WorkspaceInfo
  70. env: Record<string, string | undefined>
  71. from?: WorkspaceInfo
  72. }
  73. type RecordedAdapter = {
  74. adapter: WorkspaceAdapter
  75. calls: {
  76. configure: WorkspaceInfo[]
  77. create: RecordedCreate[]
  78. list: number
  79. remove: WorkspaceInfo[]
  80. target: WorkspaceInfo[]
  81. }
  82. }
  83. type FetchCall = {
  84. url: URL
  85. method: string
  86. headers: Headers
  87. bodyText?: string
  88. json?: unknown
  89. }
  90. function unique(prefix: string) {
  91. return `${prefix}-${Math.random().toString(36).slice(2)}`
  92. }
  93. function restoreEnv() {
  94. Object.entries(originalEnv).forEach(([key, value]) => {
  95. if (value === undefined) {
  96. delete process.env[key]
  97. return
  98. }
  99. process.env[key] = value
  100. })
  101. }
  102. beforeEach(() => {
  103. Database.close()
  104. restoreEnv()
  105. process.env.OPENCODE_EXPERIMENTAL_WORKSPACES = "true"
  106. })
  107. afterEach(async () => {
  108. mock.restore()
  109. await disposeAllInstances()
  110. restoreEnv()
  111. await resetDatabase()
  112. })
  113. async function withInstance<T>(fn: (dir: string) => T | Promise<T>) {
  114. await using tmp = await tmpdir({ git: true })
  115. return await WithInstance.provide({
  116. directory: tmp.path,
  117. fn: () => fn(tmp.path),
  118. })
  119. }
  120. async function initGitRepo(dir: string) {
  121. await fs.mkdir(dir, { recursive: true })
  122. await $`git init`.cwd(dir).quiet()
  123. await $`git config core.fsmonitor false`.cwd(dir).quiet()
  124. await $`git config commit.gpgsign false`.cwd(dir).quiet()
  125. await $`git config user.email "test@opencode.test"`.cwd(dir).quiet()
  126. await $`git config user.name "Test"`.cwd(dir).quiet()
  127. await fs.writeFile(path.join(dir, "tracked.txt"), "base\n")
  128. await $`git add tracked.txt`.cwd(dir).quiet()
  129. await $`git commit -m "base"`.cwd(dir).quiet()
  130. }
  131. const runWorkspace = <A, E>(effect: Effect.Effect<A, E, Workspace.Service>) => AppRuntime.runPromise(effect)
  132. const createWorkspace = (input: Workspace.CreateInput) =>
  133. runWorkspace(Workspace.Service.use((workspace) => workspace.create(input)))
  134. const warpWorkspaceSession = (input: Workspace.SessionWarpInput) =>
  135. runWorkspace(Workspace.Service.use((workspace) => workspace.sessionWarp(input)))
  136. const listWorkspaces = (project: Parameters<Workspace.Interface["list"]>[0]) =>
  137. runWorkspace(Workspace.Service.use((workspace) => workspace.list(project)))
  138. const syncListWorkspaces = (project: Parameters<Workspace.Interface["syncList"]>[0]) =>
  139. runWorkspace(Workspace.Service.use((workspace) => workspace.syncList(project)))
  140. const getWorkspace = (id: WorkspaceID) => runWorkspace(Workspace.Service.use((workspace) => workspace.get(id)))
  141. const removeWorkspace = (id: WorkspaceID) => runWorkspace(Workspace.Service.use((workspace) => workspace.remove(id)))
  142. const workspaceStatus = () => runWorkspace(Workspace.Service.use((workspace) => workspace.status()))
  143. const isWorkspaceSyncing = (id: WorkspaceID) =>
  144. runWorkspace(Workspace.Service.use((workspace) => workspace.isSyncing(id)))
  145. const startWorkspaceSyncing = (projectID: ProjectID) => {
  146. void runWorkspace(Workspace.Service.use((workspace) => workspace.startWorkspaceSyncing(projectID)))
  147. }
  148. const startWorkspaceSyncingWithFlag = (projectID: ProjectID, experimentalWorkspaces: boolean) =>
  149. Effect.runPromise(
  150. Workspace.Service.use((workspace) => workspace.startWorkspaceSyncing(projectID)).pipe(
  151. Effect.provide(workspaceLayer(experimentalWorkspaces)),
  152. ),
  153. )
  154. const waitForWorkspaceSync = (workspaceID: WorkspaceID, state: Record<string, number>, signal?: AbortSignal) =>
  155. runWorkspace(Workspace.Service.use((workspace) => workspace.waitForSync(workspaceID, state, signal)))
  156. function captureGlobalEvents() {
  157. const events: GlobalEvent[] = []
  158. const handler = (event: GlobalEvent) => events.push(event)
  159. GlobalBus.on("event", handler)
  160. return {
  161. events,
  162. dispose() {
  163. GlobalBus.off("event", handler)
  164. },
  165. }
  166. }
  167. async function eventually<T>(fn: () => T | Promise<T>, timeout = 1500) {
  168. const started = Date.now()
  169. let last: unknown
  170. while (Date.now() - started < timeout) {
  171. try {
  172. return await fn()
  173. } catch (err) {
  174. last = err
  175. await delay(10)
  176. }
  177. }
  178. throw last ?? new Error("Timed out waiting for condition")
  179. }
  180. function eventuallyEffect(effect: Effect.Effect<void>, timeout = 1500) {
  181. return Effect.gen(function* () {
  182. const started = Date.now()
  183. let last: unknown
  184. while (Date.now() - started < timeout) {
  185. const exit = yield* Effect.exit(effect)
  186. if (exit._tag === "Success") return
  187. last = exit.cause
  188. yield* Effect.sleep("10 millis")
  189. }
  190. throw last ?? new Error("Timed out waiting for condition")
  191. })
  192. }
  193. function recordedAdapter(input: {
  194. target: (info: WorkspaceInfo) => Target | Promise<Target>
  195. configure?: (info: WorkspaceInfo) => WorkspaceInfo | Promise<WorkspaceInfo>
  196. create?: (info: WorkspaceInfo, env: Record<string, string | undefined>, from?: WorkspaceInfo) => Promise<void>
  197. list?: () => Omit<WorkspaceInfo, "id">[] | Promise<Omit<WorkspaceInfo, "id">[]>
  198. remove?: (info: WorkspaceInfo) => Promise<void>
  199. }): RecordedAdapter {
  200. const calls: RecordedAdapter["calls"] = {
  201. configure: [],
  202. create: [],
  203. list: 0,
  204. remove: [],
  205. target: [],
  206. }
  207. return {
  208. calls,
  209. adapter: {
  210. name: "recorded",
  211. description: "recorded",
  212. configure(info) {
  213. calls.configure.push(structuredClone(info))
  214. return input.configure?.(info) ?? info
  215. },
  216. async create(info, env, from) {
  217. calls.create.push({
  218. info: structuredClone(info),
  219. env: { ...env },
  220. from: from ? structuredClone(from) : undefined,
  221. })
  222. await input.create?.(info, env, from)
  223. },
  224. ...(input.list
  225. ? {
  226. async list() {
  227. calls.list += 1
  228. return input.list?.() ?? []
  229. },
  230. }
  231. : {}),
  232. async remove(info) {
  233. calls.remove.push(structuredClone(info))
  234. await input.remove?.(info)
  235. },
  236. target(info) {
  237. calls.target.push(structuredClone(info))
  238. return input.target(info)
  239. },
  240. },
  241. }
  242. }
  243. function localAdapter(dir: string, input?: { createDir?: boolean; remove?: (info: WorkspaceInfo) => Promise<void> }) {
  244. return recordedAdapter({
  245. configure(info) {
  246. return { ...info, directory: dir }
  247. },
  248. async create() {
  249. if (input?.createDir === false) return
  250. await fs.mkdir(dir, { recursive: true })
  251. },
  252. remove: input?.remove,
  253. target() {
  254. return { type: "local", directory: dir }
  255. },
  256. })
  257. }
  258. function remoteAdapter(url: string, input?: { directory?: string | null; headers?: HeadersInit }) {
  259. return recordedAdapter({
  260. configure(info) {
  261. return { ...info, directory: input?.directory ?? info.directory }
  262. },
  263. target() {
  264. return { type: "remote", url, headers: input?.headers }
  265. },
  266. })
  267. }
  268. function eventStreamResponse(events: unknown[] = [], keepOpen = true) {
  269. const encoder = new TextEncoder()
  270. return new Response(
  271. new ReadableStream<Uint8Array>({
  272. start(controller) {
  273. if (keepOpen) controller.enqueue(encoder.encode(":\n\n"))
  274. events.forEach((event) => controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)))
  275. if (!keepOpen) controller.close()
  276. },
  277. }),
  278. { status: 200, headers: { "content-type": "text/event-stream" } },
  279. )
  280. }
  281. function serverUrl() {
  282. return Effect.gen(function* () {
  283. return HttpServer.formatAddress((yield* HttpServer.HttpServer).address)
  284. })
  285. }
  286. function workspaceInfo(projectID: ProjectID, type: string, input?: Partial<Workspace.Info>): Workspace.Info {
  287. return {
  288. id: input?.id ?? WorkspaceID.ascending(),
  289. type,
  290. name: input?.name ?? unique("workspace"),
  291. branch: input?.branch ?? null,
  292. directory: input?.directory ?? null,
  293. extra: input?.extra ?? null,
  294. projectID,
  295. timeUsed: input?.timeUsed ?? Date.now(),
  296. }
  297. }
  298. function insertWorkspace(info: Workspace.Info) {
  299. Database.use((db) =>
  300. db
  301. .insert(WorkspaceTable)
  302. .values({
  303. id: info.id,
  304. type: info.type,
  305. branch: info.branch,
  306. name: info.name,
  307. directory: info.directory,
  308. extra: info.extra,
  309. project_id: info.projectID,
  310. time_used: info.timeUsed,
  311. })
  312. .run(),
  313. )
  314. }
  315. function insertProject(id: ProjectID, worktree: string) {
  316. Database.use((db) =>
  317. db
  318. .insert(ProjectTable)
  319. .values({
  320. id,
  321. worktree,
  322. vcs: null,
  323. name: null,
  324. time_created: Date.now(),
  325. time_updated: Date.now(),
  326. sandboxes: [],
  327. })
  328. .run(),
  329. )
  330. }
  331. function attachSessionToWorkspace(sessionID: SessionID, workspaceID: WorkspaceID) {
  332. Database.use((db) =>
  333. db.update(SessionTable).set({ workspace_id: workspaceID }).where(eq(SessionTable.id, sessionID)).run(),
  334. )
  335. }
  336. function sessionSequence(sessionID: SessionID) {
  337. return Database.use((db) =>
  338. db
  339. .select({ seq: EventSequenceTable.seq })
  340. .from(EventSequenceTable)
  341. .where(eq(EventSequenceTable.aggregate_id, sessionID))
  342. .get(),
  343. )?.seq
  344. }
  345. function sessionSequenceOwner(sessionID: SessionID) {
  346. return Database.use((db) =>
  347. db
  348. .select({ ownerID: EventSequenceTable.owner_id })
  349. .from(EventSequenceTable)
  350. .where(eq(EventSequenceTable.aggregate_id, sessionID))
  351. .get(),
  352. )?.ownerID
  353. }
  354. function sessionUpdatedType() {
  355. return SyncEvent.versionedType(SessionNs.Event.Updated.type, SessionNs.Event.Updated.version)
  356. }
  357. describe("workspace schemas and exports", () => {
  358. test("keeps the historical event type names", () => {
  359. expect(Workspace.Event.Ready.type).toBe("workspace.ready")
  360. expect(Workspace.Event.Failed.type).toBe("workspace.failed")
  361. expect(Workspace.Event.Status.type).toBe("workspace.status")
  362. })
  363. test("validates create input with workspace id, project id, branch, type, and extra", () => {
  364. const input = {
  365. id: WorkspaceID.ascending("wrk_schema_create"),
  366. type: "worktree",
  367. branch: "feature/schema",
  368. projectID: ProjectID.make("project-schema"),
  369. extra: { nested: true },
  370. }
  371. const decode = Schema.decodeUnknownSync(Workspace.CreateInput)
  372. expect(decode(input)).toEqual(input)
  373. expect(() => decode({ ...input, id: 1 })).toThrow()
  374. expect(() => decode({ ...input, branch: 1 })).toThrow()
  375. })
  376. })
  377. describe("workspace CRUD", () => {
  378. test("get returns undefined for a missing workspace", async () => {
  379. await withInstance(async () => {
  380. expect(await getWorkspace(WorkspaceID.ascending("wrk_missing_get"))).toBeUndefined()
  381. })
  382. })
  383. test("list maps database rows, filters by project, and sorts by id", async () => {
  384. await withInstance(async () => {
  385. const otherProjectID = ProjectID.make("project-other")
  386. insertProject(otherProjectID, "/tmp/other")
  387. const a = workspaceInfo(Instance.project.id, "manual", {
  388. id: WorkspaceID.ascending("wrk_a_list"),
  389. branch: "a",
  390. directory: "/a",
  391. extra: { a: true },
  392. })
  393. const b = workspaceInfo(Instance.project.id, "manual", {
  394. id: WorkspaceID.ascending("wrk_b_list"),
  395. branch: "b",
  396. directory: "/b",
  397. extra: ["b"],
  398. })
  399. const other = workspaceInfo(otherProjectID, "manual", { id: WorkspaceID.ascending("wrk_c_list") })
  400. insertWorkspace(b)
  401. insertWorkspace(other)
  402. insertWorkspace(a)
  403. expect(await listWorkspaces(Instance.project)).toEqual([a, b])
  404. })
  405. })
  406. test("create configures, persists, creates, starts local sync, and passes environment", async () => {
  407. await withInstance(async (dir) => {
  408. process.env.OPENCODE_AUTH_CONTENT = JSON.stringify({ test: { type: "api", key: "secret" } })
  409. process.env.OTEL_EXPORTER_OTLP_HEADERS = "authorization=otel"
  410. process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "https://otel.test"
  411. process.env.OTEL_RESOURCE_ATTRIBUTES = "service.name=opencode-test"
  412. const workspaceID = WorkspaceID.ascending("wrk_create_local")
  413. const type = unique("create-local")
  414. const targetDir = path.join(dir, "created-local")
  415. const recorded = recordedAdapter({
  416. configure(info) {
  417. return {
  418. ...info,
  419. branch: "configured-branch",
  420. name: "Configured Name",
  421. directory: targetDir,
  422. extra: { configured: true },
  423. }
  424. },
  425. async create() {
  426. await fs.mkdir(targetDir, { recursive: true })
  427. },
  428. target() {
  429. return { type: "local", directory: targetDir }
  430. },
  431. })
  432. registerAdapter(Instance.project.id, type, recorded.adapter)
  433. const info = await createWorkspace({
  434. id: workspaceID,
  435. type,
  436. branch: null,
  437. projectID: Instance.project.id,
  438. extra: null,
  439. })
  440. expect(info).toEqual({
  441. id: workspaceID,
  442. type,
  443. branch: "configured-branch",
  444. name: "Configured Name",
  445. directory: targetDir,
  446. extra: { configured: true },
  447. projectID: Instance.project.id,
  448. timeUsed: info.timeUsed,
  449. })
  450. expect(await getWorkspace(workspaceID)).toEqual(info)
  451. expect(await listWorkspaces(Instance.project)).toEqual([info])
  452. expect(recorded.calls.configure).toHaveLength(1)
  453. expect(recorded.calls.configure[0]).toMatchObject({ id: workspaceID, type, directory: null })
  454. expect(recorded.calls.create).toHaveLength(1)
  455. expect(recorded.calls.create[0].info).toEqual({
  456. id: workspaceID,
  457. type,
  458. branch: "configured-branch",
  459. name: "Configured Name",
  460. directory: targetDir,
  461. extra: { configured: true },
  462. projectID: Instance.project.id,
  463. })
  464. expect(JSON.parse(recorded.calls.create[0].env.OPENCODE_AUTH_CONTENT ?? "{}")).toEqual({
  465. test: { type: "api", key: "secret" },
  466. })
  467. expect(recorded.calls.create[0].env.OPENCODE_WORKSPACE_ID).toBe(workspaceID)
  468. expect(recorded.calls.create[0].env.OPENCODE_EXPERIMENTAL_WORKSPACES).toBe("true")
  469. expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_HEADERS).toBe("authorization=otel")
  470. expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_ENDPOINT).toBe("https://otel.test")
  471. expect(recorded.calls.create[0].env.OTEL_RESOURCE_ATTRIBUTES).toBe("service.name=opencode-test")
  472. expect((await workspaceStatus()).find((item) => item.workspaceID === workspaceID)?.status).toBe("connected")
  473. await removeWorkspace(workspaceID)
  474. expect((await workspaceStatus()).find((item) => item.workspaceID === workspaceID)?.status).toBeUndefined()
  475. })
  476. })
  477. test("create propagates configure failures and does not insert a workspace", async () => {
  478. await withInstance(async () => {
  479. const type = unique("configure-failure")
  480. registerAdapter(
  481. Instance.project.id,
  482. type,
  483. recordedAdapter({
  484. configure() {
  485. throw new Error("configure exploded")
  486. },
  487. target() {
  488. return { type: "local", directory: "/unused" }
  489. },
  490. }).adapter,
  491. )
  492. await expect(
  493. createWorkspace({ type, branch: null, projectID: Instance.project.id, extra: null }),
  494. ).rejects.toThrow("configure exploded")
  495. expect(await listWorkspaces(Instance.project)).toEqual([])
  496. })
  497. })
  498. test("create leaves the inserted row when adapter create fails", async () => {
  499. await withInstance(async () => {
  500. const type = unique("create-failure")
  501. const recorded = recordedAdapter({
  502. async create() {
  503. throw new Error("create exploded")
  504. },
  505. target() {
  506. return { type: "local", directory: "/unused" }
  507. },
  508. })
  509. registerAdapter(Instance.project.id, type, recorded.adapter)
  510. await expect(
  511. createWorkspace({ type, branch: "branch", projectID: Instance.project.id, extra: { x: 1 } }),
  512. ).rejects.toThrow("create exploded")
  513. const rows = await listWorkspaces(Instance.project)
  514. expect(rows).toHaveLength(1)
  515. expect(rows[0]).toMatchObject({ type, branch: "branch", extra: { x: 1 } })
  516. expect(recorded.calls.target).toHaveLength(0)
  517. await removeWorkspace(rows[0].id)
  518. })
  519. })
  520. test("create returns after a local workspace reports error", async () => {
  521. await withInstance(async (dir) => {
  522. const type = unique("local-error")
  523. const missing = path.join(dir, "missing-local-target")
  524. const recorded = localAdapter(missing, { createDir: false })
  525. registerAdapter(Instance.project.id, type, recorded.adapter)
  526. const info = await createWorkspace({ type, branch: null, projectID: Instance.project.id, extra: null })
  527. expect(info.directory).toBe(missing)
  528. expect((await workspaceStatus()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
  529. await removeWorkspace(info.id)
  530. })
  531. })
  532. test("syncList registers adapter-listed workspaces that are missing by name", async () => {
  533. await withInstance(async (dir) => {
  534. const type = unique("list-sync")
  535. const existing = workspaceInfo(Instance.project.id, type, {
  536. id: WorkspaceID.ascending("wrk_list_sync_existing"),
  537. name: "existing",
  538. directory: path.join(dir, "existing"),
  539. })
  540. insertWorkspace(existing)
  541. const discovered = {
  542. type,
  543. name: "discovered",
  544. branch: "feature/discovered",
  545. directory: path.join(dir, "discovered"),
  546. extra: { source: "adapter" },
  547. projectID: Instance.project.id,
  548. }
  549. const recorded = recordedAdapter({
  550. list() {
  551. return [
  552. {
  553. type,
  554. name: existing.name,
  555. branch: "ignored",
  556. directory: path.join(dir, "ignored"),
  557. extra: null,
  558. projectID: Instance.project.id,
  559. },
  560. discovered,
  561. ]
  562. },
  563. target(info) {
  564. return { type: "local", directory: info.directory ?? dir }
  565. },
  566. })
  567. registerAdapter(Instance.project.id, type, recorded.adapter)
  568. await syncListWorkspaces(Instance.project)
  569. const synced = (await listWorkspaces(Instance.project)).filter((item) => item.name === discovered.name)
  570. expect(synced).toHaveLength(1)
  571. expect(synced[0]).toMatchObject(discovered)
  572. expect(synced[0]?.id).toStartWith("wrk_")
  573. expect(await listWorkspaces(Instance.project)).toEqual(expect.arrayContaining([existing, synced[0]]))
  574. expect(recorded.calls.list).toBe(1)
  575. expect(recorded.calls.configure).toHaveLength(0)
  576. expect(recorded.calls.create).toHaveLength(0)
  577. expect(recorded.calls.target).toHaveLength(1)
  578. })
  579. })
  580. test("syncList calls every registered adapter with a list method", async () => {
  581. await withInstance(async (dir) => {
  582. const typeA = unique("list-sync-a")
  583. const typeB = unique("list-sync-b")
  584. const adapterA = recordedAdapter({
  585. list() {
  586. return [
  587. {
  588. type: typeA,
  589. name: "adapter-a",
  590. branch: null,
  591. directory: path.join(dir, "adapter-a"),
  592. extra: null,
  593. projectID: Instance.project.id,
  594. },
  595. ]
  596. },
  597. target(info) {
  598. return { type: "local", directory: info.directory ?? dir }
  599. },
  600. })
  601. const adapterB = recordedAdapter({
  602. list() {
  603. return [
  604. {
  605. type: typeB,
  606. name: "adapter-b",
  607. branch: null,
  608. directory: path.join(dir, "adapter-b"),
  609. extra: null,
  610. projectID: Instance.project.id,
  611. },
  612. ]
  613. },
  614. target(info) {
  615. return { type: "local", directory: info.directory ?? dir }
  616. },
  617. })
  618. const noList = recordedAdapter({
  619. target() {
  620. return { type: "local", directory: dir }
  621. },
  622. })
  623. registerAdapter(Instance.project.id, typeA, adapterA.adapter)
  624. registerAdapter(Instance.project.id, typeB, adapterB.adapter)
  625. registerAdapter(Instance.project.id, unique("list-sync-none"), noList.adapter)
  626. await syncListWorkspaces(Instance.project)
  627. const synced = await listWorkspaces(Instance.project)
  628. expect(
  629. synced
  630. .filter((item) => item.type === typeA || item.type === typeB)
  631. .map((item) => item.name)
  632. .toSorted(),
  633. ).toEqual(["adapter-a", "adapter-b"])
  634. expect(adapterA.calls.list).toBe(1)
  635. expect(adapterB.calls.list).toBe(1)
  636. expect(noList.calls.list).toBe(0)
  637. })
  638. })
  639. it.live("remote create connects to routed event and history endpoints", () => {
  640. const calls: FetchCall[] = []
  641. return Effect.gen(function* () {
  642. yield* HttpServer.serveEffect()(
  643. Effect.gen(function* () {
  644. const req = yield* HttpServerRequest.HttpServerRequest
  645. const bodyText = yield* req.text
  646. const call = {
  647. url: new URL(req.url, "http://localhost"),
  648. method: req.method,
  649. headers: new Headers(req.headers),
  650. bodyText,
  651. json: bodyText ? JSON.parse(bodyText) : undefined,
  652. }
  653. calls.push(call)
  654. if (call.url.pathname === "/base/global/event")
  655. return HttpServerResponse.fromWeb(eventStreamResponse([], false))
  656. if (call.url.pathname === "/base/sync/history") return yield* HttpServerResponse.json([])
  657. return HttpServerResponse.text("unexpected", { status: 500 })
  658. }),
  659. )
  660. const url = yield* serverUrl()
  661. yield* provideTmpdirInstance(
  662. (dir) =>
  663. Effect.gen(function* () {
  664. const workspace = yield* Workspace.Service
  665. const type = unique("remote-create")
  666. const recorded = remoteAdapter(`${url}/base/?ignored=1#hash`, { directory: dir })
  667. registerAdapter(Instance.project.id, type, recorded.adapter)
  668. const info = yield* workspace.create({ type, branch: null, projectID: Instance.project.id, extra: null })
  669. expect(
  670. calls.map((call) => `${call.method} ${call.url.pathname}${call.url.search}${call.url.hash}`),
  671. ).toEqual(["GET /base/global/event", "POST /base/sync/history"])
  672. expect(calls[1].json).toEqual({})
  673. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("connected")
  674. expect(yield* workspace.isSyncing(info.id)).toBe(true)
  675. yield* workspace.remove(info.id)
  676. expect(yield* workspace.isSyncing(info.id)).toBe(false)
  677. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
  678. }),
  679. { git: true },
  680. )
  681. })
  682. })
  683. test("remove returns undefined for a missing workspace", async () => {
  684. await withInstance(async () => {
  685. expect(await removeWorkspace(WorkspaceID.ascending("wrk_missing_remove"))).toBeUndefined()
  686. })
  687. })
  688. it.instance(
  689. "remove deletes the workspace, associated sessions, adapter resources, and status",
  690. () => {
  691. return Effect.gen(function* () {
  692. const { directory: dir } = yield* TestInstance
  693. const instance = yield* InstanceRef
  694. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  695. const workspace = yield* Workspace.Service
  696. const sessionSvc = yield* SessionNs.Service
  697. const type = unique("remove-local")
  698. const recorded = localAdapter(path.join(dir, "remove-local"))
  699. registerAdapter(instance.project.id, type, recorded.adapter)
  700. const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })
  701. const one = yield* sessionSvc.create({})
  702. const two = yield* sessionSvc.create({})
  703. attachSessionToWorkspace(one.id, info.id)
  704. attachSessionToWorkspace(two.id, info.id)
  705. const removed = yield* workspace.remove(info.id)
  706. expect(removed).toEqual(info)
  707. expect(yield* workspace.get(info.id)).toBeUndefined()
  708. expect(recorded.calls.remove).toEqual([info])
  709. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
  710. expect(
  711. Database.use((db) =>
  712. db.select({ id: SessionTable.id }).from(SessionTable).where(eq(SessionTable.workspace_id, info.id)).all(),
  713. ),
  714. ).toEqual([])
  715. })
  716. },
  717. { git: true },
  718. )
  719. test("remove still deletes the row when the adapter cannot remove resources", async () => {
  720. await withInstance(async () => {
  721. const type = unique("remove-throws")
  722. const info = workspaceInfo(Instance.project.id, type, { id: WorkspaceID.ascending("wrk_remove_throws") })
  723. registerAdapter(
  724. Instance.project.id,
  725. type,
  726. recordedAdapter({
  727. async remove() {
  728. throw new Error("remove exploded")
  729. },
  730. target() {
  731. return { type: "local", directory: "/unused" }
  732. },
  733. }).adapter,
  734. )
  735. insertWorkspace(info)
  736. expect(await removeWorkspace(info.id)).toEqual(info)
  737. expect(await getWorkspace(info.id)).toBeUndefined()
  738. })
  739. })
  740. it.instance(
  741. "sessionWarp moves a session into a local workspace and claims ownership",
  742. () => {
  743. return Effect.gen(function* () {
  744. const { directory: dir } = yield* TestInstance
  745. const instance = yield* InstanceRef
  746. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  747. const workspace = yield* Workspace.Service
  748. const sessionSvc = yield* SessionNs.Service
  749. const previousType = unique("warp-prev-local")
  750. const targetType = unique("warp-target-local")
  751. const previous = workspaceInfo(instance.project.id, previousType)
  752. const target = workspaceInfo(instance.project.id, targetType)
  753. insertWorkspace(previous)
  754. insertWorkspace(target)
  755. registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-prev-local")).adapter)
  756. registerAdapter(instance.project.id, targetType, localAdapter(path.join(dir, "warp-target-local")).adapter)
  757. const session = yield* sessionSvc.create({})
  758. attachSessionToWorkspace(session.id, previous.id)
  759. yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id })
  760. expect(
  761. Database.use((db) =>
  762. db
  763. .select({ workspaceID: SessionTable.workspace_id })
  764. .from(SessionTable)
  765. .where(eq(SessionTable.id, session.id))
  766. .get(),
  767. )?.workspaceID,
  768. ).toBe(target.id)
  769. expect(sessionSequenceOwner(session.id)).toBe(target.id)
  770. })
  771. },
  772. { git: true },
  773. )
  774. it.instance(
  775. "sessionWarp applies source workspace patch to local target workspace",
  776. () => {
  777. return Effect.gen(function* () {
  778. const { directory: dir } = yield* TestInstance
  779. const instance = yield* InstanceRef
  780. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  781. const workspace = yield* Workspace.Service
  782. const sessionSvc = yield* SessionNs.Service
  783. const previousType = unique("warp-patch-prev-local")
  784. const targetType = unique("warp-patch-target-local")
  785. const previousDir = path.join(dir, "warp-patch-prev-local")
  786. const targetDir = path.join(dir, "warp-patch-target-local")
  787. yield* Effect.promise(() => initGitRepo(previousDir))
  788. yield* Effect.promise(() => initGitRepo(targetDir))
  789. yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "tracked.txt"), "changed\n"))
  790. yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "new.txt"), "new\n"))
  791. const previous = workspaceInfo(instance.project.id, previousType)
  792. const target = workspaceInfo(instance.project.id, targetType)
  793. insertWorkspace(previous)
  794. insertWorkspace(target)
  795. registerAdapter(instance.project.id, previousType, localAdapter(previousDir, { createDir: false }).adapter)
  796. registerAdapter(instance.project.id, targetType, localAdapter(targetDir, { createDir: false }).adapter)
  797. const session = yield* sessionSvc.create({})
  798. attachSessionToWorkspace(session.id, previous.id)
  799. yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true })
  800. expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "tracked.txt"), "utf8"))).toBe("changed\n")
  801. expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "new.txt"), "utf8"))).toBe("new\n")
  802. })
  803. },
  804. { git: true },
  805. )
  806. it.instance(
  807. "sessionWarp detaches a session to the local project and claims project ownership",
  808. () => {
  809. return Effect.gen(function* () {
  810. const { directory: dir } = yield* TestInstance
  811. const instance = yield* InstanceRef
  812. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  813. const workspace = yield* Workspace.Service
  814. const sessionSvc = yield* SessionNs.Service
  815. const previousType = unique("warp-detach-local")
  816. const previous = workspaceInfo(instance.project.id, previousType)
  817. insertWorkspace(previous)
  818. registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-detach-local")).adapter)
  819. const session = yield* sessionSvc.create({})
  820. attachSessionToWorkspace(session.id, previous.id)
  821. yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id })
  822. expect(
  823. Database.use((db) =>
  824. db
  825. .select({ workspaceID: SessionTable.workspace_id })
  826. .from(SessionTable)
  827. .where(eq(SessionTable.id, session.id))
  828. .get(),
  829. )?.workspaceID,
  830. ).toBeNull()
  831. expect(sessionSequenceOwner(session.id)).toBe(instance.project.id)
  832. })
  833. },
  834. { git: true },
  835. )
  836. test("sessionWarp detaches to the source project when invoked from a workspace instance", async () => {
  837. await withInstance(async () => {
  838. const projectID = Instance.project.id
  839. await using workspaceTmp = await tmpdir({ git: true })
  840. const previousType = unique("warp-detach-workspace-instance")
  841. const previous = workspaceInfo(projectID, previousType)
  842. insertWorkspace(previous)
  843. registerAdapter(projectID, previousType, localAdapter(workspaceTmp.path, { createDir: false }).adapter)
  844. const session = await AppRuntime.runPromise(SessionNs.Service.use((svc) => svc.create({})))
  845. attachSessionToWorkspace(session.id, previous.id)
  846. const workspaceProjectID = await WithInstance.provide({
  847. directory: workspaceTmp.path,
  848. fn: async () => {
  849. const id = Instance.project.id
  850. expect(id).not.toBe(projectID)
  851. await warpWorkspaceSession({ workspaceID: null, sessionID: session.id })
  852. return id
  853. },
  854. })
  855. expect(
  856. Database.use((db) =>
  857. db
  858. .select({ workspaceID: SessionTable.workspace_id })
  859. .from(SessionTable)
  860. .where(eq(SessionTable.id, session.id))
  861. .get(),
  862. )?.workspaceID,
  863. ).toBeNull()
  864. expect(sessionSequenceOwner(session.id)).toBe(projectID)
  865. expect(sessionSequenceOwner(session.id)).not.toBe(workspaceProjectID)
  866. })
  867. })
  868. it.live("sessionWarp syncs previous remote history, replays it, steals, and claims the sequence", () => {
  869. const calls: FetchCall[] = []
  870. let historySessionID: SessionID | undefined
  871. let historyNextSeq = 0
  872. return Effect.gen(function* () {
  873. yield* HttpServer.serveEffect()(
  874. Effect.gen(function* () {
  875. const req = yield* HttpServerRequest.HttpServerRequest
  876. const bodyText = yield* req.text
  877. const call = {
  878. url: new URL(req.url, "http://localhost"),
  879. method: req.method,
  880. headers: new Headers(req.headers),
  881. bodyText,
  882. json: bodyText ? JSON.parse(bodyText) : undefined,
  883. }
  884. calls.push(call)
  885. if (call.url.pathname === "/warp-source/sync/history") {
  886. return yield* HttpServerResponse.json([
  887. {
  888. id: `evt_${unique("warp-source-history")}`,
  889. aggregate_id: historySessionID!,
  890. seq: historyNextSeq,
  891. type: sessionUpdatedType(),
  892. data: { sessionID: historySessionID!, info: { title: "from source history" } },
  893. },
  894. ])
  895. }
  896. if (call.url.pathname === "/warp-source/vcs/diff/raw") return HttpServerResponse.text("remote patch")
  897. if (call.url.pathname === "/warp-target/sync/replay")
  898. return yield* HttpServerResponse.json({ sessionID: "ok" })
  899. if (call.url.pathname === "/warp-target/sync/steal")
  900. return yield* HttpServerResponse.json({ sessionID: "ok" })
  901. if (call.url.pathname === "/warp-target/vcs/apply") return yield* HttpServerResponse.json({ applied: true })
  902. return HttpServerResponse.text("unexpected", { status: 500 })
  903. }),
  904. )
  905. const url = yield* serverUrl()
  906. yield* provideTmpdirInstance(
  907. () =>
  908. Effect.gen(function* () {
  909. const workspace = yield* Workspace.Service
  910. const sessionSvc = yield* SessionNs.Service
  911. const previousType = unique("warp-remote-source")
  912. const targetType = unique("warp-remote-target")
  913. const previous = workspaceInfo(Instance.project.id, previousType)
  914. const target = workspaceInfo(Instance.project.id, targetType, { directory: "remote-target-dir" })
  915. insertWorkspace(previous)
  916. insertWorkspace(target)
  917. registerAdapter(Instance.project.id, previousType, remoteAdapter(`${url}/warp-source`).adapter)
  918. registerAdapter(Instance.project.id, targetType, remoteAdapter(`${url}/warp-target`).adapter)
  919. const session = yield* sessionSvc.create({})
  920. attachSessionToWorkspace(session.id, previous.id)
  921. historySessionID = session.id
  922. historyNextSeq = (sessionSequence(session.id) ?? -1) + 1
  923. yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true })
  924. expect(calls.map((call) => `${call.method} ${call.url.pathname}`)).toEqual([
  925. "POST /warp-source/sync/history",
  926. "GET /warp-source/vcs/diff/raw",
  927. "POST /warp-target/vcs/apply",
  928. "POST /warp-target/sync/replay",
  929. "POST /warp-target/sync/steal",
  930. ])
  931. expect(calls[0].json).toEqual({ [session.id]: historyNextSeq - 1 })
  932. expect(calls[2].json).toEqual({ patch: "remote patch" })
  933. expect(calls[3].json).toMatchObject({
  934. directory: "remote-target-dir",
  935. events: [
  936. {
  937. aggregateID: session.id,
  938. seq: 0,
  939. type: SyncEvent.versionedType(SessionNs.Event.Created.type, SessionNs.Event.Created.version),
  940. },
  941. {
  942. aggregateID: session.id,
  943. seq: historyNextSeq,
  944. type: sessionUpdatedType(),
  945. },
  946. ],
  947. })
  948. expect(calls[4].json).toEqual({ sessionID: session.id })
  949. expect((yield* sessionSvc.get(session.id)).title).toBe("from source history")
  950. expect(sessionSequenceOwner(session.id)).toBe(target.id)
  951. }),
  952. { git: true },
  953. )
  954. })
  955. })
  956. })
  957. describe("workspace sync state", () => {
  958. it.instance(
  959. "startWorkspaceSyncing is disabled by the experimental workspace flag",
  960. () =>
  961. Effect.gen(function* () {
  962. const { directory: dir } = yield* TestInstance
  963. const instance = yield* InstanceRef
  964. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  965. const workspace = yield* Workspace.Service
  966. const sessionSvc = yield* SessionNs.Service
  967. const type = unique("flag-disabled")
  968. const info = workspaceInfo(instance.project.id, type)
  969. const session = yield* sessionSvc.create({})
  970. attachSessionToWorkspace(session.id, info.id)
  971. insertWorkspace(info)
  972. registerAdapter(instance.project.id, type, localAdapter(path.join(dir, "flag-disabled")).adapter)
  973. yield* Effect.promise(() => startWorkspaceSyncingWithFlag(instance.project.id, false))
  974. yield* Effect.sleep("25 millis")
  975. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
  976. }),
  977. { git: true },
  978. )
  979. it.instance(
  980. "startWorkspaceSyncing starts all workspaces",
  981. () =>
  982. Effect.gen(function* () {
  983. const { directory: dir } = yield* TestInstance
  984. const instance = yield* InstanceRef
  985. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  986. const workspace = yield* Workspace.Service
  987. const projectID = instance.project.id
  988. const firstType = unique("first")
  989. const secondType = unique("second")
  990. const first = workspaceInfo(projectID, firstType)
  991. const second = workspaceInfo(projectID, secondType)
  992. yield* Effect.promise(() => fs.mkdir(path.join(dir, "first"), { recursive: true }))
  993. yield* Effect.promise(() => fs.mkdir(path.join(dir, "second"), { recursive: true }))
  994. yield* Effect.sync(() => {
  995. insertWorkspace(first)
  996. insertWorkspace(second)
  997. registerAdapter(projectID, firstType, localAdapter(path.join(dir, "first")).adapter)
  998. registerAdapter(projectID, secondType, localAdapter(path.join(dir, "second")).adapter)
  999. })
  1000. yield* Effect.addFinalizer(() =>
  1001. Effect.all([workspace.remove(first.id), workspace.remove(second.id)], { discard: true }).pipe(Effect.ignore),
  1002. )
  1003. yield* workspace.startWorkspaceSyncing(projectID)
  1004. yield* eventuallyEffect(
  1005. Effect.gen(function* () {
  1006. const status = yield* workspace.status()
  1007. expect(status.find((item) => item.workspaceID === first.id)?.status).toBe("connected")
  1008. expect(status.find((item) => item.workspaceID === second.id)?.status).toBe("connected")
  1009. }),
  1010. )
  1011. }),
  1012. { git: true },
  1013. )
  1014. it.instance(
  1015. "local start reports error when the target directory is missing",
  1016. () =>
  1017. Effect.gen(function* () {
  1018. const { directory: dir } = yield* TestInstance
  1019. const instance = yield* InstanceRef
  1020. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  1021. const workspace = yield* Workspace.Service
  1022. const sessionSvc = yield* SessionNs.Service
  1023. const type = unique("missing-local")
  1024. const info = workspaceInfo(instance.project.id, type)
  1025. insertWorkspace(info)
  1026. registerAdapter(
  1027. instance.project.id,
  1028. type,
  1029. localAdapter(path.join(dir, "missing-target"), { createDir: false }).adapter,
  1030. )
  1031. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1032. yield* workspace.startWorkspaceSyncing(instance.project.id)
  1033. yield* eventuallyEffect(
  1034. Effect.gen(function* () {
  1035. const status = yield* workspace.status()
  1036. expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("error")
  1037. }),
  1038. )
  1039. expect(yield* workspace.isSyncing(info.id)).toBe(false)
  1040. yield* workspace.remove(info.id)
  1041. }),
  1042. { git: true },
  1043. )
  1044. it.instance(
  1045. "duplicate local status updates are suppressed",
  1046. () =>
  1047. Effect.gen(function* () {
  1048. const { directory: dir } = yield* TestInstance
  1049. const instance = yield* InstanceRef
  1050. if (!instance) return yield* Effect.die(new Error("missing test instance"))
  1051. const workspace = yield* Workspace.Service
  1052. const sessionSvc = yield* SessionNs.Service
  1053. const captured = captureGlobalEvents()
  1054. yield* Effect.addFinalizer(() => Effect.sync(() => captured.dispose()))
  1055. const type = unique("dedupe-local")
  1056. const info = workspaceInfo(instance.project.id, type)
  1057. const target = path.join(dir, "dedupe-local")
  1058. yield* Effect.promise(() => fs.mkdir(target, { recursive: true }))
  1059. insertWorkspace(info)
  1060. registerAdapter(instance.project.id, type, localAdapter(target).adapter)
  1061. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1062. yield* workspace.startWorkspaceSyncing(instance.project.id)
  1063. yield* workspace.startWorkspaceSyncing(instance.project.id)
  1064. yield* eventuallyEffect(
  1065. Effect.gen(function* () {
  1066. const status = yield* workspace.status()
  1067. expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("connected")
  1068. }),
  1069. )
  1070. expect(
  1071. captured.events.filter(
  1072. (event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type,
  1073. ),
  1074. ).toHaveLength(1)
  1075. yield* workspace.remove(info.id)
  1076. }),
  1077. { git: true },
  1078. )
  1079. it.live("remote start emits disconnected, connecting, and connected then refuses duplicate listeners", () => {
  1080. const calls: FetchCall[] = []
  1081. return Effect.gen(function* () {
  1082. yield* HttpServer.serveEffect()(
  1083. Effect.gen(function* () {
  1084. const req = yield* HttpServerRequest.HttpServerRequest
  1085. const bodyText = yield* req.text
  1086. const call = {
  1087. url: new URL(req.url, "http://localhost"),
  1088. method: req.method,
  1089. headers: new Headers(req.headers),
  1090. bodyText,
  1091. json: bodyText ? JSON.parse(bodyText) : undefined,
  1092. }
  1093. calls.push(call)
  1094. if (call.url.pathname === "/sync/global/event") return HttpServerResponse.fromWeb(eventStreamResponse())
  1095. if (call.url.pathname === "/sync/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
  1096. return HttpServerResponse.text("unexpected", { status: 500 })
  1097. }),
  1098. )
  1099. const url = yield* serverUrl()
  1100. yield* provideTmpdirInstance(
  1101. () =>
  1102. Effect.gen(function* () {
  1103. const workspace = yield* Workspace.Service
  1104. const sessionSvc = yield* SessionNs.Service
  1105. const captured = captureGlobalEvents()
  1106. try {
  1107. const type = unique("remote-start")
  1108. const info = workspaceInfo(Instance.project.id, type)
  1109. insertWorkspace(info)
  1110. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sync`).adapter)
  1111. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1112. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1113. yield* eventuallyEffect(
  1114. Effect.gen(function* () {
  1115. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe(
  1116. "connected",
  1117. )
  1118. }),
  1119. )
  1120. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1121. yield* Effect.sleep("25 millis")
  1122. expect(
  1123. captured.events
  1124. .filter((event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type)
  1125. .map((event) => event.payload.properties.status),
  1126. ).toEqual(["disconnected", "connecting", "connected"])
  1127. expect(calls.filter((call) => call.url.pathname === "/sync/global/event")).toHaveLength(1)
  1128. expect(calls.filter((call) => call.url.pathname === "/sync/sync/history")).toHaveLength(1)
  1129. expect(yield* workspace.isSyncing(info.id)).toBe(true)
  1130. yield* workspace.remove(info.id)
  1131. expect(yield* workspace.isSyncing(info.id)).toBe(false)
  1132. } finally {
  1133. captured.dispose()
  1134. }
  1135. }),
  1136. { git: true },
  1137. )
  1138. })
  1139. })
  1140. it.live("remote connection HTTP failures set error and clear syncing", () =>
  1141. Effect.gen(function* () {
  1142. yield* HttpServer.serveEffect()(
  1143. Effect.gen(function* () {
  1144. const req = yield* HttpServerRequest.HttpServerRequest
  1145. if (new URL(req.url, "http://localhost").pathname === "/failed/global/event")
  1146. return HttpServerResponse.text("nope", { status: 503 })
  1147. return HttpServerResponse.fromWeb(Response.json([]))
  1148. }),
  1149. )
  1150. const url = yield* serverUrl()
  1151. yield* provideTmpdirInstance(
  1152. () =>
  1153. Effect.gen(function* () {
  1154. const workspace = yield* Workspace.Service
  1155. const sessionSvc = yield* SessionNs.Service
  1156. const type = unique("remote-connect-fail")
  1157. const info = workspaceInfo(Instance.project.id, type)
  1158. insertWorkspace(info)
  1159. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/failed`).adapter)
  1160. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1161. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1162. yield* eventuallyEffect(
  1163. Effect.gen(function* () {
  1164. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
  1165. }),
  1166. )
  1167. expect(yield* workspace.isSyncing(info.id)).toBe(false)
  1168. yield* workspace.remove(info.id)
  1169. }),
  1170. { git: true },
  1171. )
  1172. }),
  1173. )
  1174. it.live("remote history HTTP failures set error", () =>
  1175. Effect.gen(function* () {
  1176. yield* HttpServer.serveEffect()(
  1177. Effect.gen(function* () {
  1178. const req = yield* HttpServerRequest.HttpServerRequest
  1179. const url = new URL(req.url, "http://localhost")
  1180. if (url.pathname === "/history-failed/global/event")
  1181. return HttpServerResponse.fromWeb(eventStreamResponse([], false))
  1182. if (url.pathname === "/history-failed/sync/history")
  1183. return HttpServerResponse.text("history failed", { status: 500 })
  1184. return HttpServerResponse.fromWeb(Response.json([]))
  1185. }),
  1186. )
  1187. const url = yield* serverUrl()
  1188. yield* provideTmpdirInstance(
  1189. () =>
  1190. Effect.gen(function* () {
  1191. const workspace = yield* Workspace.Service
  1192. const sessionSvc = yield* SessionNs.Service
  1193. const type = unique("remote-history-fail")
  1194. const info = workspaceInfo(Instance.project.id, type)
  1195. insertWorkspace(info)
  1196. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/history-failed`).adapter)
  1197. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1198. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1199. yield* eventuallyEffect(
  1200. Effect.gen(function* () {
  1201. expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
  1202. }),
  1203. )
  1204. expect(yield* workspace.isSyncing(info.id)).toBe(false)
  1205. yield* workspace.remove(info.id)
  1206. }),
  1207. { git: true },
  1208. )
  1209. }),
  1210. )
  1211. it.live("sync history sends the local sequence fence and replays returned events in workspace context", () => {
  1212. const historyBodies: unknown[] = []
  1213. let historySessionID: SessionID | undefined
  1214. let historyNextSeq = 0
  1215. return Effect.gen(function* () {
  1216. yield* HttpServer.serveEffect()(
  1217. Effect.gen(function* () {
  1218. const req = yield* HttpServerRequest.HttpServerRequest
  1219. const bodyText = yield* req.text
  1220. const url = new URL(req.url, "http://localhost")
  1221. if (url.pathname === "/history/global/event") return HttpServerResponse.fromWeb(eventStreamResponse())
  1222. if (url.pathname === "/history/sync/history") {
  1223. historyBodies.push(bodyText ? JSON.parse(bodyText) : undefined)
  1224. return HttpServerResponse.fromWeb(
  1225. Response.json([
  1226. {
  1227. id: `evt_${unique("history")}`,
  1228. aggregate_id: historySessionID!,
  1229. seq: historyNextSeq,
  1230. type: sessionUpdatedType(),
  1231. data: { sessionID: historySessionID!, info: { title: "from history" } },
  1232. },
  1233. ]),
  1234. )
  1235. }
  1236. return HttpServerResponse.text("unexpected", { status: 500 })
  1237. }),
  1238. )
  1239. const url = yield* serverUrl()
  1240. yield* provideTmpdirInstance(
  1241. () =>
  1242. Effect.gen(function* () {
  1243. const workspace = yield* Workspace.Service
  1244. const sessionSvc = yield* SessionNs.Service
  1245. const captured = captureGlobalEvents()
  1246. try {
  1247. const type = unique("history-replay")
  1248. const info = workspaceInfo(Instance.project.id, type)
  1249. insertWorkspace(info)
  1250. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/history`).adapter)
  1251. const session = yield* sessionSvc.create({ title: "before history" })
  1252. attachSessionToWorkspace(session.id, info.id)
  1253. historySessionID = session.id
  1254. historyNextSeq = (sessionSequence(session.id) ?? -1) + 1
  1255. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1256. yield* eventuallyEffect(
  1257. Effect.gen(function* () {
  1258. expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from history")
  1259. }),
  1260. )
  1261. expect(historyBodies).toEqual([{ [session.id]: historyNextSeq - 1 }])
  1262. expect(
  1263. captured.events.some(
  1264. (event) =>
  1265. event.workspace === info.id &&
  1266. event.payload.type === "sync" &&
  1267. event.payload.syncEvent.seq === historyNextSeq,
  1268. ),
  1269. ).toBe(true)
  1270. yield* workspace.remove(info.id)
  1271. } finally {
  1272. captured.dispose()
  1273. }
  1274. }),
  1275. { git: true },
  1276. )
  1277. })
  1278. })
  1279. it.live("SSE forwards non-heartbeat events and ignores heartbeats", () =>
  1280. Effect.gen(function* () {
  1281. yield* HttpServer.serveEffect()(
  1282. Effect.gen(function* () {
  1283. const req = yield* HttpServerRequest.HttpServerRequest
  1284. const url = new URL(req.url, "http://localhost")
  1285. if (url.pathname === "/sse-forward/global/event")
  1286. return HttpServerResponse.fromWeb(
  1287. eventStreamResponse(
  1288. [
  1289. { directory: "remote-dir", project: "remote-project", payload: { type: "server.heartbeat" } },
  1290. {
  1291. directory: "remote-dir",
  1292. project: "remote-project",
  1293. payload: { type: "custom.remote", properties: { ok: true } },
  1294. },
  1295. ],
  1296. false,
  1297. ),
  1298. )
  1299. if (url.pathname === "/sse-forward/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
  1300. return HttpServerResponse.text("unexpected", { status: 500 })
  1301. }),
  1302. )
  1303. const url = yield* serverUrl()
  1304. yield* provideTmpdirInstance(
  1305. () =>
  1306. Effect.gen(function* () {
  1307. const workspace = yield* Workspace.Service
  1308. const sessionSvc = yield* SessionNs.Service
  1309. const captured = captureGlobalEvents()
  1310. try {
  1311. const type = unique("sse-forward")
  1312. const info = workspaceInfo(Instance.project.id, type)
  1313. insertWorkspace(info)
  1314. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sse-forward`).adapter)
  1315. attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
  1316. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1317. yield* eventuallyEffect(
  1318. Effect.sync(() =>
  1319. expect(
  1320. captured.events.some(
  1321. (event) => event.workspace === info.id && event.payload.type === "custom.remote",
  1322. ),
  1323. ).toBe(true),
  1324. ),
  1325. )
  1326. expect(
  1327. captured.events.some(
  1328. (event) => event.workspace === info.id && event.payload.type === "server.heartbeat",
  1329. ),
  1330. ).toBe(false)
  1331. expect(
  1332. captured.events.find((event) => event.workspace === info.id && event.payload.type === "custom.remote"),
  1333. ).toMatchObject({
  1334. directory: "remote-dir",
  1335. project: "remote-project",
  1336. payload: { properties: { ok: true } },
  1337. })
  1338. yield* workspace.remove(info.id)
  1339. } finally {
  1340. captured.dispose()
  1341. }
  1342. }),
  1343. { git: true },
  1344. )
  1345. }),
  1346. )
  1347. it.live("SSE sync events are replayed and forwarded", () => {
  1348. let sseSessionID: SessionID | undefined
  1349. let sseNextSeq = 0
  1350. return Effect.gen(function* () {
  1351. yield* HttpServer.serveEffect()(
  1352. Effect.gen(function* () {
  1353. const req = yield* HttpServerRequest.HttpServerRequest
  1354. const url = new URL(req.url, "http://localhost")
  1355. if (url.pathname === "/sse-sync/global/event")
  1356. return HttpServerResponse.fromWeb(
  1357. eventStreamResponse(
  1358. [
  1359. {
  1360. directory: "remote-dir",
  1361. project: "remote-project",
  1362. payload: {
  1363. type: "sync",
  1364. syncEvent: {
  1365. id: `evt_${unique("sse")}`,
  1366. aggregateID: sseSessionID!,
  1367. seq: sseNextSeq,
  1368. type: sessionUpdatedType(),
  1369. data: { sessionID: sseSessionID!, info: { title: "from sse" } },
  1370. },
  1371. },
  1372. },
  1373. ],
  1374. false,
  1375. ),
  1376. )
  1377. if (url.pathname === "/sse-sync/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
  1378. return HttpServerResponse.text("unexpected", { status: 500 })
  1379. }),
  1380. )
  1381. const url = yield* serverUrl()
  1382. yield* provideTmpdirInstance(
  1383. () =>
  1384. Effect.gen(function* () {
  1385. const workspace = yield* Workspace.Service
  1386. const sessionSvc = yield* SessionNs.Service
  1387. const captured = captureGlobalEvents()
  1388. try {
  1389. const type = unique("sse-sync")
  1390. const info = workspaceInfo(Instance.project.id, type)
  1391. insertWorkspace(info)
  1392. registerAdapter(Instance.project.id, type, remoteAdapter(`${url}/sse-sync`).adapter)
  1393. const session = yield* sessionSvc.create({ title: "before sse" })
  1394. attachSessionToWorkspace(session.id, info.id)
  1395. sseSessionID = session.id
  1396. sseNextSeq = (sessionSequence(session.id) ?? -1) + 1
  1397. yield* workspace.startWorkspaceSyncing(Instance.project.id)
  1398. yield* eventuallyEffect(
  1399. Effect.gen(function* () {
  1400. expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from sse")
  1401. }),
  1402. )
  1403. expect(
  1404. captured.events.some(
  1405. (event) =>
  1406. event.workspace === info.id &&
  1407. event.payload.type === "sync" &&
  1408. event.payload.syncEvent.seq === sseNextSeq,
  1409. ),
  1410. ).toBe(true)
  1411. yield* workspace.remove(info.id)
  1412. } finally {
  1413. captured.dispose()
  1414. }
  1415. }),
  1416. { git: true },
  1417. )
  1418. })
  1419. })
  1420. })
  1421. describe("workspace waitForSync", () => {
  1422. test("returns immediately for an empty fence", async () => {
  1423. await withInstance(async () => {
  1424. await expect(waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_empty"), {})).resolves.toBeUndefined()
  1425. })
  1426. })
  1427. test("returns immediately when the stored sequence already satisfies the fence", async () => {
  1428. await withInstance(async () => {
  1429. const sessionID = SessionID.descending("ses_wait_done")
  1430. Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 4 }).run())
  1431. await expect(
  1432. waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_done"), { [sessionID]: 4 }),
  1433. ).resolves.toBeUndefined()
  1434. await expect(
  1435. waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_done_2"), { [sessionID]: 3 }),
  1436. ).resolves.toBeUndefined()
  1437. })
  1438. })
  1439. test("waits until the database reaches the requested sequence and a workspace event arrives", async () => {
  1440. await withInstance(async () => {
  1441. const workspaceID = WorkspaceID.ascending("wrk_wait_event")
  1442. const sessionID = SessionID.descending("ses_wait_event")
  1443. Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 1 }).run())
  1444. const waited = waitForWorkspaceSync(workspaceID, { [sessionID]: 2 })
  1445. await delay(10)
  1446. Database.use((db) =>
  1447. db.update(EventSequenceTable).set({ seq: 2 }).where(eq(EventSequenceTable.aggregate_id, sessionID)).run(),
  1448. )
  1449. GlobalBus.emit("event", { workspace: workspaceID, payload: { type: "anything" } })
  1450. await expect(waited).resolves.toBeUndefined()
  1451. })
  1452. })
  1453. test("a sync event for a different workspace can also release the fence", async () => {
  1454. await withInstance(async () => {
  1455. const workspaceID = WorkspaceID.ascending("wrk_wait_sync_any")
  1456. const sessionID = SessionID.descending("ses_wait_sync_any")
  1457. Database.use((db) => db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 0 }).run())
  1458. const waited = waitForWorkspaceSync(workspaceID, { [sessionID]: 1 })
  1459. await delay(10)
  1460. Database.use((db) =>
  1461. db.update(EventSequenceTable).set({ seq: 1 }).where(eq(EventSequenceTable.aggregate_id, sessionID)).run(),
  1462. )
  1463. GlobalBus.emit("event", {
  1464. workspace: WorkspaceID.ascending("wrk_other_workspace"),
  1465. payload: { type: "sync" },
  1466. })
  1467. await expect(waited).resolves.toBeUndefined()
  1468. })
  1469. })
  1470. test("rejects with the abort reason when aborted", async () => {
  1471. await withInstance(async () => {
  1472. const abort = new AbortController()
  1473. const reason = new Error("caller aborted")
  1474. const waited = waitForWorkspaceSync(
  1475. WorkspaceID.ascending("wrk_wait_abort"),
  1476. { [SessionID.descending("ses_wait_abort")]: 1 },
  1477. abort.signal,
  1478. )
  1479. abort.abort(reason)
  1480. await expect(waited).rejects.toMatchObject({
  1481. _tag: "WorkspaceSyncAbortedError",
  1482. message: reason.message,
  1483. cause: reason,
  1484. })
  1485. })
  1486. })
  1487. test("times out with the requested fence in the error message", async () => {
  1488. await withInstance(async () => {
  1489. const sessionID = SessionID.descending("ses_wait_timeout")
  1490. await expect(waitForWorkspaceSync(WorkspaceID.ascending("wrk_wait_timeout"), { [sessionID]: 1 })).rejects.toThrow(
  1491. `Timed out waiting for sync fence: {"${sessionID}":1}`,
  1492. )
  1493. })
  1494. }, 7000)
  1495. })