httpapi-workspace-routing.test.ts 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552
  1. import { NodeHttpServer, NodeServices } from "@effect/platform-node"
  2. import { describe, expect } from "bun:test"
  3. import { Context, Effect, Layer, Queue, Ref, Schema, Stream } from "effect"
  4. import {
  5. FetchHttpClient,
  6. HttpClient,
  7. HttpClientRequest,
  8. HttpRouter,
  9. HttpServer,
  10. HttpServerRequest,
  11. HttpServerResponse,
  12. } from "effect/unstable/http"
  13. import * as Socket from "effect/unstable/socket/Socket"
  14. import { HttpApi, HttpApiBuilder, HttpApiEndpoint, HttpApiGroup } from "effect/unstable/httpapi"
  15. import Http from "node:http"
  16. import { mkdir } from "node:fs/promises"
  17. import path from "node:path"
  18. import { registerAdapter } from "../../src/control-plane/adapters"
  19. import { WorkspaceV2 } from "@kirincode-ai/core/workspace"
  20. import type { WorkspaceAdapter } from "../../src/control-plane/types"
  21. import { Workspace } from "../../src/control-plane/workspace"
  22. import { WorkspaceTable } from "@kirincode-ai/core/control-plane/workspace.sql"
  23. import { Database } from "@kirincode-ai/core/database/database"
  24. import { Project } from "../../src/project/project"
  25. import { Session } from "../../src/session/session"
  26. import { WorkspacePaths } from "../../src/server/routes/instance/httpapi/groups/workspace"
  27. import {
  28. WorkspaceRoutingMiddleware,
  29. WorkspaceRoutingQuery,
  30. WorkspaceRouteContext,
  31. workspaceRoutingLayer,
  32. } from "../../src/server/routes/instance/httpapi/middleware/workspace-routing"
  33. import { HEADER as FenceHeader } from "../../src/server/shared/fence"
  34. import { resetDatabase } from "../fixture/db"
  35. import { workspaceLayerWithRuntimeFlags } from "../fixture/workspace"
  36. import { tmpdirScoped } from "../fixture/fixture"
  37. import { testEffect } from "../lib/effect"
  38. const testStateLayer = Layer.effectDiscard(
  39. Effect.gen(function* () {
  40. yield* Effect.promise(() => resetDatabase())
  41. yield* Effect.addFinalizer(() =>
  42. Effect.promise(async () => {
  43. await resetDatabase()
  44. }),
  45. )
  46. }),
  47. )
  48. const workspaceLayer = workspaceLayerWithRuntimeFlags({ experimentalWorkspaces: true })
  49. const it = testEffect(
  50. Layer.mergeAll(
  51. testStateLayer,
  52. NodeHttpServer.layerTest,
  53. NodeServices.layer,
  54. workspaceLayer,
  55. Socket.layerWebSocketConstructorGlobal,
  56. ),
  57. )
  58. type ProxiedRequest = {
  59. url: string
  60. method: string
  61. headers: Record<string, string>
  62. body: string
  63. }
  64. type TestHandler<E, R> = (
  65. request: HttpServerRequest.HttpServerRequest,
  66. ) => Effect.Effect<HttpServerResponse.HttpServerResponse, E, R>
  67. const workspaceRoutingTestLayer = workspaceRoutingLayer.pipe(
  68. Layer.provide([Socket.layerWebSocketConstructorGlobal, FetchHttpClient.layer]),
  69. )
  70. const serverUrl = HttpServer.HttpServer.use((server) => Effect.succeed(HttpServer.formatAddress(server.address)))
  71. const requestURL = (request: { readonly url: string }) => new URL(request.url, "http://localhost")
  72. const listenAdditionalServer = <E, R>(handler: TestHandler<E, R>) =>
  73. Effect.gen(function* () {
  74. const context = yield* Layer.build(NodeHttpServer.layer(Http.createServer, { host: "127.0.0.1", port: 0 }))
  75. const server = Context.get(context, HttpServer.HttpServer)
  76. yield* server.serve(HttpServerRequest.HttpServerRequest.use(handler))
  77. return HttpServer.formatAddress(server.address)
  78. })
  79. const localAdapter = (directory: string): WorkspaceAdapter => ({
  80. name: "Local Test",
  81. description: "Create a local test workspace",
  82. configure: (info) => ({ ...info, name: "local-test", directory }),
  83. create: async () => {
  84. await mkdir(directory, { recursive: true })
  85. },
  86. async remove() {},
  87. target: () => ({ type: "local" as const, directory }),
  88. })
  89. const remoteAdapter = (directory: string, url: string, headers?: HeadersInit): WorkspaceAdapter => ({
  90. name: "Remote Test",
  91. description: "Create a remote test workspace",
  92. configure: (info) => ({ ...info, name: "remote-test", directory }),
  93. create: async () => {
  94. await mkdir(directory, { recursive: true })
  95. },
  96. async remove() {},
  97. target: () => ({ type: "remote" as const, url, headers }),
  98. })
  99. const eventStreamResponse = () =>
  100. HttpServerResponse.text('data: {"payload":{"type":"server.connected","properties":{}}}\n\n', {
  101. contentType: "text/event-stream",
  102. })
  103. const syncResponse = (request: HttpServerRequest.HttpServerRequest) => {
  104. const url = requestURL(request)
  105. if (url.pathname === "/base/global/event") return Effect.succeed(eventStreamResponse())
  106. if (url.pathname === "/base/sync/history") return HttpServerResponse.json([])
  107. return undefined
  108. }
  109. const createWorkspace = (input: { projectID: Project.Info["id"]; type: string; adapter: WorkspaceAdapter }) =>
  110. Effect.acquireRelease(
  111. Effect.gen(function* () {
  112. registerAdapter(input.projectID, input.type, input.adapter)
  113. const workspace = yield* Workspace.Service
  114. return yield* workspace.create({
  115. type: input.type,
  116. branch: null,
  117. extra: null,
  118. projectID: input.projectID,
  119. })
  120. }),
  121. (info) => Workspace.use.remove(info.id).pipe(Effect.ignore),
  122. )
  123. const createRemoteWorkspace = (input: {
  124. dir: string
  125. projectID: Project.Info["id"]
  126. type: string
  127. url: string
  128. headers?: HeadersInit
  129. }) =>
  130. // Workspace.create starts the remote sync loop. The test upstream exposes
  131. // /global/event and /sync/history so middleware proxying sees the remote
  132. // workspace as active, just like production would.
  133. createWorkspace({
  134. projectID: input.projectID,
  135. type: input.type,
  136. adapter: remoteAdapter(path.join(input.dir, `.${input.type}`), input.url, input.headers),
  137. })
  138. const createLocalWorkspace = (input: { projectID: Project.Info["id"]; type: string; directory: string }) =>
  139. createWorkspace({
  140. projectID: input.projectID,
  141. type: input.type,
  142. adapter: localAdapter(input.directory),
  143. })
  144. const insertRemoteWorkspaceWithoutSync = (input: {
  145. dir: string
  146. projectID: Project.Info["id"]
  147. type: string
  148. url: string
  149. }) =>
  150. Effect.gen(function* () {
  151. const id = WorkspaceV2.ID.ascending()
  152. registerAdapter(input.projectID, input.type, remoteAdapter(path.join(input.dir, `.${input.type}`), input.url))
  153. const { db } = yield* Database.Service
  154. yield* db
  155. .insert(WorkspaceTable)
  156. .values({ id, type: input.type, project_id: input.projectID })
  157. .run()
  158. .pipe(Effect.orDie)
  159. return id
  160. })
  161. const startRemoteWorkspaceHttpServer = <E, R>(
  162. handler: (request: ProxiedRequest) => Effect.Effect<HttpServerResponse.HttpServerResponse, E, R>,
  163. ) =>
  164. listenAdditionalServer((request) =>
  165. Effect.gen(function* () {
  166. // Remote workspaces run a sync loop against their target server. These
  167. // bootstrap routes make Workspace.isSyncing(...) true for proxy tests;
  168. // everything else is the request being proxied by the middleware.
  169. const sync = syncResponse(request)
  170. if (sync) return yield* sync
  171. return yield* handler({
  172. url: request.url,
  173. method: request.method,
  174. headers: request.headers,
  175. body: yield* request.text,
  176. })
  177. }),
  178. )
  179. const listenRemoteWebSocket = () =>
  180. listenAdditionalServer((request) => {
  181. const sync = syncResponse(request)
  182. if (sync) return sync
  183. if (requestURL(request).pathname !== "/base/probe") return Effect.succeed(HttpServerResponse.empty({ status: 404 }))
  184. return echoWebSocket(request)
  185. })
  186. const echoWebSocket = (request: HttpServerRequest.HttpServerRequest) =>
  187. Effect.gen(function* () {
  188. const socket = yield* Effect.orDie(request.upgrade)
  189. const write = yield* socket.writer
  190. yield* socket
  191. .runRaw((message) => write(`echo:${String(message)}`), {
  192. onOpen: write(`protocol:${request.headers["sec-websocket-protocol"] ?? "none"}`).pipe(
  193. Effect.catch(() => Effect.void),
  194. ),
  195. })
  196. .pipe(Effect.catch(() => Effect.void))
  197. return HttpServerResponse.empty()
  198. })
  199. const ProbeResult = Schema.Struct({
  200. directory: Schema.String,
  201. workspaceID: Schema.optional(Schema.String),
  202. })
  203. const ProbeApi = HttpApi.make("workspace-routing-probe").add(
  204. HttpApiGroup.make("probe")
  205. .add(
  206. HttpApiEndpoint.get("get", "/probe", { query: WorkspaceRoutingQuery, success: ProbeResult }),
  207. HttpApiEndpoint.patch("patch", "/probe", { query: WorkspaceRoutingQuery, success: Schema.Boolean }),
  208. HttpApiEndpoint.get("session", "/session", { query: WorkspaceRoutingQuery, success: ProbeResult }),
  209. HttpApiEndpoint.get("workspace", WorkspacePaths.list, {
  210. query: WorkspaceRoutingQuery,
  211. success: ProbeResult,
  212. }),
  213. )
  214. .middleware(WorkspaceRoutingMiddleware),
  215. )
  216. const routeContextResponse = Effect.gen(function* () {
  217. const route = yield* WorkspaceRouteContext
  218. return { directory: route.directory, workspaceID: route.workspaceID }
  219. })
  220. const probeHandlers = HttpApiBuilder.group(ProbeApi, "probe", (handlers) =>
  221. handlers
  222. .handle("get", () => routeContextResponse)
  223. .handle("patch", () => Effect.succeed(false))
  224. .handle("session", () => routeContextResponse)
  225. .handle("workspace", () => routeContextResponse),
  226. )
  227. const serveProbe = HttpApiBuilder.layer(ProbeApi).pipe(
  228. Layer.provide(probeHandlers),
  229. Layer.provide(workspaceRoutingTestLayer),
  230. Layer.provide(Layer.mock(Session.Service)({})),
  231. HttpRouter.serve,
  232. Layer.build,
  233. )
  234. describe("HttpApi workspace routing middleware", () => {
  235. it.live("proxies remote workspace HTTP requests through the selected workspace target", () =>
  236. Effect.gen(function* () {
  237. const dir = yield* tmpdirScoped({ git: true })
  238. const project = yield* Project.use.fromDirectory(dir)
  239. let forwarded: ProxiedRequest | undefined
  240. // This starts a second HTTP server that stands in for the kirincode server
  241. // backing a remote workspace. The client below still calls the local test
  242. // server; only the middleware should call this server.
  243. const remoteUrl = yield* startRemoteWorkspaceHttpServer((request) => {
  244. forwarded = request
  245. const url = requestURL(request)
  246. return HttpServerResponse.json(
  247. {
  248. proxied: true,
  249. path: url.pathname,
  250. keep: url.searchParams.get("keep"),
  251. workspace: url.searchParams.get("workspace"),
  252. },
  253. { status: 201, headers: { "x-remote": "yes" } },
  254. )
  255. })
  256. // The adapter target tells the middleware where to proxy selected remote
  257. // workspace requests. Appending /probe to this base should produce
  258. // `${remoteUrl}/base/probe` on the fake remote server above.
  259. const workspace = yield* createRemoteWorkspace({
  260. dir,
  261. projectID: project.project.id,
  262. type: "remote-http-target",
  263. url: `${remoteUrl}/base`,
  264. headers: { "x-target-auth": "secret" },
  265. })
  266. // The local /probe handler should not run. Selecting a remote workspace
  267. // should make the middleware call HttpApiProxy.http instead.
  268. yield* serveProbe
  269. const body = '{"title":"Remote workspace request"}'
  270. const response = yield* HttpClientRequest.patch(`/probe?workspace=${workspace.id}&keep=yes`).pipe(
  271. HttpClientRequest.setHeaders({
  272. "x-opencode-directory": "/secret/path",
  273. "x-opencode-workspace": "internal",
  274. }),
  275. HttpClientRequest.bodyStream(
  276. Stream.make(new TextEncoder().encode('{"title":"Remote '), new TextEncoder().encode('workspace request"}')),
  277. { contentType: "application/json" },
  278. ),
  279. HttpClient.execute,
  280. Effect.timeout("2 seconds"),
  281. )
  282. expect(response.status).toBe(201)
  283. expect(response.headers["x-remote"]).toBe("yes")
  284. expect(yield* response.json).toEqual({ proxied: true, path: "/base/probe", keep: "yes", workspace: null })
  285. const forwardedURL = forwarded ? requestURL(forwarded) : undefined
  286. // These assertions are the routing contract: append the original path to
  287. // the remote base URL, preserve normal query params, and remove workspace.
  288. expect(forwardedURL?.pathname).toBe("/base/probe")
  289. expect(forwardedURL?.searchParams.get("keep")).toBe("yes")
  290. expect(forwardedURL?.searchParams.get("workspace")).toBeNull()
  291. expect(forwarded?.method).toBe("PATCH")
  292. expect(forwarded?.body).toBe(body)
  293. expect(forwarded?.headers["content-type"]).toBe("application/json")
  294. expect(forwarded?.headers["x-target-auth"]).toBe("secret")
  295. expect(forwarded?.headers["x-opencode-directory"]).toBeUndefined()
  296. expect(forwarded?.headers["x-opencode-workspace"]).toBeUndefined()
  297. }),
  298. )
  299. it.live("waits for sync fence headers from remote workspace HTTP responses", () =>
  300. Effect.gen(function* () {
  301. const dir = yield* tmpdirScoped({ git: true })
  302. const project = yield* Project.use.fromDirectory(dir)
  303. const workspaceID = WorkspaceV2.ID.ascending()
  304. const type = "remote-http-fence-target"
  305. const waited = yield* Ref.make<{ workspaceID: WorkspaceV2.ID; state: Record<string, number> } | undefined>(
  306. undefined,
  307. )
  308. const remoteUrl = yield* startRemoteWorkspaceHttpServer(() =>
  309. HttpServerResponse.json(
  310. { proxied: true },
  311. { status: 202, headers: { [FenceHeader]: JSON.stringify({ aggregate: 3 }) } },
  312. ),
  313. )
  314. registerAdapter(project.project.id, type, remoteAdapter(path.join(dir, `.${type}`), `${remoteUrl}/base`))
  315. const workspace = Workspace.Service.of({
  316. create: () => Effect.die("unused"),
  317. sessionWarp: () => Effect.die("unused"),
  318. list: () => Effect.die("unused"),
  319. syncList: () => Effect.die("unused"),
  320. get: (id) =>
  321. Effect.succeed(
  322. id === workspaceID
  323. ? {
  324. id: workspaceID,
  325. type,
  326. branch: null,
  327. name: "remote-http-fence-target",
  328. directory: null,
  329. extra: null,
  330. projectID: project.project.id,
  331. timeUsed: Date.now(),
  332. }
  333. : undefined,
  334. ),
  335. remove: () => Effect.die("unused"),
  336. status: () => Effect.die("unused"),
  337. isSyncing: () => Effect.succeed(true),
  338. waitForSync: (id, state) => Ref.set(waited, { workspaceID: id, state }),
  339. startWorkspaceSyncing: () => Effect.die("unused"),
  340. })
  341. yield* HttpApiBuilder.layer(ProbeApi).pipe(
  342. Layer.provide(probeHandlers),
  343. Layer.provide(workspaceRoutingTestLayer),
  344. Layer.provide(Layer.succeed(Workspace.Service, workspace)),
  345. Layer.provide(Layer.mock(Session.Service)({})),
  346. HttpRouter.serve,
  347. Layer.build,
  348. )
  349. const response = yield* HttpClientRequest.patch(`/probe?workspace=${workspaceID}`).pipe(HttpClient.execute)
  350. expect(response.status).toBe(202)
  351. expect(yield* response.json).toEqual({ proxied: true })
  352. expect(yield* Ref.get(waited)).toEqual({ workspaceID, state: { aggregate: 3 } })
  353. }),
  354. )
  355. it.live("returns 503 when a remote workspace is not actively syncing", () =>
  356. Effect.gen(function* () {
  357. const dir = yield* tmpdirScoped({ git: true })
  358. const project = yield* Project.use.fromDirectory(dir)
  359. const workspaceID = yield* insertRemoteWorkspaceWithoutSync({
  360. dir,
  361. projectID: project.project.id,
  362. type: "remote-not-syncing",
  363. url: "http://127.0.0.1:1/base",
  364. })
  365. yield* serveProbe
  366. const response = yield* HttpClient.get(`/probe?workspace=${workspaceID}`)
  367. expect(response.status).toBe(503)
  368. expect(yield* response.text).toBe(`broken sync connection for workspace: ${workspaceID}`)
  369. }),
  370. )
  371. it.live("proxies remote workspace WebSocket requests through the selected workspace target", () =>
  372. Effect.gen(function* () {
  373. const dir = yield* tmpdirScoped({ git: true })
  374. const project = yield* Project.use.fromDirectory(dir)
  375. const remoteUrl = yield* listenRemoteWebSocket()
  376. const workspace = yield* createRemoteWorkspace({
  377. dir,
  378. projectID: project.project.id,
  379. type: "remote-websocket-target",
  380. url: `${remoteUrl}/base`,
  381. })
  382. // The client connects to the local test server. The middleware should
  383. // detect the WebSocket upgrade and proxy it to the remote /base/probe.
  384. yield* serveProbe
  385. const socket = yield* Socket.makeWebSocket(
  386. `${(yield* serverUrl).replace(/^http/, "ws")}/probe?workspace=${workspace.id}`,
  387. {
  388. closeCodeIsError: () => false,
  389. protocols: "chat",
  390. },
  391. )
  392. const messages = yield* Queue.unbounded<string>()
  393. yield* socket.runRaw((message) => Queue.offer(messages, String(message))).pipe(Effect.forkScoped)
  394. const write = yield* socket.writer
  395. expect(yield* Queue.take(messages)).toBe("protocol:chat")
  396. yield* write("hello")
  397. expect(yield* Queue.take(messages)).toBe("echo:hello")
  398. }),
  399. )
  400. it.live("returns a missing workspace response for unknown workspace ids", () =>
  401. Effect.gen(function* () {
  402. const workspaceID = WorkspaceV2.ID.ascending("wrk_missing")
  403. // If the middleware resolves the workspace first, this handler is never
  404. // reached and the response should be the middleware error response.
  405. yield* serveProbe
  406. const response = yield* HttpClient.get(`/probe?workspace=${workspaceID}`)
  407. expect(response.status).toBe(500)
  408. expect(yield* response.text).toBe(`Workspace not found: ${workspaceID}`)
  409. }),
  410. )
  411. it.live("keeps control-plane routes local even when workspace is selected", () =>
  412. Effect.gen(function* () {
  413. const dir = yield* tmpdirScoped({ git: true })
  414. const project = yield* Project.use.fromDirectory(dir)
  415. const workspaceDir = path.join(dir, ".workspace-local")
  416. const workspace = yield* createLocalWorkspace({
  417. projectID: project.project.id,
  418. type: "control-plane-target",
  419. directory: workspaceDir,
  420. })
  421. // GET /session is a control-plane route: it lists sessions for the main
  422. // process and should not be redirected into the selected workspace target.
  423. yield* serveProbe
  424. const response = yield* HttpClient.get(`/session?workspace=${workspace.id}`)
  425. expect(response.status).toBe(200)
  426. expect(yield* response.json).toEqual({ directory: process.cwd(), workspaceID: workspace.id })
  427. }),
  428. )
  429. it.live("keeps workspace control routes local even when workspace is selected", () =>
  430. Effect.gen(function* () {
  431. const dir = yield* tmpdirScoped({ git: true })
  432. const project = yield* Project.use.fromDirectory(dir)
  433. const workspaceDir = path.join(dir, ".workspace-local")
  434. const workspace = yield* createLocalWorkspace({
  435. projectID: project.project.id,
  436. type: "workspace-control-plane-target",
  437. directory: workspaceDir,
  438. })
  439. // Workspace CRUD/status routes manage the control plane itself. Selecting
  440. // a workspace should preserve the selected id for handlers, but must not
  441. // swap the route context to the workspace target directory.
  442. yield* serveProbe
  443. const response = yield* HttpClient.get(`${WorkspacePaths.list}?workspace=${workspace.id}`)
  444. expect(response.status).toBe(200)
  445. expect(yield* response.json).toEqual({ directory: process.cwd(), workspaceID: workspace.id })
  446. }),
  447. )
  448. it.live("uses directory query/header fallback when no workspace is selected", () =>
  449. Effect.gen(function* () {
  450. const dir = yield* tmpdirScoped()
  451. const queryDir = path.join(dir, "query-target")
  452. const headerDir = path.join(dir, "header-target")
  453. yield* serveProbe
  454. // Without a selected workspace, the middleware falls back to request
  455. // directory hints before using the process cwd.
  456. const queryResponse = yield* HttpClient.get(`/probe?directory=${encodeURIComponent(queryDir)}`)
  457. const headerResponse = yield* HttpClientRequest.get("/probe").pipe(
  458. HttpClientRequest.setHeader("x-opencode-directory", headerDir),
  459. HttpClient.execute,
  460. )
  461. expect(queryResponse.status).toBe(200)
  462. expect(yield* queryResponse.json).toEqual({ directory: queryDir, workspaceID: null })
  463. expect(headerResponse.status).toBe(200)
  464. expect(yield* headerResponse.json).toEqual({ directory: headerDir, workspaceID: null })
  465. }),
  466. )
  467. it.live("routes local workspace requests through WorkspaceRouteContext", () =>
  468. Effect.gen(function* () {
  469. const dir = yield* tmpdirScoped({ git: true })
  470. const project = yield* Project.use.fromDirectory(dir)
  471. const workspaceDir = path.join(dir, ".workspace-local")
  472. const workspace = yield* createLocalWorkspace({
  473. projectID: project.project.id,
  474. type: "local-target",
  475. directory: workspaceDir,
  476. })
  477. yield* serveProbe
  478. // /probe is not a control-plane route, so selecting a local workspace
  479. // should swap the route context to the workspace target directory.
  480. const response = yield* HttpClient.get(`/probe?workspace=${workspace.id}`)
  481. expect(response.status).toBe(200)
  482. expect(yield* response.json).toEqual({
  483. directory: workspaceDir,
  484. workspaceID: workspace.id,
  485. })
  486. }),
  487. )
  488. })