feat(team-mode): add team runtime create with tests
This commit is contained in:
@@ -0,0 +1,375 @@
|
||||
/// <reference types="bun-types" />
|
||||
|
||||
import { afterAll, beforeEach, describe, expect, mock, test } from "bun:test"
|
||||
import { access, mkdtemp, readdir, rm } from "node:fs/promises"
|
||||
import { tmpdir } from "node:os"
|
||||
import path from "node:path"
|
||||
|
||||
import { TeamModeConfigSchema } from "../../../config/schema/team-mode"
|
||||
import type { ExecutorContext } from "../../../tools/delegate-task/executor-types"
|
||||
import type { BackgroundTask, LaunchInput } from "../../background-agent/types"
|
||||
import { BackgroundManager } from "../../background-agent/manager"
|
||||
import { loadRuntimeState } from "../team-state-store/store"
|
||||
import type { TeamSpec } from "../types"
|
||||
|
||||
const resolveMemberMock = mock(async (member: TeamSpec["members"][number]) => ({
|
||||
agentToUse: `${member.name}-agent`,
|
||||
model: { providerID: "openai", modelID: "gpt-5.4-mini" },
|
||||
fallbackChain: undefined,
|
||||
systemContent: `system:${member.name}`,
|
||||
}))
|
||||
|
||||
mock.module("./resolve-member", () => ({ resolveMember: resolveMemberMock }))
|
||||
|
||||
const { createTeamRun, TeamRunCreateError } = await import("./create")
|
||||
|
||||
function createConfig(baseDir: string, maxParallelMembers = 4) {
|
||||
return TeamModeConfigSchema.parse({ base_dir: baseDir, max_parallel_members: maxParallelMembers, max_wall_clock_minutes: 1 })
|
||||
}
|
||||
|
||||
function createSpec(memberCount: number, withWorktrees = false): TeamSpec {
|
||||
return {
|
||||
version: 1,
|
||||
name: "alpha-team",
|
||||
createdAt: Date.now(),
|
||||
leadAgentId: "member-1",
|
||||
members: Array.from({ length: memberCount }, (_, index) => ({
|
||||
kind: "category",
|
||||
name: `member-${index + 1}`,
|
||||
category: ["quick", "deep", "artistry"][index] ?? "deep",
|
||||
prompt: `prompt-${index + 1}`,
|
||||
backendType: "in-process",
|
||||
isActive: true,
|
||||
color: `color-${index + 1}`,
|
||||
...(withWorktrees ? { worktreePath: `./worktrees/member-${index + 1}` } : {}),
|
||||
})),
|
||||
}
|
||||
}
|
||||
|
||||
function createContext(baseDir: string, manager: BackgroundManager): ExecutorContext & { client: { session: { create: ReturnType<typeof mock> } } } {
|
||||
return {
|
||||
client: { session: { create: mock(async () => ({ data: { id: "forbidden" } })) } } as ExecutorContext["client"] & { session: { create: ReturnType<typeof mock> } },
|
||||
manager,
|
||||
directory: baseDir,
|
||||
}
|
||||
}
|
||||
|
||||
function createManager(
|
||||
baseDir: string,
|
||||
launchImpl: (input: LaunchInput) => Promise<BackgroundTask>,
|
||||
getTaskImpl: (taskId: string) => BackgroundTask | undefined = () => undefined,
|
||||
): { manager: BackgroundManager; launchMock: ReturnType<typeof mock>; cancelTaskMock: ReturnType<typeof mock> } {
|
||||
const manager = new BackgroundManager({ client: {} as ExecutorContext["client"], directory: baseDir } as ConstructorParameters<typeof BackgroundManager>[0])
|
||||
const launchMock = mock((input: LaunchInput) => launchImpl(input))
|
||||
const getTaskMock = mock((taskId: string) => getTaskImpl(taskId))
|
||||
const cancelTaskMock = mock(async () => true)
|
||||
manager.launch = launchMock
|
||||
manager.getTask = getTaskMock
|
||||
manager.cancelTask = cancelTaskMock
|
||||
return { manager, launchMock, cancelTaskMock }
|
||||
}
|
||||
|
||||
async function pathExists(targetPath: string): Promise<boolean> {
|
||||
try {
|
||||
await access(targetPath)
|
||||
return true
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
async function loadSingleRuntimeState(baseDir: string) {
|
||||
const [teamRunId] = await readdir(path.join(baseDir, "runtime"))
|
||||
return await loadRuntimeState(teamRunId ?? "", createConfig(baseDir))
|
||||
}
|
||||
|
||||
describe("createTeamRun", () => {
|
||||
const temporaryDirectories: string[] = []
|
||||
|
||||
beforeEach(() => {
|
||||
resolveMemberMock.mockClear()
|
||||
})
|
||||
|
||||
afterAll(async () => {
|
||||
await Promise.all(temporaryDirectories.splice(0).map(async (directoryPath) => rm(directoryPath, { recursive: true, force: true })))
|
||||
})
|
||||
|
||||
test("spawns 3 members through BackgroundManager.launch without direct session creation", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-create-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager, launchMock } = createManager(baseDir, async () => ({ id: `task-${++launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask))
|
||||
const context = createContext(baseDir, manager)
|
||||
|
||||
// when
|
||||
const runtimeState = await createTeamRun(createSpec(3), "lead-session", context, createConfig(baseDir), manager)
|
||||
|
||||
// then
|
||||
expect(launchMock).toHaveBeenCalledTimes(3)
|
||||
expect(context.client.session.create).toHaveBeenCalledTimes(0)
|
||||
expect(runtimeState.status).toBe("active")
|
||||
expect(runtimeState.members.map((member) => member.sessionId)).toEqual(["session-1", "session-2", "session-3"])
|
||||
expect((launchMock.mock.calls as Array<[LaunchInput]>).every(([input]) => input.suppressTmuxSpawn === true)).toBe(true)
|
||||
})
|
||||
|
||||
test("persists the resolved subagent_type and model on each spawned runtime member", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-subagent-type-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager } = createManager(baseDir, async () => ({ id: `task-${++launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask))
|
||||
|
||||
// when
|
||||
const runtimeState = await createTeamRun(createSpec(3), "lead-session", createContext(baseDir, manager), createConfig(baseDir), manager)
|
||||
|
||||
// then
|
||||
expect(runtimeState.members.map((member) => ({
|
||||
name: member.name,
|
||||
subagent_type: member.subagent_type,
|
||||
model: member.model,
|
||||
}))).toEqual([
|
||||
{ name: "member-1", subagent_type: "member-1-agent", model: { providerID: "openai", modelID: "gpt-5.4-mini" } },
|
||||
{ name: "member-2", subagent_type: "member-2-agent", model: { providerID: "openai", modelID: "gpt-5.4-mini" } },
|
||||
{ name: "member-3", subagent_type: "member-3-agent", model: { providerID: "openai", modelID: "gpt-5.4-mini" } },
|
||||
])
|
||||
})
|
||||
|
||||
test("member prompt only documents member-safe communication tools", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-member-prompt-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
const { manager, launchMock } = createManager(baseDir, async () => ({
|
||||
id: "task-1",
|
||||
sessionID: "session-1",
|
||||
status: "running",
|
||||
} as BackgroundTask))
|
||||
|
||||
// when
|
||||
await createTeamRun(createSpec(1), "lead-session", createContext(baseDir, manager), createConfig(baseDir), manager)
|
||||
const firstPrompt = (launchMock.mock.calls as Array<[LaunchInput]>)[0]?.[0].prompt ?? ""
|
||||
|
||||
// then
|
||||
expect(firstPrompt).toContain("Do not call lead-only lifecycle tools")
|
||||
expect(firstPrompt).not.toContain("3. Request shutdown via `team_shutdown_request`")
|
||||
expect(firstPrompt).toContain("Include `summary` and `references`")
|
||||
expect(firstPrompt).toContain("Move to `status: \"in_progress\"` when you start working")
|
||||
expect(firstPrompt).toContain("delegate-task: Do not call this")
|
||||
expect(firstPrompt).toContain("lead can decide whether to request shutdown")
|
||||
})
|
||||
|
||||
test("rolls back launched members in reverse order when a later spawn fails", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-rollback-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager, cancelTaskMock } = createManager(baseDir, async () => {
|
||||
launchCount += 1
|
||||
if (launchCount === 4) throw new Error("launch-4 failed")
|
||||
return { id: `task-${launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask
|
||||
})
|
||||
|
||||
// when
|
||||
const result = createTeamRun(createSpec(4), "lead-session", createContext(baseDir, manager), createConfig(baseDir), manager)
|
||||
|
||||
// then
|
||||
try {
|
||||
await result
|
||||
throw new Error("expected createTeamRun to reject")
|
||||
} catch (error) {
|
||||
expect(error).toBeInstanceOf(TeamRunCreateError)
|
||||
}
|
||||
expect((cancelTaskMock.mock.calls as Array<[string]>).map(([taskId]) => taskId)).toEqual(["task-3", "task-2", "task-1"])
|
||||
expect((await loadSingleRuntimeState(baseDir)).status).toBe("failed")
|
||||
})
|
||||
|
||||
test("removes all created worktrees when spawn fails after worktree creation", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-worktree-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager } = createManager(baseDir, async () => {
|
||||
launchCount += 1
|
||||
if (launchCount === 2) throw new Error("launch-2 failed")
|
||||
return { id: `task-${launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask
|
||||
})
|
||||
const spec = createSpec(2, true)
|
||||
|
||||
// when
|
||||
try {
|
||||
await createTeamRun(spec, "lead-session", createContext(baseDir, manager), createConfig(baseDir), manager)
|
||||
throw new Error("expected createTeamRun to reject")
|
||||
} catch (error) {
|
||||
expect(error).toBeInstanceOf(TeamRunCreateError)
|
||||
}
|
||||
|
||||
// then
|
||||
expect(await pathExists(path.resolve(baseDir, "./worktrees/member-1"))).toBe(false)
|
||||
expect(await pathExists(path.resolve(baseDir, "./worktrees/member-2"))).toBe(false)
|
||||
})
|
||||
|
||||
test("returns the existing runtime on repeated calls with the same spec and lead session", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-idempotent-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager, launchMock } = createManager(baseDir, async () => ({ id: `task-${++launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask))
|
||||
const spec = createSpec(2)
|
||||
const context = createContext(baseDir, manager)
|
||||
|
||||
// when
|
||||
const firstRuntime = await createTeamRun(spec, "lead-session", context, createConfig(baseDir), manager)
|
||||
const secondRuntime = await createTeamRun(spec, "lead-session", context, createConfig(baseDir), manager)
|
||||
|
||||
// then
|
||||
expect(firstRuntime.teamRunId).toBe(secondRuntime.teamRunId)
|
||||
expect(launchMock).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
test("never exceeds max_parallel_members while spawning", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-parallel-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let inFlight = 0
|
||||
let maxInFlight = 0
|
||||
let launchCount = 0
|
||||
const { manager } = createManager(baseDir, async () => {
|
||||
launchCount += 1
|
||||
inFlight += 1
|
||||
maxInFlight = Math.max(maxInFlight, inFlight)
|
||||
await new Promise((resolve) => setTimeout(resolve, 10))
|
||||
inFlight -= 1
|
||||
return { id: `task-${launchCount}`, sessionID: `session-${launchCount}`, status: "running" } as BackgroundTask
|
||||
})
|
||||
|
||||
// when
|
||||
await createTeamRun(createSpec(8), "lead-session", createContext(baseDir, manager), createConfig(baseDir, 4), manager)
|
||||
|
||||
// then
|
||||
expect(maxInFlight).toBeLessThanOrEqual(4)
|
||||
})
|
||||
|
||||
test("reuses the caller session for the lead when the lead matches the caller agent", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-caller-lead-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager, launchMock } = createManager(baseDir, async (input) => ({
|
||||
id: `task-${++launchCount}`,
|
||||
sessionID: `${input.agent}-session-${launchCount}`,
|
||||
status: "running",
|
||||
} as BackgroundTask))
|
||||
const spec: TeamSpec = {
|
||||
version: 1,
|
||||
name: "alpha-team",
|
||||
createdAt: Date.now(),
|
||||
leadAgentId: "lead",
|
||||
members: [
|
||||
{ kind: "subagent_type", name: "lead", subagent_type: "sisyphus", backendType: "in-process", isActive: true },
|
||||
{ kind: "category", name: "member-1", category: "quick", prompt: "prompt-1", backendType: "in-process", isActive: true },
|
||||
],
|
||||
}
|
||||
|
||||
// when
|
||||
const runtimeState = await createTeamRun(
|
||||
spec,
|
||||
"lead-session",
|
||||
createContext(baseDir, manager),
|
||||
createConfig(baseDir),
|
||||
manager,
|
||||
undefined,
|
||||
{ callerAgentTypeId: "sisyphus" },
|
||||
)
|
||||
|
||||
// then
|
||||
expect(launchMock).toHaveBeenCalledTimes(1)
|
||||
expect(launchMock.mock.calls[0]?.[0]).toMatchObject({ description: "Create team member alpha-team/member-1" })
|
||||
expect(resolveMemberMock).toHaveBeenCalledTimes(1)
|
||||
expect(resolveMemberMock.mock.calls[0]?.[0]).toMatchObject({ name: "member-1" })
|
||||
expect(runtimeState.members.map((member) => ({ name: member.name, sessionId: member.sessionId }))).toEqual([
|
||||
{ name: "lead", sessionId: "lead-session" },
|
||||
{ name: "member-1", sessionId: "member-1-agent-session-1" },
|
||||
])
|
||||
})
|
||||
|
||||
test("persists the reused caller lead's subagent_type so live deliveries can pin it", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-caller-lead-pin-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
const { manager } = createManager(baseDir, async (input) => ({
|
||||
id: `task-${input.agent}`,
|
||||
sessionID: `${input.agent}-session`,
|
||||
status: "running",
|
||||
} as BackgroundTask))
|
||||
const spec: TeamSpec = {
|
||||
version: 1,
|
||||
name: "alpha-team",
|
||||
createdAt: Date.now(),
|
||||
leadAgentId: "lead",
|
||||
members: [
|
||||
{ kind: "subagent_type", name: "lead", subagent_type: "sisyphus", backendType: "in-process", isActive: true },
|
||||
{ kind: "category", name: "worker", category: "quick", prompt: "work hard", backendType: "in-process", isActive: true },
|
||||
],
|
||||
}
|
||||
|
||||
// when
|
||||
const runtimeState = await createTeamRun(
|
||||
spec,
|
||||
"ses_caller_sisyphus",
|
||||
createContext(baseDir, manager),
|
||||
createConfig(baseDir),
|
||||
manager,
|
||||
undefined,
|
||||
{ callerAgentTypeId: "sisyphus" },
|
||||
)
|
||||
|
||||
// then
|
||||
const leadMember = runtimeState.members.find((member) => member.name === "lead")
|
||||
expect(leadMember?.sessionId).toBe("ses_caller_sisyphus")
|
||||
expect(leadMember?.subagent_type).toBe("sisyphus")
|
||||
expect(leadMember?.model).toBeUndefined()
|
||||
})
|
||||
|
||||
test("reuses the caller session for the lead even when the lead subagent_type differs", async () => {
|
||||
// given
|
||||
const baseDir = await mkdtemp(path.join(tmpdir(), "team-runtime-explicit-lead-"))
|
||||
temporaryDirectories.push(baseDir)
|
||||
let launchCount = 0
|
||||
const { manager, launchMock } = createManager(baseDir, async (input) => ({
|
||||
id: `task-${++launchCount}`,
|
||||
sessionID: `${input.agent}-session-${launchCount}`,
|
||||
status: "running",
|
||||
} as BackgroundTask))
|
||||
const spec: TeamSpec = {
|
||||
version: 1,
|
||||
name: "alpha-team",
|
||||
createdAt: Date.now(),
|
||||
leadAgentId: "captain",
|
||||
members: [
|
||||
{ kind: "subagent_type", name: "captain", subagent_type: "atlas", backendType: "in-process", isActive: true },
|
||||
{ kind: "category", name: "member-1", category: "quick", prompt: "prompt-1", backendType: "in-process", isActive: true },
|
||||
],
|
||||
}
|
||||
|
||||
// when
|
||||
const runtimeState = await createTeamRun(
|
||||
spec,
|
||||
"lead-session",
|
||||
createContext(baseDir, manager),
|
||||
createConfig(baseDir),
|
||||
manager,
|
||||
undefined,
|
||||
{ callerAgentTypeId: "sisyphus" },
|
||||
)
|
||||
|
||||
// then
|
||||
expect(launchMock).toHaveBeenCalledTimes(1)
|
||||
expect(launchMock.mock.calls.map(([input]) => input.description)).toEqual([
|
||||
"Create team member alpha-team/member-1",
|
||||
])
|
||||
expect(runtimeState.members.map((member) => ({ name: member.name, sessionId: member.sessionId }))).toEqual([
|
||||
{ name: "captain", sessionId: "lead-session" },
|
||||
{ name: "member-1", sessionId: "member-1-agent-session-1" },
|
||||
])
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,260 @@
|
||||
import { access, mkdir } from "node:fs/promises"
|
||||
import path from "node:path"
|
||||
|
||||
import type { TeamModeConfig } from "../../../config/schema/team-mode"
|
||||
import { QUESTION_DENIED_SESSION_PERMISSION } from "../../../shared/question-denied-session-permission"
|
||||
import type { ExecutorContext } from "../../../tools/delegate-task/executor-types"
|
||||
import type { BackgroundTask } from "../../background-agent/types"
|
||||
import type { BackgroundManager } from "../../background-agent/manager"
|
||||
import type { TmuxSessionManager } from "../../tmux-subagent/manager"
|
||||
import { ensureBaseDirs, getInboxDir, getTeamSpecPath, resolveBaseDir } from "../team-registry/paths"
|
||||
import { createRuntimeState, listActiveTeams, loadRuntimeState, transitionRuntimeState } from "../team-state-store/store"
|
||||
import { registerTeamSession } from "../team-session-registry"
|
||||
import type { RuntimeState, TeamSpec } from "../types"
|
||||
import { activateTeamLayout } from "./activate-team-layout"
|
||||
import { cleanupTeamRunResources } from "./cleanup-team-run-resources"
|
||||
import { buildTeammateCommunicationAddendum } from "../member-guidance"
|
||||
import { resolveMember } from "./resolve-member"
|
||||
import { shouldReuseCallerLeadSession } from "../resolve-caller-team-lead"
|
||||
import { sweepStaleTeamSessions } from "../team-layout-tmux/sweep-stale-team-sessions"
|
||||
|
||||
const SESSION_ID_POLL_MS = 25
|
||||
|
||||
type SpawnedMemberResource = {
|
||||
taskId?: string
|
||||
worktreePath?: string
|
||||
}
|
||||
|
||||
type CreateTeamRunOptions = {
|
||||
callerAgentTypeId?: string
|
||||
parentMessageID?: string
|
||||
}
|
||||
|
||||
export class TeamRunCreateError extends Error {
|
||||
constructor(
|
||||
message: string,
|
||||
public readonly cleanupReport: {
|
||||
cancelledTaskIds: string[]
|
||||
removedLayout: boolean
|
||||
removedWorktrees: string[]
|
||||
errors: string[]
|
||||
},
|
||||
cause: Error,
|
||||
) {
|
||||
super(`${message}: ${cause.message}`)
|
||||
this.name = "TeamRunCreateError"
|
||||
this.cause = cause
|
||||
}
|
||||
}
|
||||
|
||||
function normalizeError(error: unknown): Error {
|
||||
return error instanceof Error ? error : new Error(String(error))
|
||||
}
|
||||
|
||||
async function pathExists(filePath: string): Promise<boolean> {
|
||||
try {
|
||||
await access(filePath)
|
||||
return true
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
async function resolveSpecSource(spec: TeamSpec, ctx: ExecutorContext, config: TeamModeConfig): Promise<"project" | "user"> {
|
||||
const baseDir = resolveBaseDir(config)
|
||||
if (await pathExists(getTeamSpecPath(baseDir, spec.name, "project", ctx.directory))) return "project"
|
||||
if (await pathExists(getTeamSpecPath(baseDir, spec.name, "user"))) return "user"
|
||||
return "project"
|
||||
}
|
||||
|
||||
async function findExistingRuntime(spec: TeamSpec, leadSessionId: string, config: TeamModeConfig): Promise<RuntimeState | undefined> {
|
||||
for (const candidate of await listActiveTeams(config)) {
|
||||
if (candidate.teamName !== spec.name || (candidate.status !== "creating" && candidate.status !== "active")) continue
|
||||
const runtimeState = await loadRuntimeState(candidate.teamRunId, config).catch(() => undefined)
|
||||
if (runtimeState?.leadSessionId === leadSessionId) return runtimeState
|
||||
}
|
||||
}
|
||||
|
||||
async function createMemberWorktree(memberWorktreePath: string, projectRoot: string): Promise<string> {
|
||||
const absolutePath = path.isAbsolute(memberWorktreePath) ? memberWorktreePath : path.resolve(projectRoot, memberWorktreePath)
|
||||
await mkdir(absolutePath, { recursive: true })
|
||||
return absolutePath
|
||||
}
|
||||
|
||||
async function waitForTaskSessionId(bgMgr: BackgroundManager, task: BackgroundTask, deadlineAt: number): Promise<string> {
|
||||
let sessionId = task.sessionID
|
||||
while (!sessionId) {
|
||||
if (Date.now() > deadlineAt) throw new Error(`timed out waiting for child session for task ${task.id}`)
|
||||
const updatedTask = bgMgr.getTask(task.id)
|
||||
if (updatedTask?.status === "error" || updatedTask?.status === "cancelled" || updatedTask?.status === "interrupt") {
|
||||
throw new Error(updatedTask.error ?? `task ${task.id} failed before session creation`)
|
||||
}
|
||||
sessionId = updatedTask?.sessionID
|
||||
if (!sessionId) await new Promise((resolve) => setTimeout(resolve, SESSION_ID_POLL_MS))
|
||||
}
|
||||
return sessionId
|
||||
}
|
||||
|
||||
function buildMemberPrompt(
|
||||
spec: TeamSpec,
|
||||
member: TeamSpec["members"][number],
|
||||
teamRunId: string,
|
||||
config: TeamModeConfig,
|
||||
worktreePath?: string,
|
||||
): string {
|
||||
const promptLines = [`Team: ${spec.name}`, `TeamRunId: ${teamRunId}`, `Member: ${member.name}`]
|
||||
if (worktreePath) promptLines.push(`Worktree: ${worktreePath}`)
|
||||
if (member.prompt) promptLines.push(member.prompt)
|
||||
promptLines.push(buildTeammateCommunicationAddendum(config))
|
||||
return promptLines.join("\n")
|
||||
}
|
||||
|
||||
export async function createTeamRun(
|
||||
spec: TeamSpec,
|
||||
leadSessionId: string,
|
||||
ctx: ExecutorContext,
|
||||
config: TeamModeConfig,
|
||||
bgMgr: BackgroundManager,
|
||||
tmuxMgr?: TmuxSessionManager,
|
||||
options?: CreateTeamRunOptions,
|
||||
): Promise<RuntimeState> {
|
||||
const existingRuntime = await findExistingRuntime(spec, leadSessionId, config)
|
||||
if (existingRuntime) return existingRuntime
|
||||
|
||||
const activeTeams = await listActiveTeams(config)
|
||||
const activeRunIds = new Set(activeTeams.map((t) => t.teamRunId))
|
||||
sweepStaleTeamSessions(activeRunIds).catch(() => {})
|
||||
|
||||
const baseDir = resolveBaseDir(config)
|
||||
await ensureBaseDirs(baseDir)
|
||||
const reusesCallerLeadSession = shouldReuseCallerLeadSession(spec, options?.callerAgentTypeId)
|
||||
let runtimeState = await createRuntimeState(spec, leadSessionId, await resolveSpecSource(spec, ctx, config), config)
|
||||
if (reusesCallerLeadSession && spec.leadAgentId) {
|
||||
const callerLeadSubagentType = options?.callerAgentTypeId
|
||||
registerTeamSession(leadSessionId, {
|
||||
teamRunId: runtimeState.teamRunId,
|
||||
memberName: spec.leadAgentId,
|
||||
role: "lead",
|
||||
})
|
||||
runtimeState = await transitionRuntimeState(runtimeState.teamRunId, (currentState) => ({
|
||||
...currentState,
|
||||
members: currentState.members.map((member) => member.name === spec.leadAgentId
|
||||
? {
|
||||
...member,
|
||||
sessionId: leadSessionId,
|
||||
status: "running",
|
||||
...(callerLeadSubagentType ? { subagent_type: callerLeadSubagentType } : {}),
|
||||
}
|
||||
: member),
|
||||
}), config)
|
||||
}
|
||||
await Promise.all(spec.members.map((member) => mkdir(getInboxDir(baseDir, runtimeState.teamRunId, member.name), { recursive: true })))
|
||||
|
||||
const deadlineAt = Date.now() + (config.max_wall_clock_minutes * 60_000)
|
||||
const resources: SpawnedMemberResource[] = spec.members.map(() => ({}))
|
||||
let createdLayout = false
|
||||
|
||||
try {
|
||||
let nextMemberIndex = 0
|
||||
let failure: Error | undefined
|
||||
const workerCount = Math.min(config.max_parallel_members, spec.members.length)
|
||||
const categoryExamples = Object.keys(ctx.userCategories ?? {}).join(", ")
|
||||
|
||||
await Promise.all(Array.from({ length: workerCount }, async () => {
|
||||
while (!failure) {
|
||||
if (Date.now() > deadlineAt) {
|
||||
failure = new Error("team creation exceeded max_wall_clock_minutes")
|
||||
return
|
||||
}
|
||||
const memberIndex = nextMemberIndex++
|
||||
const member = spec.members[memberIndex]
|
||||
if (!member) return
|
||||
const resource = resources[memberIndex]
|
||||
if (!resource) return
|
||||
|
||||
try {
|
||||
if (member.worktreePath) resource.worktreePath = await createMemberWorktree(member.worktreePath, ctx.directory)
|
||||
if (reusesCallerLeadSession && member.name === spec.leadAgentId) {
|
||||
if (resource.worktreePath) {
|
||||
await transitionRuntimeState(runtimeState.teamRunId, (currentState) => ({
|
||||
...currentState,
|
||||
members: currentState.members.map((currentMember, currentIndex) => currentIndex === memberIndex
|
||||
? { ...currentMember, worktreePath: resource.worktreePath }
|
||||
: currentMember),
|
||||
}), config)
|
||||
}
|
||||
continue
|
||||
}
|
||||
const resolvedMember = await resolveMember(member, ctx, categoryExamples, spec.leadAgentId)
|
||||
const task = await bgMgr.launch({
|
||||
description: `Create team member ${spec.name}/${member.name}`,
|
||||
prompt: buildMemberPrompt(spec, member, runtimeState.teamRunId, config, resource.worktreePath),
|
||||
agent: resolvedMember.agentToUse,
|
||||
parentSessionID: leadSessionId,
|
||||
parentMessageID: options?.parentMessageID ?? `team-create:${runtimeState.teamRunId}:${member.name}`,
|
||||
teamRunId: runtimeState.teamRunId,
|
||||
suppressTmuxSpawn: true,
|
||||
model: resolvedMember.model,
|
||||
fallbackChain: resolvedMember.fallbackChain,
|
||||
skillContent: resolvedMember.systemContent,
|
||||
category: member.kind === "category" ? member.category : undefined,
|
||||
sessionPermission: QUESTION_DENIED_SESSION_PERMISSION,
|
||||
})
|
||||
resource.taskId = task.id
|
||||
const sessionId = await waitForTaskSessionId(bgMgr, task, deadlineAt)
|
||||
registerTeamSession(sessionId, {
|
||||
teamRunId: runtimeState.teamRunId,
|
||||
memberName: member.name,
|
||||
role: member.name === spec.leadAgentId ? "lead" : "member",
|
||||
})
|
||||
const persistedModel = resolvedMember.model
|
||||
? {
|
||||
providerID: resolvedMember.model.providerID,
|
||||
modelID: resolvedMember.model.modelID,
|
||||
...(resolvedMember.model.variant ? { variant: resolvedMember.model.variant } : {}),
|
||||
...(resolvedMember.model.reasoningEffort ? { reasoningEffort: resolvedMember.model.reasoningEffort } : {}),
|
||||
...(resolvedMember.model.temperature !== undefined ? { temperature: resolvedMember.model.temperature } : {}),
|
||||
...(resolvedMember.model.top_p !== undefined ? { top_p: resolvedMember.model.top_p } : {}),
|
||||
...(resolvedMember.model.maxTokens !== undefined ? { maxTokens: resolvedMember.model.maxTokens } : {}),
|
||||
...(resolvedMember.model.thinking ? { thinking: resolvedMember.model.thinking } : {}),
|
||||
}
|
||||
: undefined
|
||||
await transitionRuntimeState(runtimeState.teamRunId, (currentState) => ({
|
||||
...currentState,
|
||||
members: currentState.members.map((currentMember, currentIndex) => currentIndex === memberIndex
|
||||
? {
|
||||
...currentMember,
|
||||
sessionId,
|
||||
status: "running",
|
||||
worktreePath: resource.worktreePath,
|
||||
subagent_type: resolvedMember.agentToUse,
|
||||
...(member.kind === "category" ? { category: member.category } : {}),
|
||||
...(persistedModel ? { model: persistedModel } : {}),
|
||||
}
|
||||
: currentMember),
|
||||
}), config)
|
||||
} catch (error) {
|
||||
failure = normalizeError(error)
|
||||
return
|
||||
}
|
||||
}
|
||||
}))
|
||||
|
||||
if (failure) throw failure
|
||||
|
||||
const launchedRuntimeState = await loadRuntimeState(runtimeState.teamRunId, config)
|
||||
createdLayout = await activateTeamLayout(launchedRuntimeState, config, ctx.directory, tmuxMgr)
|
||||
|
||||
return await transitionRuntimeState(runtimeState.teamRunId, (currentState) => ({ ...currentState, status: "active" }), config)
|
||||
} catch (error) {
|
||||
const cleanupReport = await cleanupTeamRunResources({
|
||||
teamRunId: runtimeState.teamRunId,
|
||||
config,
|
||||
resources,
|
||||
bgMgr,
|
||||
tmuxMgr,
|
||||
createdLayout,
|
||||
})
|
||||
throw new TeamRunCreateError(`Failed to create team run '${spec.name}'`, cleanupReport, normalizeError(error))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user