已发布 上游基线 bf46254 原文 ↗ 在 GitHub 编辑

Supervisor

Effect 的 Supervisor 负责管理 Fiber 的生命周期,让你能够跟踪、监控并控制应用内 Fiber 的行为。

Supervisor<A> 是 Effect 中用于管理 Fiber 的工具,它让你能够跟踪 Fiber 的生命周期(创建与终止),并产出一个类型为 A 的值来反映这种监督。当你需要洞察或控制应用中 Fiber 的行为时,Supervisor 会很有用。

要创建一个 supervisor,可以使用 Supervisor.track 函数。它会生成一个新的 supervisor,用来跟踪其子 Fiber,并把它们维护在一个集合中。这样你就可以在执行过程中观察和监控它们的状态。

你可以使用 Effect.supervised 函数来监督一个 effect。该函数接收一个 supervisor 作为参数,并返回一个 effect,其中在该 effect 内 fork 出来的所有子 Fiber 都由所提供的 supervisor 监督。由此,你可以通过 supervisor 捕获这些子 Fiber 的详细信息,例如它们的状态。

示例(监控 Fiber 数量)

在这个示例中,我们将使用一个 supervisor 定期监控应用中正在运行的 Fiber 数量。程序会计算一个斐波那契数,在此过程中生成多个 Fiber,同时另有一个监控器跟踪 Fiber 的数量。

import { Effect, Supervisor, Schedule, Fiber, FiberStatus } from "effect"

// Main program that monitors fibers while calculating a Fibonacci number
const program = Effect.gen(function* () {
  // Create a supervisor to track child fibers
  const supervisor = yield* Supervisor.track

  // Start a Fibonacci calculation, supervised by the supervisor
  const fibFiber = yield* fib(20).pipe(
    Effect.supervised(supervisor),
    // Fork the Fibonacci effect into a fiber
    Effect.fork,
  )

  // Define a schedule to periodically monitor the fiber count every 500ms
  const policy = Schedule.spaced("500 millis").pipe(
    Schedule.whileInputEffect((_) =>
      Fiber.status(fibFiber).pipe(
        // Continue while the Fibonacci fiber is not done
        Effect.andThen((status) => status !== FiberStatus.done),
      ),
    ),
  )

  // Start monitoring the fibers, using the supervisor to track the count
  const monitorFiber = yield* monitorFibers(supervisor).pipe(
    // Repeat the monitoring according to the schedule
    Effect.repeat(policy),
    // Fork the monitoring into its own fiber
    Effect.fork,
  )

  // Join the monitor and Fibonacci fibers to ensure they complete
  yield* Fiber.join(monitorFiber)
  const result = yield* Fiber.join(fibFiber)

  console.log(`fibonacci result: ${result}`)
})

// Function to monitor and log the number of active fibers
const monitorFibers = (
  supervisor: Supervisor.Supervisor<Array<Fiber.RuntimeFiber<any, any>>>,
): Effect.Effect<void> =>
  Effect.gen(function* () {
    const fibers = yield* supervisor.value // Get the current set of fibers
    console.log(`number of fibers: ${fibers.length}`)
  })

// Recursive Fibonacci calculation, spawning fibers for each recursive step
const fib = (n: number): Effect.Effect<number> =>
  Effect.gen(function* () {
    if (n <= 1) {
      return 1
    }
    yield* Effect.sleep("500 millis") // Simulate work by delaying

    // Fork two fibers for the recursive Fibonacci calls
    const fiber1 = yield* Effect.fork(fib(n - 2))
    const fiber2 = yield* Effect.fork(fib(n - 1))

    // Join the fibers to retrieve their results
    const v1 = yield* Fiber.join(fiber1)
    const v2 = yield* Fiber.join(fiber2)

    return v1 + v2 // Combine the results
  })

Effect.runPromise(program)
/*
Output:
number of fibers: 0
number of fibers: 2
number of fibers: 6
number of fibers: 14
number of fibers: 30
number of fibers: 62
number of fibers: 126
number of fibers: 254
number of fibers: 510
number of fibers: 1022
number of fibers: 2034
number of fibers: 3795
number of fibers: 5810
number of fibers: 6474
number of fibers: 4942
number of fibers: 2515
number of fibers: 832
number of fibers: 170
number of fibers: 18
number of fibers: 0
fibonacci result: 10946
*/