Skip to content

API reference

Public exports from pydantic_team.

pydantic_team

Type-safe team orchestration for pydantic-ai Agents.

__version__ module-attribute

__version__ = '0.4.0'

BaseTeam

Bases: ABC, Generic[OutputT]

Abstract base for team orchestration strategies.

members abstractmethod property

members: Sequence[TeamMember]

Agents or nested teams that participate in this team.

run abstractmethod async

run(
    user_prompt: str, *, usage: RunUsage | None = None
) -> TeamResult[OutputT]

Execute the team against user_prompt.

Parameters:

Name Type Description Default
user_prompt str

User input for the team run.

required
usage RunUsage | None

Optional usage accumulator shared with nested agent runs.

None

Returns:

Type Description
TeamResult[OutputT]

A TeamResult with final data and aggregated usage.

Source code in pydantic_team/base.py
@abstractmethod
async def run(self, user_prompt: str, *, usage: RunUsage | None = None) -> TeamResult[OutputT]:
    """Execute the team against `user_prompt`.

    Args:
        user_prompt: User input for the team run.
        usage: Optional usage accumulator shared with nested agent runs.

    Returns:
        A [`TeamResult`][pydantic_team.base.TeamResult] with final data and aggregated usage.
    """

TeamResult dataclass

TeamResult(data: OutputT, usage: RunUsage)

Bases: Generic[OutputT]

Outcome of a team run.

Attributes:

Name Type Description
data OutputT

Final output produced by the team (leader output for hierarchical teams).

usage RunUsage

Aggregated token/request usage for all agents involved in the run.

HierarchicalTeam

HierarchicalTeam(
    *,
    members: Sequence[TeamMember],
    leader_agent: Agent[object, object] | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
)

Bases: BaseTeam[object]

Leader-driven team that registers each member as a delegation tool.

Follows pydantic-ai agent delegation: nested member runs receive usage=ctx.usage so tokens aggregate on the leader run.

Create a hierarchical team.

Parameters:

Name Type Description Default
members Sequence[TeamMember]

Specialist agents or nested teams (at least one required).

required
leader_agent Agent[object, object] | None

Existing leader agent. Mutually exclusive with leader_model.

None
leader_model str | None

Model string used to construct a leader agent when leader_agent is not provided.

None
system_prompt_override str | None

Optional leader instructions / extra system prompt.

None
name str | None

Optional team name (used when this team is nested as a member tool).

None
Source code in pydantic_team/hierarchical.py
def __init__(
    self,
    *,
    members: Sequence[TeamMember],
    leader_agent: Agent[object, object] | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
) -> None:
    """Create a hierarchical team.

    Args:
        members: Specialist agents or nested teams (at least one required).
        leader_agent: Existing leader agent. Mutually exclusive with `leader_model`.
        leader_model: Model string used to construct a leader agent when
            `leader_agent` is not provided.
        system_prompt_override: Optional leader instructions / extra system prompt.
        name: Optional team name (used when this team is nested as a member tool).
    """
    if leader_agent is not None and leader_model is not None:
        raise ValueError('Provide leader_agent or leader_model, not both')
    if leader_agent is None and leader_model is None:
        raise ValueError('Provide leader_agent or leader_model')
    if not members:
        raise ValueError('members must be a non-empty sequence')

    self.name = name
    self._members: list[TeamMember] = list(members)

    if leader_agent is not None:
        self._leader = leader_agent
        if system_prompt_override is not None:

            def _leader_override_prompt() -> str:
                return system_prompt_override

            self._leader.system_prompt(_leader_override_prompt)

    else:
        assert leader_model is not None
        self._leader = Agent(
            leader_model,
            instructions=system_prompt_override or _DEFAULT_LEADER_INSTRUCTIONS,
        )

    self._register_member_tools()

leader property

leader: Agent[object, object]

The coordinating leader agent.

CollaborativeTeam

CollaborativeTeam(
    *,
    members: Sequence[AnyAgent],
    leader_agent: AnyAgent | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
    max_rounds: int = 3,
    max_replans: int = 0,
    max_assignments_per_tick: int | None = None,
    dispatch_mode: DispatchMode = 'phased',
    require_review: bool = False,
)

Bases: BaseTeam[object]

Team that coordinates work through a shared TaskBoard.

The leader creates and assigns tasks by role; members complete their assigned work in parallel (phased rounds or streaming dispatch). When the board remains incomplete, the leader may replan (up to max_replans) before synthesizing. Members may message each other directly via board tools (send_message / list_messages).

Create a collaborative team.

Parameters:

Name Type Description Default
members Sequence[AnyAgent]

Teammate agents (at least one). Nested teams are not supported here.

required
leader_agent AnyAgent | None

Existing leader. Mutually exclusive with leader_model.

None
leader_model str | None

Model string used to build the leader when leader_agent is omitted.

None
system_prompt_override str | None

Optional leader instructions / extra system prompt.

None
name str | None

Optional team name.

None
max_rounds int

In phased mode, parallel member ticks per phase (after seed and after each replan). In streaming mode, max ticks per member per phase.

3
max_replans int

How many times the leader may replan after incomplete member phases. 0 preserves seed → work → synthesize with no replan.

0
max_assignments_per_tick int | None

If set, each member tick only lists this many incomplete assignments (caps work per round without relying on soft prompt wording).

None
dispatch_mode DispatchMode

phased (default) runs seed then member rounds; streaming starts member ticks as soon as tasks are assigned (overlap with seed/replan).

'phased'
require_review bool

When True, new tasks get reviewer=leader so complete goes to pending_review until approve/reject. When False, tasks complete to done unless assign_reviewer sets a reviewer.

False
Source code in pydantic_team/collaborative.py
def __init__(
    self,
    *,
    members: Sequence[AnyAgent],
    leader_agent: AnyAgent | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
    max_rounds: int = 3,
    max_replans: int = 0,
    max_assignments_per_tick: int | None = None,
    dispatch_mode: DispatchMode = 'phased',
    require_review: bool = False,
) -> None:
    """Create a collaborative team.

    Args:
        members: Teammate agents (at least one). Nested teams are not supported here.
        leader_agent: Existing leader. Mutually exclusive with `leader_model`.
        leader_model: Model string used to build the leader when `leader_agent` is omitted.
        system_prompt_override: Optional leader instructions / extra system prompt.
        name: Optional team name.
        max_rounds: In ``phased`` mode, parallel member ticks per phase (after seed and
            after each replan). In ``streaming`` mode, max ticks per member per phase.
        max_replans: How many times the leader may replan after incomplete member phases.
            ``0`` preserves seed → work → synthesize with no replan.
        max_assignments_per_tick: If set, each member tick only lists this many incomplete
            assignments (caps work per round without relying on soft prompt wording).
        dispatch_mode: ``phased`` (default) runs seed then member rounds; ``streaming``
            starts member ticks as soon as tasks are assigned (overlap with seed/replan).
        require_review: When True, new tasks get ``reviewer=leader`` so ``complete`` goes
            to ``pending_review`` until approve/reject. When False, tasks complete to
            ``done`` unless ``assign_reviewer`` sets a reviewer.
    """
    if leader_agent is not None and leader_model is not None:
        raise ValueError('Provide leader_agent or leader_model, not both')
    if leader_agent is None and leader_model is None:
        raise ValueError('Provide leader_agent or leader_model')
    if not members:
        raise ValueError('members must be a non-empty sequence')
    if max_rounds < 1:
        raise ValueError('max_rounds must be >= 1')
    if max_replans < 0:
        raise ValueError('max_replans must be >= 0')
    if max_assignments_per_tick is not None and max_assignments_per_tick < 1:
        raise ValueError('max_assignments_per_tick must be >= 1')
    if dispatch_mode not in ('phased', 'streaming'):
        raise ValueError("dispatch_mode must be 'phased' or 'streaming'")

    self.name = name
    self._max_rounds = max_rounds
    self._max_replans = max_replans
    self._max_assignments_per_tick = max_assignments_per_tick
    self._dispatch_mode: DispatchMode = dispatch_mode
    self._require_review = require_review
    self._members: list[AnyAgent] = list(members)
    self._member_ids: list[str] = [
        _agent_id(member, fallback=f'member-{index}') for index, member in enumerate(self._members)
    ]

    if leader_agent is not None:
        self._leader: AnyAgent = leader_agent
        if system_prompt_override is not None:

            def _leader_override_prompt() -> str:
                return system_prompt_override

            self._leader.system_prompt(_leader_override_prompt)
    else:
        assert leader_model is not None
        self._leader = cast(
            AnyAgent,
            Agent(
                leader_model,
                name='leader',
                instructions=system_prompt_override or default_leader_instructions(self._member_ids),
                deps_type=BoardDeps,
            ),
        )

    _register_leader_tools(self._leader)
    for member in self._members:
        _register_member_tools(member)

leader property

leader: AnyAgent

The team lead agent.

dispatch_mode property

dispatch_mode: DispatchMode

Scheduling mode: phased or streaming.

member_ids property

member_ids: Sequence[str]

Stable teammate ids used with assign_task (agent names).

iter

iter(
    user_prompt: str, *, usage: RunUsage | None = None
) -> CollaborativeRun

Start an observable collaborative run (async context + async iterator).

Source code in pydantic_team/collaborative.py
def iter(self, user_prompt: str, *, usage: RunUsage | None = None) -> CollaborativeRun:
    """Start an observable collaborative run (async context + async iterator)."""
    return CollaborativeRun(
        leader=self._leader,
        members=self._members,
        member_ids=self._member_ids,
        max_rounds=self._max_rounds,
        max_replans=self._max_replans,
        max_assignments_per_tick=self._max_assignments_per_tick,
        dispatch_mode=self._dispatch_mode,
        run_member=self._run_member,
        user_prompt=user_prompt,
        require_review=self._require_review,
        _usage=usage or RunUsage(),
    )

CollaborativeRun dataclass

CollaborativeRun(
    leader: AnyAgent,
    members: Sequence[AnyAgent],
    member_ids: Sequence[str],
    max_rounds: int,
    max_replans: int,
    max_assignments_per_tick: int | None,
    dispatch_mode: DispatchMode,
    run_member: RunMemberFn,
    user_prompt: str,
    require_review: bool,
    _usage: RunUsage,
    _board: TaskBoard = TaskBoard(),
    _events: Queue[TeamEvent | None] = (
        lambda: asyncio.Queue[TeamEvent | None]()
    )(),
    _result: TeamResult[object] | None = None,
    _driver: Task[None] | None = None,
    _error: BaseException | None = None,
)

Step-by-step collaborative run (inspired by pydantic-graph GraphRun).

board property

board: TaskBoard

Live task board for this run.

result property

result: TeamResult[object] | None

Final result once the run has ended; otherwise None.

usage property

usage: RunUsage

Aggregated usage accumulated during the run.

TeamTask dataclass

TeamTask(
    kind: TaskKind,
    agent_id: str | None = None,
    task_ids: tuple[str, ...] = (),
)

A unit of scheduled collaborative work (inspired by pydantic-graph GraphTask).

TasksScheduled dataclass

TasksScheduled(tasks: tuple[TeamTask, ...])

One or more tasks are about to run (or have just been spawned).

TaskCompleted dataclass

TaskCompleted(task: TeamTask)

A scheduled task finished (board may have been mutated via tools).

PhaseJoined dataclass

PhaseJoined(phase: PhaseName, incomplete: bool)

No inflight work remains for this phase (join / barrier).

MessagePosted dataclass

MessagePosted(message: BoardMessage)

A teammate posted a peer message on the board.

TaskReviewDecided dataclass

TaskReviewDecided(task: Task, decision: ReviewDecision)

A reviewer approved or rejected a board task.

RunEnded dataclass

RunEnded(result: TeamResult[object])

The collaborative run finished with a final result.

TaskBoard dataclass

TaskBoard(
    _tasks: dict[str, Task] = (lambda: {})(),
    _messages: list[BoardMessage] = (
        lambda: list[BoardMessage]()
    )(),
    _lock: Lock = asyncio.Lock(),
    _counter: int = 0,
    _message_counter: int = 0,
    _wakeup: Event = asyncio.Event(),
    _wakeup_agents: set[str] = (lambda: set[str]())(),
)

Thread-safe in-process task list shared by a collaborative team.

add_task async

add_task(
    title: str,
    description: str = '',
    *,
    reviewer: str | None = None,
) -> Task

Create an open task and return its snapshot.

Source code in pydantic_team/board.py
async def add_task(
    self,
    title: str,
    description: str = '',
    *,
    reviewer: str | None = None,
) -> Task:
    """Create an open task and return its snapshot."""
    async with self._lock:
        self._counter += 1
        task = Task(
            id=f'task-{self._counter}',
            title=title,
            description=description,
            reviewer=reviewer,
        )
        self._tasks[task.id] = task
        return task

list_tasks async

list_tasks(status: TaskStatus | None = None) -> list[Task]

Return task snapshots, optionally filtered by status.

Source code in pydantic_team/board.py
async def list_tasks(self, status: TaskStatus | None = None) -> list[Task]:
    """Return task snapshots, optionally filtered by status."""
    async with self._lock:
        tasks = list(self._tasks.values())
    if status is not None:
        tasks = [task for task in tasks if task.status is status]
    return tasks

post_message async

post_message(
    sender: str,
    to: str,
    body: str,
    *,
    task_id: str | None = None,
) -> BoardMessage

Append a peer message; optionally link to an existing task.

Source code in pydantic_team/board.py
async def post_message(
    self,
    sender: str,
    to: str,
    body: str,
    *,
    task_id: str | None = None,
) -> BoardMessage:
    """Append a peer message; optionally link to an existing task."""
    async with self._lock:
        if task_id is not None:
            self._require(task_id)
        self._message_counter += 1
        message = BoardMessage(
            id=f'msg-{self._message_counter}',
            sender=sender,
            to=to,
            body=body,
            task_id=task_id,
        )
        self._messages.append(message)
    if to == '*':
        self.signal_wakeup()
    else:
        self.signal_wakeup(to)
    return message

list_messages async

list_messages(
    *, agent_id: str | None = None
) -> list[BoardMessage]

Return messages; when agent_id is set, only visible ones for that agent.

Source code in pydantic_team/board.py
async def list_messages(self, *, agent_id: str | None = None) -> list[BoardMessage]:
    """Return messages; when ``agent_id`` is set, only visible ones for that agent."""
    async with self._lock:
        messages = list(self._messages)
    if agent_id is None:
        return messages
    return [
        message for message in messages if message.to == agent_id or message.to == '*' or message.sender == agent_id
    ]

messages_snapshot

messages_snapshot() -> list[BoardMessage]

Return a stable list copy of all messages (asyncio-safe between awaits).

Source code in pydantic_team/board.py
def messages_snapshot(self) -> list[BoardMessage]:
    """Return a stable list copy of all messages (asyncio-safe between awaits)."""
    return list(self._messages)

claim async

claim(task_id: str, agent_id: str) -> Task

Atomically claim an open task for agent_id.

Source code in pydantic_team/board.py
async def claim(self, task_id: str, agent_id: str) -> Task:
    """Atomically claim an open task for `agent_id`."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.OPEN:
            raise TaskClaimError(f'task {task_id!r} is not open (status={task.status})')
        updated = replace(task, status=TaskStatus.CLAIMED, assignee=agent_id)
        self._tasks[task_id] = updated
    self.signal_wakeup(agent_id)
    return updated

assign async

assign(task_id: str, agent_id: str) -> Task

Force-assign a non-done, non-pending-review task to agent_id (lead operation).

Source code in pydantic_team/board.py
async def assign(self, task_id: str, agent_id: str) -> Task:
    """Force-assign a non-done, non-pending-review task to `agent_id` (lead operation)."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is TaskStatus.DONE:
            raise TaskClaimError(f'task {task_id!r} is already done')
        if task.status is TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is pending review')
        updated = replace(task, status=TaskStatus.CLAIMED, assignee=agent_id)
        self._tasks[task_id] = updated
    self.signal_wakeup(agent_id)
    return updated

assign_reviewer async

assign_reviewer(task_id: str, reviewer_id: str) -> Task

Set or replace the reviewer on a non-done task.

Source code in pydantic_team/board.py
async def assign_reviewer(self, task_id: str, reviewer_id: str) -> Task:
    """Set or replace the reviewer on a non-done task."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is TaskStatus.DONE:
            raise TaskClaimError(f'task {task_id!r} is already done')
        updated = replace(task, reviewer=reviewer_id)
        self._tasks[task_id] = updated
    if updated.status is TaskStatus.PENDING_REVIEW:
        self.signal_wakeup(reviewer_id)
    return updated

complete async

complete(
    task_id: str, *, result: str, agent_id: str
) -> Task

Mark work submitted; gated tasks become pending_review, others done.

Source code in pydantic_team/board.py
async def complete(self, task_id: str, *, result: str, agent_id: str) -> Task:
    """Mark work submitted; gated tasks become pending_review, others done."""
    async with self._lock:
        task = self._require(task_id)
        if task.status not in (TaskStatus.CLAIMED, TaskStatus.NEEDS_REVISION):
            raise TaskClaimError(f'task {task_id!r} is not completable (status={task.status})')
        if task.assignee != agent_id:
            raise TaskClaimError(f'task {task_id!r} is assigned to {task.assignee!r}, not {agent_id!r}')
        if task.reviewer is not None:
            updated = replace(
                task,
                status=TaskStatus.PENDING_REVIEW,
                result=result,
                rejection_reason=None,
            )
            self._tasks[task_id] = updated
            wakeup_agent = task.reviewer
        else:
            updated = replace(
                task,
                status=TaskStatus.DONE,
                result=result,
                rejection_reason=None,
            )
            self._tasks[task_id] = updated
            wakeup_agent = None
    if wakeup_agent is not None:
        self.signal_wakeup(wakeup_agent)
    else:
        self.signal_wakeup()
    return updated

approve async

approve(task_id: str, *, agent_id: str) -> Task

Accept a pending_review task; only the reviewer may approve.

Source code in pydantic_team/board.py
async def approve(self, task_id: str, *, agent_id: str) -> Task:
    """Accept a pending_review task; only the reviewer may approve."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is not pending review (status={task.status})')
        if task.reviewer != agent_id:
            raise TaskClaimError(f'task {task_id!r} reviewer is {task.reviewer!r}, not {agent_id!r}')
        updated = replace(task, status=TaskStatus.DONE, rejection_reason=None)
        self._tasks[task_id] = updated
    self.signal_wakeup()
    return updated

reject async

reject(task_id: str, *, reason: str, agent_id: str) -> Task

Reject a pending_review task back to needs_revision; reason required.

Source code in pydantic_team/board.py
async def reject(self, task_id: str, *, reason: str, agent_id: str) -> Task:
    """Reject a pending_review task back to needs_revision; reason required."""
    cleaned = reason.strip()
    if not cleaned:
        raise TaskClaimError('rejection reason must be non-empty')
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is not pending review (status={task.status})')
        if task.reviewer != agent_id:
            raise TaskClaimError(f'task {task_id!r} reviewer is {task.reviewer!r}, not {agent_id!r}')
        updated = replace(
            task,
            status=TaskStatus.NEEDS_REVISION,
            rejection_reason=cleaned,
        )
        self._tasks[task_id] = updated
        assignee = task.assignee
    self.signal_wakeup(assignee)
    return updated

is_complete

is_complete() -> bool

Return True when there are no tasks or every task is done.

Source code in pydantic_team/board.py
def is_complete(self) -> bool:
    """Return True when there are no tasks or every task is done."""
    if not self._tasks:
        return True
    return all(task.status is TaskStatus.DONE for task in self._tasks.values())

snapshot

snapshot() -> list[Task]

Return a stable list copy of all tasks (asyncio-safe between awaits).

Source code in pydantic_team/board.py
def snapshot(self) -> list[Task]:
    """Return a stable list copy of all tasks (asyncio-safe between awaits)."""
    return list(self._tasks.values())

signal_wakeup

signal_wakeup(agent_id: str | None = None) -> None

Wake waiters; optionally record which agent gained work.

Source code in pydantic_team/board.py
def signal_wakeup(self, agent_id: str | None = None) -> None:
    """Wake waiters; optionally record which agent gained work."""
    if agent_id is not None:
        self._wakeup_agents.add(agent_id)
    self._wakeup.set()

wait_wakeup async

wait_wakeup() -> set[str]

Block until signal_wakeup; return agent ids recorded since last wait.

Source code in pydantic_team/board.py
async def wait_wakeup(self) -> set[str]:
    """Block until ``signal_wakeup``; return agent ids recorded since last wait."""
    await self._wakeup.wait()
    self._wakeup.clear()
    agents = set(self._wakeup_agents)
    self._wakeup_agents.clear()
    return agents

Task dataclass

Task(
    id: str,
    title: str,
    description: str = '',
    status: TaskStatus = TaskStatus.OPEN,
    assignee: str | None = None,
    result: str | None = None,
    reviewer: str | None = None,
    rejection_reason: str | None = None,
)

A unit of work on the shared board.

BoardMessage dataclass

BoardMessage(
    id: str,
    sender: str,
    to: str,
    body: str,
    task_id: str | None = None,
)

A peer message posted on the shared board.

TaskStatus

Bases: str, Enum

Lifecycle status of a board task.

BoardDeps dataclass

BoardDeps(
    board: TaskBoard,
    agent_id: str,
    member_ids: tuple[str, ...],
    leader_id: str = 'leader',
    emit: EmitFn | None = None,
    require_review: bool = False,
)

Dependencies injected into leader and member agent runs.

instrument_pydantic_team

instrument_pydantic_team(*, enabled: bool = True) -> None

Enable or disable OpenTelemetry spans for team orchestration (idempotent).

Source code in pydantic_team/_instrumentation.py
def instrument_pydantic_team(*, enabled: bool = True) -> None:
    """Enable or disable OpenTelemetry spans for team orchestration (idempotent)."""
    global _instrumentation_enabled
    _instrumentation_enabled = enabled

is_instrumented

is_instrumented() -> bool

Return whether team orchestration spans are currently enabled.

Source code in pydantic_team/_instrumentation.py
def is_instrumented() -> bool:
    """Return whether team orchestration spans are currently enabled."""
    return _instrumentation_enabled

Base types

pydantic_team.base.TeamResult dataclass

TeamResult(data: OutputT, usage: RunUsage)

Bases: Generic[OutputT]

Outcome of a team run.

Attributes:

Name Type Description
data OutputT

Final output produced by the team (leader output for hierarchical teams).

usage RunUsage

Aggregated token/request usage for all agents involved in the run.

pydantic_team.base.BaseTeam

Bases: ABC, Generic[OutputT]

Abstract base for team orchestration strategies.

members abstractmethod property

members: Sequence[TeamMember]

Agents or nested teams that participate in this team.

run abstractmethod async

run(
    user_prompt: str, *, usage: RunUsage | None = None
) -> TeamResult[OutputT]

Execute the team against user_prompt.

Parameters:

Name Type Description Default
user_prompt str

User input for the team run.

required
usage RunUsage | None

Optional usage accumulator shared with nested agent runs.

None

Returns:

Type Description
TeamResult[OutputT]

A TeamResult with final data and aggregated usage.

Source code in pydantic_team/base.py
@abstractmethod
async def run(self, user_prompt: str, *, usage: RunUsage | None = None) -> TeamResult[OutputT]:
    """Execute the team against `user_prompt`.

    Args:
        user_prompt: User input for the team run.
        usage: Optional usage accumulator shared with nested agent runs.

    Returns:
        A [`TeamResult`][pydantic_team.base.TeamResult] with final data and aggregated usage.
    """

Hierarchical team

pydantic_team.hierarchical.HierarchicalTeam

HierarchicalTeam(
    *,
    members: Sequence[TeamMember],
    leader_agent: Agent[object, object] | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
)

Bases: BaseTeam[object]

Leader-driven team that registers each member as a delegation tool.

Follows pydantic-ai agent delegation: nested member runs receive usage=ctx.usage so tokens aggregate on the leader run.

Create a hierarchical team.

Parameters:

Name Type Description Default
members Sequence[TeamMember]

Specialist agents or nested teams (at least one required).

required
leader_agent Agent[object, object] | None

Existing leader agent. Mutually exclusive with leader_model.

None
leader_model str | None

Model string used to construct a leader agent when leader_agent is not provided.

None
system_prompt_override str | None

Optional leader instructions / extra system prompt.

None
name str | None

Optional team name (used when this team is nested as a member tool).

None
Source code in pydantic_team/hierarchical.py
def __init__(
    self,
    *,
    members: Sequence[TeamMember],
    leader_agent: Agent[object, object] | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
) -> None:
    """Create a hierarchical team.

    Args:
        members: Specialist agents or nested teams (at least one required).
        leader_agent: Existing leader agent. Mutually exclusive with `leader_model`.
        leader_model: Model string used to construct a leader agent when
            `leader_agent` is not provided.
        system_prompt_override: Optional leader instructions / extra system prompt.
        name: Optional team name (used when this team is nested as a member tool).
    """
    if leader_agent is not None and leader_model is not None:
        raise ValueError('Provide leader_agent or leader_model, not both')
    if leader_agent is None and leader_model is None:
        raise ValueError('Provide leader_agent or leader_model')
    if not members:
        raise ValueError('members must be a non-empty sequence')

    self.name = name
    self._members: list[TeamMember] = list(members)

    if leader_agent is not None:
        self._leader = leader_agent
        if system_prompt_override is not None:

            def _leader_override_prompt() -> str:
                return system_prompt_override

            self._leader.system_prompt(_leader_override_prompt)

    else:
        assert leader_model is not None
        self._leader = Agent(
            leader_model,
            instructions=system_prompt_override or _DEFAULT_LEADER_INSTRUCTIONS,
        )

    self._register_member_tools()

leader property

leader: Agent[object, object]

The coordinating leader agent.

Collaborative team / board

pydantic_team.board.TaskStatus

Bases: str, Enum

Lifecycle status of a board task.

pydantic_team.board.Task dataclass

Task(
    id: str,
    title: str,
    description: str = '',
    status: TaskStatus = TaskStatus.OPEN,
    assignee: str | None = None,
    result: str | None = None,
    reviewer: str | None = None,
    rejection_reason: str | None = None,
)

A unit of work on the shared board.

pydantic_team.board.BoardMessage dataclass

BoardMessage(
    id: str,
    sender: str,
    to: str,
    body: str,
    task_id: str | None = None,
)

A peer message posted on the shared board.

pydantic_team.board.TaskBoard dataclass

TaskBoard(
    _tasks: dict[str, Task] = (lambda: {})(),
    _messages: list[BoardMessage] = (
        lambda: list[BoardMessage]()
    )(),
    _lock: Lock = asyncio.Lock(),
    _counter: int = 0,
    _message_counter: int = 0,
    _wakeup: Event = asyncio.Event(),
    _wakeup_agents: set[str] = (lambda: set[str]())(),
)

Thread-safe in-process task list shared by a collaborative team.

add_task async

add_task(
    title: str,
    description: str = '',
    *,
    reviewer: str | None = None,
) -> Task

Create an open task and return its snapshot.

Source code in pydantic_team/board.py
async def add_task(
    self,
    title: str,
    description: str = '',
    *,
    reviewer: str | None = None,
) -> Task:
    """Create an open task and return its snapshot."""
    async with self._lock:
        self._counter += 1
        task = Task(
            id=f'task-{self._counter}',
            title=title,
            description=description,
            reviewer=reviewer,
        )
        self._tasks[task.id] = task
        return task

list_tasks async

list_tasks(status: TaskStatus | None = None) -> list[Task]

Return task snapshots, optionally filtered by status.

Source code in pydantic_team/board.py
async def list_tasks(self, status: TaskStatus | None = None) -> list[Task]:
    """Return task snapshots, optionally filtered by status."""
    async with self._lock:
        tasks = list(self._tasks.values())
    if status is not None:
        tasks = [task for task in tasks if task.status is status]
    return tasks

post_message async

post_message(
    sender: str,
    to: str,
    body: str,
    *,
    task_id: str | None = None,
) -> BoardMessage

Append a peer message; optionally link to an existing task.

Source code in pydantic_team/board.py
async def post_message(
    self,
    sender: str,
    to: str,
    body: str,
    *,
    task_id: str | None = None,
) -> BoardMessage:
    """Append a peer message; optionally link to an existing task."""
    async with self._lock:
        if task_id is not None:
            self._require(task_id)
        self._message_counter += 1
        message = BoardMessage(
            id=f'msg-{self._message_counter}',
            sender=sender,
            to=to,
            body=body,
            task_id=task_id,
        )
        self._messages.append(message)
    if to == '*':
        self.signal_wakeup()
    else:
        self.signal_wakeup(to)
    return message

list_messages async

list_messages(
    *, agent_id: str | None = None
) -> list[BoardMessage]

Return messages; when agent_id is set, only visible ones for that agent.

Source code in pydantic_team/board.py
async def list_messages(self, *, agent_id: str | None = None) -> list[BoardMessage]:
    """Return messages; when ``agent_id`` is set, only visible ones for that agent."""
    async with self._lock:
        messages = list(self._messages)
    if agent_id is None:
        return messages
    return [
        message for message in messages if message.to == agent_id or message.to == '*' or message.sender == agent_id
    ]

messages_snapshot

messages_snapshot() -> list[BoardMessage]

Return a stable list copy of all messages (asyncio-safe between awaits).

Source code in pydantic_team/board.py
def messages_snapshot(self) -> list[BoardMessage]:
    """Return a stable list copy of all messages (asyncio-safe between awaits)."""
    return list(self._messages)

claim async

claim(task_id: str, agent_id: str) -> Task

Atomically claim an open task for agent_id.

Source code in pydantic_team/board.py
async def claim(self, task_id: str, agent_id: str) -> Task:
    """Atomically claim an open task for `agent_id`."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.OPEN:
            raise TaskClaimError(f'task {task_id!r} is not open (status={task.status})')
        updated = replace(task, status=TaskStatus.CLAIMED, assignee=agent_id)
        self._tasks[task_id] = updated
    self.signal_wakeup(agent_id)
    return updated

assign async

assign(task_id: str, agent_id: str) -> Task

Force-assign a non-done, non-pending-review task to agent_id (lead operation).

Source code in pydantic_team/board.py
async def assign(self, task_id: str, agent_id: str) -> Task:
    """Force-assign a non-done, non-pending-review task to `agent_id` (lead operation)."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is TaskStatus.DONE:
            raise TaskClaimError(f'task {task_id!r} is already done')
        if task.status is TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is pending review')
        updated = replace(task, status=TaskStatus.CLAIMED, assignee=agent_id)
        self._tasks[task_id] = updated
    self.signal_wakeup(agent_id)
    return updated

assign_reviewer async

assign_reviewer(task_id: str, reviewer_id: str) -> Task

Set or replace the reviewer on a non-done task.

Source code in pydantic_team/board.py
async def assign_reviewer(self, task_id: str, reviewer_id: str) -> Task:
    """Set or replace the reviewer on a non-done task."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is TaskStatus.DONE:
            raise TaskClaimError(f'task {task_id!r} is already done')
        updated = replace(task, reviewer=reviewer_id)
        self._tasks[task_id] = updated
    if updated.status is TaskStatus.PENDING_REVIEW:
        self.signal_wakeup(reviewer_id)
    return updated

complete async

complete(
    task_id: str, *, result: str, agent_id: str
) -> Task

Mark work submitted; gated tasks become pending_review, others done.

Source code in pydantic_team/board.py
async def complete(self, task_id: str, *, result: str, agent_id: str) -> Task:
    """Mark work submitted; gated tasks become pending_review, others done."""
    async with self._lock:
        task = self._require(task_id)
        if task.status not in (TaskStatus.CLAIMED, TaskStatus.NEEDS_REVISION):
            raise TaskClaimError(f'task {task_id!r} is not completable (status={task.status})')
        if task.assignee != agent_id:
            raise TaskClaimError(f'task {task_id!r} is assigned to {task.assignee!r}, not {agent_id!r}')
        if task.reviewer is not None:
            updated = replace(
                task,
                status=TaskStatus.PENDING_REVIEW,
                result=result,
                rejection_reason=None,
            )
            self._tasks[task_id] = updated
            wakeup_agent = task.reviewer
        else:
            updated = replace(
                task,
                status=TaskStatus.DONE,
                result=result,
                rejection_reason=None,
            )
            self._tasks[task_id] = updated
            wakeup_agent = None
    if wakeup_agent is not None:
        self.signal_wakeup(wakeup_agent)
    else:
        self.signal_wakeup()
    return updated

approve async

approve(task_id: str, *, agent_id: str) -> Task

Accept a pending_review task; only the reviewer may approve.

Source code in pydantic_team/board.py
async def approve(self, task_id: str, *, agent_id: str) -> Task:
    """Accept a pending_review task; only the reviewer may approve."""
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is not pending review (status={task.status})')
        if task.reviewer != agent_id:
            raise TaskClaimError(f'task {task_id!r} reviewer is {task.reviewer!r}, not {agent_id!r}')
        updated = replace(task, status=TaskStatus.DONE, rejection_reason=None)
        self._tasks[task_id] = updated
    self.signal_wakeup()
    return updated

reject async

reject(task_id: str, *, reason: str, agent_id: str) -> Task

Reject a pending_review task back to needs_revision; reason required.

Source code in pydantic_team/board.py
async def reject(self, task_id: str, *, reason: str, agent_id: str) -> Task:
    """Reject a pending_review task back to needs_revision; reason required."""
    cleaned = reason.strip()
    if not cleaned:
        raise TaskClaimError('rejection reason must be non-empty')
    async with self._lock:
        task = self._require(task_id)
        if task.status is not TaskStatus.PENDING_REVIEW:
            raise TaskClaimError(f'task {task_id!r} is not pending review (status={task.status})')
        if task.reviewer != agent_id:
            raise TaskClaimError(f'task {task_id!r} reviewer is {task.reviewer!r}, not {agent_id!r}')
        updated = replace(
            task,
            status=TaskStatus.NEEDS_REVISION,
            rejection_reason=cleaned,
        )
        self._tasks[task_id] = updated
        assignee = task.assignee
    self.signal_wakeup(assignee)
    return updated

is_complete

is_complete() -> bool

Return True when there are no tasks or every task is done.

Source code in pydantic_team/board.py
def is_complete(self) -> bool:
    """Return True when there are no tasks or every task is done."""
    if not self._tasks:
        return True
    return all(task.status is TaskStatus.DONE for task in self._tasks.values())

snapshot

snapshot() -> list[Task]

Return a stable list copy of all tasks (asyncio-safe between awaits).

Source code in pydantic_team/board.py
def snapshot(self) -> list[Task]:
    """Return a stable list copy of all tasks (asyncio-safe between awaits)."""
    return list(self._tasks.values())

signal_wakeup

signal_wakeup(agent_id: str | None = None) -> None

Wake waiters; optionally record which agent gained work.

Source code in pydantic_team/board.py
def signal_wakeup(self, agent_id: str | None = None) -> None:
    """Wake waiters; optionally record which agent gained work."""
    if agent_id is not None:
        self._wakeup_agents.add(agent_id)
    self._wakeup.set()

wait_wakeup async

wait_wakeup() -> set[str]

Block until signal_wakeup; return agent ids recorded since last wait.

Source code in pydantic_team/board.py
async def wait_wakeup(self) -> set[str]:
    """Block until ``signal_wakeup``; return agent ids recorded since last wait."""
    await self._wakeup.wait()
    self._wakeup.clear()
    agents = set(self._wakeup_agents)
    self._wakeup_agents.clear()
    return agents

pydantic_team.collaborative.BoardDeps dataclass

BoardDeps(
    board: TaskBoard,
    agent_id: str,
    member_ids: tuple[str, ...],
    leader_id: str = 'leader',
    emit: EmitFn | None = None,
    require_review: bool = False,
)

Dependencies injected into leader and member agent runs.

pydantic_team.collaborative.CollaborativeTeam

CollaborativeTeam(
    *,
    members: Sequence[AnyAgent],
    leader_agent: AnyAgent | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
    max_rounds: int = 3,
    max_replans: int = 0,
    max_assignments_per_tick: int | None = None,
    dispatch_mode: DispatchMode = 'phased',
    require_review: bool = False,
)

Bases: BaseTeam[object]

Team that coordinates work through a shared TaskBoard.

The leader creates and assigns tasks by role; members complete their assigned work in parallel (phased rounds or streaming dispatch). When the board remains incomplete, the leader may replan (up to max_replans) before synthesizing. Members may message each other directly via board tools (send_message / list_messages).

Create a collaborative team.

Parameters:

Name Type Description Default
members Sequence[AnyAgent]

Teammate agents (at least one). Nested teams are not supported here.

required
leader_agent AnyAgent | None

Existing leader. Mutually exclusive with leader_model.

None
leader_model str | None

Model string used to build the leader when leader_agent is omitted.

None
system_prompt_override str | None

Optional leader instructions / extra system prompt.

None
name str | None

Optional team name.

None
max_rounds int

In phased mode, parallel member ticks per phase (after seed and after each replan). In streaming mode, max ticks per member per phase.

3
max_replans int

How many times the leader may replan after incomplete member phases. 0 preserves seed → work → synthesize with no replan.

0
max_assignments_per_tick int | None

If set, each member tick only lists this many incomplete assignments (caps work per round without relying on soft prompt wording).

None
dispatch_mode DispatchMode

phased (default) runs seed then member rounds; streaming starts member ticks as soon as tasks are assigned (overlap with seed/replan).

'phased'
require_review bool

When True, new tasks get reviewer=leader so complete goes to pending_review until approve/reject. When False, tasks complete to done unless assign_reviewer sets a reviewer.

False
Source code in pydantic_team/collaborative.py
def __init__(
    self,
    *,
    members: Sequence[AnyAgent],
    leader_agent: AnyAgent | None = None,
    leader_model: str | None = None,
    system_prompt_override: str | None = None,
    name: str | None = None,
    max_rounds: int = 3,
    max_replans: int = 0,
    max_assignments_per_tick: int | None = None,
    dispatch_mode: DispatchMode = 'phased',
    require_review: bool = False,
) -> None:
    """Create a collaborative team.

    Args:
        members: Teammate agents (at least one). Nested teams are not supported here.
        leader_agent: Existing leader. Mutually exclusive with `leader_model`.
        leader_model: Model string used to build the leader when `leader_agent` is omitted.
        system_prompt_override: Optional leader instructions / extra system prompt.
        name: Optional team name.
        max_rounds: In ``phased`` mode, parallel member ticks per phase (after seed and
            after each replan). In ``streaming`` mode, max ticks per member per phase.
        max_replans: How many times the leader may replan after incomplete member phases.
            ``0`` preserves seed → work → synthesize with no replan.
        max_assignments_per_tick: If set, each member tick only lists this many incomplete
            assignments (caps work per round without relying on soft prompt wording).
        dispatch_mode: ``phased`` (default) runs seed then member rounds; ``streaming``
            starts member ticks as soon as tasks are assigned (overlap with seed/replan).
        require_review: When True, new tasks get ``reviewer=leader`` so ``complete`` goes
            to ``pending_review`` until approve/reject. When False, tasks complete to
            ``done`` unless ``assign_reviewer`` sets a reviewer.
    """
    if leader_agent is not None and leader_model is not None:
        raise ValueError('Provide leader_agent or leader_model, not both')
    if leader_agent is None and leader_model is None:
        raise ValueError('Provide leader_agent or leader_model')
    if not members:
        raise ValueError('members must be a non-empty sequence')
    if max_rounds < 1:
        raise ValueError('max_rounds must be >= 1')
    if max_replans < 0:
        raise ValueError('max_replans must be >= 0')
    if max_assignments_per_tick is not None and max_assignments_per_tick < 1:
        raise ValueError('max_assignments_per_tick must be >= 1')
    if dispatch_mode not in ('phased', 'streaming'):
        raise ValueError("dispatch_mode must be 'phased' or 'streaming'")

    self.name = name
    self._max_rounds = max_rounds
    self._max_replans = max_replans
    self._max_assignments_per_tick = max_assignments_per_tick
    self._dispatch_mode: DispatchMode = dispatch_mode
    self._require_review = require_review
    self._members: list[AnyAgent] = list(members)
    self._member_ids: list[str] = [
        _agent_id(member, fallback=f'member-{index}') for index, member in enumerate(self._members)
    ]

    if leader_agent is not None:
        self._leader: AnyAgent = leader_agent
        if system_prompt_override is not None:

            def _leader_override_prompt() -> str:
                return system_prompt_override

            self._leader.system_prompt(_leader_override_prompt)
    else:
        assert leader_model is not None
        self._leader = cast(
            AnyAgent,
            Agent(
                leader_model,
                name='leader',
                instructions=system_prompt_override or default_leader_instructions(self._member_ids),
                deps_type=BoardDeps,
            ),
        )

    _register_leader_tools(self._leader)
    for member in self._members:
        _register_member_tools(member)

leader property

leader: AnyAgent

The team lead agent.

dispatch_mode property

dispatch_mode: DispatchMode

Scheduling mode: phased or streaming.

member_ids property

member_ids: Sequence[str]

Stable teammate ids used with assign_task (agent names).

iter

iter(
    user_prompt: str, *, usage: RunUsage | None = None
) -> CollaborativeRun

Start an observable collaborative run (async context + async iterator).

Source code in pydantic_team/collaborative.py
def iter(self, user_prompt: str, *, usage: RunUsage | None = None) -> CollaborativeRun:
    """Start an observable collaborative run (async context + async iterator)."""
    return CollaborativeRun(
        leader=self._leader,
        members=self._members,
        member_ids=self._member_ids,
        max_rounds=self._max_rounds,
        max_replans=self._max_replans,
        max_assignments_per_tick=self._max_assignments_per_tick,
        dispatch_mode=self._dispatch_mode,
        run_member=self._run_member,
        user_prompt=user_prompt,
        require_review=self._require_review,
        _usage=usage or RunUsage(),
    )

pydantic_team.collaborative.CollaborativeRun dataclass

CollaborativeRun(
    leader: AnyAgent,
    members: Sequence[AnyAgent],
    member_ids: Sequence[str],
    max_rounds: int,
    max_replans: int,
    max_assignments_per_tick: int | None,
    dispatch_mode: DispatchMode,
    run_member: RunMemberFn,
    user_prompt: str,
    require_review: bool,
    _usage: RunUsage,
    _board: TaskBoard = TaskBoard(),
    _events: Queue[TeamEvent | None] = (
        lambda: asyncio.Queue[TeamEvent | None]()
    )(),
    _result: TeamResult[object] | None = None,
    _driver: Task[None] | None = None,
    _error: BaseException | None = None,
)

Step-by-step collaborative run (inspired by pydantic-graph GraphRun).

board property

board: TaskBoard

Live task board for this run.

result property

result: TeamResult[object] | None

Final result once the run has ended; otherwise None.

usage property

usage: RunUsage

Aggregated usage accumulated during the run.

Run events

pydantic_team.events.TeamTask dataclass

TeamTask(
    kind: TaskKind,
    agent_id: str | None = None,
    task_ids: tuple[str, ...] = (),
)

A unit of scheduled collaborative work (inspired by pydantic-graph GraphTask).

pydantic_team.events.TasksScheduled dataclass

TasksScheduled(tasks: tuple[TeamTask, ...])

One or more tasks are about to run (or have just been spawned).

pydantic_team.events.TaskCompleted dataclass

TaskCompleted(task: TeamTask)

A scheduled task finished (board may have been mutated via tools).

pydantic_team.events.PhaseJoined dataclass

PhaseJoined(phase: PhaseName, incomplete: bool)

No inflight work remains for this phase (join / barrier).

pydantic_team.events.MessagePosted dataclass

MessagePosted(message: BoardMessage)

A teammate posted a peer message on the board.

pydantic_team.events.TaskReviewDecided dataclass

TaskReviewDecided(task: Task, decision: ReviewDecision)

A reviewer approved or rejected a board task.

pydantic_team.events.RunEnded dataclass

RunEnded(result: TeamResult[object])

The collaborative run finished with a final result.