run-service.ts 2.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647
  1. import { Effect, Fiber, Layer, ManagedRuntime } from "effect"
  2. import * as Context from "effect/Context"
  3. import { InstanceRef, WorkspaceRef } from "./instance-ref"
  4. import * as Observability from "@kirincode-ai/core/observability"
  5. import { WorkspaceContext } from "@/control-plane/workspace-context"
  6. import type { InstanceContext } from "@/project/instance-context"
  7. import { memoMap } from "@kirincode-ai/core/effect/memo-map"
  8. type Refs = {
  9. instance?: InstanceContext
  10. workspace?: string
  11. }
  12. export function attachWith<A, E, R>(effect: Effect.Effect<A, E, R>, refs: Refs): Effect.Effect<A, E, R> {
  13. if (!refs.instance && !refs.workspace) return effect
  14. if (!refs.instance) return effect.pipe(Effect.provideService(WorkspaceRef, refs.workspace))
  15. if (!refs.workspace) return effect.pipe(Effect.provideService(InstanceRef, refs.instance))
  16. return effect.pipe(
  17. Effect.provideService(InstanceRef, refs.instance),
  18. Effect.provideService(WorkspaceRef, refs.workspace),
  19. )
  20. }
  21. export function attach<A, E, R>(effect: Effect.Effect<A, E, R>): Effect.Effect<A, E, R> {
  22. const workspace = WorkspaceContext.workspaceID
  23. const fiber = Fiber.getCurrent()
  24. return attachWith(effect, {
  25. instance: fiber ? Context.getReferenceUnsafe(fiber.context, InstanceRef) : undefined,
  26. workspace: workspace ?? (fiber ? Context.getReferenceUnsafe(fiber.context, WorkspaceRef) : undefined),
  27. })
  28. }
  29. export function makeRuntime<I, S, E>(service: Context.Service<I, S>, layer: Layer.Layer<I, E>) {
  30. let rt: ManagedRuntime.ManagedRuntime<I, E> | undefined
  31. const getRuntime = () => (rt ??= ManagedRuntime.make(Layer.provideMerge(layer, Observability.layer), { memoMap }))
  32. return {
  33. runSync: <A, Err>(fn: (svc: S) => Effect.Effect<A, Err, I>) => getRuntime().runSync(attach(service.use(fn))),
  34. runPromiseExit: <A, Err>(fn: (svc: S) => Effect.Effect<A, Err, I>, options?: Effect.RunOptions) =>
  35. getRuntime().runPromiseExit(attach(service.use(fn)), options),
  36. runPromise: <A, Err>(fn: (svc: S) => Effect.Effect<A, Err, I>, options?: Effect.RunOptions) =>
  37. getRuntime().runPromise(attach(service.use(fn)), options),
  38. runFork: <A, Err>(fn: (svc: S) => Effect.Effect<A, Err, I>) => getRuntime().runFork(attach(service.use(fn))),
  39. runCallback: <A, Err>(fn: (svc: S) => Effect.Effect<A, Err, I>) =>
  40. getRuntime().runCallback(attach(service.use(fn))),
  41. }
  42. }