---
slug: examples/spawn-patterns
title: Spawn patterns
kind: howto
surface: agents
summary: Three Python and Node patterns for Aetherfy task agents that spawn other tasks — an independent assessor that checks an output against a rubric, a manager that fans out one task per account and aggregates, and a child whose progress its parent reads from the run's logs.
sources:
  - aetherfy-vectors-python-sdk:aetherfy_agent/__init__.py
  - aetherfy-vectors-python-sdk:aetherfy_agent/models.py
  - aetherfy-vectors-python-sdk:aetherfy_agent/exceptions.py
  - aetherfy-vectors-python-sdk:aetherfy_memory/scope.py
  - aetherfy-vectors-js-sdk:src/agent/index.ts
  - aetherfy-vectors-js-sdk:src/agent/models.ts
  - aetherfy-vectors-js-sdk:src/agent/errors.ts
  - aetherfy-vectors-js-sdk:src/memory/scope.ts
  - aetherfy-control-plane:shared/job_runs.py
  - aetherfy-control-plane:api/routes/agents.py
  - aetherfy-control-plane:orchestrator/fly_manager.py
  - aetherfy-control-plane:supervisor/logs.go
  - dashboard/migrations/119_plans_max_in_flight_runs.sql
---

# Spawn patterns

Three patterns for an Aetherfy task agent that starts other task agents. Every one
of them is built from calls [The task contract](/agents/task-contract) already
documents — spawn with a payload, wait, read the result, fan out inside the
machine — plus the per-run logs on [Runs and logs](/agents/runs-and-logs). None
of them needs anything new from the platform.

Every file is given in Python, on `aetherfy_agent`, and in Node, on the
`aetherfy-vectors/agent` helper. The two helpers make the same calls under each
language's naming — `write_result` is `writeResult`, `fan_out` is `fanOut` — and
the Node calls return promises. `call_your_model`, `judge_with_your_model`,
`process_page` and `list_pages_to_crawl`, and their camel-case Node spellings,
stand for your own code.

## What every Aetherfy spawn pattern assumes

A parent can spawn a child only when both sides of the link are declared. The
parent enables spawning and names its workers, and each worker is a `type: job`
agent that is deployed and live before the parent is deployed. The full list of
conditions is on [afy spawn](/cli/spawn).

```yaml
name: drafter
type: job
runtime: python3.12
entrypoint: main.py
spawn:
  enabled: true
  workers:
    - assessor
```

```yaml
name: assessor
type: job
runtime: python3.12
entrypoint: main.py
```

The Node version of an agent declares a Node runtime and its entrypoint, and its
`package.json` sets `"type": "module"`, as on
[An agent that uses memory](/examples/agent-with-memory):

```yaml
name: assessor
type: job
runtime: node22
entrypoint: index.js
```

Four facts about Aetherfy spawns shape all three patterns:

| Fact | What it means for the pattern |
|---|---|
| `spawn` returns when the run is recorded, not when it finishes | Always wait for the child and read its `state` |
| `wait` holds a request open for 1 to 60 seconds, and a wait that times out returns the run as it stands | Waiting longer is another call, in a loop |
| Every task run counts toward the account's runs-in-flight limit, **the parent's own run included** | On a plan whose limit is 1, a task cannot spawn while it runs: the spawn is refused with `429 AGENT_RUN_CONCURRENCY_LIMIT_EXCEEDED`, which the SDK raises as `TooManyRunsInFlight`. The per-plan values are on [Limits](/platform/limits) |
| A run is never queued | A refused spawn is not retried for you; your code decides what to do instead |

Every child is a machine that is awake, and billed, for as long as it runs, and a
parent that waits is awake too. Both are held to the 60-minute backstop. When the
work does not need a separate machine, fan out inside the parent's machine instead;
the trade-off is on [The task contract](/agents/task-contract).

### runs.py and runs.js

The waiting loop the first two patterns share. Save it next to the parent's
entrypoint.

```python
"""Wait for an Aetherfy run to finish, one bounded wait at a time."""

import time

from aetherfy_agent import wait

# The terminal states in the state table on /agents/runs-and-logs. A run in
# any other state is still going.
FINISHED = {"completed", "failed", "superseded", "rolled_back"}


def finish(run_id, deadline_seconds=600):
    """Return the run once it is finished, or as it stands once the deadline
    has passed. Either way, read `state` before trusting `result`."""
    started = time.monotonic()
    run = wait(run_id, timeout_seconds=60)
    while run.state not in FINISHED and time.monotonic() - started < deadline_seconds:
        run = wait(run_id, timeout_seconds=60)
    return run
```

```javascript
// Wait for an Aetherfy run to finish, one bounded wait at a time.

import { wait } from 'aetherfy-vectors/agent';

// The terminal states in the state table on /agents/runs-and-logs. A run in
// any other state is still going.
export const FINISHED = new Set(['completed', 'failed', 'superseded', 'rolled_back']);

// Resolves to the run once it is finished, or as it stands once the deadline
// has passed. Either way, read `state` before trusting `result`.
export async function finish(runId, deadlineSeconds = 600) {
  const started = Date.now();
  let run = await wait(runId, 60);
  while (!FINISHED.has(run.state) && (Date.now() - started) / 1000 < deadlineSeconds) {
    run = await wait(runId, 60);
  }
  return run;
}
```

## An independent assessor on Aetherfy

The parent drafts an answer, and a separate task decides whether it is good enough
to return. The assessor runs on its own machine and sees only what is in its
payload, the output and the rubric, not the prompt or reasoning that produced
the draft. The parent **fails closed**: an assessment that failed, ran past its
deadline or returned nothing readable is treated as a fail, so the parent returns
nothing rather than something wrong.

### The parent: drafter, in Python

```python
"""An Aetherfy task that returns its answer only if an assessor passed it."""

import json

from aetherfy_agent import payload, spawn, write_result

from runs import finish

RUBRIC = [
    "The answer addresses the question that was asked",
    "Every figure in the answer is stated in the question or the source",
]


def main():
    data = payload()
    draft = call_your_model(data.get("question", ""))

    run = spawn("assessor", {"output": draft, "rubric": RUBRIC})
    verdict = finish(run.spawn_id, deadline_seconds=300)

    # Fail closed: only a completed run that explicitly passed counts as a pass.
    assessment = verdict.result if isinstance(verdict.result, dict) else {}
    if verdict.state == "completed" and assessment.get("passed") is True:
        write_result({"answer": draft})
        return

    print(json.dumps({
        "event": "answer_withheld",
        "assessor_state": verdict.state,
        "reasons": assessment.get("reasons"),
    }), flush=True)
    write_result({"answer": None})


if __name__ == "__main__":
    main()
```

### The parent: drafter, in Node

```javascript
// An Aetherfy task that returns its answer only if an assessor passed it.

import { payload, spawn, writeResult } from 'aetherfy-vectors/agent';

import { finish } from './runs.js';

const RUBRIC = [
  'The answer addresses the question that was asked',
  'Every figure in the answer is stated in the question or the source',
];

async function main() {
  const data = await payload();
  const draft = await callYourModel(data.question ?? '');

  const run = await spawn('assessor', { output: draft, rubric: RUBRIC });
  const verdict = await finish(run.spawn_id, 300);

  // Fail closed: only a completed run that explicitly passed counts as a pass.
  const assessment =
    verdict.result !== null && typeof verdict.result === 'object' ? verdict.result : {};
  if (verdict.state === 'completed' && assessment.passed === true) {
    await writeResult({ answer: draft });
    return;
  }

  console.log(JSON.stringify({
    event: 'answer_withheld',
    assessor_state: verdict.state,
    reasons: assessment.reasons ?? null,
  }));
  await writeResult({ answer: null });
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

### The child: assessor, in Python

```python
"""An Aetherfy task that judges one output against the rubric it was handed."""

from aetherfy_agent import payload, write_result


def main():
    data = payload()
    output = data.get("output")
    rubric = data.get("rubric") or []

    # Nothing to judge is a fail, never a pass.
    if not output or not rubric:
        write_result({"passed": False, "reasons": ["no output or no rubric in the payload"]})
        return

    reasons = []
    for criterion in rubric:
        met, why = judge_with_your_model(output, criterion)
        if not met:
            reasons.append(f"{criterion}: {why}")

    write_result({"passed": not reasons, "reasons": reasons})


if __name__ == "__main__":
    main()
```

### The child: assessor, in Node

```javascript
// An Aetherfy task that judges one output against the rubric it was handed.

import { payload, writeResult } from 'aetherfy-vectors/agent';

async function main() {
  const data = await payload();
  const output = data.output;
  const rubric = data.rubric ?? [];

  // Nothing to judge is a fail, never a pass.
  if (!output || rubric.length === 0) {
    await writeResult({ passed: false, reasons: ['no output or no rubric in the payload'] });
    return;
  }

  const reasons = [];
  for (const criterion of rubric) {
    const { met, why } = await judgeWithYourModel(output, criterion);
    if (!met) reasons.push(`${criterion}: ${why}`);
  }

  await writeResult({ passed: reasons.length === 0, reasons });
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

If the assessor raises or rejects, its process exits non-zero, the run is `failed`, and the
parent withholds the answer. The payload holds the draft itself, so it is subject to
the 256 KB inline cap. A longer output belongs in a collection, with its id in the
payload.

## A manager that fans out on Aetherfy

The manager reads shared definitions from an Aetherfy memory namespace, spawns one
task per account with that account's definition in the payload, and aggregates
what the tasks return. It spawns the same child once per account; how those runs
are placed on machines is on [The task contract](/agents/task-contract).

The namespace `account-playbooks` holds one item per account: the definition as
the item's `text`, and the account as `metadata.account`. The manager's payload
names the accounts to run, as `{"accounts": ["acme", "globex"]}`.

### The parent: manager, in Python

```python
"""An Aetherfy task that runs one worker per account and aggregates the results."""

from aetherfy_agent import TooManyRunsInFlight, fan_out, payload, spawn, write_result
from aetherfy_memory import MemoryClient

from runs import finish


def load_definitions():
    """Every definition in the namespace, keyed by the account it is for."""
    memory = MemoryClient()
    ns = memory.namespace("account-playbooks")
    definitions = {}
    for point in ns.iter():
        item = point["payload"]
        account = (item.get("metadata") or {}).get("account")
        if account:
            definitions[account] = item["text"]
    return definitions


def main():
    accounts = payload().get("accounts") or []
    definitions = load_definitions()

    started, not_started = {}, {}
    for account in accounts:
        if account not in definitions:
            not_started[account] = "no definition in account-playbooks"
            continue
        try:
            run = spawn("account-worker", {"account": account, "definition": definitions[account]})
        except TooManyRunsInFlight:
            # The account's runs-in-flight limit is full, this manager's own
            # run included. Nothing is queued, so record it and carry on.
            not_started[account] = "runs-in-flight limit reached"
            continue
        started[account] = run.spawn_id

    # Waiting is I/O, so wait on every child at once, on threads in this machine.
    finished = fan_out(finish, list(started.values()))

    report = {"results": {}, "failed": {}, "not_started": not_started}
    for account, run in zip(started, finished):
        if run.state == "completed" and run.has_result:
            report["results"][account] = run.result
        else:
            report["failed"][account] = run.error_message or run.result_error or run.state
    write_result(report)


if __name__ == "__main__":
    main()
```

### The parent: manager, in Node

```javascript
// An Aetherfy task that runs one worker per account and aggregates the results.

import { MemoryClient } from 'aetherfy-vectors';
import { TooManyRunsInFlight, fanOut, payload, spawn, writeResult } from 'aetherfy-vectors/agent';

import { finish } from './runs.js';

// Every definition in the namespace, keyed by the account it is for.
async function loadDefinitions() {
  const memory = new MemoryClient();
  const ns = await memory.namespace('account-playbooks');
  const definitions = new Map();
  for await (const point of ns.iter()) {
    const item = point.payload ?? {};
    const account = item.metadata?.account;
    if (account) definitions.set(account, item.text);
  }
  memory.dispose();
  return definitions;
}

async function main() {
  const accounts = (await payload()).accounts ?? [];
  const definitions = await loadDefinitions();

  const started = new Map();
  const notStarted = {};
  for (const account of accounts) {
    if (!definitions.has(account)) {
      notStarted[account] = 'no definition in account-playbooks';
      continue;
    }
    try {
      const run = await spawn('account-worker', {
        account,
        definition: definitions.get(account),
      });
      started.set(account, run.spawn_id);
    } catch (error) {
      if (!(error instanceof TooManyRunsInFlight)) throw error;
      // The account's runs-in-flight limit is full, this manager's own run
      // included. Nothing is queued, so record it and carry on.
      notStarted[account] = 'runs-in-flight limit reached';
    }
  }

  // Waiting is I/O, so wait on every child at once, as concurrent promises.
  const finished = await fanOut(finish, [...started.values()]);

  const report = { results: {}, failed: {}, not_started: notStarted };
  [...started.keys()].forEach((account, index) => {
    const run = finished[index];
    if (run.state === 'completed' && run.has_result) {
      report.results[account] = run.result;
    } else {
      report.failed[account] = run.error_message ?? run.result_error ?? run.state;
    }
  });
  await writeResult(report);
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

### The child: account-worker, in Python

```python
"""An Aetherfy task that applies one account's definition and returns a summary."""

from aetherfy_agent import payload, write_result


def main():
    data = payload()
    summary = call_your_model(f"Apply this playbook to {data['account']}: {data['definition']}")
    write_result({"summary": summary})


if __name__ == "__main__":
    main()
```

### The child: account-worker, in Node

```javascript
// An Aetherfy task that applies one account's definition and returns a summary.

import { payload, writeResult } from 'aetherfy-vectors/agent';

async function main() {
  const data = await payload();
  const summary = await callYourModel(
    `Apply this playbook to ${data.account}: ${data.definition}`
  );
  await writeResult({ summary });
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

The manager's own result is subject to the same inline cap as each child's, and it
holds every child's result at once. Keep each worker's result to a short summary,
and write anything longer to a collection, returning its id.

`fan_out` and `fanOut` return results in input order, which is what lets the
manager put them back onto the accounts. A child that fails is not an exception
here: `finish` returns its run, and the manager records it. A wait that fails,
because the control plane could not be reached, is raised again by the helper once
every other wait has returned, and the manager's run fails.

## Progress from a child on Aetherfy

A long child can report progress by printing one JSON line per milestone, with an
`event` key. Its parent reads the child's logs for that run, filtered by the run's
id, until the run is finished. A person following the same run reads the same
lines with the CLI:

```bash
afy logs page-crawler --run 9f3a1c72-5b8e-4d61-a0f4-7c2e9b5d3a18
```

### Staying inside the Aetherfy log caps

Progress lines are ordinary log lines, so they share the caps on
[Runs and logs](/agents/runs-and-logs) with everything else the child prints:

| Limit | Value |
|---|---|
| Per line | 4 KB |
| Chunks per minute, per agent | 60 |

The per-minute cap is per **agent**, so every run of the same child, including
siblings from a fan-out, shares it. Print a milestone, not a line per item, and do
not use logs as a high-frequency stream. A child that logs past the caps loses
lines, and the stream shows `[SYSTEM] N log line(s) dropped` where they went.

Lines also reach the log store in batches, a few seconds behind the child and
longer while the child is over the caps. Treat
progress as a view of the run, not a record of it: the child's result is what it
finished, and the parent reads that from the run.

### The child: page-crawler, in Python

```python
"""An Aetherfy task that prints one JSON line per milestone."""

import json

from aetherfy_agent import payload, write_result


def milestone(event, **fields):
    print(json.dumps({"event": event, **fields}), flush=True)


def main():
    pages = payload().get("pages") or []
    milestone("started", total=len(pages))

    for done, page in enumerate(pages, start=1):
        process_page(page)
        if done % 25 == 0:
            milestone("progress", done=done, total=len(pages))

    milestone("finished", total=len(pages))
    write_result({"pages": len(pages)})


if __name__ == "__main__":
    main()
```

### The child: page-crawler, in Node

```javascript
// An Aetherfy task that prints one JSON line per milestone.

import { payload, writeResult } from 'aetherfy-vectors/agent';

function milestone(event, fields = {}) {
  console.log(JSON.stringify({ event, ...fields }));
}

async function main() {
  const pages = (await payload()).pages ?? [];
  milestone('started', { total: pages.length });

  for (const [index, page] of pages.entries()) {
    await processPage(page);
    const done = index + 1;
    if (done % 25 === 0) milestone('progress', { done, total: pages.length });
  }

  milestone('finished', { total: pages.length });
  await writeResult({ pages: pages.length });
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

### The parent: following the child's logs

The parent reads the logs endpoint documented on
[the runs API](/agents/api-runs) with `deployment_id` set to the run's id. Passing
`after_id` returns lines oldest first, so the id of the last line read is the cursor
for the next read. Each read follows a bounded `wait`, so the parent reads at most
once every 15 seconds and stops as soon as the run is finished.

In Python:

```python
"""An Aetherfy task that spawns a child and follows its progress."""

import json
import os
import urllib.parse
import urllib.request

from aetherfy_agent import spawn, wait, write_result

CHILD = "page-crawler"
FINISHED = {"completed", "failed", "superseded", "rolled_back"}
PAGE = 1000  # the most lines one read returns


def read_lines(run_id, after_id):
    """This run's stdout lines with an id above after_id, oldest first."""
    query = urllib.parse.urlencode(
        {"deployment_id": run_id, "after_id": after_id, "stream": "stdout", "tail": PAGE}
    )
    request = urllib.request.Request(
        f"{os.environ['AETHERFY_API_URL']}/agents/{CHILD}/logs?{query}",
        headers={
            # A task's key is issued per run, so read it from the environment.
            "Authorization": f"Bearer {os.environ['AETHERFY_API_KEY']}",
            # Always set a User-Agent: urllib's default is refused at the edge.
            "User-Agent": "progress-follower/1.0",
        },
    )
    with urllib.request.urlopen(request, timeout=30) as response:
        return json.load(response)


def report(lines):
    for line in lines:
        try:
            record = json.loads(line["message"])
        except ValueError:
            continue  # an ordinary log line, not a milestone
        if isinstance(record, dict) and "event" in record:
            print(json.dumps({"child_event": record}), flush=True)


def main():
    run = spawn(CHILD, {"pages": list_pages_to_crawl()})

    after_id = 0
    while True:
        current = wait(run.spawn_id, timeout_seconds=15)
        lines = read_lines(run.spawn_id, after_id)
        report(lines)
        if lines:
            after_id = lines[-1]["id"]
        # A full page means more lines are waiting: read again before stopping.
        if current.state in FINISHED and len(lines) < PAGE:
            break

    write_result({"child_state": current.state, "child_result": current.result})


if __name__ == "__main__":
    main()
```

In Node, with the built-in `fetch`:

```javascript
// An Aetherfy task that spawns a child and follows its progress.

import { spawn, wait, writeResult } from 'aetherfy-vectors/agent';

const CHILD = 'page-crawler';
const FINISHED = new Set(['completed', 'failed', 'superseded', 'rolled_back']);
const PAGE = 1000; // the most lines one read returns

// This run's stdout lines with an id above afterId, oldest first.
async function readLines(runId, afterId) {
  const query = new URLSearchParams({
    deployment_id: runId,
    after_id: String(afterId),
    stream: 'stdout',
    tail: String(PAGE),
  });
  const response = await fetch(`${process.env.AETHERFY_API_URL}/agents/${CHILD}/logs?${query}`, {
    headers: {
      // A task's key is issued per run, so read it from the environment.
      Authorization: `Bearer ${process.env.AETHERFY_API_KEY}`,
      // Always set a User-Agent of your own for requests to Aetherfy's API hosts.
      'User-Agent': 'progress-follower/1.0',
    },
  });
  if (!response.ok) throw new Error(`reading the logs failed with ${response.status}`);
  return response.json();
}

function report(lines) {
  for (const line of lines) {
    let record;
    try {
      record = JSON.parse(line.message);
    } catch {
      continue; // an ordinary log line, not a milestone
    }
    if (record !== null && typeof record === 'object' && 'event' in record) {
      console.log(JSON.stringify({ child_event: record }));
    }
  }
}

async function main() {
  const run = await spawn(CHILD, { pages: await listPagesToCrawl() });

  let afterId = 0;
  let current;
  for (;;) {
    current = await wait(run.spawn_id, 15);
    const lines = await readLines(run.spawn_id, afterId);
    report(lines);
    if (lines.length > 0) afterId = lines[lines.length - 1].id;
    // A full page means more lines are waiting: read again before stopping.
    if (FINISHED.has(current.state) && lines.length < PAGE) break;
  }

  await writeResult({ child_state: current.state, child_result: current.result });
}

main().catch((error) => {
  console.error(error);
  process.exit(1);
});
```

The loop reads once more after the run finishes, but a milestone printed in the
child's last seconds can still be in a batch on its way. The parent's result is
built from the child's run, not from the last line it saw.
