bridge.ts 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384
  1. import { Context, Effect, Exit, Fiber } from "effect"
  2. import { WorkspaceContext } from "@/control-plane/workspace-context"
  3. import type { WorkspaceV2 } from "@kirincode-ai/core/workspace"
  4. import { InstanceRef, WorkspaceRef } from "./instance-ref"
  5. import { attachWith } from "./run-service"
  6. export interface Shape {
  7. readonly promise: <A, E, R>(effect: Effect.Effect<A, E, R>) => Promise<A>
  8. readonly fork: <A, E, R>(effect: Effect.Effect<A, E, R>) => Fiber.Fiber<A, E>
  9. readonly run: <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E>
  10. readonly bind: <Args extends readonly unknown[], Result>(fn: (...args: Args) => Result) => (...args: Args) => Result
  11. }
  12. function restoreWorkspace<R>(workspace: WorkspaceV2.ID | undefined, fn: () => R): R {
  13. if (workspace !== undefined) return WorkspaceContext.restore(workspace, fn)
  14. return fn()
  15. }
  16. function captureSync() {
  17. const fiber = Fiber.getCurrent()
  18. const instance = fiber ? Context.getReferenceUnsafe(fiber.context, InstanceRef) : undefined
  19. const workspace =
  20. (fiber ? Context.getReferenceUnsafe(fiber.context, WorkspaceRef) : undefined) ?? WorkspaceContext.workspaceID
  21. return { instance, workspace }
  22. }
  23. export const bind = <Args extends readonly unknown[], Result>(fn: (...args: Args) => Result) => {
  24. const captured = captureSync()
  25. return (...args: Args) =>
  26. restoreWorkspace(captured.workspace, () =>
  27. Effect.runSync(
  28. attachWith(
  29. Effect.sync(() => fn(...args)),
  30. captured,
  31. ),
  32. ),
  33. )
  34. }
  35. /**
  36. * Bridge from Effect into a Promise-returning JS callback while preserving
  37. * `WorkspaceContext` AsyncLocalStorage for callback code that still reads it.
  38. * `InstanceRef` is captured for effects run through the returned bridge APIs;
  39. * plain JS callbacks that need it should receive the ref explicitly.
  40. *
  41. * Mirrors `Effect.promise` but restores workspace ALS first.
  42. */
  43. export const fromPromise = <T>(fn: () => Promise<T> | T): Effect.Effect<T> =>
  44. Effect.gen(function* () {
  45. const workspace = yield* WorkspaceRef
  46. return yield* Effect.promise(() => Promise.resolve(restoreWorkspace(workspace, () => fn())))
  47. })
  48. export function make(): Effect.Effect<Shape> {
  49. return Effect.gen(function* () {
  50. const ctx = yield* Effect.context()
  51. const captured = captureSync()
  52. const instance = (yield* InstanceRef) ?? captured.instance
  53. const workspace = (yield* WorkspaceRef) ?? captured.workspace
  54. const wrap = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
  55. attachWith(effect.pipe(Effect.provide(ctx)) as Effect.Effect<A, E, never>, { instance, workspace })
  56. return {
  57. promise: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
  58. restoreWorkspace(workspace, () => Effect.runPromise(wrap(effect))),
  59. fork: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
  60. restoreWorkspace(workspace, () => Effect.runFork(wrap(effect))),
  61. run: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
  62. Effect.callback<A, E>((resume) => {
  63. restoreWorkspace(workspace, () =>
  64. Effect.runPromiseExit(wrap(effect)).then((exit) =>
  65. resume(Exit.isSuccess(exit) ? Effect.succeed(exit.value) : Effect.failCause(exit.cause)),
  66. ),
  67. )
  68. }),
  69. bind:
  70. <Args extends readonly unknown[], Result>(fn: (...args: Args) => Result) =>
  71. (...args: Args) =>
  72. restoreWorkspace(workspace, () => Effect.runSync(wrap(Effect.sync(() => fn(...args))))),
  73. } satisfies Shape
  74. })
  75. }
  76. export * as EffectBridge from "./bridge"