Effect v4Module 7 of 8: Bounded parallel work
Module 7 of 8Concurrency policysrc/course/examples/concurrency.ts

What happens when 500 records start at once?

Limit active concurrent work before the provider, browser, or database becomes the bottleneck.

First principle

Concurrency is a resource budget. Starting more work than the system can absorb turns latency into failure.

07 / 08module position
01

The concrete friction

Promise.all unleashes unbounded concurrent execution

In JavaScript, Promise.all(records.map(processRecord)) invokes processRecord on every item immediately. If a batch contains 500 items, 500 concurrent HTTP connections or database transactions are opened simultaneously. This triggers connection pool exhaustion, file descriptor limits (EMFILE), upstream 429 Too Many Requests rate-limiting, and severe event loop lag.

First-principles consequence

Hardware has finite physical bandwidth and queue depth (Little's Law: L = λW). Unbounded concurrency converts throughput into latency and crashes downstream services. Furthermore, Promise.all cannot cancel in-flight work if one item fails.

baseline-friction.tsbaseline problem
const issueIds = ["ISSUE-101", "ISSUE-102", "ISSUE-103"];

const loadIssue = async (issueId: string): Promise<string> => {
  const response = await fetch("/api/issues/" + issueId);
  if (!response.ok) throw new Error("HTTP " + response.status);
  return "loaded " + issueId;
};

// map calls loadIssue for every ID before Promise.all begins waiting.
const issues = await Promise.all(issueIds.map(loadIssue));
02

The mental model

Structured concurrency enforces bounded worker budgets

Effect.all(effects, { concurrency: N }) executes a collection of Effects under a strict concurrency budget. If concurrency is 2 and there are 3 tasks, exactly 2 fibers execute while the third waits in queue. When any fiber completes, the next queued task is scheduled. The parent fiber supervises all child fibers: if one fails or the parent is interrupted, pending tasks are cancelled automatically.

Issue IDs: [101, 102, 103] with { concurrency: 2 }

t=0ms:   [101: Active]   [102: Active]   [103: Queued]
         (max active = 2)

t=50ms:  [101: Complete] ──> 103 becomes Active
         [102: Active]   [103: Active]

t=100ms: [102: Complete] [103: Complete]
03

Minimal working example

Executable source

This is the checked source file used by the course. Read it before you run it.

src/course/examples/concurrency.tsConcurrency policy
import { Console, Effect, Ref } from "effect"

const loadIssue = (
  issueId: string,
  activeRef: Ref.Ref<number>,
  maxActiveRef: Ref.Ref<number>
) =>
  Effect.acquireUseRelease(
    Effect.gen(function* () {
      const active = yield* Ref.modify(activeRef, (count) => [count + 1, count + 1])
      yield* Ref.update(maxActiveRef, (maximum) => Math.max(maximum, active))
    }),
    () => Effect.sleep("50 millis").pipe(Effect.as(`loaded ${issueId}`)),
    () => Ref.update(activeRef, (count) => count - 1)
  )

export const boundedIssueImport = Effect.gen(function* () {
  const activeRef = yield* Ref.make(0)
  const maxActiveRef = yield* Ref.make(0)
  const issueIds = ["ISSUE-101", "ISSUE-102", "ISSUE-103"]
  const results = yield* Effect.all(
    issueIds.map((issueId) => loadIssue(issueId, activeRef, maxActiveRef)),
    { concurrency: 2 }
  )
  const maxActive = yield* Ref.get(maxActiveRef)
  return { results, maxActive }
})

await Effect.runPromise(
  Effect.gen(function* () {
    const completed = yield* boundedIssueImport
    yield* Console.log(`max active: ${completed.maxActive}`)
    yield* Console.log(`results: ${completed.results.join(", ")}`)
  })
)
04

Execute and verify

Predict the result, run the command, and compare the output with the model.

Run this command
pnpm exec tsx src/course/examples/concurrency.ts
Before you run it

What happens to the third issue when Effect.all runs three detail requests with concurrency: 2?

Reveal expected output
max active: 2
results: loaded ISSUE-101, loaded ISSUE-102, loaded ISSUE-103
05

Error anatomy and edge cases

The interruption boundary: local fiber vs. remote server execution

Assuming that interrupting an Effect cancels an HTTP mutation that the external server has already processed.

Compiler or runtime diagnostic

Local fiber cancellation terminates local processing and closes local sockets, but an external API that already received the payload will commit the change. Subsequent retries may create duplicate records if the endpoint is not idempotent.

Remedy

Use bounded concurrency for local pressure control, wire request cancellation (AbortSignal) where supported, and design external mutations to be idempotent.

06

Active lab exercise

src/exercises/exercise7/starter.ts

Exercise: bound issue loading concurrency

Limit active issue loads to two while allowing the complete batch to finish.

ObjectiveFix starter.ts so the maximum active request count is two or less.
pnpm exec tsx src/exercises/exercise7/exercise.test.ts starter

The harness should fail against the starter. Inspect the assertion, change the starter implementation, and run it again.

Inspect starter code
Broken starter implementation
src/exercises/exercise7/starter.tsstarter
import { Effect, Ref } from "effect"

const loadIssue = (
  issueId: string,
  activeRef: Ref.Ref<number>,
  maxActiveRef: Ref.Ref<number>
) =>
  Effect.acquireUseRelease(
    Effect.gen(function* () {
      const active = yield* Ref.modify(activeRef, (count) => [count + 1, count + 1])
      yield* Ref.update(maxActiveRef, (maximum) => Math.max(maximum, active))
    }),
    () => Effect.sleep("20 millis").pipe(Effect.as(`loaded ${issueId}`)),
    () => Ref.update(activeRef, (count) => count - 1)
  )

export const boundedIssueImport = Effect.gen(function* () {
  const activeRef = yield* Ref.make(0)
  const maxActiveRef = yield* Ref.make(0)
  const issueIds = ["ISSUE-101", "ISSUE-102", "ISSUE-103"]
  const results = yield* Effect.all(
    issueIds.map((issueId) => loadIssue(issueId, activeRef, maxActiveRef)),
    { concurrency: "unbounded" }
  )
  const maxActive = yield* Ref.get(maxActiveRef)
  return { results, maxActive }
})
Reveal reference solution
Target reference solution
solution.tsverified solution
import { Effect, Ref } from "effect"

const loadIssue = (
  issueId: string,
  activeRef: Ref.Ref<number>,
  maxActiveRef: Ref.Ref<number>
) =>
  Effect.acquireUseRelease(
    Effect.gen(function* () {
      const active = yield* Ref.modify(activeRef, (count) => [count + 1, count + 1])
      yield* Ref.update(maxActiveRef, (maximum) => Math.max(maximum, active))
    }),
    () => Effect.sleep("20 millis").pipe(Effect.as(`loaded ${issueId}`)),
    () => Ref.update(activeRef, (count) => count - 1)
  )

export const boundedIssueImport = Effect.gen(function* () {
  const activeRef = yield* Ref.make(0)
  const maxActiveRef = yield* Ref.make(0)
  const issueIds = ["ISSUE-101", "ISSUE-102", "ISSUE-103"]
  const results = yield* Effect.all(
    issueIds.map((issueId) => loadIssue(issueId, activeRef, maxActiveRef)),
    { concurrency: 2 }
  )
  const maxActive = yield* Ref.get(maxActiveRef)
  return { results, maxActive }
})
Why this solution works

Pass { concurrency: 2 } to Effect.all. The third issue waits for one of the first two loads to release its active-work slot.