# Concurrent work needs an owner and a bound

> For agents: start with the [agent guide](https://ratstack.sh/llms.txt). Every page is Markdown by default; add `Accept: text/html` for HTML.

Free workshop: [how to burn a trillion tokens and get good results](/tokenmaxx#interested).

Choose how many operations may run at once. Give child work an owner that can stop it.
A completed caller must not silently leave its work running.

## Bound a traversal before running it

`Effect.forEach` describes one Effect per input and collects results in input order.
Its concurrency option bounds active operations, not the input collection's size.
The checked fixture measures active work with a Ref.

packages/capability/test/fixtures/wiki/bounded-traversal.ts at 4019cffbcbdb282d6f006854ff5b1f26d3b407e0; lines 8-30; highlighted none.

```typescript
    const activity = yield* Ref.make({ active: 0, maximum: 0 });
    const input = Array.from({ length: options.count }, (_, index) => index);

    const results = yield* Effect.forEach(
      input,
      (index) =>
        Effect.acquireUseRelease(
          Ref.update(activity, (state) => ({
            active: state.active + 1,
            maximum: Math.max(state.maximum, state.active + 1),
          })),
          () => Effect.yieldNow.pipe(Effect.andThen(Effect.succeed(index))),
          () =>
            Ref.update(activity, (state) => ({
              ...state,
              active: state.active - 1,
            }))
        ),
      { concurrency: options.concurrency }
    );

    return { ...(yield* Ref.get(activity)), results };
  }
```

The generated test varies both the input count and the bound.
It requires completed work, no remaining activity, input-order results, and a peak that does not exceed the supplied bound.

packages/capability/test/concurrency-ownership.test.ts at 4019cffbcbdb282d6f006854ff5b1f26d3b407e0; lines 7-26; highlighted none.

```typescript
const Traversal = Schema.Struct({
  concurrency: Schema.Literals([1, 2, 3, 4]),
  count: Schema.Literals([1, 2, 3, 4, 5, 6, 7, 8]),
});

it.effect.prop(
  "traversal never exceeds its supplied bound and preserves input order",
  { scenario: Arbitrary.schema(Traversal) },
  ({ scenario }) =>
    Effect.gen(function* traversalBound() {
      const observed = yield* boundedTraversal(scenario);
      expect(observed.maximum).toBeGreaterThan(0);
      expect(observed.maximum).toBeLessThanOrEqual(scenario.concurrency);
      expect(observed.active).toBe(0);
      expect(observed.results).toEqual(
        Array.from({ length: scenario.count }, (_, index) => index)
      );
    })
);

```

The expectation uses the supplied bound, not a copy of the traversal's scheduling algorithm.

A planted violation replaced the bound with `unbounded`.
The property failed with seven active operations against a bound of one. The restored implementation passes.
Measure production limits separately from these generated test counts.

## Stop child work with its Scope

The sandbox writes queued input to a child process through a scoped fiber:

packages/capability/src/sandbox-subprocess.ts at 65e9465f38092e24486392f22c08f45d61230c20; lines 224-229; highlighted none.

```typescript
      const outbox = yield* Queue.unbounded<Uint8Array, Cause.Done>();
      yield* Stream.fromQueue(outbox).pipe(
        Stream.run(handle.stdin),
        Effect.ignore,
        Effect.forkScoped
      );
```

The checked ownership test starts a child that waits indefinitely.
When the owner closes, the child must observe interruption before the enclosing Effect returns.

packages/capability/test/concurrency-ownership.test.ts at 4019cffbcbdb282d6f006854ff5b1f26d3b407e0; lines 30-46; highlighted none.

```typescript
    Effect.gen(function* childOwnership() {
      const started = yield* Deferred.make<boolean>();
      const interrupted = yield* Ref.make(0);
      yield* Effect.scoped(
        Effect.gen(function* ownChild() {
          yield* Deferred.succeed(started, true).pipe(
            Effect.andThen(Effect.never),
            Effect.onInterrupt(() =>
              Ref.update(interrupted, (count) => count + 1)
            ),
            Effect.forkScoped
          );
          yield* Deferred.await(started);
        })
      );
      expect(yield* Ref.get(interrupted)).toBe(1);
    })
```

A Deferred confirms that work started. The test does not depend on an arbitrary sleep.

Use [scopes](/lore/scopes-own-resources) to decide the lifetime before choosing a fork operation.
Different fork operations carry different ownership rules; do not treat their names as interchangeable.

## Consume a stream without promising a bounded queue

The sandbox decodes stdout, splits it into lines, and interprets each protocol message with an Effect.
Its consumer stops after the first completed outcome.

packages/capability/src/sandbox-subprocess.ts at 65e9465f38092e24486392f22c08f45d61230c20; lines 241-253,276-280; highlighted none.

```typescript
      const outcome = yield* Stream.decodeText(handle.stdout).pipe(
        Stream.splitLines,
        Stream.filter((line) => line.trim() !== ""),
        Stream.mapEffect((line) =>
          decodeChildMessage(line).pipe(
            Effect.mapError(
              (error) =>
                new SandboxError({
                  logs: [],
                  message: `Unreadable sandbox message: ${error.message}`,
                  reason: "protocol",
                })
            ),
        Stream.filter(Option.isSome),
        Stream.map((option) => option.value),
        Stream.runHead,
        Effect.mapError((error) =>
          isSandboxError(error)
```

`Stream.mapEffect` keeps decoding and dispatch in the stream description. The consumer does not launch eager Promises for each line.

Backpressure lets downstream demand limit upstream progress. An unbounded queue can decouple the producer from that demand.

The producer uses `Queue.unbounded`. This code does not prove a bounded producer queue or bounded total memory.
A bounded consumer is not enough when producers can accumulate data faster than it is consumed.

For a growing input, choose a stream and inspect each buffer's policy.
For a known finite collection, a bounded traversal can be enough.
Measure the service's limits before selecting an operational bound.

## Common mistakes

- Create eager Promises before building the owned Effect traversal.
- Detach work that should stop with the caller.
- Bound worker concurrency and claim every upstream buffer is bounded too.
- Copy test counts into production without measuring the provider and host.

## Effect idiom and house rule

- Fibers, traversal bounds and streams compose owned work without replacing the typed Effect description.
- In rat-stack, measure operational limits instead of inventing caps.
- In rat-stack, [machines own finite domain retries and cancellation](/lore/lifecycles-are-machines). Effect's native fibers and schedules still exist.

Next, [validate HTTP responses](/lore/http-responses-need-validation) before treating concurrent provider work as successful.
Return to the [Effect reading order](/lore/effect-basics#read-next).

## Sources

1. [Effect contributors. 2026. Effect. Effect 4.0.0.](https://github.com/Effect-TS/effect/blob/67ba4e46a11ccda0b6761578bfd22c04ae00167d/packages/effect/src/Effect.ts)
   Effect-TS. forEach limits traversal concurrency and retains input order; forkScoped attaches child interruption to the current Scope. Accessed 2026-10-06.

2. [Effect contributors. 2026. Stream.](https://github.com/Effect-TS/effect/blob/67ba4e46a11ccda0b6761578bfd22c04ae00167d/packages/effect/src/Stream.ts)
   Effect-TS. mapEffect performs effectful transformation; splitLines and runHead consume line-oriented output. Accessed 2026-10-06.

3. [Effect contributors. 2026. Consuming and transforming streams.](https://github.com/Effect-TS/effect/blob/67ba4e46a11ccda0b6761578bfd22c04ae00167d/ai-docs/src/03_stream/20_consuming-streams.ts)
   Effect-TS. Effectful stream transformation and consumption with explicit concurrency choices. Accessed 2026-10-06.

4. [rat-stack contributors. 2026. Subprocess sandbox ownership.](https://github.com/joelhooks/rat-stack/blob/65e9465f38092e24486392f22c08f45d61230c20/packages/capability/src/sandbox-subprocess.ts)
   rat-stack. A scoped stdin writer and decoded stdout stream; its producer queue is explicitly unbounded. Accessed 2026-10-06.

5. [rat-stack contributors. 2026. Traversal and child-ownership fixtures.](https://github.com/joelhooks/rat-stack/blob/4019cffbcbdb282d6f006854ff5b1f26d3b407e0/packages/capability/test/concurrency-ownership.test.ts)
   rat-stack. Generated traversals check peak active work, completion and order. A separate seam proves scoped child interruption. The bound property catches a planted unbounded traversal. Accessed 2026-10-06.

6. [Langton, Kit. 2026. Testing. Effect Solutions.](https://github.com/kitlangton/effect-solutions/blob/09f82e6c5c928e7232cd32daf04d7c6a830b63f7/packages/website/docs/08-testing.md)
   Kit Langton. Fiber examples frame the ownership question; the snapshot has no complete concurrency guide. Independent fixtures establish this page's claims. Accessed 2026-10-06.

## Linked from

- Agent guide → An Effect stack so pure (aspirational) Kit Langton will blush. → [Read page](https://ratstack.sh/llms.txt)

- Change log → What changed in the files served here, newest first. → [Read page](https://ratstack.sh/log)

- Effect basics → An Effect describes work, expected failures, and required services before a runtime executes it. → [Read page](https://ratstack.sh/lore/effect-basics)

- Full agent guide → An Effect stack so pure (aspirational) Kit Langton will blush. → [Read page](https://ratstack.sh/llms-full.txt)

- Rat Stack lore | rat-stack → Short, source-grounded notes on the ideas and decisions behind rat-stack. → [Read page](https://ratstack.sh/lore)

- Scopes release resources when work ends → Acquire resources lazily, keep their use inside the owner, and release them after success, failure, or interruption. → [Read page](https://ratstack.sh/lore/scopes-own-resources)

- Source change history (log.md) → The generated change log as Markdown. → [Read page](https://ratstack.sh/log.md)
