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 already documents — spawn with a payload, wait, read the result, fan out inside the machine — plus the per-run logs on 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.
name: drafter
type: job
runtime: python3.12
entrypoint: main.py
spawn:
enabled: true
workers:
- assessorname: assessor
type: job
runtime: python3.12
entrypoint: main.pyThe 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:
name: assessor
type: job
runtime: node22
entrypoint: index.jsFour 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 |
| 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.
runs.py and runs.js
The waiting loop the first two patterns share. Save it next to the parent’s entrypoint.
"""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// 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
"""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
// 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
"""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
// 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.
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
"""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
// 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
"""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
// 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:
afy logs page-crawler --run 9f3a1c72-5b8e-4d61-a0f4-7c2e9b5d3a18Staying inside the Aetherfy log caps
Progress lines are ordinary log lines, so they share the caps on 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
"""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
// 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 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:
"""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:
// 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.