| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701 |
- import { afterEach, beforeEach, describe, expect, mock, test } from "bun:test"
- import { $ } from "bun"
- import fs from "node:fs/promises"
- import Http from "node:http"
- import path from "node:path"
- import { NodeHttpServer } from "@effect/platform-node"
- import { Effect, Exit, Fiber, Layer, Schema } from "effect"
- import { HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
- import { eq } from "drizzle-orm"
- import { GlobalBus, type GlobalEvent } from "@/bus/global"
- import { Database } from "@kirincode-ai/core/database/database"
- import { ProjectV2 } from "@kirincode-ai/core/project"
- import { ProjectTable } from "@kirincode-ai/core/project/sql"
- import { AbsolutePath } from "@kirincode-ai/core/schema"
- import { Session as SessionNs } from "@/session/session"
- import { SessionID } from "@/session/schema"
- import { SessionTable } from "@kirincode-ai/core/session/sql"
- import { SessionProjector } from "@kirincode-ai/core/session/projector"
- import { EventSequenceTable } from "@kirincode-ai/core/event/sql"
- import { resetDatabase } from "../fixture/db"
- import { disposeAllInstances, provideTmpdirInstance, requireInstance, TestInstance } from "../fixture/fixture"
- import { testEffect } from "../lib/effect"
- import { registerAdapter } from "../../src/control-plane/adapters"
- import { WorkspaceV2 } from "@kirincode-ai/core/workspace"
- import { WorkspaceTable } from "@kirincode-ai/core/control-plane/workspace.sql"
- import type { Target, WorkspaceAdapter, WorkspaceInfo } from "../../src/control-plane/types"
- import * as Workspace from "../../src/control-plane/workspace"
- import { InstanceStore } from "@/project/instance-store"
- import { InstanceBootstrap } from "@/project/bootstrap"
- import { RuntimeFlags } from "@/effect/runtime-flags"
- import { Ripgrep } from "@kirincode-ai/core/ripgrep"
- import { AppNodeBuilder } from "@kirincode-ai/core/effect/app-node-builder"
- import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
- const originalEnv = {
- KIRINCODE_AUTH_CONTENT: process.env.KIRINCODE_AUTH_CONTENT,
- KIRINCODE_EXPERIMENTAL_WORKSPACES: process.env.KIRINCODE_EXPERIMENTAL_WORKSPACES,
- OTEL_EXPORTER_OTLP_HEADERS: process.env.OTEL_EXPORTER_OTLP_HEADERS,
- OTEL_EXPORTER_OTLP_ENDPOINT: process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
- OTEL_RESOURCE_ATTRIBUTES: process.env.OTEL_RESOURCE_ATTRIBUTES,
- }
- const workspaceLayer = (experimentalWorkspaces: boolean) =>
- AppNodeBuilder.build(
- LayerNode.group([
- Workspace.node,
- SessionNs.node,
- SessionProjector.node,
- Database.node,
- InstanceStore.node,
- Ripgrep.node,
- ]),
- [
- [RuntimeFlags.node, RuntimeFlags.layer({ experimentalWorkspaces })],
- [
- InstanceStore.bootstrapNode,
- Layer.succeed(InstanceBootstrap.Service, InstanceBootstrap.Service.of({ run: Effect.void })),
- ],
- ],
- )
- const testServerLayer = Layer.mergeAll(
- NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }),
- workspaceLayer(true),
- )
- const it = testEffect(testServerLayer)
- type RecordedCreate = {
- info: WorkspaceInfo
- env: Record<string, string | undefined>
- from?: WorkspaceInfo
- }
- type RecordedAdapter = {
- adapter: WorkspaceAdapter
- calls: {
- configure: WorkspaceInfo[]
- create: RecordedCreate[]
- list: number
- remove: WorkspaceInfo[]
- target: WorkspaceInfo[]
- }
- }
- type FetchCall = {
- url: URL
- method: string
- headers: Headers
- bodyText?: string
- json?: unknown
- }
- function unique(prefix: string) {
- return `${prefix}-${Math.random().toString(36).slice(2)}`
- }
- function restoreEnv() {
- Object.entries(originalEnv).forEach(([key, value]) => {
- if (value === undefined) {
- delete process.env[key]
- return
- }
- process.env[key] = value
- })
- }
- beforeEach(() => {
- restoreEnv()
- process.env.KIRINCODE_EXPERIMENTAL_WORKSPACES = "true"
- })
- afterEach(async () => {
- mock.restore()
- await disposeAllInstances()
- restoreEnv()
- await resetDatabase()
- })
- async function initGitRepo(dir: string) {
- await fs.mkdir(dir, { recursive: true })
- await $`git init`.cwd(dir).quiet()
- await $`git config core.fsmonitor false`.cwd(dir).quiet()
- await $`git config commit.gpgsign false`.cwd(dir).quiet()
- await $`git config user.email "test@opencode.test"`.cwd(dir).quiet()
- await $`git config user.name "Test"`.cwd(dir).quiet()
- await fs.writeFile(path.join(dir, "tracked.txt"), "base\n")
- await $`git add tracked.txt`.cwd(dir).quiet()
- await $`git commit -m "base"`.cwd(dir).quiet()
- }
- const startWorkspaceSyncingWithFlag = (projectID: ProjectV2.ID, experimentalWorkspaces: boolean) =>
- Effect.runPromise(
- Workspace.use.startWorkspaceSyncing(projectID).pipe(Effect.provide(workspaceLayer(experimentalWorkspaces))),
- )
- function captureGlobalEvents() {
- const events: GlobalEvent[] = []
- const handler = (event: GlobalEvent) => events.push(event)
- GlobalBus.on("event", handler)
- return {
- events,
- dispose() {
- GlobalBus.off("event", handler)
- },
- }
- }
- function expectExitContains(exit: Exit.Exit<unknown, unknown>, ...messages: string[]) {
- expect(Exit.isFailure(exit)).toBe(true)
- if (!Exit.isFailure(exit)) return
- for (const message of messages) expect(String(exit.cause)).toContain(message)
- }
- function eventuallyEffect(effect: Effect.Effect<void>, timeout = 1500) {
- return Effect.gen(function* () {
- const started = Date.now()
- let last: unknown
- while (Date.now() - started < timeout) {
- const exit = yield* Effect.exit(effect)
- if (exit._tag === "Success") return
- last = exit.cause
- yield* Effect.sleep("10 millis")
- }
- throw last ?? new Error("Timed out waiting for condition")
- })
- }
- function recordedAdapter(input: {
- target: (info: WorkspaceInfo) => Target | Promise<Target>
- configure?: (info: WorkspaceInfo) => WorkspaceInfo | Promise<WorkspaceInfo>
- create?: (info: WorkspaceInfo, env: Record<string, string | undefined>, from?: WorkspaceInfo) => Promise<void>
- list?: () => Omit<WorkspaceInfo, "id">[] | Promise<Omit<WorkspaceInfo, "id">[]>
- remove?: (info: WorkspaceInfo) => Promise<void>
- }): RecordedAdapter {
- const calls: RecordedAdapter["calls"] = {
- configure: [],
- create: [],
- list: 0,
- remove: [],
- target: [],
- }
- return {
- calls,
- adapter: {
- name: "recorded",
- description: "recorded",
- configure(info) {
- calls.configure.push(structuredClone(info))
- return input.configure?.(info) ?? info
- },
- async create(info, env, from) {
- calls.create.push({
- info: structuredClone(info),
- env: { ...env },
- from: from ? structuredClone(from) : undefined,
- })
- await input.create?.(info, env, from)
- },
- ...(input.list
- ? {
- async list() {
- calls.list += 1
- return input.list?.() ?? []
- },
- }
- : {}),
- async remove(info) {
- calls.remove.push(structuredClone(info))
- await input.remove?.(info)
- },
- target(info) {
- calls.target.push(structuredClone(info))
- return input.target(info)
- },
- },
- }
- }
- function localAdapter(dir: string, input?: { createDir?: boolean; remove?: (info: WorkspaceInfo) => Promise<void> }) {
- return recordedAdapter({
- configure(info) {
- return { ...info, directory: dir }
- },
- async create() {
- if (input?.createDir === false) return
- await fs.mkdir(dir, { recursive: true })
- },
- remove: input?.remove,
- target() {
- return { type: "local", directory: dir }
- },
- })
- }
- function remoteAdapter(url: string, input?: { directory?: string | null; headers?: HeadersInit }) {
- return recordedAdapter({
- configure(info) {
- return { ...info, directory: input?.directory ?? info.directory }
- },
- target() {
- return { type: "remote", url, headers: input?.headers }
- },
- })
- }
- function eventStreamResponse(events: unknown[] = [], keepOpen = true) {
- const encoder = new TextEncoder()
- return new Response(
- new ReadableStream<Uint8Array>({
- start(controller) {
- if (keepOpen) controller.enqueue(encoder.encode(":\n\n"))
- events.forEach((event) => controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)))
- if (!keepOpen) controller.close()
- },
- }),
- { status: 200, headers: { "content-type": "text/event-stream" } },
- )
- }
- function serverUrl() {
- return Effect.gen(function* () {
- return HttpServer.formatAddress((yield* HttpServer.HttpServer).address)
- })
- }
- function workspaceInfo(projectID: ProjectV2.ID, type: string, input?: Partial<Workspace.Info>): Workspace.Info {
- return {
- id: input?.id ?? WorkspaceV2.ID.ascending(),
- type,
- name: input?.name ?? unique("workspace"),
- branch: input?.branch ?? null,
- directory: input?.directory ?? null,
- extra: input?.extra ?? null,
- projectID,
- timeUsed: input?.timeUsed ?? Date.now(),
- }
- }
- function insertWorkspace(info: Workspace.Info) {
- return Database.Service.use(({ db }) =>
- db
- .insert(WorkspaceTable)
- .values({
- id: info.id,
- type: info.type,
- branch: info.branch,
- name: info.name,
- directory: info.directory,
- extra: info.extra,
- project_id: info.projectID,
- time_used: info.timeUsed,
- })
- .run()
- .pipe(Effect.orDie),
- )
- }
- function insertProject(id: ProjectV2.ID, worktree: string) {
- return Database.Service.use(({ db }) =>
- db
- .insert(ProjectTable)
- .values({
- id,
- worktree: AbsolutePath.make(worktree),
- vcs: null,
- name: null,
- time_created: Date.now(),
- time_updated: Date.now(),
- sandboxes: [],
- })
- .run()
- .pipe(Effect.orDie),
- )
- }
- function attachSessionToWorkspace(sessionID: SessionID, workspaceID: WorkspaceV2.ID) {
- return Database.Service.use(({ db }) =>
- db
- .update(SessionTable)
- .set({ workspace_id: workspaceID })
- .where(eq(SessionTable.id, sessionID))
- .run()
- .pipe(Effect.orDie),
- )
- }
- function sessionSequence(sessionID: SessionID) {
- return Database.Service.use(({ db }) =>
- db
- .select({ seq: EventSequenceTable.seq })
- .from(EventSequenceTable)
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .get()
- .pipe(
- Effect.orDie,
- Effect.map((row) => row?.seq),
- ),
- )
- }
- function sessionSequenceOwner(sessionID: SessionID) {
- return Database.Service.use(({ db }) =>
- db
- .select({ ownerID: EventSequenceTable.owner_id })
- .from(EventSequenceTable)
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .get()
- .pipe(
- Effect.orDie,
- Effect.map((row) => row?.ownerID),
- ),
- )
- }
- describe("workspace schemas and exports", () => {
- test("keeps the historical event type names", () => {
- expect(Workspace.Event.Ready.type).toBe("workspace.ready")
- expect(Workspace.Event.Failed.type).toBe("workspace.failed")
- expect(Workspace.Event.Status.type).toBe("workspace.status")
- })
- test("validates create input with workspace id, project id, branch, type, and extra", () => {
- const input = {
- id: WorkspaceV2.ID.ascending("wrk_schema_create"),
- type: "worktree",
- branch: "feature/schema",
- projectID: ProjectV2.ID.make("project-schema"),
- extra: { nested: true },
- }
- const decode = Schema.decodeUnknownSync(Workspace.CreateInput)
- expect(decode(input)).toEqual(input)
- expect(() => decode({ ...input, id: 1 })).toThrow()
- expect(() => decode({ ...input, branch: 1 })).toThrow()
- })
- })
- describe("workspace CRUD", () => {
- it.instance(
- "get returns undefined for a missing workspace",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- expect(yield* workspace.get(WorkspaceV2.ID.ascending("wrk_missing_get"))).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "list maps database rows, filters by project, and sorts by id",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const otherProjectID = ProjectV2.ID.make("project-other")
- yield* insertProject(otherProjectID, "/tmp/other")
- const a = workspaceInfo(instance.project.id, "manual", {
- id: WorkspaceV2.ID.ascending("wrk_a_list"),
- branch: "a",
- directory: "/a",
- extra: { a: true },
- })
- const b = workspaceInfo(instance.project.id, "manual", {
- id: WorkspaceV2.ID.ascending("wrk_b_list"),
- branch: "b",
- directory: "/b",
- extra: ["b"],
- })
- const other = workspaceInfo(otherProjectID, "manual", { id: WorkspaceV2.ID.ascending("wrk_c_list") })
- yield* insertWorkspace(b)
- yield* insertWorkspace(other)
- yield* insertWorkspace(a)
- expect(yield* workspace.list(instance.project)).toEqual([a, b])
- }),
- { git: true },
- )
- it.instance(
- "create configures, persists, creates, starts local sync, and passes environment",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- process.env.KIRINCODE_AUTH_CONTENT = JSON.stringify({ test: { type: "api", key: "secret" } })
- process.env.OTEL_EXPORTER_OTLP_HEADERS = "authorization=otel"
- process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "https://otel.test"
- process.env.OTEL_RESOURCE_ATTRIBUTES = "service.name=opencode-test"
- const workspaceID = WorkspaceV2.ID.ascending("wrk_create_local")
- const type = unique("create-local")
- const targetDir = path.join(instance.directory, "created-local")
- const recorded = recordedAdapter({
- configure(info) {
- return {
- ...info,
- branch: "configured-branch",
- name: "Configured Name",
- directory: targetDir,
- extra: { configured: true },
- }
- },
- async create() {
- await fs.mkdir(targetDir, { recursive: true })
- },
- target() {
- return { type: "local", directory: targetDir }
- },
- })
- registerAdapter(instance.project.id, type, recorded.adapter)
- const info = yield* workspace.create({
- id: workspaceID,
- type,
- branch: null,
- projectID: instance.project.id,
- extra: null,
- })
- expect(info).toEqual({
- id: workspaceID,
- type,
- branch: "configured-branch",
- name: "Configured Name",
- directory: targetDir,
- extra: { configured: true },
- projectID: instance.project.id,
- timeUsed: info.timeUsed,
- })
- expect(yield* workspace.get(workspaceID)).toEqual(info)
- expect(yield* workspace.list(instance.project)).toEqual([info])
- expect(recorded.calls.configure).toHaveLength(1)
- expect(recorded.calls.configure[0]).toMatchObject({ id: workspaceID, type, directory: null })
- expect(recorded.calls.create).toHaveLength(1)
- expect(recorded.calls.create[0].info).toEqual({
- id: workspaceID,
- type,
- branch: "configured-branch",
- name: "Configured Name",
- directory: targetDir,
- extra: { configured: true },
- projectID: instance.project.id,
- })
- expect(JSON.parse(recorded.calls.create[0].env.KIRINCODE_AUTH_CONTENT ?? "{}")).toEqual({
- test: { type: "api", key: "secret" },
- })
- expect(recorded.calls.create[0].env.KIRINCODE_WORKSPACE_ID).toBe(workspaceID)
- expect(recorded.calls.create[0].env.KIRINCODE_EXPERIMENTAL_WORKSPACES).toBe("true")
- expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_HEADERS).toBe("authorization=otel")
- expect(recorded.calls.create[0].env.OTEL_EXPORTER_OTLP_ENDPOINT).toBe("https://otel.test")
- expect(recorded.calls.create[0].env.OTEL_RESOURCE_ATTRIBUTES).toBe("service.name=opencode-test")
- expect((yield* workspace.status()).find((item) => item.workspaceID === workspaceID)?.status).toBe("connected")
- yield* workspace.remove(workspaceID)
- expect((yield* workspace.status()).find((item) => item.workspaceID === workspaceID)?.status).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "create propagates configure failures and does not insert a workspace",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const type = unique("configure-failure")
- registerAdapter(
- instance.project.id,
- type,
- recordedAdapter({
- configure() {
- throw new Error("configure exploded")
- },
- target() {
- return { type: "local", directory: "/unused" }
- },
- }).adapter,
- )
- expectExitContains(
- yield* Effect.exit(workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })),
- "configure exploded",
- )
- expect(yield* workspace.list(instance.project)).toEqual([])
- }),
- { git: true },
- )
- it.instance(
- "create leaves the inserted row when adapter create fails",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const type = unique("create-failure")
- const recorded = recordedAdapter({
- async create() {
- throw new Error("create exploded")
- },
- target() {
- return { type: "local", directory: "/unused" }
- },
- })
- registerAdapter(instance.project.id, type, recorded.adapter)
- expectExitContains(
- yield* Effect.exit(
- workspace.create({ type, branch: "branch", projectID: instance.project.id, extra: { x: 1 } }),
- ),
- "create exploded",
- )
- const rows = yield* workspace.list(instance.project)
- expect(rows).toHaveLength(1)
- expect(rows[0]).toMatchObject({ type, branch: "branch", extra: { x: 1 } })
- expect(recorded.calls.target).toHaveLength(0)
- yield* workspace.remove(rows[0].id)
- }),
- { git: true },
- )
- it.instance(
- "create returns after a local workspace reports error",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const type = unique("local-error")
- const missing = path.join(instance.directory, "missing-local-target")
- const recorded = localAdapter(missing, { createDir: false })
- registerAdapter(instance.project.id, type, recorded.adapter)
- const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })
- expect(info.directory).toBe(missing)
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- it.instance(
- "syncList registers adapter-listed workspaces that are missing by name",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const type = unique("list-sync")
- const existing = workspaceInfo(instance.project.id, type, {
- id: WorkspaceV2.ID.ascending("wrk_list_sync_existing"),
- name: "existing",
- directory: path.join(instance.directory, "existing"),
- })
- yield* insertWorkspace(existing)
- const discovered = {
- type,
- name: "discovered",
- branch: "feature/discovered",
- directory: path.join(instance.directory, "discovered"),
- extra: { source: "adapter" },
- projectID: instance.project.id,
- }
- const recorded = recordedAdapter({
- list() {
- return [
- {
- type,
- name: existing.name,
- branch: "ignored",
- directory: path.join(instance.directory, "ignored"),
- extra: null,
- projectID: instance.project.id,
- },
- discovered,
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? instance.directory }
- },
- })
- registerAdapter(instance.project.id, type, recorded.adapter)
- yield* workspace.syncList(instance.project)
- const synced = (yield* workspace.list(instance.project)).filter((item) => item.name === discovered.name)
- expect(synced).toHaveLength(1)
- expect(synced[0]).toMatchObject(discovered)
- expect(synced[0]?.id).toStartWith("wrk_")
- expect(yield* workspace.list(instance.project)).toEqual(expect.arrayContaining([existing, synced[0]]))
- expect(recorded.calls.list).toBe(1)
- expect(recorded.calls.configure).toHaveLength(0)
- expect(recorded.calls.create).toHaveLength(0)
- expect(recorded.calls.target).toHaveLength(1)
- }),
- { git: true },
- )
- it.instance(
- "syncList calls every registered adapter with a list method",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const typeA = unique("list-sync-a")
- const typeB = unique("list-sync-b")
- const adapterA = recordedAdapter({
- list() {
- return [
- {
- type: typeA,
- name: "adapter-a",
- branch: null,
- directory: path.join(instance.directory, "adapter-a"),
- extra: null,
- projectID: instance.project.id,
- },
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? instance.directory }
- },
- })
- const adapterB = recordedAdapter({
- list() {
- return [
- {
- type: typeB,
- name: "adapter-b",
- branch: null,
- directory: path.join(instance.directory, "adapter-b"),
- extra: null,
- projectID: instance.project.id,
- },
- ]
- },
- target(info) {
- return { type: "local", directory: info.directory ?? instance.directory }
- },
- })
- const noList = recordedAdapter({
- target() {
- return { type: "local", directory: instance.directory }
- },
- })
- registerAdapter(instance.project.id, typeA, adapterA.adapter)
- registerAdapter(instance.project.id, typeB, adapterB.adapter)
- registerAdapter(instance.project.id, unique("list-sync-none"), noList.adapter)
- yield* workspace.syncList(instance.project)
- const synced = yield* workspace.list(instance.project)
- expect(
- synced
- .filter((item) => item.type === typeA || item.type === typeB)
- .map((item) => item.name)
- .toSorted(),
- ).toEqual(["adapter-a", "adapter-b"])
- expect(adapterA.calls.list).toBe(1)
- expect(adapterB.calls.list).toBe(1)
- expect(noList.calls.list).toBe(0)
- }),
- { git: true },
- )
- it.live("remote create connects to routed event and history endpoints", () => {
- const calls: FetchCall[] = []
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- const call = {
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- }
- calls.push(call)
- if (call.url.pathname === "/base/global/event")
- return HttpServerResponse.fromWeb(eventStreamResponse([], false))
- if (call.url.pathname === "/base/sync/history") return yield* HttpServerResponse.json([])
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const instance = yield* requireInstance
- const type = unique("remote-create")
- const recorded = remoteAdapter(`${url}/base/?ignored=1#hash`, { directory: dir })
- registerAdapter(instance.project.id, type, recorded.adapter)
- const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })
- expect(
- calls.map((call) => `${call.method} ${call.url.pathname}${call.url.search}${call.url.hash}`),
- ).toEqual(["GET /base/global/event", "POST /base/sync/history"])
- expect(calls[1].json).toEqual({})
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("connected")
- expect(yield* workspace.isSyncing(info.id)).toBe(true)
- yield* workspace.remove(info.id)
- expect(yield* workspace.isSyncing(info.id)).toBe(false)
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
- }),
- { git: true },
- )
- })
- })
- it.instance(
- "remove returns undefined for a missing workspace",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- expect(yield* workspace.remove(WorkspaceV2.ID.ascending("wrk_missing_remove"))).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "remove deletes the workspace, associated sessions, adapter resources, and status",
- () => {
- return Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("remove-local")
- const recorded = localAdapter(path.join(dir, "remove-local"))
- registerAdapter(instance.project.id, type, recorded.adapter)
- const info = yield* workspace.create({ type, branch: null, projectID: instance.project.id, extra: null })
- const one = yield* sessionSvc.create({})
- const two = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(one.id, info.id)
- yield* attachSessionToWorkspace(two.id, info.id)
- const removed = yield* workspace.remove(info.id)
- expect(removed).toEqual(info)
- expect(yield* workspace.get(info.id)).toBeUndefined()
- expect(recorded.calls.remove).toEqual([info])
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
- const { db } = yield* Database.Service
- expect(
- yield* db
- .select({ id: SessionTable.id })
- .from(SessionTable)
- .where(eq(SessionTable.workspace_id, info.id))
- .all()
- .pipe(Effect.orDie),
- ).toEqual([])
- })
- },
- { git: true },
- )
- it.instance(
- "remove still deletes the row when the adapter cannot remove resources",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const type = unique("remove-throws")
- const info = workspaceInfo(instance.project.id, type, { id: WorkspaceV2.ID.ascending("wrk_remove_throws") })
- registerAdapter(
- instance.project.id,
- type,
- recordedAdapter({
- async remove() {
- throw new Error("remove exploded")
- },
- target() {
- return { type: "local", directory: "/unused" }
- },
- }).adapter,
- )
- yield* insertWorkspace(info)
- expect(yield* workspace.remove(info.id)).toEqual(info)
- expect(yield* workspace.get(info.id)).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "sessionWarp moves a session into a local workspace and claims ownership",
- () => {
- return Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const previousType = unique("warp-prev-local")
- const targetType = unique("warp-target-local")
- const previous = workspaceInfo(instance.project.id, previousType)
- const target = workspaceInfo(instance.project.id, targetType)
- yield* insertWorkspace(previous)
- yield* insertWorkspace(target)
- registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-prev-local")).adapter)
- registerAdapter(instance.project.id, targetType, localAdapter(path.join(dir, "warp-target-local")).adapter)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, previous.id)
- yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id })
- const { db } = yield* Database.Service
- expect(
- (yield* db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get()
- .pipe(Effect.orDie))?.workspaceID,
- ).toBe(target.id)
- expect(yield* sessionSequenceOwner(session.id)).toBe(target.id)
- })
- },
- { git: true },
- )
- it.instance(
- "sessionWarp applies source workspace patch to local target workspace",
- () => {
- return Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const previousType = unique("warp-patch-prev-local")
- const targetType = unique("warp-patch-target-local")
- const previousDir = path.join(dir, "warp-patch-prev-local")
- const targetDir = path.join(dir, "warp-patch-target-local")
- yield* Effect.promise(() => initGitRepo(previousDir))
- yield* Effect.promise(() => initGitRepo(targetDir))
- yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "tracked.txt"), "changed\n"))
- yield* Effect.promise(() => fs.writeFile(path.join(previousDir, "new.txt"), "new\n"))
- const previous = workspaceInfo(instance.project.id, previousType)
- const target = workspaceInfo(instance.project.id, targetType)
- yield* insertWorkspace(previous)
- yield* insertWorkspace(target)
- registerAdapter(instance.project.id, previousType, localAdapter(previousDir, { createDir: false }).adapter)
- registerAdapter(instance.project.id, targetType, localAdapter(targetDir, { createDir: false }).adapter)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, previous.id)
- yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true })
- expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "tracked.txt"), "utf8"))).toBe("changed\n")
- expect(yield* Effect.promise(() => fs.readFile(path.join(targetDir, "new.txt"), "utf8"))).toBe("new\n")
- })
- },
- { git: true },
- )
- it.instance(
- "sessionWarp detaches a session to the local project and claims project ownership",
- () => {
- return Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const previousType = unique("warp-detach-local")
- const previous = workspaceInfo(instance.project.id, previousType)
- yield* insertWorkspace(previous)
- registerAdapter(instance.project.id, previousType, localAdapter(path.join(dir, "warp-detach-local")).adapter)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, previous.id)
- yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id })
- const { db } = yield* Database.Service
- expect(
- (yield* db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get()
- .pipe(Effect.orDie))?.workspaceID,
- ).toBeNull()
- expect(yield* sessionSequenceOwner(session.id)).toBe(instance.project.id)
- })
- },
- { git: true },
- )
- const itCrossInstance = process.platform === "win32" ? it.instance.skip : it.instance
- itCrossInstance(
- "sessionWarp detaches to the source project when invoked from a workspace instance",
- () =>
- Effect.gen(function* () {
- const instance = yield* requireInstance
- const projectID = instance.project.id
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const previousType = unique("warp-detach-workspace-instance")
- const previous = workspaceInfo(projectID, previousType)
- yield* insertWorkspace(previous)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, previous.id)
- const workspaceProjectID = yield* provideTmpdirInstance(
- (workspaceDir) =>
- Effect.gen(function* () {
- registerAdapter(projectID, previousType, localAdapter(workspaceDir, { createDir: false }).adapter)
- const workspaceCtx = yield* requireInstance
- expect(workspaceCtx.project.id).not.toBe(projectID)
- yield* workspace.sessionWarp({ workspaceID: null, sessionID: session.id })
- return workspaceCtx.project.id
- }),
- { git: true },
- )
- const { db } = yield* Database.Service
- expect(
- (yield* db
- .select({ workspaceID: SessionTable.workspace_id })
- .from(SessionTable)
- .where(eq(SessionTable.id, session.id))
- .get()
- .pipe(Effect.orDie))?.workspaceID,
- ).toBeNull()
- expect(yield* sessionSequenceOwner(session.id)).toBe(projectID)
- expect(yield* sessionSequenceOwner(session.id)).not.toBe(workspaceProjectID)
- }),
- { git: true },
- )
- it.live("sessionWarp syncs previous remote history, replays it, steals, and claims the sequence", () => {
- const calls: FetchCall[] = []
- let historySessionID: SessionID | undefined
- let historySession: SessionNs.Info | undefined
- let historyNextSeq = 0
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- const call = {
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- }
- calls.push(call)
- if (call.url.pathname === "/warp-source/sync/history") {
- return yield* HttpServerResponse.json([
- {
- id: `evt_${unique("warp-source-history")}`,
- aggregate_id: historySessionID!,
- seq: historyNextSeq,
- type: "session.updated.1",
- data: { sessionID: historySessionID!, info: historySession! },
- },
- ])
- }
- if (call.url.pathname === "/warp-source/vcs/diff/raw") return HttpServerResponse.text("remote patch")
- if (call.url.pathname === "/warp-target/sync/replay")
- return yield* HttpServerResponse.json({ sessionID: "ok" })
- if (call.url.pathname === "/warp-target/sync/steal")
- return yield* HttpServerResponse.json({ sessionID: "ok" })
- if (call.url.pathname === "/warp-target/vcs/apply") return yield* HttpServerResponse.json({ applied: true })
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const previousType = unique("warp-remote-source")
- const targetType = unique("warp-remote-target")
- const previous = workspaceInfo(instance.project.id, previousType)
- const target = workspaceInfo(instance.project.id, targetType, { directory: "remote-target-dir" })
- yield* insertWorkspace(previous)
- yield* insertWorkspace(target)
- registerAdapter(instance.project.id, previousType, remoteAdapter(`${url}/warp-source`).adapter)
- registerAdapter(instance.project.id, targetType, remoteAdapter(`${url}/warp-target`).adapter)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, previous.id)
- historySessionID = session.id
- historySession = { ...session, workspaceID: previous.id, title: "from source history" }
- historyNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1
- yield* workspace.sessionWarp({ workspaceID: target.id, sessionID: session.id, copyChanges: true })
- expect(calls.map((call) => `${call.method} ${call.url.pathname}`)).toEqual([
- "POST /warp-source/sync/history",
- "GET /warp-source/vcs/diff/raw",
- "POST /warp-target/vcs/apply",
- "POST /warp-target/sync/replay",
- "POST /warp-target/sync/steal",
- ])
- expect(calls[0].json).toEqual({ [session.id]: historyNextSeq - 1 })
- expect(calls[2].json).toEqual({ patch: "remote patch" })
- expect(calls[3].json).toMatchObject({
- directory: "remote-target-dir",
- events: [
- {
- aggregateID: session.id,
- seq: 0,
- type: "session.created.1",
- },
- {
- aggregateID: session.id,
- seq: historyNextSeq,
- type: "session.updated.1",
- },
- ],
- })
- expect(calls[4].json).toEqual({ sessionID: session.id })
- expect((yield* sessionSvc.get(session.id)).title).toBe("from source history")
- expect(yield* sessionSequenceOwner(session.id)).toBe(target.id)
- }),
- { git: true },
- )
- })
- })
- })
- describe("workspace sync state", () => {
- it.instance(
- "startWorkspaceSyncing is disabled by the experimental workspace flag",
- () =>
- Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("flag-disabled")
- const info = workspaceInfo(instance.project.id, type)
- const session = yield* sessionSvc.create({})
- yield* attachSessionToWorkspace(session.id, info.id)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, localAdapter(path.join(dir, "flag-disabled")).adapter)
- yield* Effect.promise(() => startWorkspaceSyncingWithFlag(instance.project.id, false))
- yield* Effect.sleep("25 millis")
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "startWorkspaceSyncing starts all workspaces",
- () =>
- Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const projectID = instance.project.id
- const firstType = unique("first")
- const secondType = unique("second")
- const first = workspaceInfo(projectID, firstType)
- const second = workspaceInfo(projectID, secondType)
- yield* Effect.promise(() => fs.mkdir(path.join(dir, "first"), { recursive: true }))
- yield* Effect.promise(() => fs.mkdir(path.join(dir, "second"), { recursive: true }))
- yield* insertWorkspace(first)
- yield* insertWorkspace(second)
- registerAdapter(projectID, firstType, localAdapter(path.join(dir, "first")).adapter)
- registerAdapter(projectID, secondType, localAdapter(path.join(dir, "second")).adapter)
- yield* Effect.addFinalizer(() =>
- Effect.all([workspace.remove(first.id), workspace.remove(second.id)], { discard: true }).pipe(Effect.ignore),
- )
- yield* workspace.startWorkspaceSyncing(projectID)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- const status = yield* workspace.status()
- expect(status.find((item) => item.workspaceID === first.id)?.status).toBe("connected")
- expect(status.find((item) => item.workspaceID === second.id)?.status).toBe("connected")
- }),
- )
- }),
- { git: true },
- )
- it.instance(
- "local start reports error when the target directory is missing",
- () =>
- Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const type = unique("missing-local")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(
- instance.project.id,
- type,
- localAdapter(path.join(dir, "missing-target"), { createDir: false }).adapter,
- )
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- const status = yield* workspace.status()
- expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("error")
- }),
- )
- expect(yield* workspace.isSyncing(info.id)).toBe(false)
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- it.instance(
- "duplicate local status updates are suppressed",
- () =>
- Effect.gen(function* () {
- const { directory: dir } = yield* TestInstance
- const instance = yield* requireInstance
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const captured = captureGlobalEvents()
- yield* Effect.addFinalizer(() => Effect.sync(() => captured.dispose()))
- const type = unique("dedupe-local")
- const info = workspaceInfo(instance.project.id, type)
- const target = path.join(dir, "dedupe-local")
- yield* Effect.promise(() => fs.mkdir(target, { recursive: true }))
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, localAdapter(target).adapter)
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- const status = yield* workspace.status()
- expect(status.find((item) => item.workspaceID === info.id)?.status).toBe("connected")
- }),
- )
- expect(
- captured.events.filter(
- (event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type,
- ),
- ).toHaveLength(1)
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- it.live("remote start emits disconnected, connecting, and connected then refuses duplicate listeners", () => {
- const calls: FetchCall[] = []
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- const call = {
- url: new URL(req.url, "http://localhost"),
- method: req.method,
- headers: new Headers(req.headers),
- bodyText,
- json: bodyText ? JSON.parse(bodyText) : undefined,
- }
- calls.push(call)
- if (call.url.pathname === "/sync/global/event") return HttpServerResponse.fromWeb(eventStreamResponse())
- if (call.url.pathname === "/sync/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const captured = captureGlobalEvents()
- try {
- const type = unique("remote-start")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sync`).adapter)
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe(
- "connected",
- )
- }),
- )
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* Effect.sleep("25 millis")
- expect(
- captured.events
- .filter((event) => event.workspace === info.id && event.payload.type === Workspace.Event.Status.type)
- .map((event) => event.payload.properties.status),
- ).toEqual(["disconnected", "connecting", "connected"])
- expect(calls.filter((call) => call.url.pathname === "/sync/global/event")).toHaveLength(1)
- expect(calls.filter((call) => call.url.pathname === "/sync/sync/history")).toHaveLength(1)
- expect(yield* workspace.isSyncing(info.id)).toBe(true)
- yield* workspace.remove(info.id)
- expect(yield* workspace.isSyncing(info.id)).toBe(false)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- })
- })
- it.live("remote connection HTTP failures set error and clear syncing", () =>
- Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- if (new URL(req.url, "http://localhost").pathname === "/failed/global/event")
- return HttpServerResponse.text("nope", { status: 503 })
- return HttpServerResponse.fromWeb(Response.json([]))
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const type = unique("remote-connect-fail")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/failed`).adapter)
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
- }),
- )
- expect(yield* workspace.isSyncing(info.id)).toBe(false)
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- }),
- )
- it.live("remote history HTTP failures set error", () =>
- Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const url = new URL(req.url, "http://localhost")
- if (url.pathname === "/history-failed/global/event")
- return HttpServerResponse.fromWeb(eventStreamResponse([], false))
- if (url.pathname === "/history-failed/sync/history")
- return HttpServerResponse.text("history failed", { status: 500 })
- return HttpServerResponse.fromWeb(Response.json([]))
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const type = unique("remote-history-fail")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/history-failed`).adapter)
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- expect((yield* workspace.status()).find((item) => item.workspaceID === info.id)?.status).toBe("error")
- }),
- )
- expect(yield* workspace.isSyncing(info.id)).toBe(false)
- yield* workspace.remove(info.id)
- }),
- { git: true },
- )
- }),
- )
- it.live("sync history sends the local sequence fence and replays returned events in workspace context", () => {
- const historyBodies: unknown[] = []
- let historySessionID: SessionID | undefined
- let historySession: SessionNs.Info | undefined
- let historyNextSeq = 0
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const bodyText = yield* req.text
- const url = new URL(req.url, "http://localhost")
- if (url.pathname === "/history/global/event") return HttpServerResponse.fromWeb(eventStreamResponse())
- if (url.pathname === "/history/sync/history") {
- historyBodies.push(bodyText ? JSON.parse(bodyText) : undefined)
- return HttpServerResponse.fromWeb(
- Response.json([
- {
- id: `evt_${unique("history")}`,
- aggregate_id: historySessionID!,
- seq: historyNextSeq,
- type: "session.updated.1",
- data: { sessionID: historySessionID!, info: historySession! },
- },
- ]),
- )
- }
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const captured = captureGlobalEvents()
- try {
- const type = unique("history-replay")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/history`).adapter)
- const session = yield* sessionSvc.create({ title: "before history" })
- yield* attachSessionToWorkspace(session.id, info.id)
- historySessionID = session.id
- historySession = { ...session, workspaceID: info.id, title: "from history" }
- historyNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from history")
- }),
- )
- expect(historyBodies).toEqual([{ [session.id]: historyNextSeq - 1 }])
- expect(
- captured.events.some(
- (event) =>
- event.workspace === info.id &&
- event.payload.type === "session.updated" &&
- event.payload.properties.sessionID === session.id &&
- event.payload.properties.info.title === "from history",
- ),
- ).toBe(true)
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- })
- })
- it.live("SSE forwards non-heartbeat events and ignores heartbeats", () =>
- Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const url = new URL(req.url, "http://localhost")
- if (url.pathname === "/sse-forward/global/event")
- return HttpServerResponse.fromWeb(
- eventStreamResponse(
- [
- { directory: "remote-dir", project: "remote-project", payload: { type: "server.heartbeat" } },
- {
- directory: "remote-dir",
- project: "remote-project",
- payload: { type: "custom.remote", properties: { ok: true } },
- },
- ],
- false,
- ),
- )
- if (url.pathname === "/sse-forward/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const captured = captureGlobalEvents()
- try {
- const type = unique("sse-forward")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sse-forward`).adapter)
- yield* attachSessionToWorkspace((yield* sessionSvc.create({})).id, info.id)
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.sync(() =>
- expect(
- captured.events.some(
- (event) => event.workspace === info.id && event.payload.type === "custom.remote",
- ),
- ).toBe(true),
- ),
- )
- expect(
- captured.events.some(
- (event) => event.workspace === info.id && event.payload.type === "server.heartbeat",
- ),
- ).toBe(false)
- expect(
- captured.events.find((event) => event.workspace === info.id && event.payload.type === "custom.remote"),
- ).toMatchObject({
- directory: "remote-dir",
- project: "remote-project",
- payload: { properties: { ok: true } },
- })
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- }),
- )
- it.live("SSE sync events are replayed and forwarded", () => {
- let sseSessionID: SessionID | undefined
- let sseSession: SessionNs.Info | undefined
- let sseNextSeq = 0
- return Effect.gen(function* () {
- yield* HttpServer.serveEffect()(
- Effect.gen(function* () {
- const req = yield* HttpServerRequest.HttpServerRequest
- const url = new URL(req.url, "http://localhost")
- if (url.pathname === "/sse-sync/global/event")
- return HttpServerResponse.fromWeb(
- eventStreamResponse(
- [
- {
- directory: "remote-dir",
- project: "remote-project",
- payload: {
- type: "sync",
- syncEvent: {
- id: `evt_${unique("sse")}`,
- aggregateID: sseSessionID!,
- seq: sseNextSeq,
- type: "session.updated.1",
- data: { sessionID: sseSessionID!, info: sseSession! },
- },
- },
- },
- ],
- false,
- ),
- )
- if (url.pathname === "/sse-sync/sync/history") return HttpServerResponse.fromWeb(Response.json([]))
- return HttpServerResponse.text("unexpected", { status: 500 })
- }),
- )
- const url = yield* serverUrl()
- yield* provideTmpdirInstance(
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionSvc = yield* SessionNs.Service
- const instance = yield* requireInstance
- const captured = captureGlobalEvents()
- try {
- const type = unique("sse-sync")
- const info = workspaceInfo(instance.project.id, type)
- yield* insertWorkspace(info)
- registerAdapter(instance.project.id, type, remoteAdapter(`${url}/sse-sync`).adapter)
- const session = yield* sessionSvc.create({ title: "before sse" })
- yield* attachSessionToWorkspace(session.id, info.id)
- sseSessionID = session.id
- sseSession = { ...session, workspaceID: info.id, title: "from sse" }
- sseNextSeq = ((yield* sessionSequence(session.id)) ?? -1) + 1
- yield* workspace.startWorkspaceSyncing(instance.project.id)
- yield* eventuallyEffect(
- Effect.gen(function* () {
- expect((yield* sessionSvc.get(session.id).pipe(Effect.orDie)).title).toBe("from sse")
- }),
- )
- expect(
- captured.events.some(
- (event) =>
- event.workspace === info.id &&
- event.payload.type === "sync" &&
- event.payload.syncEvent.seq === sseNextSeq,
- ),
- ).toBe(true)
- yield* workspace.remove(info.id)
- } finally {
- captured.dispose()
- }
- }),
- { git: true },
- )
- })
- })
- })
- describe("workspace waitForSync", () => {
- it.instance(
- "returns immediately for an empty fence",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- expect(yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_empty"), {})).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "returns immediately when the stored sequence already satisfies the fence",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionID = SessionID.descending("ses_wait_done")
- const { db } = yield* Database.Service
- yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 4 }).run().pipe(Effect.orDie)
- expect(
- yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_done"), { [sessionID]: 4 }),
- ).toBeUndefined()
- expect(
- yield* workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_done_2"), { [sessionID]: 3 }),
- ).toBeUndefined()
- }),
- { git: true },
- )
- it.instance(
- "waits until the database reaches the requested sequence and a workspace event arrives",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const workspaceID = WorkspaceV2.ID.ascending("wrk_wait_event")
- const sessionID = SessionID.descending("ses_wait_event")
- const { db } = yield* Database.Service
- yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 1 }).run().pipe(Effect.orDie)
- yield* Effect.all(
- [
- workspace.waitForSync(workspaceID, { [sessionID]: 2 }),
- Effect.gen(function* () {
- yield* Effect.sleep("10 millis")
- yield* db
- .update(EventSequenceTable)
- .set({ seq: 2 })
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .run()
- .pipe(Effect.orDie)
- GlobalBus.emit("event", { workspace: workspaceID, payload: { type: "anything" } })
- }),
- ],
- { concurrency: "unbounded" },
- )
- }),
- { git: true },
- )
- it.instance(
- "a sync event for a different workspace can also release the fence",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const workspaceID = WorkspaceV2.ID.ascending("wrk_wait_sync_any")
- const sessionID = SessionID.descending("ses_wait_sync_any")
- const { db } = yield* Database.Service
- yield* db.insert(EventSequenceTable).values({ aggregate_id: sessionID, seq: 0 }).run().pipe(Effect.orDie)
- yield* Effect.all(
- [
- workspace.waitForSync(workspaceID, { [sessionID]: 1 }),
- Effect.gen(function* () {
- yield* Effect.sleep("10 millis")
- yield* db
- .update(EventSequenceTable)
- .set({ seq: 1 })
- .where(eq(EventSequenceTable.aggregate_id, sessionID))
- .run()
- .pipe(Effect.orDie)
- GlobalBus.emit("event", {
- workspace: WorkspaceV2.ID.ascending("wrk_other_workspace"),
- payload: { type: "sync" },
- })
- }),
- ],
- { concurrency: "unbounded" },
- )
- }),
- { git: true },
- )
- it.instance(
- "rejects with the abort reason when aborted",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const abort = new AbortController()
- const reason = new Error("caller aborted")
- const fiber = yield* Effect.forkChild(
- workspace.waitForSync(
- WorkspaceV2.ID.ascending("wrk_wait_abort"),
- { [SessionID.descending("ses_wait_abort")]: 1 },
- abort.signal,
- ),
- )
- abort.abort(reason)
- expectExitContains(yield* Fiber.await(fiber), "WorkspaceSyncAbortedError", reason.message)
- }),
- { git: true },
- )
- it.instance(
- "times out with the requested fence in the error message",
- () =>
- Effect.gen(function* () {
- const workspace = yield* Workspace.Service
- const sessionID = SessionID.descending("ses_wait_timeout")
- expectExitContains(
- yield* Effect.exit(
- workspace.waitForSync(WorkspaceV2.ID.ascending("wrk_wait_timeout"), { [sessionID]: 1 }, undefined, 25),
- ),
- `Timed out waiting for sync fence: {"${sessionID}":1}`,
- )
- }),
- { git: true },
- 7000,
- )
- })
|