Skip to content

Delegation

Agents can delegate work to other agents over RabbitMQ. A DelegationNode sends a DelegationRequest to a target queue and waits for a DelegationResponse, correlated by ID.

Delegating from an agent

from sofias_sdk_lite import RabbitMQDelegationTransport, RabbitMQClient, RabbitMQConfig
from sofias_sdk_lite.nodes import DelegationNodeConfig

async with RabbitMQClient(RabbitMQConfig(host="localhost")) as client:
    transport = RabbitMQDelegationTransport(client)
    await transport.start()

    builder.with_delegation_transport(transport)
    builder.add_delegation_node(
        "translate",
        DelegationNodeConfig(target_queue="translator-agent-tasks", timeout_seconds=30),
        contract,
    )

DelegationTransport is a protocol (sofias_sdk_lite.nodes.delegation_transport) — RabbitMQDelegationTransport is the RabbitMQ implementation, NullDelegationTransport is an in-memory stand-in for tests.

Responding to a delegation request

On the receiving side, task.reply_to and task.correlation_id will be set. Use BaseDelegationResponseWorkflow instead of ChatResponseWorkflow so the reply goes back over the delegation transport, not a chat stream:

from sofias_sdk_lite import BaseDelegationResponseWorkflow

def create_response_workflow(self, client, task):
    if task.reply_to:
        async def publish(payload: dict) -> None:
            await client.publish_reply(task.reply_to, payload, correlation_id=task.correlation_id)
        return BaseDelegationResponseWorkflow(publish_fn=publish)
    return super().create_response_workflow(client, task)

Override AgentRunner.create_response_workflow() with logic like the above in your runner subclass.

Fan-out and aggregation

To delegate to several agents in parallel and merge results, use FanOutStrategy to route to multiple delegation nodes, and an AggregatorNode with a resolution policy (all, any, majority) to collect responses:

from sofias_sdk_lite import FanOutStrategy
from sofias_sdk_lite.nodes import AggregatorNodeConfig

builder.add_route("dispatch", FanOutStrategy(targets=["translate", "summarize"], join_node="merge"))
builder.add_aggregator_node("merge", AggregatorNodeConfig(resolution_policy="all", timeout_seconds=60), contract)