workspace.test.ts 64 KB

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