Effect v4Module 8 of 8: Resilient production operation
Module 8 of 8Full operationsrc/course/engine.ts

Can the operation survive a failure halfway through?

Combine validation, ownership, retry, and concurrency policies in one unified pipeline.

First principle

A reliable operation makes each boundary explicit: input, outcome (A), failure (E), dependency (R), lifetime, time, and parallel work.

08 / 08module position
01

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.

First-principles consequence

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.

baseline-friction.tsbaseline problem
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" });
});
02

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 }
03

Minimal working example

Executable source

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

src/course/engine.tsFull operation
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 }
  )
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/engine.test.ts
Before you run it

How 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.
05

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.

Compiler or runtime diagnostic

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.

Remedy

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).

06

Scripted pipeline trace

Scripted pipeline trace • Conceptual model

Issue Import Engine: Follow the execution order

Illustrative trace

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.

Issue Import Runtime Execution Console
// Ready. Select an execution trial above to inspect the pipeline trace.