Планируйте ready set, а не очередь разговоров
DAG-scheduler запускает только узлы с завершёнными зависимостями, ограничивает concurrency и делает critical path видимым. В каждой итерации он вычисляет ready set - pending work items, у которых все зависимости завершены, - и берёт из него не больше concurrency limit. После завершения batch появляются новые кандидаты. Это проще и надёжнее, чем просить supervisor "следить, кто уже закончил" в длинном transcript.
import type { TaskStatus, WorkItem } from "./types.js";
export function readyTasks(
tasks: readonly WorkItem[],
status: Readonly<Record<string, TaskStatus>>
): readonly WorkItem[] {
return tasks.filter((task) =>
status[task.id] === "pending" &&
task.dependsOn.every((dependency) => status[dependency] === "completed")
);
}
export function createInitialStatus(tasks: readonly WorkItem[]): Record<string, TaskStatus> {
return Object.fromEntries(tasks.map((task) => [task.id, "pending" as const]));
}
export function hasUnfinishedTasks(status: Readonly<Record<string, TaskStatus>>): boolean {
return Object.values(status).some((value) => value === "pending" || value === "running");
}Сам батч исполняется с ограниченной параллельностью, а результаты проверяются на объявленный write set до записи.
const store = new ArtifactStore(handlers.map((handler) => handler.contract));
let failureReason: string | undefined;
eventLog.append({ type: "run.started", runId });
while (hasUnfinishedTasks(status) && !failureReason) {
const batch = readyTasks(this.#tasks, status).slice(0, this.#maxConcurrency);
if (batch.length === 0) {
failureReason = "No runnable task remains";
break;
}
for (const task of batch) {
status[task.id] = "running";
}
const outcomes = await Promise.all(
batch.map(async (task) => {
const agentId = routeTask(task);
const handler = this.#handlers.get(agentId);
if (!handler) {
return { task, reason: `No handler for ${agentId}` };
}
assertTaskAccepted(handler.contract, task);
for (let attempt = 1; attempt <= task.maxAttempts; attempt += 1) {
eventLog.append({ type: "task.started", taskId: task.id, agentId, attempt });
const openSpan = trace.start(task.id, agentId, attempt);
try {
const result = await handler.run({
runId,
task,
inputArtifacts: store.list(),
attempt
});
budget.consume(result.usage);
for (const draft of result.artifacts) {
if (!task.writeSet.includes(draft.key)) {
throw new Error(`Task ${task.id} produced undeclared artifact ${draft.key}`);
}
const artifact = store.write(agentId, draft);Оценить эффективный параллелизм графа и увидеть critical path помогает лаборатория.
Интерактивная лаборатория 3
Оцените эффективный параллелизм
Теоретически одновременно: 3. Эффективная оценка после coordination penalties: 1.
Нужны оба лимита. Concurrency limit ограничивает одновременно активных workers, а total worker limit - весь run. Иначе завершившийся worker может породить новый fan-out и обойти ограничение одновременности через несколько волн. И считайте wall-clock правильно: для параллельного batch latency - это самая медленная ветвь плюс fan-out/fan-in overhead, а cost - сумма потребления всех ветвей.
Расписание есть. Но часть узлов упадёт, и систему определяет то, как она классифицирует отказ до повтора.