跳转至

jianmu.swarm

适用对象:多 Agent 应用开发者 / 协作运行时维护者 是否必读:按需 相关模块:jianmu.message, jianmu.tool, jianmu.model, jianmu.mcp

1. 模块职责

jianmu.swarm 提供多 Agent 角色、事件、消息封装、运行时、宿主和工具 provider 的公开 API。

如果你需要的不再是单个节点或单条工具调用,而是多个 Agent 之间的协作,这个模块就是主入口。

2. 适合查什么

  • 配置与领域对象:AgentProfile、GroupProfile、BudgetSpec
  • 运行时:AgentRuntime、DefaultAgentState
  • 事件与消息:AgentEvent、MessageEnvelope、InMemoryEventBus
  • 消息存储:EnvelopeStoreProtocol、InMemoryMessageStore
  • 宿主与角色:AgentHost、SwarmRole、role
  • provider:SwarmToolProvider、MCPToolProvider

3. 使用建议

  • 多 Agent 编排优先从这里的 profile / role / runtime 组合开始
  • 普通单 Agent 流程不要过早引入 swarm 抽象
  • 更底层的调度细节属于内部 runtime 子模块,不应默认直接依赖
  • 如果你要自定义消息持久化或投递语义,优先实现 EnvelopeStoreProtocol,而不是直接耦合 runtime 内部细节

4. 最小示例

from jianmu.swarm import AgentProfile, role


@role
def researcher() -> str:
    return "Research the question and report back."


profile = AgentProfile(name="researcher")

5. 常见入口

  • 想定义角色:看 role、SwarmRole
  • 想配置 Agent / Group:看 AgentProfile、GroupProfile
  • 想接消息总线:看 InMemoryEventBus
  • 想替换消息存储:看 EnvelopeStoreProtocol、InMemoryMessageStore
  • 想启动协作运行时:看 AgentRuntime、AgentHost

6. 注意事项

  • 单 Agent 流程不必为了后续扩展提前引入 swarm 抽象
  • 更底层的调度细节仍属于内部 runtime 子模块,不建议直接依赖
  • 角色、事件、runtime 的职责是分层的,建议先按 profile / role / runtime 入口理解
  • context_builder 是统一入口:既可参与 mailbox/history 合并,也会注入到 prompt runtime 供 LLM 上下文装配使用

7. API 参考

swarm

Multi-agent collaboration runtime, message routing, roles, and event bus for Jianmu swarms.

BudgetSpec dataclass

BudgetSpec(
    total_tokens: Optional[int] = None,
    prompt_tokens: Optional[int] = None,
    completion_tokens: Optional[int] = None,
    max_tool_calls: Optional[int] = None,
)

Token and call budget limits applied to one agent run.

属性:

名称 类型 描述
total_tokens Optional[int]

Overall token ceiling for the run.

prompt_tokens Optional[int]

Prompt-token ceiling across all model calls.

completion_tokens Optional[int]

Completion-token ceiling across all model calls.

max_tool_calls Optional[int]

Maximum number of tool invocations allowed.

AgentProfile dataclass

AgentProfile(
    id: str,
    role: str,
    parent_id: Optional[str] = None,
    task: Optional[str] = None,
    budget: Optional[BudgetSpec] = None,
    metadata: dict[str, Any] = dict(),
)

Identity and configuration snapshot for one agent in the swarm.

属性:

名称 类型 描述
id str

Stable agent identifier used by the swarm runtime.

role str

Registered role name used to build the agent.

parent_id Optional[str]

Optional parent agent ID when this agent was spawned.

task Optional[str]

Optional task statement currently assigned to the agent.

budget Optional[BudgetSpec]

Optional run budget associated with the agent.

metadata dict[str, Any]

Arbitrary application metadata attached to the profile.

GroupProfile dataclass

GroupProfile(
    id: str,
    members: list[str] = list(),
    metadata: dict[str, Any] = dict(),
)

Membership record for a broadcast group channel.

属性:

名称 类型 描述
id str

Group identifier used for group message routing.

members list[str]

Agent IDs currently subscribed to the group.

metadata dict[str, Any]

Arbitrary metadata associated with the group.

AgentEvent dataclass

AgentEvent(
    event_type: str,
    agent_id: Optional[str] = None,
    payload: dict[str, Any] = dict(),
    created_at: Optional[float] = None,
    event_id: Optional[str] = None,
)

Runtime event emitted by an agent or the scheduler.

属性:

名称 类型 描述
event_type str

Event name such as lifecycle, routing, or mailbox actions.

agent_id Optional[str]

Optional agent ID associated with the event.

payload dict[str, Any]

Structured event-specific data.

created_at Optional[float]

Optional event timestamp in Unix seconds.

event_id Optional[str]

Optional unique event identifier.

MessageEnvelope dataclass

MessageEnvelope(
    seq: int,
    message: Message,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    created_at: float = 0.0,
)

Stored message row with monotonic sequence for cursor-based reads.

属性:

名称 类型 描述
seq int

Monotonic sequence number used for cursor-based consumption.

message Message

The stored message payload.

to_agent_id str | None

Optional direct-recipient agent ID.

group_id str | None

Optional group recipient ID.

topic str | None

Optional topic label used for selective delivery.

created_at float

Message creation timestamp in Unix seconds.

AgentRole

Bases: Protocol

Protocol implemented by role factories that build agent trees.

build_tree

build_tree(profile: AgentProfile) -> Any

Build the root tree for an agent profile.

参数:

名称 类型 描述 默认
profile AgentProfile

Agent profile used to configure the role-specific tree.

必需

返回:

类型 描述
Any

Root node or tree-like object that the runtime can execute.

源代码位于: jianmu/swarm/base.py
def build_tree(self, profile: AgentProfile) -> Any:
    """Build the root tree for an agent profile.

    Args:
        profile: Agent profile used to configure the role-specific tree.

    Returns:
        Root node or tree-like object that the runtime can execute.
    """
    ...

AgentRuntimeProtocol

Bases: Protocol

Public runtime interface for spawning and coordinating swarm agents.

spawn

spawn(
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
) -> str

Create an agent host and return its runtime identifier.

参数:

名称 类型 描述 默认
role str | AgentRole

Registered role name or role factory instance.

必需
task str | None

Optional initial task for the spawned agent.

None
parent_id str | None

Optional parent agent identifier for nested spawns.

None
mode str

Spawn mode understood by the runtime implementation.

'detached'
constraints Any | None

Optional runtime-specific execution constraints.

None

返回:

类型 描述
str

Runtime agent identifier.

源代码位于: jianmu/swarm/base.py
def spawn(
    self,
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
) -> str:
    """Create an agent host and return its runtime identifier.

    Args:
        role: Registered role name or role factory instance.
        task: Optional initial task for the spawned agent.
        parent_id: Optional parent agent identifier for nested spawns.
        mode: Spawn mode understood by the runtime implementation.
        constraints: Optional runtime-specific execution constraints.

    Returns:
        Runtime agent identifier.
    """
    ...

kill

kill(agent_id: str) -> None

Terminate one agent synchronously.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
源代码位于: jianmu/swarm/base.py
def kill(self, agent_id: str) -> None:
    """Terminate one agent synchronously.

    Args:
        agent_id: Target agent identifier.
    """
    ...

kill_async async

kill_async(agent_id: str) -> None

Terminate one agent asynchronously.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
源代码位于: jianmu/swarm/base.py
async def kill_async(self, agent_id: str) -> None:
    """Terminate one agent asynchronously.

    Args:
        agent_id: Target agent identifier.
    """
    ...

pause

pause(agent_id: str) -> None

Pause one agent so it stops taking new work.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
源代码位于: jianmu/swarm/base.py
def pause(self, agent_id: str) -> None:
    """Pause one agent so it stops taking new work.

    Args:
        agent_id: Target agent identifier.
    """
    ...

resume

resume(agent_id: str) -> None

Resume a previously paused agent.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
源代码位于: jianmu/swarm/base.py
def resume(self, agent_id: str) -> None:
    """Resume a previously paused agent.

    Args:
        agent_id: Target agent identifier.
    """
    ...

wake

wake(agent_id: str, reason: str | None = None) -> None

Signal an agent that new work is available.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
reason str | None

Optional explanation or trigger reason for wake-up.

None
源代码位于: jianmu/swarm/base.py
def wake(self, agent_id: str, reason: str | None = None) -> None:
    """Signal an agent that new work is available.

    Args:
        agent_id: Target agent identifier.
        reason: Optional explanation or trigger reason for wake-up.
    """
    ...

preempt_agent async

preempt_agent(
    agent_id: str, reason: str | None = None
) -> bool

Interrupt one agent and ask it to yield control.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
reason str | None

Optional reason for the preemption.

None

返回:

类型 描述
bool

True if the agent yielded successfully, False otherwise.

源代码位于: jianmu/swarm/base.py
async def preempt_agent(self, agent_id: str, reason: str | None = None) -> bool:
    """Interrupt one agent and ask it to yield control.

    Args:
        agent_id: Target agent identifier.
        reason: Optional reason for the preemption.

    Returns:
        True if the agent yielded successfully, False otherwise.
    """
    ...

send_message async

send_message(
    *,
    sender_id: str | None,
    content: str,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: Optional[dict[str, Any]] = None,
) -> int

Send a message into the swarm message store.

参数:

名称 类型 描述 默认
sender_id str | None

Originating agent identifier, or None for external senders.

必需
content str

Message body.

必需
to_agent_id str | None

Optional direct-recipient agent identifier.

None
group_id str | None

Optional group recipient.

None
topic str | None

Optional topic recipient.

None
content_type str

Message content type label.

'text'
metadata Optional[dict[str, Any]]

Optional structured metadata persisted with the message.

None

返回:

类型 描述
int

Monotonic store sequence number or equivalent delivery token.

源代码位于: jianmu/swarm/base.py
async def send_message(
    self,
    *,
    sender_id: str | None,
    content: str,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: Optional[dict[str, Any]] = None,
) -> int:
    """Send a message into the swarm message store.

    Args:
        sender_id: Originating agent identifier, or ``None`` for external senders.
        content: Message body.
        to_agent_id: Optional direct-recipient agent identifier.
        group_id: Optional group recipient.
        topic: Optional topic recipient.
        content_type: Message content type label.
        metadata: Optional structured metadata persisted with the message.

    Returns:
        Monotonic store sequence number or equivalent delivery token.
    """
    ...

emit

emit(event: AgentEvent) -> None

Publish one runtime event to subscribers.

参数:

名称 类型 描述 默认
event AgentEvent

The AgentEvent instance to emit.

必需
源代码位于: jianmu/swarm/base.py
def emit(self, event: AgentEvent) -> None:
    """Publish one runtime event to subscribers.

    Args:
        event: The AgentEvent instance to emit.
    """
    ...

subscribe_events

subscribe_events(
    callback: Callable[[AgentEvent], None],
) -> Callable[[], None]

Subscribe to runtime events.

参数:

名称 类型 描述 默认
callback Callable[[AgentEvent], None]

Consumer invoked for each emitted AgentEvent.

必需

返回:

类型 描述
Callable[[], None]

Unsubscribe callback.

源代码位于: jianmu/swarm/base.py
def subscribe_events(self, callback: Callable[[AgentEvent], None]) -> Callable[[], None]:
    """Subscribe to runtime events.

    Args:
        callback: Consumer invoked for each emitted ``AgentEvent``.

    Returns:
        Unsubscribe callback.
    """
    ...

subscribe_observability

subscribe_observability(
    callback: Callable[[str, dict[str, Any]], None],
    *,
    runtime_events: bool = True,
    trace_events: bool = False,
) -> Callable[[], None]

Subscribe to runtime and trace observability streams.

参数:

名称 类型 描述 默认
callback Callable[[str, dict[str, Any]], None]

Consumer invoked with an event name and structured payload.

必需
runtime_events bool

Whether runtime-level events should be forwarded.

True
trace_events bool

Whether lower-level trace events should be forwarded.

False

返回:

类型 描述
Callable[[], None]

Unsubscribe callback.

源代码位于: jianmu/swarm/base.py
def subscribe_observability(
    self,
    callback: Callable[[str, dict[str, Any]], None],
    *,
    runtime_events: bool = True,
    trace_events: bool = False,
) -> Callable[[], None]:
    """Subscribe to runtime and trace observability streams.

    Args:
        callback: Consumer invoked with an event name and structured payload.
        runtime_events: Whether runtime-level events should be forwarded.
        trace_events: Whether lower-level trace events should be forwarded.

    Returns:
        Unsubscribe callback.
    """
    ...

get_profile

get_profile(agent_id: str) -> Optional[AgentProfile]

Return the runtime profile for one agent.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需

返回:

类型 描述
Optional[AgentProfile]

AgentProfile if found, else None.

源代码位于: jianmu/swarm/base.py
def get_profile(self, agent_id: str) -> Optional[AgentProfile]:
    """Return the runtime profile for one agent.

    Args:
        agent_id: Target agent identifier.

    Returns:
        AgentProfile if found, else None.
    """
    ...

list_agents

list_agents() -> list[AgentProfile]

Return profiles for all known agents.

返回:

类型 描述
list[AgentProfile]

List of AgentProfile objects.

源代码位于: jianmu/swarm/base.py
def list_agents(self) -> list[AgentProfile]:
    """Return profiles for all known agents.

    Returns:
        List of AgentProfile objects.
    """
    ...

save_snapshot async

save_snapshot(
    *, thread_id: str | None = None, step: int | None = None
) -> dict[str, Any]

Persist a runtime snapshot.

参数:

名称 类型 描述 默认
thread_id str | None

Optional checkpoint thread identifier.

None
step int | None

Optional explicit step number for the snapshot.

None

返回:

类型 描述
dict[str, Any]

Serializable snapshot payload.

源代码位于: jianmu/swarm/base.py
async def save_snapshot(self, *, thread_id: str | None = None, step: int | None = None) -> dict[str, Any]:
    """Persist a runtime snapshot.

    Args:
        thread_id: Optional checkpoint thread identifier.
        step: Optional explicit step number for the snapshot.

    Returns:
        Serializable snapshot payload.
    """
    ...

restore_snapshot async

restore_snapshot(*, thread_id: str | None = None) -> bool

Restore runtime state from a previously saved snapshot.

参数:

名称 类型 描述 默认
thread_id str | None

Optional checkpoint thread identifier.

None

返回:

类型 描述
bool

True if a snapshot was found and restored, otherwise False.

源代码位于: jianmu/swarm/base.py
async def restore_snapshot(self, *, thread_id: str | None = None) -> bool:
    """Restore runtime state from a previously saved snapshot.

    Args:
        thread_id: Optional checkpoint thread identifier.

    Returns:
        ``True`` if a snapshot was found and restored, otherwise ``False``.
    """
    ...

start_all

start_all() -> None

Start background execution for all root agents.

源代码位于: jianmu/swarm/base.py
def start_all(self) -> None:
    """Start background execution for all root agents."""
    ...

stop_all

stop_all() -> None

Stop background execution for all managed agents.

源代码位于: jianmu/swarm/base.py
def stop_all(self) -> None:
    """Stop background execution for all managed agents."""
    ...

step_agent async

step_agent(
    agent_id: str, obs: Optional[dict[str, Any]] = None
) -> Any

Advance a single agent by one step.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
obs Optional[dict[str, Any]]

Optional observation payload injected before stepping.

None

返回:

类型 描述
Any

Runtime-defined step result.

源代码位于: jianmu/swarm/base.py
async def step_agent(self, agent_id: str, obs: Optional[dict[str, Any]] = None) -> Any:
    """Advance a single agent by one step.

    Args:
        agent_id: Target agent identifier.
        obs: Optional observation payload injected before stepping.

    Returns:
        Runtime-defined step result.
    """
    ...

step_agent_with_actions async

step_agent_with_actions(
    agent_id: str, obs: Optional[dict[str, Any]] = None
) -> tuple[Any, dict[str, Any]]

Advance a single agent and return both status and emitted actions.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
obs Optional[dict[str, Any]]

Optional observation payload injected before stepping.

None

返回:

类型 描述
tuple[Any, dict[str, Any]]

Tuple of runtime-defined status and action mapping.

源代码位于: jianmu/swarm/base.py
async def step_agent_with_actions(
    self,
    agent_id: str,
    obs: Optional[dict[str, Any]] = None,
) -> tuple[Any, dict[str, Any]]:
    """Advance a single agent and return both status and emitted actions.

    Args:
        agent_id: Target agent identifier.
        obs: Optional observation payload injected before stepping.

    Returns:
        Tuple of runtime-defined status and action mapping.
    """
    ...

step_all async

step_all() -> dict[str, Any]

Advance all agents by one scheduler step.

返回:

类型 描述
dict[str, Any]

Mapping of agent IDs to their step execution status/results.

源代码位于: jianmu/swarm/base.py
async def step_all(self) -> dict[str, Any]:
    """Advance all agents by one scheduler step.

    Returns:
        Mapping of agent IDs to their step execution status/results.
    """
    ...

step_all_concurrent async

step_all_concurrent(
    *, max_concurrency: int | None = None
) -> dict[str, Any]

Advance all agents concurrently.

参数:

名称 类型 描述 默认
max_concurrency int | None

Optional cap for concurrent stepping tasks.

None

返回:

类型 描述
dict[str, Any]

Mapping from agent identifier to runtime-defined step result.

源代码位于: jianmu/swarm/base.py
async def step_all_concurrent(self, *, max_concurrency: int | None = None) -> dict[str, Any]:
    """Advance all agents concurrently.

    Args:
        max_concurrency: Optional cap for concurrent stepping tasks.

    Returns:
        Mapping from agent identifier to runtime-defined step result.
    """
    ...

run_until_idle async

run_until_idle(max_steps: int = 200) -> dict[str, Any]

Drive the runtime until no agents have active work.

参数:

名称 类型 描述 默认
max_steps int

Maximum scheduler iterations before aborting.

200

返回:

类型 描述
dict[str, Any]

Final per-agent status mapping.

源代码位于: jianmu/swarm/base.py
async def run_until_idle(self, max_steps: int = 200) -> dict[str, Any]:
    """Drive the runtime until no agents have active work.

    Args:
        max_steps: Maximum scheduler iterations before aborting.

    Returns:
        Final per-agent status mapping.
    """
    ...

join_group

join_group(agent_id: str, group_id: str) -> None

Subscribe an agent to a broadcast group.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
group_id str

Group identifier to subscribe to.

必需
源代码位于: jianmu/swarm/base.py
def join_group(self, agent_id: str, group_id: str) -> None:
    """Subscribe an agent to a broadcast group.

    Args:
        agent_id: Target agent identifier.
        group_id: Group identifier to subscribe to.
    """
    ...

leave_group

leave_group(agent_id: str, group_id: str) -> None

Remove an agent from a broadcast group.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
group_id str

Group identifier to leave.

必需
源代码位于: jianmu/swarm/base.py
def leave_group(self, agent_id: str, group_id: str) -> None:
    """Remove an agent from a broadcast group.

    Args:
        agent_id: Target agent identifier.
        group_id: Group identifier to leave.
    """
    ...

subscribe_topic

subscribe_topic(agent_id: str, topic: str) -> None

Subscribe an agent to a topic stream.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
topic str

Topic identifier to subscribe to.

必需
源代码位于: jianmu/swarm/base.py
def subscribe_topic(self, agent_id: str, topic: str) -> None:
    """Subscribe an agent to a topic stream.

    Args:
        agent_id: Target agent identifier.
        topic: Topic identifier to subscribe to.
    """
    ...

unsubscribe_topic

unsubscribe_topic(agent_id: str, topic: str) -> None

Unsubscribe an agent from a topic stream.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需
topic str

Topic identifier to unsubscribe from.

必需
源代码位于: jianmu/swarm/base.py
def unsubscribe_topic(self, agent_id: str, topic: str) -> None:
    """Unsubscribe an agent from a topic stream.

    Args:
        agent_id: Target agent identifier.
        topic: Topic identifier to unsubscribe from.
    """
    ...

get_all_tools

get_all_tools(agent_id: str) -> list

Get all available tools for an agent.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需

返回:

类型 描述
list

Effective tool list after provider resolution.

源代码位于: jianmu/swarm/base.py
def get_all_tools(self, agent_id: str) -> list:
    """Get all available tools for an agent.

    Args:
        agent_id: Target agent identifier.

    Returns:
        Effective tool list after provider resolution.
    """
    ...

explain_tools

explain_tools(agent_id: str) -> dict[str, Any]

Explain tool resolution for an agent.

参数:

名称 类型 描述 默认
agent_id str

Target agent identifier.

必需

返回:

类型 描述
dict[str, Any]

Diagnostic payload describing providers, overrides, and final tools.

源代码位于: jianmu/swarm/base.py
def explain_tools(self, agent_id: str) -> dict[str, Any]:
    """Explain tool resolution for an agent.

    Args:
        agent_id: Target agent identifier.

    Returns:
        Diagnostic payload describing providers, overrides, and final tools.
    """
    ...

get_dead_letters async

get_dead_letters(
    *, agent_id: str | None = None, limit: int = 100
) -> list[dict[str, Any]]

Return failed-delivery envelopes retained by the runtime.

参数:

名称 类型 描述 默认
agent_id str | None

Optional filter for one agent's dead letters.

None
limit int

Maximum number of envelopes to return.

100

返回:

类型 描述
list[dict[str, Any]]

Dead-letter payloads in runtime-defined dictionary form.

源代码位于: jianmu/swarm/base.py
async def get_dead_letters(self, *, agent_id: str | None = None, limit: int = 100) -> list[dict[str, Any]]:
    """Return failed-delivery envelopes retained by the runtime.

    Args:
        agent_id: Optional filter for one agent's dead letters.
        limit: Maximum number of envelopes to return.

    Returns:
        Dead-letter payloads in runtime-defined dictionary form.
    """
    ...

DefaultAgentState

Bases: BaseModel

Default schema for runtime host state (extra fields allowed).

属性:

名称 类型 描述
done Annotated[bool, Ephemeral(scope='run')]

Completion flag produced during the current run.

final_answer Annotated[str, Ephemeral(scope='run')]

Final answer produced by the agent during the current run.

rounds Annotated[int, Ephemeral(scope='run')]

Current round counter for the agent loop.

tool_effects Annotated[dict[str, int], Ephemeral(scope='run')]

Aggregated tool-side-effect counters.

actions Annotated[list[Any], Ephemeral(scope='step')]

Per-step action trace for the current step scope.

streaming_output Annotated[str, Ephemeral(scope='call')]

Incremental streamed text for the current call scope.

AgentRuntime

AgentRuntime(
    *,
    roles: Optional[Dict[str, AgentRole]] = None,
    event_bus: Optional[InMemoryEventBus] = None,
    message_store: Optional[EnvelopeStoreProtocol] = None,
    state_factory: Optional[
        Callable[[AgentProfile], StateManager]
    ] = None,
    runner_factory: Optional[
        Callable[
            [Behaviour, StateManager, RunContext],
            ReactiveRunner,
        ]
    ] = None,
    tool_providers: Optional[list[ToolProvider]] = None,
    default_model_client: Any = None,
    auto_start: bool = False,
    session_id: Optional[str] = None,
    mcp_config_path: Optional[str] = None,
    skills_dir: Optional[str] = None,
    mailbox_drain_max_iterations: int = 100,
    message_store_gc_min_advance: int = 1000,
    delivery_max_failures: int = 3,
    restart_policy: str = "on_failure",
    max_restarts: int = 3,
    restart_backoff_s: float = 0.2,
    same_run_no_send_retries: int = 1,
    checkpointer: Optional[CheckpointerProtocol] = None,
    checkpoint_thread_id: Optional[str] = None,
    context_builder: Optional[
        ContextBuilderProtocol
    ] = None,
    runtime_event_bus: Optional[RuntimeEventBus] = None,
)

Unified multi-agent runtime with host/message-store/tool-provider composition.

AgentRuntime APIs include deterministic stepping: - step_agent(agent_id) - step_all() - run_until_idle()

属性:

名称 类型 描述
default_model_client

Default model client injected into spawned hosts.

session_id

Optional session identifier propagated into runtime events.

_event_bus

In-memory bus carrying runtime lifecycle events.

_message_store

Envelope store used for inter-agent message delivery.

_tool_providers list[ToolProvider]

Ordered tool providers merged for each host.

_checkpointer

Optional snapshot backend for runtime state persistence.

_checkpoint_thread_id

Default checkpoint thread identifier override.

_context_builder

Optional shared context builder for mailbox and prompt use.

_runtime_event_bus

Optional cross-module event bus for swarm lifecycle events.

Create a multi-agent runtime.

参数:

名称 类型 描述 默认
roles Optional[Dict[str, AgentRole]]

Optional mapping of role names to role factories.

None
event_bus Optional[InMemoryEventBus]

Optional event bus implementation.

None
message_store Optional[EnvelopeStoreProtocol]

Optional message store implementation.

None
state_factory Optional[Callable[[AgentProfile], StateManager]]

Factory used to build a state manager for each agent.

None
runner_factory Optional[Callable[[Behaviour, StateManager, RunContext], ReactiveRunner]]

Factory used to build a runner for each agent tree.

None
tool_providers Optional[list[ToolProvider]]

Additional tool providers merged after builtin swarm tools.

None
default_model_client Any

Default model client injected into spawned agents.

None
auto_start bool

Whether spawned agents should begin running immediately.

False
session_id Optional[str]

Optional session identifier propagated into events and state.

None
mcp_config_path Optional[str]

Optional MCP provider configuration path.

None
skills_dir Optional[str]

Optional skills directory used to bootstrap prompt configuration.

None
mailbox_drain_max_iterations int

Max envelopes to drain in one mailbox pass.

100
message_store_gc_min_advance int

Minimum sequence advance before store pruning.

1000
delivery_max_failures int

Maximum failed deliveries before dead-letter handling.

3
restart_policy str

Host restart policy: never, on_failure, or always.

'on_failure'
max_restarts int

Maximum restart attempts per host.

3
restart_backoff_s float

Delay before retrying a failed host.

0.2
same_run_no_send_retries int

Immediate retries when a run completes without sending.

1
checkpointer Optional[CheckpointerProtocol]

Optional runtime checkpoint backend.

None
checkpoint_thread_id Optional[str]

Optional checkpoint thread identifier override.

None
context_builder Optional[ContextBuilderProtocol]

Optional context builder used both when merging incoming history and when injecting prompt-time context.

None
runtime_event_bus Optional[RuntimeEventBus]

Optional cross-module runtime event bus used to expose swarm lifecycle events alongside single-agent flows.

None

引发:

类型 描述
ValueError

If numeric guardrails or restart-policy values are invalid.

源代码位于: jianmu/swarm/runtime/core.py
def __init__(
    self,
    *,
    roles: Optional[Dict[str, AgentRole]] = None,
    event_bus: Optional[InMemoryEventBus] = None,
    message_store: Optional[EnvelopeStoreProtocol] = None,
    state_factory: Optional[Callable[[AgentProfile], StateManager]] = None,
    runner_factory: Optional[Callable[[Behaviour, StateManager, RunContext], ReactiveRunner]] = None,
    tool_providers: Optional[list[ToolProvider]] = None,
    default_model_client: Any = None,
    auto_start: bool = False,
    session_id: Optional[str] = None,
    mcp_config_path: Optional[str] = None,
    skills_dir: Optional[str] = None,
    mailbox_drain_max_iterations: int = 100,
    message_store_gc_min_advance: int = 1000,
    delivery_max_failures: int = 3,
    restart_policy: str = "on_failure",
    max_restarts: int = 3,
    restart_backoff_s: float = 0.2,
    same_run_no_send_retries: int = 1,
    checkpointer: Optional[CheckpointerProtocol] = None,
    checkpoint_thread_id: Optional[str] = None,
    context_builder: Optional[ContextBuilderProtocol] = None,
    runtime_event_bus: Optional[RuntimeEventBus] = None,
):
    """Create a multi-agent runtime.

    Args:
        roles: Optional mapping of role names to role factories.
        event_bus: Optional event bus implementation.
        message_store: Optional message store implementation.
        state_factory: Factory used to build a state manager for each agent.
        runner_factory: Factory used to build a runner for each agent tree.
        tool_providers: Additional tool providers merged after builtin swarm tools.
        default_model_client: Default model client injected into spawned agents.
        auto_start: Whether spawned agents should begin running immediately.
        session_id: Optional session identifier propagated into events and state.
        mcp_config_path: Optional MCP provider configuration path.
        skills_dir: Optional skills directory used to bootstrap prompt configuration.
        mailbox_drain_max_iterations: Max envelopes to drain in one mailbox pass.
        message_store_gc_min_advance: Minimum sequence advance before store pruning.
        delivery_max_failures: Maximum failed deliveries before dead-letter handling.
        restart_policy: Host restart policy: ``never``, ``on_failure``, or ``always``.
        max_restarts: Maximum restart attempts per host.
        restart_backoff_s: Delay before retrying a failed host.
        same_run_no_send_retries: Immediate retries when a run completes without sending.
        checkpointer: Optional runtime checkpoint backend.
        checkpoint_thread_id: Optional checkpoint thread identifier override.
        context_builder: Optional context builder used both when merging
            incoming history and when injecting prompt-time context.
        runtime_event_bus: Optional cross-module runtime event bus used to
            expose swarm lifecycle events alongside single-agent flows.

    Raises:
        ValueError: If numeric guardrails or restart-policy values are invalid.
    """
    if mailbox_drain_max_iterations <= 0:
        raise ValueError("mailbox_drain_max_iterations must be > 0")
    if message_store_gc_min_advance <= 0:
        raise ValueError("message_store_gc_min_advance must be > 0")
    if delivery_max_failures <= 0:
        raise ValueError("delivery_max_failures must be > 0")
    if restart_policy not in {"never", "on_failure", "always"}:
        raise ValueError("restart_policy must be one of: never, on_failure, always")
    if max_restarts < 0:
        raise ValueError("max_restarts must be >= 0")
    if restart_backoff_s < 0:
        raise ValueError("restart_backoff_s must be >= 0")
    if same_run_no_send_retries < 0:
        raise ValueError("same_run_no_send_retries must be >= 0")
    self._roles: Dict[str, AgentRole] = roles or {}
    self._event_bus = event_bus or InMemoryEventBus()
    self._message_store = message_store or InMemoryMessageStore()
    self._state_factory = state_factory or (lambda profile: default_state_factory(DefaultAgentState, profile))
    self._runner_factory = runner_factory or (lambda root, state, ctx: ReactiveRunner(root, state, ctx=ctx))
    self.default_model_client = default_model_client
    self._auto_start = auto_start
    self.session_id = session_id
    self._mailbox_drain_max_iterations = mailbox_drain_max_iterations
    self._message_store_gc_min_advance = message_store_gc_min_advance
    self._last_pruned_seq = 0
    self._delivery_max_failures = delivery_max_failures
    self._restart_policy = restart_policy
    self._max_restarts = max_restarts
    self._restart_backoff_s = restart_backoff_s
    self._same_run_no_send_retries = same_run_no_send_retries
    self._checkpointer = checkpointer
    self._checkpoint_thread_id = checkpoint_thread_id
    self._context_builder = context_builder
    self._runtime_event_bus = runtime_event_bus
    self._mcp_config_path = mcp_config_path
    self._prompt_runtime = build_prompt_runtime(
        skills_dir=skills_dir,
        context_builder=context_builder,
    )

    self._hosts: Dict[str, AgentHost] = {}
    self._attachments: Dict[str, set[str]] = {}
    self._group_members: Dict[str, set[str]] = {}
    self._topic_subscribers: Dict[str, set[str]] = {}
    self._pending_profiles: Dict[str, AgentProfile] = {}
    self._managed_unsubscribers: list[Callable[[], None]] = []
    self._initialized = False

    # Ordering + merge semantics:
    # - get_all_tools uses last-wins by tool name.
    # - Put builtin first so any following provider can override it.
    providers: list[ToolProvider] = [SwarmToolProvider(self)]
    if mcp_config_path:
        providers.append(MCPToolProvider(config_path=mcp_config_path))
    providers.extend(tool_providers or [])
    self._tool_providers: list[ToolProvider] = providers

event_bus property

event_bus: InMemoryEventBus

Return the runtime event bus.

返回:

类型 描述
InMemoryEventBus

The event bus instance.

message_store property

message_store: EnvelopeStoreProtocol

Return the runtime message store.

返回:

类型 描述
EnvelopeStoreProtocol

The message store instance.

hosts property

hosts: Dict[str, AgentHost]

Return the live host registry.

返回:

类型 描述
Dict[str, AgentHost]

A dictionary mapping agent IDs to AgentHost instances.

initialize_async async

initialize_async() -> None

Initialize runtime-owned providers and resources.

返回:

类型 描述
None

None.

源代码位于: jianmu/swarm/runtime/core.py
async def initialize_async(self) -> None:
    """Initialize runtime-owned providers and resources.

    Returns:
        ``None``.
    """
    if self._initialized:
        return
    for provider in self._tool_providers:
        await provider.initialize()
    self._initialized = True

close_async async

close_async() -> None

Stop active agents and close runtime-owned providers.

返回:

类型 描述
None

None.

源代码位于: jianmu/swarm/runtime/core.py
async def close_async(self) -> None:
    """Stop active agents and close runtime-owned providers.

    Returns:
        ``None``.
    """
    root_ids = self._root_host_ids()
    if root_ids:
        await asyncio.gather(*(self.kill_async(host_id) for host_id in root_ids), return_exceptions=True)
    self._clear_managed_unsubscribers()
    for provider in reversed(self._tool_providers):
        try:
            await provider.close()
        except Exception as exc:
            logger.warning("⚠️ Failed to close tool provider: {}", exc)
    self._initialized = False

save_snapshot async

save_snapshot(
    *, thread_id: str | None = None, step: int | None = None
) -> dict[str, Any]

Serialize hosts and message-store state into a runtime snapshot.

参数:

名称 类型 描述 默认
thread_id str | None

Optional checkpoint thread identifier override.

None
step int | None

Optional explicit checkpoint step.

None

返回:

类型 描述
dict[str, Any]

Serializable runtime snapshot payload.

源代码位于: jianmu/swarm/runtime/core.py
async def save_snapshot(
    self,
    *,
    thread_id: str | None = None,
    step: int | None = None,
) -> dict[str, Any]:
    """Serialize hosts and message-store state into a runtime snapshot.

    Args:
        thread_id: Optional checkpoint thread identifier override.
        step: Optional explicit checkpoint step.

    Returns:
        Serializable runtime snapshot payload.
    """
    return await snapshot_impl.save_snapshot(self, thread_id=thread_id, step=step)

restore_snapshot async

restore_snapshot(*, thread_id: str | None = None) -> bool

Restore runtime hosts and message store from a saved snapshot.

参数:

名称 类型 描述 默认
thread_id str | None

Optional checkpoint thread identifier override.

None

返回:

类型 描述
bool

True when a snapshot is found and restored, otherwise False.

源代码位于: jianmu/swarm/runtime/core.py
async def restore_snapshot(self, *, thread_id: str | None = None) -> bool:
    """Restore runtime hosts and message store from a saved snapshot.

    Args:
        thread_id: Optional checkpoint thread identifier override.

    Returns:
        ``True`` when a snapshot is found and restored, otherwise ``False``.
    """
    if self._checkpointer is None:
        return False
    target_thread = thread_id or self._checkpoint_thread_id or self.session_id or "runtime"
    checkpoint = self._checkpointer.get_checkpoint(target_thread)
    if not checkpoint:
        return False
    snapshot = checkpoint.get("state") if isinstance(checkpoint, dict) else None
    if snapshot is None and isinstance(checkpoint, dict):
        snapshot = checkpoint
    if not isinstance(snapshot, dict):
        return False

    await snapshot_impl.kill_root_hosts(self)
    await snapshot_impl.restore_message_store_snapshot(self, snapshot)

    hosts_payload = list(snapshot.get("hosts") or [])
    self._hosts.clear()

    for item in hosts_payload:
        host = snapshot_impl.rebuild_host_from_snapshot(self, item)
        self._hosts[host.profile.id] = host

    snapshot_impl.restore_routing_indexes(self)
    await snapshot_impl.restore_auto_tasks(self)
    snapshot_impl.emit_restored(self)
    return True

register_tool_provider

register_tool_provider(provider: ToolProvider) -> None

Register a tool provider before runtime initialization.

参数:

名称 类型 描述 默认
provider ToolProvider

Provider appended after the builtin swarm provider chain.

必需

引发:

类型 描述
RuntimeError

If called after the runtime has already initialized.

源代码位于: jianmu/swarm/runtime/core.py
def register_tool_provider(self, provider: ToolProvider) -> None:
    """Register a tool provider before runtime initialization.

    Args:
        provider: Provider appended after the builtin swarm provider chain.

    Raises:
        RuntimeError: If called after the runtime has already initialized.
    """
    tools_impl.register_tool_provider(self, provider)

register_role

register_role(role: AgentRole) -> None

Register a named role factory.

参数:

名称 类型 描述 默认
role AgentRole

Role object keyed by its name attribute.

必需
源代码位于: jianmu/swarm/runtime/core.py
def register_role(self, role: AgentRole) -> None:
    """Register a named role factory.

    Args:
        role: Role object keyed by its ``name`` attribute.
    """
    self._roles[role.name] = role

spawn

spawn(
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
    agent_id: str | None = None,
) -> str

Create a host for a role and return its agent identifier.

参数:

名称 类型 描述 默认
role str | AgentRole

Registered role name or role object.

必需
task str | None

Optional initial task for the agent.

None
parent_id str | None

Optional parent host identifier for attached agents.

None
mode str

Attachment mode, typically detached or attached.

'detached'
constraints Any | None

Optional execution constraints applied to the new host.

None
agent_id str | None

Optional stable identifier to reuse while reconstructing a checkpointed host.

None

返回:

类型 描述
str

Newly created agent identifier.

引发:

类型 描述
ValueError

If the role builds an invalid root tree.

源代码位于: jianmu/swarm/runtime/core.py
def spawn(
    self,
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
    agent_id: str | None = None,
) -> str:
    """Create a host for a role and return its agent identifier.

    Args:
        role: Registered role name or role object.
        task: Optional initial task for the agent.
        parent_id: Optional parent host identifier for attached agents.
        mode: Attachment mode, typically ``detached`` or ``attached``.
        constraints: Optional execution constraints applied to the new host.
        agent_id: Optional stable identifier to reuse while reconstructing
            a checkpointed host.

    Returns:
        Newly created agent identifier.

    Raises:
        ValueError: If the role builds an invalid root tree.
    """
    role_obj = self._resolve_role(role)
    agent_id = str(agent_id or "").strip() or str(uuid.uuid4())
    if agent_id in self._hosts or agent_id in self._pending_profiles:
        raise ValueError(f"Agent id already exists: {agent_id}")
    constraints_obj = coerce_constraints(constraints)
    metadata: Dict[str, Any] = {}

    serialized_constraints = constraints_payload(constraints, constraints_obj)
    if serialized_constraints:
        metadata["constraints"] = serialized_constraints

    profile = AgentProfile(
        id=agent_id,
        role=role_obj.name,
        parent_id=parent_id,
        task=task,
        budget=budget_spec_from_constraints(constraints_obj),
        metadata=metadata,
    )

    self._pending_profiles[agent_id] = profile
    try:
        root = role_obj.build_tree(profile)
        if not isinstance(root, Behaviour):
            raise ValueError("AgentRole.build_tree() must return a py_trees Behaviour root")
    finally:
        self._pending_profiles.pop(agent_id, None)

    state = self._state_factory(profile)
    initial_state = {"agent_id": agent_id, "role": role_obj.name, "task": task}
    if self.session_id:
        initial_state["session_id"] = self.session_id
    state.initialize(initial_state)

    ctx = RunContext(
        model_client=self.default_model_client,
        runtime=self,
        constraints=constraints_obj,
        sandbox=constraints_obj.execution if constraints_obj else None,
        prompt_runtime=self._prompt_runtime,
        runtime_event_bus=self._runtime_event_bus,
    )
    runner = self._runner_factory(root, state, ctx)

    host = AgentHost(profile=profile, role=role_obj, runner=runner, state=state)
    if mode == "attached" and parent_id:
        host.attached_to = parent_id
        self._attachments.setdefault(parent_id, set()).add(agent_id)

    self._hosts[agent_id] = host
    self.emit(AgentEvent(event_type="agent_created", agent_id=agent_id, payload={"role": role_obj.name}))

    if self._auto_start:
        self._start_agent(agent_id)
    return agent_id

kill

kill(agent_id: str) -> None

Stop a host synchronously and schedule cleanup of its tasks.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
源代码位于: jianmu/swarm/runtime/core.py
def kill(self, agent_id: str) -> None:
    """Stop a host synchronously and schedule cleanup of its tasks.

    Args:
        agent_id: Target host identifier.
    """
    host, attached = self._detach_host(agent_id)
    if not host:
        return
    tasks_to_cleanup: list[asyncio.Task] = []
    if host.sync_task and not host.sync_task.done():
        host.sync_task.cancel()
        tasks_to_cleanup.append(host.sync_task)
    if host.task and not host.task.done():
        host.task.cancel()
        tasks_to_cleanup.append(host.task)
    for child in list(attached):
        self.kill(child)
    self._schedule_task_cleanup(tasks_to_cleanup)

    self.emit(AgentEvent(event_type="agent_stopped", agent_id=agent_id))
    self.emit(AgentEvent(event_type="agent_killed", agent_id=agent_id))

kill_async async

kill_async(agent_id: str) -> None

Stop a host asynchronously and await task cancellation.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
源代码位于: jianmu/swarm/runtime/core.py
async def kill_async(self, agent_id: str) -> None:
    """Stop a host asynchronously and await task cancellation.

    Args:
        agent_id: Target host identifier.
    """
    host, attached = self._detach_host(agent_id)
    if not host:
        return

    if attached:
        await asyncio.gather(*(self.kill_async(child_id) for child_id in attached), return_exceptions=True)

    await self._cancel_task(host.sync_task)
    await self._cancel_task(host.task)

    self.emit(AgentEvent(event_type="agent_stopped", agent_id=agent_id))
    self.emit(AgentEvent(event_type="agent_killed", agent_id=agent_id))

pause

pause(agent_id: str) -> None

Pause a host so it no longer consumes scheduler work.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
源代码位于: jianmu/swarm/runtime/core.py
def pause(self, agent_id: str) -> None:
    """Pause a host so it no longer consumes scheduler work.

    Args:
        agent_id: Target host identifier.
    """
    host = self._hosts.get(agent_id)
    if not host:
        return
    host.paused = True
    self.emit(AgentEvent(event_type="agent_paused", agent_id=agent_id))

resume

resume(agent_id: str) -> None

Resume a paused host and schedule new work.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
源代码位于: jianmu/swarm/runtime/core.py
def resume(self, agent_id: str) -> None:
    """Resume a paused host and schedule new work.

    Args:
        agent_id: Target host identifier.
    """
    host = self._hosts.get(agent_id)
    if not host:
        return
    host.paused = False
    self.emit(AgentEvent(event_type="agent_resumed", agent_id=agent_id))
    self.wake(agent_id, reason="resume")

wake

wake(agent_id: str, reason: str | None = None) -> None

Mark a host dirty and ensure its runner task exists.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
reason str | None

Optional observability reason attached to the wake event.

None
源代码位于: jianmu/swarm/runtime/core.py
def wake(self, agent_id: str, reason: str | None = None) -> None:
    """Mark a host dirty and ensure its runner task exists.

    Args:
        agent_id: Target host identifier.
        reason: Optional observability reason attached to the wake event.
    """
    self._mark_dirty(agent_id, force_signal=True)
    self._ensure_runner_task(agent_id)
    if reason:
        self.emit(AgentEvent(event_type="agent_wake", agent_id=agent_id, payload={"reason": reason}))

preempt_agent async

preempt_agent(
    agent_id: str, reason: str | None = None
) -> bool

Interrupt a running host and schedule it for a fresh restart.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
reason str | None

Optional reason emitted in the preemption event.

None

返回:

类型 描述
bool

True if a running host was preempted, otherwise False.

源代码位于: jianmu/swarm/runtime/core.py
async def preempt_agent(self, agent_id: str, reason: str | None = None) -> bool:
    """Interrupt a running host and schedule it for a fresh restart.

    Args:
        agent_id: Target host identifier.
        reason: Optional reason emitted in the preemption event.

    Returns:
        ``True`` if a running host was preempted, otherwise ``False``.
    """
    host = self._hosts.get(agent_id)
    if not host or host.paused:
        return False

    was_running = bool(host.task and not host.task.done())
    if host.sync_task and not host.sync_task.done():
        await self._cancel_task(host.sync_task)
        host.sync_task = None
    if host.task and not host.task.done():
        await self._cancel_task(host.task)
        host.task = None

    if reason:
        self.emit(
            AgentEvent(
                event_type="agent_preempted",
                agent_id=agent_id,
                payload={"reason": reason, "was_running": was_running},
            )
        )

    await self._reset_runner_for_restart(agent_id)
    # The replacement runner owns its single prime step.  A drain may
    # already have prepared fresh incoming context before this preemption;
    # forcing and priming here would consume or clear that batch before the
    # replacement runner can execute it.
    self._ensure_runner_task(agent_id)
    return True

send_message async

send_message(
    *,
    sender_id: str | None,
    content: str,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: Optional[dict[str, Any]] = None,
) -> int

Append a routed message and wake matching recipients.

参数:

名称 类型 描述 默认
sender_id str | None

Originating agent identifier, or None for external senders.

必需
content str

Message body.

必需
to_agent_id str | None

Optional direct recipient identifier.

None
group_id str | None

Optional group recipient identifier.

None
topic str | None

Optional topic recipient.

None
content_type str

Message content-type label stored in metadata.

'text'
metadata Optional[dict[str, Any]]

Optional additional metadata persisted with the message.

None

返回:

类型 描述
int

Monotonic message sequence number assigned by the store.

引发:

类型 描述
ValueError

If zero or multiple routing targets are supplied.

源代码位于: jianmu/swarm/runtime/core.py
async def send_message(
    self,
    *,
    sender_id: str | None,
    content: str,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: Optional[dict[str, Any]] = None,
) -> int:
    """Append a routed message and wake matching recipients.

    Args:
        sender_id: Originating agent identifier, or ``None`` for external senders.
        content: Message body.
        to_agent_id: Optional direct recipient identifier.
        group_id: Optional group recipient identifier.
        topic: Optional topic recipient.
        content_type: Message content-type label stored in metadata.
        metadata: Optional additional metadata persisted with the message.

    Returns:
        Monotonic message sequence number assigned by the store.

    Raises:
        ValueError: If zero or multiple routing targets are supplied.
    """
    return await routing_impl.send_message(
        self,
        sender_id=sender_id,
        content=content,
        to_agent_id=to_agent_id,
        group_id=group_id,
        topic=topic,
        content_type=content_type,
        metadata=metadata,
    )

enqueue_internal_message

enqueue_internal_message(
    *,
    agent_id: str,
    message: Message,
    source: str | None = None,
    wake: bool = False,
) -> bool

Deliver a runtime-local message through the mailbox handoff.

Unlike direct host.inbox mutation, this always advances the mailbox version and marks the next execution input as pending.

参数:

名称 类型 描述 默认
agent_id str

Identifier of the local target host.

必需
message Message

Message to append to the target host's mailbox.

必需
source str | None

Optional provenance label stored in message metadata.

None
wake bool

Whether to schedule and signal immediate auto-mode handling.

False

返回:

类型 描述
bool

True when the target host exists and accepts the message.

源代码位于: jianmu/swarm/runtime/core.py
def enqueue_internal_message(
    self,
    *,
    agent_id: str,
    message: Message,
    source: str | None = None,
    wake: bool = False,
) -> bool:
    """Deliver a runtime-local message through the mailbox handoff.

    Unlike direct ``host.inbox`` mutation, this always advances the
    mailbox version and marks the next execution input as pending.

    Args:
        agent_id: Identifier of the local target host.
        message: Message to append to the target host's mailbox.
        source: Optional provenance label stored in message metadata.
        wake: Whether to schedule and signal immediate auto-mode handling.

    Returns:
        ``True`` when the target host exists and accepts the message.
    """
    return mailbox_impl.enqueue_internal_message(
        self,
        agent_id,
        message,
        source=source,
        wake=wake,
    )

step_agent async

step_agent(
    agent_id: str, obs: Optional[dict[str, Any]] = None
) -> Status

Advance one agent in manual mode.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
obs Optional[dict[str, Any]]

Optional observation payload injected before the step.

None

返回:

类型 描述
Status

Root status after the step completes.

源代码位于: jianmu/swarm/runtime/core.py
async def step_agent(self, agent_id: str, obs: Optional[dict[str, Any]] = None) -> Status:
    """Advance one agent in manual mode.

    Args:
        agent_id: Target host identifier.
        obs: Optional observation payload injected before the step.

    Returns:
        Root status after the step completes.
    """
    return await scheduler_impl.step_agent(self, agent_id, obs=obs)

step_agent_with_actions async

step_agent_with_actions(
    agent_id: str, obs: Optional[dict[str, Any]] = None
) -> tuple[Status, dict[str, Any]]

Advance one agent and return both status and emitted actions.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需
obs Optional[dict[str, Any]]

Optional observation payload injected before the step.

None

返回:

类型 描述
tuple[Status, dict[str, Any]]

Tuple of root status and the action mapping returned by the runner.

引发:

类型 描述
RuntimeError

If the target host is in auto mode.

源代码位于: jianmu/swarm/runtime/core.py
async def step_agent_with_actions(
    self,
    agent_id: str,
    obs: Optional[dict[str, Any]] = None,
) -> tuple[Status, dict[str, Any]]:
    """Advance one agent and return both status and emitted actions.

    Args:
        agent_id: Target host identifier.
        obs: Optional observation payload injected before the step.

    Returns:
        Tuple of root status and the action mapping returned by the runner.

    Raises:
        RuntimeError: If the target host is in auto mode.
    """
    return await scheduler_impl.step_agent_with_actions(self, agent_id, obs=obs)

step_all async

step_all() -> dict[str, Status]

Advance all hosts sequentially in manual mode.

返回:

类型 描述
dict[str, Status]

Mapping from agent identifier to root status.

源代码位于: jianmu/swarm/runtime/core.py
async def step_all(self) -> dict[str, Status]:
    """Advance all hosts sequentially in manual mode.

    Returns:
        Mapping from agent identifier to root status.
    """
    return await self._step_all()

step_all_concurrent async

step_all_concurrent(
    *, max_concurrency: int | None = None
) -> dict[str, Status]

Advance all hosts concurrently in manual mode.

参数:

名称 类型 描述 默认
max_concurrency int | None

Optional cap for concurrent stepping tasks.

None

返回:

类型 描述
dict[str, Status]

Mapping from agent identifier to root status.

引发:

类型 描述
ValueError

If max_concurrency is provided but not positive.

源代码位于: jianmu/swarm/runtime/core.py
async def step_all_concurrent(self, *, max_concurrency: int | None = None) -> dict[str, Status]:
    """Advance all hosts concurrently in manual mode.

    Args:
        max_concurrency: Optional cap for concurrent stepping tasks.

    Returns:
        Mapping from agent identifier to root status.

    Raises:
        ValueError: If ``max_concurrency`` is provided but not positive.
    """
    return await scheduler_impl.step_all_concurrent(self, max_concurrency=max_concurrency)

run_until_idle async

run_until_idle(max_steps: int = 200) -> dict[str, Status]

Advance hosts until the runtime becomes idle.

参数:

名称 类型 描述 默认
max_steps int

Maximum scheduler iterations before returning.

200

返回:

类型 描述
dict[str, Status]

Last per-agent root-status mapping observed during the run.

引发:

类型 描述
RuntimeError

If any host is running in auto mode.

源代码位于: jianmu/swarm/runtime/core.py
async def run_until_idle(self, max_steps: int = 200) -> dict[str, Status]:
    """Advance hosts until the runtime becomes idle.

    Args:
        max_steps: Maximum scheduler iterations before returning.

    Returns:
        Last per-agent root-status mapping observed during the run.

    Raises:
        RuntimeError: If any host is running in auto mode.
    """
    return await scheduler_impl.run_until_idle(self, max_steps=max_steps)

join_group

join_group(agent_id: str, group_id: str) -> None

Subscribe a host to a broadcast group.

参数:

名称 类型 描述 默认
agent_id str

The ID of the agent/host joining the group.

必需
group_id str

The ID of the broadcast group to join.

必需
源代码位于: jianmu/swarm/runtime/core.py
def join_group(self, agent_id: str, group_id: str) -> None:
    """Subscribe a host to a broadcast group.

    Args:
        agent_id: The ID of the agent/host joining the group.
        group_id: The ID of the broadcast group to join.
    """
    routing_impl.join_group(self, agent_id, group_id)

leave_group

leave_group(agent_id: str, group_id: str) -> None

Remove a host from a broadcast group.

参数:

名称 类型 描述 默认
agent_id str

The ID of the agent/host leaving the group.

必需
group_id str

The ID of the broadcast group to leave.

必需
源代码位于: jianmu/swarm/runtime/core.py
def leave_group(self, agent_id: str, group_id: str) -> None:
    """Remove a host from a broadcast group.

    Args:
        agent_id: The ID of the agent/host leaving the group.
        group_id: The ID of the broadcast group to leave.
    """
    routing_impl.leave_group(self, agent_id, group_id)

subscribe_topic

subscribe_topic(agent_id: str, topic: str) -> None

Subscribe a host to a topic.

参数:

名称 类型 描述 默认
agent_id str

The ID of the agent/host subscribing.

必需
topic str

The topic name to subscribe to.

必需
源代码位于: jianmu/swarm/runtime/core.py
def subscribe_topic(self, agent_id: str, topic: str) -> None:
    """Subscribe a host to a topic.

    Args:
        agent_id: The ID of the agent/host subscribing.
        topic: The topic name to subscribe to.
    """
    routing_impl.subscribe_topic(self, agent_id, topic)

unsubscribe_topic

unsubscribe_topic(agent_id: str, topic: str) -> None

Unsubscribe a host from a topic.

参数:

名称 类型 描述 默认
agent_id str

The ID of the agent/host unsubscribing.

必需
topic str

The topic name to unsubscribe from.

必需
源代码位于: jianmu/swarm/runtime/core.py
def unsubscribe_topic(self, agent_id: str, topic: str) -> None:
    """Unsubscribe a host from a topic.

    Args:
        agent_id: The ID of the agent/host unsubscribing.
        topic: The topic name to unsubscribe from.
    """
    routing_impl.unsubscribe_topic(self, agent_id, topic)

emit

emit(event: AgentEvent) -> None

Normalize and publish one runtime event.

参数:

名称 类型 描述 默认
event AgentEvent

The event to publish.

必需
源代码位于: jianmu/swarm/runtime/core.py
def emit(self, event: AgentEvent) -> None:
    """Normalize and publish one runtime event.

    Args:
        event: The event to publish.
    """
    payload = dict(event.payload or {})
    payload.setdefault("session_id", self.session_id)
    payload.setdefault("agent_id", event.agent_id)
    event.payload = payload
    self._event_bus.emit(event)
    self._emit_runtime_bridge_event(event)

subscribe_events

subscribe_events(
    callback: Callable[[AgentEvent], None],
) -> Callable[[], None]

Subscribe to normalized runtime events.

参数:

名称 类型 描述 默认
callback Callable[[AgentEvent], None]

Consumer invoked for each emitted AgentEvent.

必需

返回:

类型 描述
Callable[[], None]

Unsubscribe callback.

源代码位于: jianmu/swarm/runtime/core.py
def subscribe_events(self, callback: Callable[[AgentEvent], None]) -> Callable[[], None]:
    """Subscribe to normalized runtime events.

    Args:
        callback: Consumer invoked for each emitted ``AgentEvent``.

    Returns:
        Unsubscribe callback.
    """
    return self._track_unsubscriber(self._event_bus.subscribe(callback))

subscribe_observability

subscribe_observability(
    callback: Callable[[str, dict[str, Any]], None],
    *,
    runtime_events: bool = True,
    trace_events: bool = False,
) -> Callable[[], None]

Subscribe to runtime and trace observability streams.

参数:

名称 类型 描述 默认
callback Callable[[str, dict[str, Any]], None]

Consumer invoked as callback(event_name, payload).

必需
runtime_events bool

Whether normalized runtime events should be forwarded.

True
trace_events bool

Whether low-level trace events should be forwarded.

False

返回:

类型 描述
Callable[[], None]

Unsubscribe callback that detaches all registered listeners.

源代码位于: jianmu/swarm/runtime/core.py
def subscribe_observability(
    self,
    callback: Callable[[str, dict[str, Any]], None],
    *,
    runtime_events: bool = True,
    trace_events: bool = False,
) -> Callable[[], None]:
    """Subscribe to runtime and trace observability streams.

    Args:
        callback: Consumer invoked as ``callback(event_name, payload)``.
        runtime_events: Whether normalized runtime events should be forwarded.
        trace_events: Whether low-level trace events should be forwarded.

    Returns:
        Unsubscribe callback that detaches all registered listeners.
    """
    unsubscribers: list[Callable[[], None]] = []

    if runtime_events:
        def _runtime_cb(event: AgentEvent) -> None:
            """Forward runtime events into the observability callback."""
            payload = dict(event.payload or {})
            payload.setdefault("session_id", self.session_id)
            payload.setdefault("agent_id", event.agent_id)
            payload["source"] = "runtime"
            callback(event.event_type, payload)

        unsubscribers.append(self._event_bus.subscribe(_runtime_cb))

    if trace_events:
        def _trace_cb(event) -> None:
            """Forward trace events into the observability callback."""
            body = dict(event.payload or {})
            payload = {
                "event": event.name,
                "name": event.name,
                "type": "trace",
                "ts": event.ts,
                "trace_id": event.trace_id,
                "span_id": event.span_id,
                "payload": body,
                **body,
            }
            payload.setdefault("session_id", self.session_id)
            payload["source"] = "trace"
            callback(event.name, payload)

        subscribed_trace_hub = None

        def _attach_trace_hub(next_hub) -> None:
            nonlocal subscribed_trace_hub
            if next_hub is subscribed_trace_hub:
                return
            if subscribed_trace_hub is not None:
                try:
                    subscribed_trace_hub.unsubscribe(_trace_cb)
                except Exception:
                    pass
            next_hub.subscribe(_trace_cb)
            subscribed_trace_hub = next_hub

        _attach_trace_hub(get_default_hub())
        unsubscribe_hub_changes = subscribe_default_hub_changes(
            lambda _previous_hub, next_hub: _attach_trace_hub(next_hub)
        )

        def _unsub_trace() -> None:
            """Detach the temporary trace subscriber."""
            try:
                unsubscribe_hub_changes()
            finally:
                if subscribed_trace_hub is not None:
                    try:
                        subscribed_trace_hub.unsubscribe(_trace_cb)
                    except Exception:
                        pass

        unsubscribers.append(_unsub_trace)

    def _unsubscribe_all() -> None:
        """Detach all observability subscribers created by this call."""
        for unsub in unsubscribers:
            try:
                unsub()
            except Exception:
                continue

    return self._track_unsubscriber(_unsubscribe_all)

get_profile

get_profile(agent_id: str) -> Optional[AgentProfile]

Return the profile for a host, if it exists.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需

返回:

类型 描述
Optional[AgentProfile]

Agent profile, or None if the host is unknown.

源代码位于: jianmu/swarm/runtime/core.py
def get_profile(self, agent_id: str) -> Optional[AgentProfile]:
    """Return the profile for a host, if it exists.

    Args:
        agent_id: Target host identifier.

    Returns:
        Agent profile, or ``None`` if the host is unknown.
    """
    host = self._hosts.get(agent_id)
    return host.profile if host else None

list_agents

list_agents() -> list[AgentProfile]

Return profiles for all active hosts.

返回:

类型 描述
list[AgentProfile]

A list of profiles for all active hosts.

源代码位于: jianmu/swarm/runtime/core.py
def list_agents(self) -> list[AgentProfile]:
    """Return profiles for all active hosts.

    Returns:
        A list of profiles for all active hosts.
    """
    return [host.profile for host in self._hosts.values()]

get_host

get_host(agent_id: str) -> Optional[AgentHost]

Return the live host object for one agent.

参数:

名称 类型 描述 默认
agent_id str

The ID of the target agent.

必需

返回:

类型 描述
Optional[AgentHost]

The live host object, or None if not found.

源代码位于: jianmu/swarm/runtime/core.py
def get_host(self, agent_id: str) -> Optional[AgentHost]:
    """Return the live host object for one agent.

    Args:
        agent_id: The ID of the target agent.

    Returns:
        The live host object, or None if not found.
    """
    return self._hosts.get(agent_id)

resolve_suspension async

resolve_suspension(
    *,
    agent_id: str,
    thread_id: str,
    request_id: str,
    payload: dict[str, Any],
    suspension_mode: str
    | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpoint_interval: int = 1,
    max_fps: float = 60.0,
) -> ResumeInteractionResult

Resolve one host-owned suspension through the runtime host registry.

This facade only serves runtime-host scenarios where the runtime already owns the agent_id -> host -> runner relationship. It does not replace standalone ReactiveRunner.resume_and_continue() or Agent.resume_interaction() entrypoints.

参数:

名称 类型 描述 默认
agent_id str

Runtime host identifier used to locate the target runner.

必需
thread_id str

Checkpoint thread identifier routed to the target runner.

必需
request_id str

Active suspension request identifier to resume.

必需
payload dict[str, Any]

Host-provided resume payload.

必需
suspension_mode str | SuspensionMode

Host-facing mode for any later suspension.

YIELD
reset_tree bool

Whether to reset tree execution before resuming.

True
reset_data bool

Whether to reset the backing state first.

False
max_ticks int | None

Optional hard tick limit.

None
timeout_s float | None

Optional wall-clock timeout for the resumed run.

None
checkpoint_interval int

Checkpoint save interval for the resumed run.

1
max_fps float

Runner FPS cap for the resumed run.

60.0

返回:

类型 描述
ResumeInteractionResult

Structured resume status plus resumed run result when available.

源代码位于: jianmu/swarm/runtime/core.py
async def resolve_suspension(
    self,
    *,
    agent_id: str,
    thread_id: str,
    request_id: str,
    payload: dict[str, Any],
    suspension_mode: str | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpoint_interval: int = 1,
    max_fps: float = 60.0,
) -> ResumeInteractionResult:
    """Resolve one host-owned suspension through the runtime host registry.

    This facade only serves runtime-host scenarios where the runtime
    already owns the ``agent_id -> host -> runner`` relationship. It does
    not replace standalone ``ReactiveRunner.resume_and_continue()`` or
    ``Agent.resume_interaction()`` entrypoints.

    Args:
        agent_id: Runtime host identifier used to locate the target runner.
        thread_id: Checkpoint thread identifier routed to the target runner.
        request_id: Active suspension request identifier to resume.
        payload: Host-provided resume payload.
        suspension_mode: Host-facing mode for any later suspension.
        reset_tree: Whether to reset tree execution before resuming.
        reset_data: Whether to reset the backing state first.
        max_ticks: Optional hard tick limit.
        timeout_s: Optional wall-clock timeout for the resumed run.
        checkpoint_interval: Checkpoint save interval for the resumed run.
        max_fps: Runner FPS cap for the resumed run.

    Returns:
        Structured resume status plus resumed run result when available.
    """
    host = self.get_host(agent_id)
    if host is None:
        return ResumeInteractionResult(
            resume=ResumeSuspensionResult(
                status=ResumeStatus.INVALID_REQUEST,
                request_id=str(request_id or "").strip(),
                reason="unknown_agent",
                message=f"No runtime host matched agent_id={agent_id!r}.",
                payload={"agent_id": agent_id, "thread_id": thread_id},
                retryable=False,
            ),
            run=None,
        )
    try:
        run_result = await host.runner.resume_and_continue(
            thread_id=thread_id,
            request_id=request_id,
            payload=payload,
            checkpointer=self._checkpointer,
            suspension_mode=suspension_mode,
            reset_tree=reset_tree,
            reset_data=reset_data,
            max_ticks=max_ticks,
            timeout_s=timeout_s,
            checkpoint_interval=checkpoint_interval,
            max_fps=max_fps,
        )
    except Exception as exc:
        result = getattr(exc, "result", None)
        if isinstance(result, ResumeSuspensionResult):
            return ResumeInteractionResult(resume=result, run=None)
        raise
    return ResumeInteractionResult(
        resume=ResumeSuspensionResult(
            status=ResumeStatus.SUCCESS,
            request_id=request_id,
            message="Resume request accepted.",
            payload=dict(payload or {}),
            retryable=False,
        ),
        run=run_result,
    )

get_dead_letters async

get_dead_letters(
    *, agent_id: str | None = None, limit: int = 100
) -> list[dict[str, Any]]

Return dead-letter payloads from the message store.

参数:

名称 类型 描述 默认
agent_id str | None

Optional filter for one agent's dead letters.

None
limit int

Maximum number of dead-letter records to return.

100

返回:

类型 描述
list[dict[str, Any]]

Dead-letter payloads, or an empty list when unsupported.

源代码位于: jianmu/swarm/runtime/core.py
async def get_dead_letters(self, *, agent_id: str | None = None, limit: int = 100) -> list[dict[str, Any]]:
    """Return dead-letter payloads from the message store.

    Args:
        agent_id: Optional filter for one agent's dead letters.
        limit: Maximum number of dead-letter records to return.

    Returns:
        Dead-letter payloads, or an empty list when unsupported.
    """
    get_dlq = getattr(self._message_store, "get_dlq", None)
    if not callable(get_dlq):
        return []
    try:
        return await get_dlq(agent_id=agent_id, limit=limit)
    except Exception:
        return []

get_all_tools

get_all_tools(agent_id: str) -> list

Return the merged tool list available to a host.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需

返回:

类型 描述
list

Effective merged tool list.

源代码位于: jianmu/swarm/runtime/core.py
def get_all_tools(self, agent_id: str) -> list:
    """Return the merged tool list available to a host.

    Args:
        agent_id: Target host identifier.

    Returns:
        Effective merged tool list.
    """
    return tools_impl.get_all_tools(self, agent_id)

explain_tools

explain_tools(agent_id: str) -> dict[str, Any]

Return diagnostics for tool resolution on a host.

参数:

名称 类型 描述 默认
agent_id str

Target host identifier.

必需

返回:

类型 描述
dict[str, Any]

Diagnostic payload describing providers, overrides, and effective tools.

源代码位于: jianmu/swarm/runtime/core.py
def explain_tools(self, agent_id: str) -> dict[str, Any]:
    """Return diagnostics for tool resolution on a host.

    Args:
        agent_id: Target host identifier.

    Returns:
        Diagnostic payload describing providers, overrides, and effective tools.
    """
    return tools_impl.explain_tools(self, agent_id)

start_all

start_all() -> None

Start all registered hosts.

源代码位于: jianmu/swarm/runtime/core.py
def start_all(self) -> None:
    """Start all registered hosts."""
    for agent_id in list(self._hosts.keys()):
        self._start_agent(agent_id)

stop_all

stop_all() -> None

Stop all registered hosts.

源代码位于: jianmu/swarm/runtime/core.py
def stop_all(self) -> None:
    """Stop all registered hosts."""
    for agent_id in list(self._hosts.keys()):
        self.kill(agent_id)

AgentTeam

AgentTeam(
    *,
    roles: dict[str, AgentRole] | None = None,
    default_model_client: Any = None,
    tool_providers: Sequence[ToolProvider] | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_thread_id: str | None = None,
    runtime_event_bus: RuntimeEventBus | None = None,
    context_builder: ContextBuilderProtocol | None = None,
    session_id: str | None = None,
    auto_start: bool = False,
)

User-facing facade that wraps one AgentRuntime instance.

AgentTeam centralizes runtime assembly, lifecycle management, and snapshot entrypoints for Jianmu's multi-agent runtime without replacing the underlying swarm engine.

属性:

名称 类型 描述
_defaults

Stable defaults reused when spawning or configuring hosts.

_runtime

Wrapped low-level swarm runtime instance.

_checkpointer

Resolved checkpoint backend shared with the runtime.

_checkpoint_thread_id

Default checkpoint thread identifier.

Create a high-level facade around one assembled AgentRuntime.

参数:

名称 类型 描述 默认
roles dict[str, AgentRole] | None

Optional mapping of role names to role objects.

None
default_model_client Any

Optional default model provider shared by hosts.

None
tool_providers Sequence[ToolProvider] | None

Optional additional runtime tool providers.

None
checkpointer CheckpointerProtocol | None

Default runtime checkpointer, if any.

None
checkpoint_thread_id str | None

Default checkpoint thread identifier, if any.

None
runtime_event_bus RuntimeEventBus | None

Optional runtime-semantic event bus shared across spawned agents and bridged swarm lifecycle events.

None
context_builder ContextBuilderProtocol | None

Optional context builder shared by mailbox-history merge and prompt-time context assembly.

None
session_id str | None

Optional runtime session identifier.

None
auto_start bool

Whether spawned hosts should begin running automatically.

False
源代码位于: jianmu/swarm/team.py
def __init__(
    self,
    *,
    roles: dict[str, AgentRole] | None = None,
    default_model_client: Any = None,
    tool_providers: Sequence[ToolProvider] | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_thread_id: str | None = None,
    runtime_event_bus: RuntimeEventBus | None = None,
    context_builder: ContextBuilderProtocol | None = None,
    session_id: str | None = None,
    auto_start: bool = False,
):
    """Create a high-level facade around one assembled ``AgentRuntime``.

    Args:
        roles: Optional mapping of role names to role objects.
        default_model_client: Optional default model provider shared by hosts.
        tool_providers: Optional additional runtime tool providers.
        checkpointer: Default runtime checkpointer, if any.
        checkpoint_thread_id: Default checkpoint thread identifier, if any.
        runtime_event_bus: Optional runtime-semantic event bus shared across
            spawned agents and bridged swarm lifecycle events.
        context_builder: Optional context builder shared by mailbox-history
            merge and prompt-time context assembly.
        session_id: Optional runtime session identifier.
        auto_start: Whether spawned hosts should begin running automatically.
    """
    resolved_checkpointer = checkpointer if checkpointer is not None else _resolve_default_checkpointer()
    self._defaults = AgentTeamDefaults(
        default_model_client=default_model_client,
        context_builder=context_builder,
        tool_providers=tuple(tool_providers or ()),
    )
    self._runtime = AgentRuntime(
        roles=roles,
        default_model_client=default_model_client,
        tool_providers=list(tool_providers or []),
        auto_start=auto_start,
        session_id=session_id,
        checkpointer=resolved_checkpointer,
        checkpoint_thread_id=checkpoint_thread_id,
        runtime_event_bus=runtime_event_bus,
        context_builder=context_builder,
    )
    self._checkpointer = resolved_checkpointer
    self._checkpoint_thread_id = checkpoint_thread_id

runtime property

runtime: AgentRuntime

Return the wrapped low-level runtime.

返回:

类型 描述
AgentRuntime

The resulting AgentRuntime value.

defaults property

defaults: AgentTeamDefaults

Return the team-level defaults captured during assembly.

返回:

类型 描述
AgentTeamDefaults

The resulting AgentTeamDefaults value.

start async

start() -> None

Initialize runtime-owned providers and resources.

源代码位于: jianmu/swarm/team.py
async def start(self) -> None:
    """Initialize runtime-owned providers and resources."""
    await self._runtime.initialize_async()

close async

close() -> None

Close runtime-owned providers and stop managed agents.

源代码位于: jianmu/swarm/team.py
async def close(self) -> None:
    """Close runtime-owned providers and stop managed agents."""
    await self._runtime.close_async()

spawn

spawn(
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
    agent_id: str | None = None,
) -> str

Create one runtime host via the wrapped runtime.

参数:

名称 类型 描述 默认
role str | AgentRole

Role name to resolve or execute.

必需
task str | None

The task value.

None
parent_id str | None

Identifier for parent.

None
mode str

The mode value.

'detached'
constraints Any | None

Collection of constraint values.

None
agent_id str | None

Optional stable identifier for checkpointed host rebuilds.

None

返回:

类型 描述
str

The resulting string value.

源代码位于: jianmu/swarm/team.py
def spawn(
    self,
    role: str | AgentRole,
    task: str | None = None,
    parent_id: str | None = None,
    mode: str = "detached",
    constraints: Any | None = None,
    agent_id: str | None = None,
) -> str:
    """Create one runtime host via the wrapped runtime.

    Args:
        role: Role name to resolve or execute.
        task: The `task` value.
        parent_id: Identifier for parent.
        mode: The `mode` value.
        constraints: Collection of constraint values.
        agent_id: Optional stable identifier for checkpointed host rebuilds.

    Returns:
        The resulting string value.
    """
    return self._runtime.spawn(
        role=role,
        task=task,
        parent_id=parent_id,
        mode=mode,
        constraints=constraints,
        agent_id=agent_id,
    )

send async

send(
    *,
    content: str,
    sender_id: str | None = None,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: dict[str, Any] | None = None,
) -> int

Send one message through the wrapped runtime.

参数:

名称 类型 描述 默认
content str

The content value.

必需
sender_id str | None

Identifier for sender.

None
to_agent_id str | None

Identifier for to agent.

None
group_id str | None

Identifier for group.

None
topic str | None

The topic value.

None
content_type str

The content_type value.

'text'
metadata dict[str, Any] | None

The metadata value.

None

返回:

类型 描述
int

The resulting integer value.

源代码位于: jianmu/swarm/team.py
async def send(
    self,
    *,
    content: str,
    sender_id: str | None = None,
    to_agent_id: str | None = None,
    group_id: str | None = None,
    topic: str | None = None,
    content_type: str = "text",
    metadata: dict[str, Any] | None = None,
) -> int:
    """Send one message through the wrapped runtime.

    Args:
        content: The `content` value.
        sender_id: Identifier for sender.
        to_agent_id: Identifier for to agent.
        group_id: Identifier for group.
        topic: The `topic` value.
        content_type: The `content_type` value.
        metadata: The `metadata` value.

    Returns:
        The resulting integer value.
    """
    return await self._runtime.send_message(
        sender_id=sender_id,
        content=content,
        to_agent_id=to_agent_id,
        group_id=group_id,
        topic=topic,
        content_type=content_type,
        metadata=metadata,
    )

save_snapshot async

save_snapshot(
    *, thread_id: str | None = None, step: int | None = None
) -> dict[str, Any]

Persist one runtime snapshot through the wrapped runtime.

参数:

名称 类型 描述 默认
thread_id str | None

Identifier for thread.

None
step int | None

The step value.

None

返回:

类型 描述
dict[str, Any]

The resulting mapping value.

源代码位于: jianmu/swarm/team.py
async def save_snapshot(
    self,
    *,
    thread_id: str | None = None,
    step: int | None = None,
) -> dict[str, Any]:
    """Persist one runtime snapshot through the wrapped runtime.

    Args:
        thread_id: Identifier for thread.
        step: The `step` value.

    Returns:
        The resulting mapping value.
    """
    target_thread = thread_id or self._checkpoint_thread_id
    return await self._runtime.save_snapshot(thread_id=target_thread, step=step)

restore_snapshot async

restore_snapshot(*, thread_id: str | None = None) -> bool

Restore one runtime snapshot through the wrapped runtime.

参数:

名称 类型 描述 默认
thread_id str | None

Identifier for thread.

None

返回:

类型 描述
bool

True if the operation succeeds; otherwise False.

源代码位于: jianmu/swarm/team.py
async def restore_snapshot(self, *, thread_id: str | None = None) -> bool:
    """Restore one runtime snapshot through the wrapped runtime.

    Args:
        thread_id: Identifier for thread.

    Returns:
        True if the operation succeeds; otherwise False.
    """
    target_thread = thread_id or self._checkpoint_thread_id
    return await self._runtime.restore_snapshot(thread_id=target_thread)

AgentTeamDefaults dataclass

AgentTeamDefaults(
    default_model_client: Any = None,
    context_builder: ContextBuilderProtocol | None = None,
    tool_providers: tuple[ToolProvider, ...] = tuple(),
)

Stable team-level defaults propagated into newly spawned hosts.

属性:

名称 类型 描述
default_model_client Any

Default model provider injected into agent hosts.

context_builder ContextBuilderProtocol | None

Optional context builder shared by mailbox-history merge and prompt-time context assembly.

tool_providers tuple[ToolProvider, ...]

Additional tool providers registered for the runtime.

MCPToolProvider

MCPToolProvider(config_path: str = 'mcp.json')

Bases: BaseToolProvider

Load MCP tools from config and expose them as Jianmu tools.

属性:

名称 类型 描述
config_path

Filesystem path to the MCP config file.

_clients dict[str, MCPClient]

Connected MCP clients keyed by configured server name.

_tools list[Tool]

Cached wrapped Tool instances exported by all connected servers.

Create a provider backed by an MCP config file.

源代码位于: jianmu/mcp/provider.py
def __init__(self, config_path: str = "mcp.json"):
    """Create a provider backed by an MCP config file."""
    self.config_path = config_path
    self._clients: dict[str, MCPClient] = {}
    self._tools: list[Tool] = []

initialize async

initialize() -> None

Connect configured MCP clients eagerly.

源代码位于: jianmu/mcp/provider.py
async def initialize(self) -> None:
    """Connect configured MCP clients eagerly."""
    config = self._load_config()
    if not config:
        return
    servers = config.get("mcpServers", {}) if isinstance(config, dict) else {}
    if not isinstance(servers, dict):
        return
    for server_name, server_config in servers.items():
        if not isinstance(server_config, dict):
            continue
        try:
            await self._connect_server(server_name, server_config)
        except Exception as exc:
            logger.warning("⚠️ Failed to connect MCP server {}: {}", server_name, exc)

get_tools

get_tools(**_: Any) -> list[Tool]

Return cached MCP-backed Jianmu tools.

参数:

名称 类型 描述 默认
_ Any

Unused keyword arguments.

{}

返回:

类型 描述
list[Tool]

A list of cached Tool instances.

源代码位于: jianmu/mcp/provider.py
def get_tools(self, **_: Any) -> list[Tool]:
    """Return cached MCP-backed Jianmu tools.

    Args:
        _: Unused keyword arguments.

    Returns:
        A list of cached Tool instances.
    """
    return list(self._tools)

close async

close() -> None

Close all managed MCP clients.

源代码位于: jianmu/mcp/provider.py
async def close(self) -> None:
    """Close all managed MCP clients."""
    for client in self._clients.values():
        try:
            await client.close()
        except Exception as exc:
            logger.warning("⚠️ Failed to close MCP client: {}", exc)
    self._clients.clear()
    self._tools.clear()

SwarmRole

SwarmRole(
    *,
    name: str,
    runtime: AgentRuntimeProtocol,
    system_prompt: str = "You are a helpful assistant.",
    model_client: Optional[ModelClient] = None,
    model_client_factory: Optional[
        Callable[[AgentProfile], ModelClient]
    ] = None,
    tools: Optional[list[Tool]] = None,
    runtime_tool_names: Optional[list[str]] = None,
    max_iterations: int = 15,
    constraints: Optional[Constraints] = None,
    skills_dir: str | Path | None = None,
    enabled_skills: Optional[list[str]] = None,
    explicit_enabled_skills: bool = False,
    skill_files: Optional[list[str | Path]] = None,
    skill_prompt_mode: str = "summary",
)

Ready-to-use role implementation for liquid-topology collaboration.

属性:

名称 类型 描述
name

Public role name used for lookup and runtime diagnostics.

runtime

Swarm runtime used to resolve tools and spawn profiles.

system_prompt

Base system prompt prepended to each task prompt.

model_client

Shared model client used when no factory is provided.

model_client_factory

Optional factory that resolves a model client per agent profile.

tools

Additional static tools attached to each spawned agent.

runtime_tool_names

Optional allowlist for runtime-provided tool names.

max_iterations

Max inner ReAct rounds for the spawned agent tree.

constraints

Optional execution constraints forwarded into the node config.

skills_dir

Optional resolved skill directory path.

enabled_skills

Skill names enabled for the role by default.

skill_files

Additional explicit skill file paths for the role.

skill_prompt_mode

Skill prompt rendering mode for the node config.

Configure a reusable swarm agent role template.

源代码位于: jianmu/swarm/role.py
def __init__(
    self,
    *,
    name: str,
    runtime: AgentRuntimeProtocol,
    system_prompt: str = "You are a helpful assistant.",
    model_client: Optional[ModelClient] = None,
    model_client_factory: Optional[Callable[[AgentProfile], ModelClient]] = None,
    tools: Optional[list[Tool]] = None,
    runtime_tool_names: Optional[list[str]] = None,
    max_iterations: int = 15,
    constraints: Optional[Constraints] = None,
    skills_dir: str | Path | None = None,
    enabled_skills: Optional[list[str]] = None,
    explicit_enabled_skills: bool = False,
    skill_files: Optional[list[str | Path]] = None,
    skill_prompt_mode: str = "summary",
):
    """Configure a reusable swarm agent role template."""
    self.name = name
    self.runtime = runtime
    self.system_prompt = system_prompt
    self.model_client = model_client
    self.model_client_factory = model_client_factory
    self.tools = tools or []
    self.runtime_tool_names = set(runtime_tool_names or []) if runtime_tool_names else None
    # Swarm roles forward this as the inner skill/react agent-loop round cap.
    self.max_iterations = max_iterations
    self.constraints = constraints
    self.skills_dir = str(Path(skills_dir).expanduser().resolve()) if skills_dir else None
    self.enabled_skills = list(enabled_skills or [])
    self.explicit_enabled_skills = bool(explicit_enabled_skills)
    self.skill_files = [str(Path(item).expanduser().resolve()) for item in (skill_files or [])]
    self.skill_prompt_mode = skill_prompt_mode

build_tree

build_tree(profile: AgentProfile) -> Any

Build the swarm-oriented skill tree for one agent profile.

参数:

名称 类型 描述 默认
profile AgentProfile

AgentProfile used to build the tree.

必需

返回:

类型 描述
Any

The configured SwarmNode.

源代码位于: jianmu/swarm/role.py
def build_tree(self, profile: AgentProfile) -> Any:
    """Build the swarm-oriented skill tree for one agent profile.

    Args:
        profile: AgentProfile used to build the tree.

    Returns:
        The configured SwarmNode.
    """
    model_client = self._resolve_model_client(profile)

    # Prefer runtime-provided tool registry; fallback to builtin swarm tools.
    runtime_tools = self._resolve_runtime_tools(profile)
    all_tools = [*runtime_tools, *self.tools]

    prompt = f"{self.system_prompt}\nTask: {profile.task or ''}".strip()
    config = SkillNodeConfig(
        push_to_chat=True,
        use_history=True,
        react_config=ReActConfig(
            system_prompt=prompt,
            max_iterations=self.max_iterations,
        ),
        constraints=self.constraints,
        skill_prompt_mode=self.skill_prompt_mode if self.skill_prompt_mode in {"summary", "full"} else "summary",
    )
    return SwarmNode(
        name=f"{self.name}-{profile.id[:8]}",
        model_client=model_client,
        skills_dir=self.skills_dir,
        enabled_skills=self.enabled_skills,
        explicit_enabled_skills=self.explicit_enabled_skills,
        skill_files=self.skill_files,
        messages=[human(profile.task)] if profile.task else None,
        tools=all_tools,
        config=config,
    )

FunctionRole

FunctionRole(
    *,
    name: str,
    description: str,
    handler: Callable[..., Any],
)

Bases: FunctionalRole

Wrap a plain function as a FunctionalRole with metadata.

Create a functional role with explicit display description.

源代码位于: jianmu/swarm/decorator.py
def __init__(
    self,
    *,
    name: str,
    description: str,
    handler: Callable[..., Any],
):
    """Create a functional role with explicit display description."""
    super().__init__(name=name, handler=handler)
    self.description = description

FunctionalRole

FunctionalRole(
    *,
    name: str,
    handler: Callable[[AgentProfile, Any, Any], Any],
)

Role adapter: define agent logic via a plain callable.

Create a role wrapper around a plain handler function.

源代码位于: jianmu/swarm/decorator.py
def __init__(self, *, name: str, handler: Callable[[AgentProfile, Any, Any], Any]):
    """Create a role wrapper around a plain handler function."""
    self.name = name
    self._handler = handler

build_tree

build_tree(profile: AgentProfile) -> Any

Build a single functional node for one agent profile.

参数:

名称 类型 描述 默认
profile AgentProfile

AgentProfile instance to build the tree for.

必需

返回:

类型 描述
Any

An instance of _FunctionalNode.

源代码位于: jianmu/swarm/decorator.py
def build_tree(self, profile: AgentProfile) -> Any:
    """Build a single functional node for one agent profile.

    Args:
        profile: AgentProfile instance to build the tree for.

    Returns:
        An instance of _FunctionalNode.
    """
    return _FunctionalNode(
        name=f"{self.name}-{profile.id[:8]}",
        profile=profile,
        handler=self._handler,
    )