httpapi-workspace-routing.test.ts 21 KB

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