workspace.test.ts 64 KB

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