The concrete friction
A production request traverses every operational boundary simultaneously
A production endpoint cannot treat validation, locking, dependency management, transient recovery, cleanup, and concurrency as disconnected concerns. In a naive Express or Fastify handler, unvalidated input crashes the server, concurrent duplicate requests cause race conditions, unhandled 503s fail client workflows, and failed operations leak in-memory locks.
Partial failure in production is catastrophic if boundaries are ad-hoc. Without guaranteed finalizers, a crash midway leaves resources locked forever; without selective retries, transient network blips become permanent user-facing 500 errors.
type Request = { readonly body: unknown };
type Response = { readonly json: (body: unknown) => void };
const app = {
post: (_path: string, _handler: (req: Request, res: Response) => void) => undefined
};
const processIssueImport = async (_payload: unknown) => undefined;
// The naive API route has no operational policy
app.post("/imports/issues", async (req, res) => {
// Unvalidated input reaches the domain directly.
// Duplicate requests can process the same entities twice.
// Provider outages are retried without limits or error policies.
await processIssueImport(req.body);
res.json({ status: "ok" });
});The mental model
The unified end-to-end operational pipeline
The capstone engine composes all seven primitives: 1. Schema.decodeUnknownEffect validates raw input. 2. InFlightImportStore acquires a scoped deduplication lock. 3. IssueProvider executes with selective exponential retries. 4. ImportReportStore records results. 5. MetricsCollector tracks throughput and retries. 6. Scoped finalizers guarantee lock release on success, failure, or cancellation. 7. Batch processing runs under a bounded concurrency budget.
Issue Import Payload
│
▼ 1. Schema.decodeUnknownEffect (IssueImportValidationError)
Scoped Context Opens
│
├─► 2. InFlightImportStore.acquireLock (DuplicateIssueImportError)
│
├─► 3. IssueProvider.importIssues with Effect.retry
│ (IssueImportRejected fails immediately without retry)
│
├─► 4. ImportReportStore.recordImport
│
├─► 5. MetricsCollector.incrementProcessed
│
Scoped Context Closes
│
└─► 6. Finalizer removes importId from in-flight set
│
▼ Return { importId, issueCount }Minimal working example
Executable source
This is the checked source file used by the course. Read it before you run it.
import { Effect, Schedule, Schema, pipe } from "effect"
import { IssueImportPayload } from "./schema.js"
import {
DuplicateIssueImportError,
IssueImportRejected,
IssueImportValidationError,
IssueProviderUnavailable
} from "./errors.js"
import {
ImportReportStore,
InFlightImportStore,
IssueProvider,
MetricsCollector
} from "./services.js"
export const validateImport = (
rawPayload: unknown
): Effect.Effect<IssueImportPayload, IssueImportValidationError> =>
pipe(
Schema.decodeUnknownEffect(IssueImportPayload)(rawPayload),
Effect.mapError((schemaError) =>
IssueImportValidationError({ message: schemaError.message })
)
)
export const retryPolicy = Schedule.exponential("20 millis", 2).pipe(
Schedule.upTo({ times: 2 })
)
export const processIssueImport = (
rawPayload: unknown
): Effect.Effect<
{ readonly importId: string; readonly issueCount: number },
| IssueImportValidationError
| DuplicateIssueImportError
| IssueProviderUnavailable
| IssueImportRejected,
InFlightImportStore | IssueProvider | ImportReportStore | MetricsCollector
> =>
Effect.scoped(
Effect.gen(function* () {
const inFlightImports = yield* InFlightImportStore
const provider = yield* IssueProvider
const report = yield* ImportReportStore
const metrics = yield* MetricsCollector
const payload = yield* validateImport(rawPayload)
yield* inFlightImports.acquireLock(payload.importId)
const imported = yield* pipe(
provider.importIssues({
importId: payload.importId,
projectId: payload.projectId,
query: payload.query
}),
Effect.tapError((error) =>
error._tag === "IssueProviderUnavailable"
? metrics.incrementRetries()
: Effect.void
),
Effect.retry({
schedule: retryPolicy,
while: (error) => error._tag === "IssueProviderUnavailable"
})
)
yield* report.recordImport({
importId: payload.importId,
projectId: payload.projectId,
issueCount: imported.issueCount
})
yield* metrics.incrementProcessed()
return {
importId: payload.importId,
issueCount: imported.issueCount
}
})
)
export const processIssueImportBatch = (
rawPayloads: ReadonlyArray<unknown>,
concurrency = 4
): Effect.Effect<
ReadonlyArray<{ readonly importId: string; readonly issueCount: number }>,
| IssueImportValidationError
| DuplicateIssueImportError
| IssueProviderUnavailable
| IssueImportRejected,
InFlightImportStore | IssueProvider | ImportReportStore | MetricsCollector
> =>
Effect.all(
rawPayloads.map((rawPayload) => processIssueImport(rawPayload)),
{ concurrency }
)Execute and verify
Predict the result, run the command, and compare the output with the model.
pnpm exec tsx src/course/engine.test.tsHow do the concepts from Modules 1 through 7 combine in the processIssueImport pipeline? Which guarantees are local to this process, and which require durable coordination?
Reveal expected output
▶ Running Issue Import Engine Test Suite (Effect v4)...
✔ Test 1: Invalid payload rejected by Schema
✔ Test 2: Temporary provider errors retried and recorded
✔ Test 3: Bounded concurrent import batch completed
✔ Test 4: Retry policy stops after two retries
✔ Test 5: Permanent provider errors are not retried
✔ Test 6: Concurrent duplicate imports share one scoped lock
All Issue Import tests passed.Error anatomy and edge cases
Confusing local process guarantees with distributed system invariants
Assuming an in-memory Ref lock protects against duplicate requests in a clustered, multi-instance deployment.
An in-memory lock (HashSet in Ref) coordinates fibers within a single Node.js process. When multiple server instances run behind a load balancer, concurrent requests routed to different containers can process the same importId concurrently.
Use the in-process Scope and lock for local queue coordination, but enforce distributed idempotency at the database layer using unique constraints or a distributed lock (e.g. Redis/PostgreSQL advisory lock).
Scripted pipeline trace
Issue Import Engine: Follow the execution order
Select an import trial to follow the intended execution order. This panel replays a scripted trace; run the engine test command to execute the real implementation. Compare schema validation, scoped locks, selective retries, and bounded work with the source code.