-- subagents/workflow.lua: the callback API's awaiting, failure-as-value rule, -- structured workers, the concurrency gate, and the job budget. -- -- The fake host settles jobs in the order each outcome's `settle` key asks for, -- which is what makes the ordering cases meaningful: input order and settle -- order are deliberately different everywhere below. local fake = require("spec.fake_ext") local workflow = require("subagents.workflow") local function has(text, needle) assert(type(text) == "string", "expected a string, got " .. type(text)) assert(text:find(needle, 1, true), "expected to find " .. needle .. " in:\n" .. tostring(text)) end local function profile_set() local set = { list = {}, by_name = {}, warnings = {} } for _, name in ipairs({ "alpha", "beta", "gamma" }) do local profile = { name = name, description = name, body = "You are " .. name .. ".\n" } set.list[#set.list + 1] = profile set.by_name[name] = profile end return set end local function with_host(fn, opts) local handle = fake.install(opts) local ok, err = pcall(fn, handle, profile_set()) handle.restore() if not ok then error(err, 0) end end local function three_children(ctx) local handles = {} for index, name in ipairs({ "alpha", "beta", "gamma" }) do handles[index] = ctx:agent({ name = name, agent = name, prompt = "work on " .. name }) end return handles end local function json_available() return pcall(require, "dkjson") end local ITEMS_SCHEMA = { type = "object", required = { "items" }, properties = { items = { type = "array", items = { type = "string" } }, }, } local GHOST = { unknown_models = { ["openai:ghost"] = true } } return { { "await all returns results in input order, not settle order", function() with_host(function(handle, profiles) handle.queue_for("alpha", { output = "A", settle = 3 }) handle.queue_for("beta", { output = "B", settle = 1 }) handle.queue_for("gamma", { output = "C", settle = 2 }) local results = workflow.execute(workflow.workflow(function(ctx) return ctx:await(three_children(ctx), "all") end), "input", { profiles = profiles }) assert(#results == 3, "expected three results") assert(results[1].output == "A", tostring(results[1].output)) assert(results[2].output == "B", tostring(results[2].output)) assert(results[3].output == "C", tostring(results[3].output)) assert(#handle.spawns == 3, "one child per ctx:agent") assert(handle.spawns[1].prompt == "work on alpha") end) end }, { "await first returns the earliest settler and the remaining handles", function() with_host(function(handle, profiles) handle.queue_for("alpha", { output = "A", settle = 3 }) handle.queue_for("beta", { output = "B", settle = 1 }) handle.queue_for("gamma", { output = "C", settle = 2 }) local order = workflow.execute(workflow.workflow(function(ctx) local handles = three_children(ctx) local seen = {} while #handles > 0 do local result, remaining = ctx:await(handles, "first") seen[#seen + 1] = result.output assert(#remaining == #handles - 1, "one handle settles per await") handles = remaining end return seen end), "input", { profiles = profiles }) assert(table.concat(order, ",") == "B,C,A", table.concat(order, ",")) end) end }, { "handle:await settles one child", function() with_host(function(handle, profiles) handle.queue_for("alpha", { output = "just me" }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "only", agent = "alpha", prompt = "go" }):await() end), "input", { profiles = profiles }) assert(result.status == "completed", tostring(result.status)) assert(result.output == "just me", tostring(result.output)) end) end }, { "the gate keeps five children in flight and queues the rest", function() with_host(function(handle, profiles) local results = workflow.execute(workflow.workflow(function(ctx) local handles = {} for index = 1, 6 do handles[index] = ctx:agent({ name = "job-" .. index, agent = "alpha", prompt = "task " .. index }) end assert(#handle.runs == 5, "five turns run at once, saw " .. #handle.runs) return ctx:await(handles, "all") end), "input", { profiles = profiles }) assert(#results == 6, "every child reports") assert(#handle.runs == 6, "the queued child starts once a slot frees up") assert(handle.max_live == 5, "never more than five at once, peaked at " .. handle.max_live) end) end }, { "a rejected child is a failed result, not an error", function() with_host(function(handle, profiles) handle.queue_for("beta", { output = "B" }) local results = workflow.execute(workflow.workflow(function(ctx) local first = ctx:agent({ name = "first", agent = "alpha", model = "openai:ghost", prompt = "a" }) local second = ctx:agent({ name = "second", agent = "beta", prompt = "b" }) return ctx:await({ first, second }, "all") end), "input", { profiles = profiles }) assert(results[1].status == "failed", tostring(results[1].status)) has(results[1].error, "unknown model 'openai:ghost'") assert(results[1].resumable == false, "a child that never allocated is not resumable") assert(results[2].status == "completed", "a sibling failure must not disturb this one") assert(#handle.runs == 1, "only the sibling was ever started") end, GHOST) end }, { "a failed child is reported as a value", function() with_host(function(handle, profiles) handle.queue_for("alpha", { status = "failed", error = "provider refused", resumable = false }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "failed", agent = "alpha", prompt = "a" }):await() end), "input", { profiles = profiles }) assert(result.status == "failed") assert(result.error == "provider refused", tostring(result.error)) end) end }, { "an id with a turn already in flight is refused", function() with_host(function(handle, profiles) handle.add_session("0198-child", { { role = "system", text = "You are alpha.\n", metadata = { subagents = { owner = "0198-primary", agent = "alpha" } } }, }) handle.queue({ output = "the first turn wins" }) local results = workflow.execute(workflow.workflow(function(ctx) local first = ctx:agent({ name = "first", id = "0198-child", prompt = "a" }) local second = ctx:agent({ name = "second", id = "0198-child", prompt = "b" }) return ctx:await({ first, second }, "all") end), "input", { profiles = profiles }) assert(results[1].status == "completed", tostring(results[1].error)) assert(results[2].status == "failed", "one turn per child at a time") has(results[2].error, "already has a turn in flight") assert(#handle.runs == 1, "the second call never starts a turn") end) end }, { "a structured worker decodes and validates its output", function() if not json_available() then return "skip", "dkjson is not installed" end with_host(function(handle, profiles) handle.queue_for("alpha", { structured_json = '{"items":["x","y"]}' }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "worker", agent = "alpha", prompt = "split it", output = { description = "Return the work items.", schema = ITEMS_SCHEMA }, }):await() end), "input", { profiles = profiles }) assert(result.status == "completed", tostring(result.error)) assert(type(result.output) == "table", "structured output decodes to a table") assert(result.output.items[1] == "x" and result.output.items[2] == "y") assert(result.resumable == false, "a one-shot child is not resumable") assert(result.id == nil, "and has no durable session to name") assert(#handle.stores == 0, "a structured worker never touches the child catalog") local child = handle.spawns[1] assert(type(child.tool_choice) == "table" and child.tool_choice.name == "emit_result", "the output tool is the only choice the child has") assert(child.run.dispatch_tools == false, "a one-shot child never dispatches a tool call") assert(#child.tool_decls == 1, "the output tool replaces the inherited set, saw " .. table.concat(child.tools, ",")) local decl = child.tool_decls[1] assert(decl.name == "emit_result", "the synthetic tool is named for the host") assert(decl.description == "Return the work items.") assert(decl.schema == ITEMS_SCHEMA, "the schema is passed through untouched") assert(decl._source == nil and decl._ctx == nil and decl._vt == nil, "the output tool is declaration-only: nothing may dispatch it") end) end }, { "structured output that violates the schema fails the result", function() if not json_available() then return "skip", "dkjson is not installed" end with_host(function(handle, profiles) handle.queue_for("alpha", { structured_json = '{"nope":1}' }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "worker", agent = "alpha", prompt = "split it", output = { schema = ITEMS_SCHEMA }, }):await() end), "input", { profiles = profiles }) assert(result.status == "failed", "a schema violation is never a success") has(result.error, "structured output failed validation") assert(result.output == nil, "invalid output is not handed to the caller") end) end }, -- One validator, always the built-in one: a schema must not be judged by -- whether an undeclared rock happens to be installed on this machine. { "validation never consults the jsonschema rock", function() if not json_available() then return "skip", "dkjson is not installed" end local restore = package.loaded.jsonschema package.loaded.jsonschema = { generate_validator = function() return function() return false, "the rock must not be consulted" end end, } local ok, err = pcall(with_host, function(handle, profiles) handle.queue_for("alpha", { structured_json = '{"items":["x"]}' }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "worker", agent = "alpha", prompt = "split it", output = { schema = ITEMS_SCHEMA }, }):await() end), "input", { profiles = profiles }) assert(result.status == "completed", tostring(result.error)) assert(result.output.items[1] == "x", "the built-in validator accepted it") end) package.loaded.jsonschema = restore if not ok then error(err, 0) end end }, { "a structured worker whose output tool carried no arguments fails as missing", function() with_host(function(handle, profiles) handle.queue_for("alpha", { structured_json = "" }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "worker", agent = "alpha", prompt = "split it", output = { schema = ITEMS_SCHEMA }, }):await() end), "input", { profiles = profiles }) assert(result.status == "failed", tostring(result.status)) has(result.error, "structured output missing") end) end }, { "a structured worker that answers in prose fails", function() with_host(function(handle, profiles) handle.queue_for("alpha", { output = "prose, not a tool call" }) local result = workflow.execute(workflow.workflow(function(ctx) return ctx:agent({ name = "worker", agent = "alpha", prompt = "split it", output = { schema = ITEMS_SCHEMA }, }):await() end), "input", { profiles = profiles }) assert(result.status == "failed", tostring(result.status)) has(result.error, "did not call the required 'emit_result' output tool") end) end }, { "an already settled handle is served without polling again", function() with_host(function(handle, profiles) handle.queue_for("beta", { output = "B" }) handle.queue_for("gamma", { output = "C" }) local first_status = workflow.execute(workflow.workflow(function(ctx) local rejected = ctx:agent({ name = "rejected", agent = "alpha", model = "openai:ghost", prompt = "a" }) local live = { rejected, ctx:agent({ name = "live", agent = "beta", prompt = "b" }) } local result, remaining = ctx:await(live, "first") assert(#remaining == 1, "the pending sibling stays outstanding") assert(handle.polls == 0, "a cached result must not reach the job machinery") return result.status end), "input", { profiles = profiles }) assert(first_status == "failed", tostring(first_status)) assert(#handle.runs == 1, "only the live sibling was ever started") end, GHOST) end }, { "max_jobs caps how many children one workflow starts", function() with_host(function(handle, profiles) local ok, err = pcall(workflow.execute, workflow.workflow(function(ctx) three_children(ctx) end), "input", { profiles = profiles, max_jobs = 2 }) assert(not ok, "the third ctx:agent must be refused") has(tostring(err), "job limit exceeded") assert(#handle.runs == 2, "the capped call never reaches the host") end) end }, { "handles the callback never awaited are settled before returning", function() with_host(function(handle, profiles) workflow.execute(workflow.workflow(function(ctx) ctx:agent({ name = "orphan", agent = "alpha", prompt = "orphan" }) return "done" end), "input", { profiles = profiles }) assert(#handle.jobs == 1, "one child ran") assert(handle.jobs[1]._settled, "no child is left running") end) end }, { "an unknown agent inside a workflow is a workflow error", function() with_host(function(handle, profiles) local ok, err = pcall(workflow.execute, workflow.workflow(function(ctx) ctx:agent({ name = "ghost", agent = "ghost", prompt = "go" }) end), "input", { profiles = profiles }) assert(not ok, "an unknown profile is a programmer error, not a child failure") has(tostring(err), "unknown agent 'ghost'") end) end }, }