workspace.test.ts 58 KB

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