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.
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.
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));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]Minimal working example
Executable source
This is the checked source file used by the course. Read it before you run it.
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(", ")}`)
})
)Execute and verify
Predict the result, run the command, and compare the output with the model.
pnpm exec tsx src/course/examples/concurrency.tsWhat 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-103Error 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.
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.
Use bounded concurrency for local pressure control, wire request cancellation (AbortSignal) where supported, and design external mutations to be idempotent.
Active lab exercise
Exercise: bound issue loading concurrency
Limit active issue loads to two while allowing the complete batch to finish.
pnpm exec tsx src/exercises/exercise7/exercise.test.ts starterThe harness should fail against the starter. Inspect the assertion, change the starter implementation, and run it again.
Inspect starter code
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
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 }
})Pass { concurrency: 2 } to Effect.all. The third issue waits for one of the first two loads to release its active-work slot.