workspace.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617
  1. import { Schema } from "effect"
  2. import { setTimeout as sleep } from "node:timers/promises"
  3. import { fn } from "@/util/fn"
  4. import { Database, asc, eq, inArray } from "@/storage"
  5. import { Project } from "@/project"
  6. import { BusEvent } from "@/bus/bus-event"
  7. import { GlobalBus } from "@/bus/global"
  8. import { Auth } from "@/auth"
  9. import { SyncEvent } from "@/sync"
  10. import { EventSequenceTable, EventTable } from "@/sync/event.sql"
  11. import { Flag } from "@opencode-ai/core/flag/flag"
  12. import { Log } from "@/util"
  13. import { Filesystem } from "@/util"
  14. import { ProjectID } from "@/project/schema"
  15. import { Slug } from "@opencode-ai/core/util/slug"
  16. import { WorkspaceTable } from "./workspace.sql"
  17. import { getAdaptor } from "./adaptors"
  18. import { type WorkspaceInfo, WorkspaceInfo as WorkspaceInfoSchema } from "./types"
  19. import { WorkspaceID } from "./schema"
  20. import { parseSSE } from "./sse"
  21. import { Session } from "@/session"
  22. import { SessionTable } from "@/session/session.sql"
  23. import { SessionID } from "@/session/schema"
  24. import { errorData } from "@/util/error"
  25. import { AppRuntime } from "@/effect/app-runtime"
  26. import { waitEvent } from "./util"
  27. import { WorkspaceContext } from "./workspace-context"
  28. import { NonNegativeInt, withStatics } from "@/util/schema"
  29. import { zod as effectZod, zodObject } from "@/util/effect-zod"
  30. export const Info = WorkspaceInfoSchema
  31. export type Info = WorkspaceInfo
  32. export const ConnectionStatus = Schema.Struct({
  33. workspaceID: WorkspaceID,
  34. status: Schema.Literals(["connected", "connecting", "disconnected", "error"]),
  35. })
  36. export type ConnectionStatus = Schema.Schema.Type<typeof ConnectionStatus>
  37. const Restore = Schema.Struct({
  38. workspaceID: WorkspaceID,
  39. sessionID: SessionID,
  40. total: NonNegativeInt,
  41. step: NonNegativeInt,
  42. })
  43. export const Event = {
  44. Ready: BusEvent.define(
  45. "workspace.ready",
  46. Schema.Struct({
  47. name: Schema.String,
  48. }),
  49. ),
  50. Failed: BusEvent.define(
  51. "workspace.failed",
  52. Schema.Struct({
  53. message: Schema.String,
  54. }),
  55. ),
  56. Restore: BusEvent.define("workspace.restore", Restore),
  57. Status: BusEvent.define("workspace.status", ConnectionStatus),
  58. }
  59. function fromRow(row: typeof WorkspaceTable.$inferSelect): Info {
  60. return {
  61. id: row.id,
  62. type: row.type,
  63. branch: row.branch,
  64. name: row.name,
  65. directory: row.directory,
  66. extra: row.extra,
  67. projectID: row.project_id,
  68. }
  69. }
  70. export const CreateInput = Schema.Struct({
  71. id: Schema.optional(WorkspaceID),
  72. type: Info.fields.type,
  73. branch: Info.fields.branch,
  74. projectID: ProjectID,
  75. extra: Info.fields.extra,
  76. }).pipe(withStatics((s) => ({ zod: effectZod(s), zodObject: zodObject(s) })))
  77. export type CreateInput = Schema.Schema.Type<typeof CreateInput>
  78. export const create = fn(CreateInput.zod, async (input) => {
  79. const id = WorkspaceID.ascending(input.id)
  80. const adaptor = await getAdaptor(input.projectID, input.type)
  81. const config = await adaptor.configure({ ...input, id, name: Slug.create(), directory: null })
  82. const info: Info = {
  83. id,
  84. type: config.type,
  85. branch: config.branch ?? null,
  86. name: config.name ?? null,
  87. directory: config.directory ?? null,
  88. extra: config.extra ?? null,
  89. projectID: input.projectID,
  90. }
  91. Database.use((db) => {
  92. db.insert(WorkspaceTable)
  93. .values({
  94. id: info.id,
  95. type: info.type,
  96. branch: info.branch,
  97. name: info.name,
  98. directory: info.directory,
  99. extra: info.extra,
  100. project_id: info.projectID,
  101. })
  102. .run()
  103. })
  104. const env = {
  105. OPENCODE_AUTH_CONTENT: JSON.stringify(await AppRuntime.runPromise(Auth.Service.use((auth) => auth.all()))),
  106. OPENCODE_WORKSPACE_ID: config.id,
  107. OPENCODE_EXPERIMENTAL_WORKSPACES: "true",
  108. OTEL_EXPORTER_OTLP_HEADERS: process.env.OTEL_EXPORTER_OTLP_HEADERS,
  109. OTEL_EXPORTER_OTLP_ENDPOINT: process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
  110. OTEL_RESOURCE_ATTRIBUTES: process.env.OTEL_RESOURCE_ATTRIBUTES,
  111. }
  112. await adaptor.create(config, env)
  113. startSync(info)
  114. await waitEvent({
  115. timeout: TIMEOUT,
  116. fn(event) {
  117. if (event.workspace === info.id && event.payload.type === Event.Status.type) {
  118. const { status } = event.payload.properties
  119. return status === "error" || status === "connected"
  120. }
  121. return false
  122. },
  123. })
  124. return info
  125. })
  126. export const SessionRestoreInput = Schema.Struct({
  127. workspaceID: WorkspaceID,
  128. sessionID: SessionID,
  129. }).pipe(withStatics((s) => ({ zod: effectZod(s), zodObject: zodObject(s) })))
  130. export type SessionRestoreInput = Schema.Schema.Type<typeof SessionRestoreInput>
  131. export const sessionRestore = fn(SessionRestoreInput.zod, async (input) => {
  132. log.info("session restore requested", {
  133. workspaceID: input.workspaceID,
  134. sessionID: input.sessionID,
  135. })
  136. try {
  137. const space = await get(input.workspaceID)
  138. if (!space) throw new Error(`Workspace not found: ${input.workspaceID}`)
  139. const adaptor = await getAdaptor(space.projectID, space.type)
  140. const target = await adaptor.target(space)
  141. // Need to switch the workspace of the session
  142. SyncEvent.run(Session.Event.Updated, {
  143. sessionID: input.sessionID,
  144. info: {
  145. workspaceID: input.workspaceID,
  146. },
  147. })
  148. const rows = Database.use((db) =>
  149. db
  150. .select({
  151. id: EventTable.id,
  152. aggregateID: EventTable.aggregate_id,
  153. seq: EventTable.seq,
  154. type: EventTable.type,
  155. data: EventTable.data,
  156. })
  157. .from(EventTable)
  158. .where(eq(EventTable.aggregate_id, input.sessionID))
  159. .orderBy(asc(EventTable.seq))
  160. .all(),
  161. )
  162. if (rows.length === 0) throw new Error(`No events found for session: ${input.sessionID}`)
  163. const all = rows
  164. const size = 10
  165. const sets = Array.from({ length: Math.ceil(all.length / size) }, (_, i) => all.slice(i * size, (i + 1) * size))
  166. const total = sets.length
  167. log.info("session restore prepared", {
  168. workspaceID: input.workspaceID,
  169. sessionID: input.sessionID,
  170. workspaceType: space.type,
  171. directory: space.directory,
  172. target: target.type === "remote" ? String(route(target.url, "/sync/replay")) : target.directory,
  173. events: all.length,
  174. batches: total,
  175. first: all[0]?.seq,
  176. last: all.at(-1)?.seq,
  177. })
  178. GlobalBus.emit("event", {
  179. directory: "global",
  180. workspace: input.workspaceID,
  181. payload: {
  182. type: Event.Restore.type,
  183. properties: {
  184. workspaceID: input.workspaceID,
  185. sessionID: input.sessionID,
  186. total,
  187. step: 0,
  188. },
  189. },
  190. })
  191. for (const [i, events] of sets.entries()) {
  192. log.info("session restore batch starting", {
  193. workspaceID: input.workspaceID,
  194. sessionID: input.sessionID,
  195. step: i + 1,
  196. total,
  197. events: events.length,
  198. first: events[0]?.seq,
  199. last: events.at(-1)?.seq,
  200. target: target.type === "remote" ? String(route(target.url, "/sync/replay")) : target.directory,
  201. })
  202. if (target.type === "local") {
  203. SyncEvent.replayAll(events)
  204. log.info("session restore batch replayed locally", {
  205. workspaceID: input.workspaceID,
  206. sessionID: input.sessionID,
  207. step: i + 1,
  208. total,
  209. events: events.length,
  210. })
  211. } else {
  212. const url = route(target.url, "/sync/replay")
  213. const headers = new Headers(target.headers)
  214. headers.set("content-type", "application/json")
  215. const res = await fetch(url, {
  216. method: "POST",
  217. headers,
  218. body: JSON.stringify({
  219. directory: space.directory ?? "",
  220. events,
  221. }),
  222. })
  223. if (!res.ok) {
  224. const body = await res.text()
  225. log.error("session restore batch failed", {
  226. workspaceID: input.workspaceID,
  227. sessionID: input.sessionID,
  228. step: i + 1,
  229. total,
  230. status: res.status,
  231. body,
  232. })
  233. throw new Error(
  234. `Failed to replay session ${input.sessionID} into workspace ${input.workspaceID}: HTTP ${res.status} ${body}`,
  235. )
  236. }
  237. log.info("session restore batch posted", {
  238. workspaceID: input.workspaceID,
  239. sessionID: input.sessionID,
  240. step: i + 1,
  241. total,
  242. status: res.status,
  243. })
  244. }
  245. GlobalBus.emit("event", {
  246. directory: "global",
  247. workspace: input.workspaceID,
  248. payload: {
  249. type: Event.Restore.type,
  250. properties: {
  251. workspaceID: input.workspaceID,
  252. sessionID: input.sessionID,
  253. total,
  254. step: i + 1,
  255. },
  256. },
  257. })
  258. }
  259. log.info("session restore complete", {
  260. workspaceID: input.workspaceID,
  261. sessionID: input.sessionID,
  262. batches: total,
  263. })
  264. return {
  265. total,
  266. }
  267. } catch (err) {
  268. log.error("session restore failed", {
  269. workspaceID: input.workspaceID,
  270. sessionID: input.sessionID,
  271. error: errorData(err),
  272. })
  273. throw err
  274. }
  275. })
  276. export function list(project: Project.Info) {
  277. const rows = Database.use((db) =>
  278. db.select().from(WorkspaceTable).where(eq(WorkspaceTable.project_id, project.id)).all(),
  279. )
  280. const spaces = rows.map(fromRow).sort((a, b) => a.id.localeCompare(b.id))
  281. return spaces
  282. }
  283. export const get = fn(WorkspaceID.zod, async (id) => {
  284. const row = Database.use((db) => db.select().from(WorkspaceTable).where(eq(WorkspaceTable.id, id)).get())
  285. if (!row) return
  286. return fromRow(row)
  287. })
  288. export const remove = fn(WorkspaceID.zod, async (id) => {
  289. const sessions = Database.use((db) =>
  290. db.select({ id: SessionTable.id }).from(SessionTable).where(eq(SessionTable.workspace_id, id)).all(),
  291. )
  292. for (const session of sessions) {
  293. await AppRuntime.runPromise(Session.Service.use((svc) => svc.remove(session.id)))
  294. }
  295. const row = Database.use((db) => db.select().from(WorkspaceTable).where(eq(WorkspaceTable.id, id)).get())
  296. if (row) {
  297. stopSync(id)
  298. const info = fromRow(row)
  299. try {
  300. const adaptor = await getAdaptor(info.projectID, row.type)
  301. await adaptor.remove(info)
  302. } catch {
  303. log.error("adaptor not available when removing workspace", { type: row.type })
  304. }
  305. Database.use((db) => db.delete(WorkspaceTable).where(eq(WorkspaceTable.id, id)).run())
  306. return info
  307. }
  308. })
  309. const connections = new Map<WorkspaceID, ConnectionStatus>()
  310. const aborts = new Map<WorkspaceID, AbortController>()
  311. const TIMEOUT = 5000
  312. function setStatus(id: WorkspaceID, status: ConnectionStatus["status"]) {
  313. const prev = connections.get(id)
  314. if (prev?.status === status) return
  315. const next = { workspaceID: id, status }
  316. connections.set(id, next)
  317. if (status === "error") {
  318. aborts.delete(id)
  319. }
  320. GlobalBus.emit("event", {
  321. directory: "global",
  322. workspace: id,
  323. payload: {
  324. type: Event.Status.type,
  325. properties: next,
  326. },
  327. })
  328. }
  329. export function status(): ConnectionStatus[] {
  330. return [...connections.values()]
  331. }
  332. function synced(state: Record<string, number>) {
  333. const ids = Object.keys(state)
  334. if (ids.length === 0) return true
  335. const done = Object.fromEntries(
  336. Database.use((db) =>
  337. db
  338. .select({
  339. id: EventSequenceTable.aggregate_id,
  340. seq: EventSequenceTable.seq,
  341. })
  342. .from(EventSequenceTable)
  343. .where(inArray(EventSequenceTable.aggregate_id, ids))
  344. .all(),
  345. ).map((row) => [row.id, row.seq]),
  346. ) as Record<string, number>
  347. return ids.every((id) => {
  348. return (done[id] ?? -1) >= state[id]
  349. })
  350. }
  351. export async function isSyncing(workspaceID: WorkspaceID) {
  352. return aborts.has(workspaceID)
  353. }
  354. export async function waitForSync(workspaceID: WorkspaceID, state: Record<string, number>, signal?: AbortSignal) {
  355. if (synced(state)) return
  356. try {
  357. await waitEvent({
  358. timeout: TIMEOUT,
  359. signal,
  360. fn(event) {
  361. if (event.workspace !== workspaceID && event.payload.type !== "sync") {
  362. return false
  363. }
  364. return synced(state)
  365. },
  366. })
  367. } catch {
  368. if (signal?.aborted) throw signal.reason ?? new Error("Request aborted")
  369. throw new Error(`Timed out waiting for sync fence: ${JSON.stringify(state)}`)
  370. }
  371. }
  372. const log = Log.create({ service: "workspace-sync" })
  373. function route(url: string | URL, path: string) {
  374. const next = new URL(url)
  375. next.pathname = `${next.pathname.replace(/\/$/, "")}${path}`
  376. next.search = ""
  377. next.hash = ""
  378. return next
  379. }
  380. async function connectSSE(url: URL | string, headers: HeadersInit | undefined, signal: AbortSignal) {
  381. const res = await fetch(route(url, "/global/event"), {
  382. method: "GET",
  383. headers,
  384. signal,
  385. })
  386. if (!res.ok) throw new Error(`Workspace sync HTTP failure: ${res.status}`)
  387. if (!res.body) throw new Error("No response body from global sync")
  388. return res.body
  389. }
  390. async function syncHistory(space: Info, url: URL | string, headers: HeadersInit | undefined, signal: AbortSignal) {
  391. const sessionIDs = Database.use((db) =>
  392. db
  393. .select({ id: SessionTable.id })
  394. .from(SessionTable)
  395. .where(eq(SessionTable.workspace_id, space.id))
  396. .all()
  397. .map((row) => row.id),
  398. )
  399. const state = sessionIDs.length
  400. ? Object.fromEntries(
  401. Database.use((db) =>
  402. db.select().from(EventSequenceTable).where(inArray(EventSequenceTable.aggregate_id, sessionIDs)).all(),
  403. ).map((row) => [row.aggregate_id, row.seq]),
  404. )
  405. : {}
  406. log.info("syncing workspace history", {
  407. workspaceID: space.id,
  408. sessions: sessionIDs.length,
  409. known: Object.keys(state).length,
  410. })
  411. const requestHeaders = new Headers(headers)
  412. requestHeaders.set("content-type", "application/json")
  413. const res = await fetch(route(url, "/sync/history"), {
  414. method: "POST",
  415. headers: requestHeaders,
  416. body: JSON.stringify(state),
  417. signal,
  418. })
  419. if (!res.ok) {
  420. const body = await res.text()
  421. throw new Error(`Workspace history HTTP failure: ${res.status} ${body}`)
  422. }
  423. const events = await res.json()
  424. return WorkspaceContext.provide({
  425. workspaceID: space.id,
  426. fn: () => {
  427. for (const event of events) {
  428. SyncEvent.replay(
  429. {
  430. id: event.id,
  431. aggregateID: event.aggregate_id,
  432. seq: event.seq,
  433. type: event.type,
  434. data: event.data,
  435. },
  436. { publish: true },
  437. )
  438. }
  439. },
  440. })
  441. log.info("workspace history synced", {
  442. workspaceID: space.id,
  443. events: events.length,
  444. })
  445. }
  446. async function syncWorkspaceLoop(space: Info, signal: AbortSignal) {
  447. const adaptor = await getAdaptor(space.projectID, space.type)
  448. const target = await adaptor.target(space)
  449. if (target.type === "local") return null
  450. let attempt = 0
  451. while (!signal.aborted) {
  452. log.info("connecting to global sync", { workspace: space.name })
  453. setStatus(space.id, "connecting")
  454. let stream
  455. try {
  456. stream = await connectSSE(target.url, target.headers, signal)
  457. await syncHistory(space, target.url, target.headers, signal)
  458. } catch (err) {
  459. stream = null
  460. setStatus(space.id, "error")
  461. log.info("failed to connect to global sync", {
  462. workspace: space.name,
  463. err,
  464. })
  465. }
  466. if (stream) {
  467. attempt = 0
  468. log.info("global sync connected", { workspace: space.name })
  469. setStatus(space.id, "connected")
  470. await parseSSE(stream, signal, (evt: any) => {
  471. try {
  472. if (!("payload" in evt)) return
  473. if (evt.payload.type === "server.heartbeat") return
  474. if (evt.payload.type === "sync") {
  475. SyncEvent.replay(evt.payload.syncEvent as SyncEvent.SerializedEvent)
  476. }
  477. GlobalBus.emit("event", {
  478. directory: evt.directory,
  479. project: evt.project,
  480. workspace: space.id,
  481. payload: evt.payload,
  482. })
  483. } catch (err) {
  484. log.info("failed to replay global event", {
  485. workspaceID: space.id,
  486. error: err,
  487. })
  488. }
  489. })
  490. log.info("disconnected from global sync: " + space.id)
  491. setStatus(space.id, "disconnected")
  492. }
  493. // Back off reconnect attempts up to 2 minutes while the workspace
  494. // stays unavailable.
  495. await sleep(Math.min(120_000, 1_000 * 2 ** attempt))
  496. attempt += 1
  497. }
  498. }
  499. async function startSync(space: Info) {
  500. if (!Flag.OPENCODE_EXPERIMENTAL_WORKSPACES) return
  501. const adaptor = await getAdaptor(space.projectID, space.type)
  502. const target = await adaptor.target(space)
  503. if (target.type === "local") {
  504. void Filesystem.exists(target.directory).then((exists) => {
  505. setStatus(space.id, exists ? "connected" : "error")
  506. })
  507. return
  508. }
  509. if (aborts.has(space.id)) return true
  510. setStatus(space.id, "disconnected")
  511. const abort = new AbortController()
  512. aborts.set(space.id, abort)
  513. void syncWorkspaceLoop(space, abort.signal).catch((error) => {
  514. aborts.delete(space.id)
  515. setStatus(space.id, "error")
  516. log.warn("workspace listener failed", {
  517. workspaceID: space.id,
  518. error,
  519. })
  520. })
  521. }
  522. function stopSync(id: WorkspaceID) {
  523. aborts.get(id)?.abort()
  524. aborts.delete(id)
  525. connections.delete(id)
  526. }
  527. export function startWorkspaceSyncing(projectID: ProjectID) {
  528. const spaces = Database.use((db) =>
  529. db
  530. .select({ workspace: WorkspaceTable })
  531. .from(WorkspaceTable)
  532. .innerJoin(SessionTable, eq(SessionTable.workspace_id, WorkspaceTable.id))
  533. .where(eq(WorkspaceTable.project_id, projectID))
  534. .all(),
  535. )
  536. for (const row of new Map(spaces.map((row) => [row.workspace.id, row.workspace])).values()) {
  537. void startSync(fromRow(row))
  538. }
  539. }
  540. export * as Workspace from "./workspace"