Implementation:Microsoft Autogen ChatAgentContainer
| Key | Value |
|---|---|
| id | Microsoft_Autogen_ChatAgentContainer |
| source | Microsoft_Autogen |
| category | Group Chat |
Overview
Description
The ChatAgentContainer is a core agent wrapper class that bridges individual ChatAgent or Team instances with the group chat event system. It extends SequentialRoutedAgent to enable ChatAgent and Team instances to participate in distributed group chat scenarios managed by orchestrators.
The container acts as an adapter that:
- Message Buffering: Maintains a buffer of chat messages from the group conversation
- Event Handling: Processes group chat events (start, request, reset, pause, resume)
- Delegation: Forwards buffered messages to the wrapped agent/team for processing
- Response Publishing: Publishes agent responses or team results back to the group chat
- Error Management: Catches and publishes errors as GroupChatError events
- State Management: Saves and restores both agent state and message buffer
- Message Logging: Publishes intermediate messages and events to output topic
The container ensures that agents process messages sequentially using a FIFO lock, preventing race conditions in distributed group chat scenarios.
Usage
The ChatAgentContainer is used internally by group chat teams (RoundRobinGroupChat, SelectorGroupChat, MagenticOneGroupChat) to wrap participant agents. It is not typically instantiated directly by users.
Key responsibilities:
- Buffer messages from group chat events (GroupChatStart, GroupChatAgentResponse, GroupChatTeamResponse)
- Process GroupChatRequestPublish events by invoking the wrapped agent/team
- Handle lifecycle events (reset, pause, resume)
- Maintain sequential processing guarantees for group chat messages
- Publish responses and errors back to the orchestrator
Code Reference
Source Location
- Repository: https://github.com/microsoft/autogen
- File Path: /tmp/kapso_repo_2mr4n2g4/python/packages/autogen-agentchat/src/autogen_agentchat/teams/_group_chat/_chat_agent_container.py
- Lines: 1-214
Signature
class ChatAgentContainer(SequentialRoutedAgent):
def __init__(
self,
parent_topic_type: str,
output_topic_type: str,
agent: ChatAgent | Team,
message_factory: MessageFactory
) -> None:
...
Import
from autogen_agentchat.teams._group_chat import ChatAgentContainer
I/O Contract
Inputs
| Parameter | Type | Required | Description |
|---|---|---|---|
| parent_topic_type | str | Yes | The topic type of the parent orchestrator for publishing responses |
| output_topic_type | str | Yes | The topic type for publishing intermediate messages and events |
| agent | Team | Yes | The agent or team to delegate message handling to |
| message_factory | MessageFactory | Yes | Factory for creating messages from JSON data and validating message types |
Event Handlers
| Method | Event Type | Description |
|---|---|---|
| handle_start | GroupChatStart | Buffers initial messages when group chat starts |
| handle_agent_response | GroupChatAgentResponse | Buffers agent response messages from other participants |
| handle_team_response | GroupChatTeamResponse | Buffers team result messages from other participants |
| handle_request | GroupChatRequestPublish | Processes buffered messages through wrapped agent/team and publishes response |
| handle_reset | GroupChatReset | Clears message buffer and resets wrapped agent/team |
| handle_pause | GroupChatPause | Pauses the wrapped agent/team |
| handle_resume | GroupChatResume | Resumes the wrapped agent/team |
Outputs
| Published Event | When | Description |
|---|---|---|
| GroupChatAgentResponse | After agent processing | Contains Response from ChatAgent with chat_message |
| GroupChatTeamResponse | After team processing | Contains TaskResult from Team with messages |
| GroupChatMessage | During processing | Intermediate messages and events for logging/monitoring |
| GroupChatError | On exception | Serialized exception with traceback |
Usage Examples
Container Creation (Internal)
from autogen_agentchat.teams._group_chat import ChatAgentContainer
from autogen_agentchat.agents import AssistantAgent
from autogen_agentchat.messages import MessageFactory
# This is typically done internally by group chat teams
message_factory = MessageFactory()
agent = AssistantAgent("assistant", model_client=model_client)
container = ChatAgentContainer(
parent_topic_type="orchestrator_topic",
output_topic_type="output_topic",
agent=agent,
message_factory=message_factory
)
Message Buffering Flow
from autogen_agentchat.teams._group_chat._events import (
GroupChatStart, GroupChatAgentResponse, GroupChatRequestPublish
)
from autogen_agentchat.messages import TextMessage
from autogen_agentchat.base import Response
# 1. Start event adds messages to buffer
start_event = GroupChatStart(messages=[
TextMessage(content="Hello", source="user")
])
# Container buffers the message
# 2. Other agent responses are buffered
agent_response = GroupChatAgentResponse(
response=Response(chat_message=TextMessage(content="Hi", source="other_agent")),
name="other_agent"
)
# Container buffers the response message
# 3. Request publish triggers processing
request_event = GroupChatRequestPublish()
# Container passes buffered messages to wrapped agent
# Agent processes and returns response
# Container publishes GroupChatAgentResponse to orchestrator
Team Container Processing
from autogen_agentchat.teams import RoundRobinGroupChat
from autogen_agentchat.teams._group_chat import ChatAgentContainer
async def team_container_example():
# Create inner team
inner_team = RoundRobinGroupChat(
participants=[agent1, agent2],
termination_condition=max_turns_condition
)
# Wrap in container
container = ChatAgentContainer(
parent_topic_type="parent_topic",
output_topic_type="output_topic",
agent=inner_team,
message_factory=message_factory
)
# When request is handled:
# 1. Container calls inner_team.run_stream(task=buffered_messages)
# 2. Streams TaskResult and events from inner team
# 3. Publishes GroupChatTeamResponse with final TaskResult
Error Handling
from autogen_agentchat.teams._group_chat._events import GroupChatError
async def error_handling_example():
# If wrapped agent throws exception during processing:
try:
response = await agent.on_messages(messages, cancellation_token)
except Exception as e:
# Container catches exception
error_message = SerializableException.from_exception(e)
# Publishes error to orchestrator
error_event = GroupChatError(error=error_message)
# Then re-raises for runtime
State Management
from autogen_agentchat.state import ChatAgentContainerState
async def state_management_example(container: ChatAgentContainer):
# Save container state
state_dict = await container.save_state()
# State includes:
# 1. Wrapped agent/team state
# 2. Message buffer (serialized messages)
container_state = ChatAgentContainerState.model_validate(state_dict)
print(f"Agent state: {container_state.agent_state}")
print(f"Buffered messages: {len(container_state.message_buffer)}")
# Restore state
await container.load_state(state_dict)
# Both agent state and message buffer are restored
Message Logging
from autogen_agentchat.messages import BaseAgentEvent, BaseChatMessage
async def message_logging_example():
# During agent processing, container logs intermediate messages
async for msg in agent.on_messages_stream(messages, cancellation_token):
if isinstance(msg, Response):
# Log final response message
await container._log_message(msg.chat_message)
else:
# Log intermediate event or message
await container._log_message(msg)
# Each logged message is published as GroupChatMessage to output_topic
# This enables monitoring and debugging of group chat conversations
Reset and Lifecycle
from autogen_agentchat.teams._group_chat._events import (
GroupChatReset, GroupChatPause, GroupChatResume
)
async def lifecycle_example(container: ChatAgentContainer):
# Reset container
reset_event = GroupChatReset()
await container.handle_reset(reset_event, ctx)
# Clears message buffer and resets wrapped agent/team
# Pause during processing
pause_event = GroupChatPause()
await container.handle_pause(pause_event, ctx)
# Wrapped agent/team is paused
# Resume later
resume_event = GroupChatResume()
await container.handle_resume(resume_event, ctx)
# Wrapped agent/team is resumed
Related Pages
- Microsoft_Autogen_SequentialRoutedAgent - Base class providing sequential message processing
- Microsoft_Autogen_ChatAgent_Protocol - Protocol implemented by wrapped agents
- Microsoft_Autogen_Team - Base class for wrapped teams
- Microsoft_Autogen_GroupChat_Events - Events handled by the container
- Microsoft_Autogen_MessageFactory - Factory for message serialization/deserialization
- Microsoft_Autogen_ChatAgentContainerState - State model for container persistence
- Microsoft_Autogen_BaseGroupChat - Base class for group chat teams that use containers
- Microsoft_Autogen_RoundRobinGroupChat - Group chat team that uses containers
- Microsoft_Autogen_SelectorGroupChat - Group chat team that uses containers
- Microsoft_Autogen_MagenticOneGroupChat - Group chat team that uses containers