Victor Dibia 0e985d4b40
v1 of AutoGen Studio on AgentChat (#4097)
* add skeleton worflow manager

* add test notebook

* update test nb

* add sample team spec

* refactor requirements to agentchat and ext

* add base provider to return agentchat agents from json spec

* initial api refactor, update dbmanager

* api refactor

* refactor tests

* ags api tutorial update

* ui refactor

* general refactor

* minor refactor updates

* backend api refaactor

* ui refactor and update

* implement v1 for streaming connection with ui updates

* backend refactor

* ui refactor

* minor ui tweak

* minor refactor and tweaks

* general refactor

* update tests

* sync uv.lock with main

* uv lock update
2024-11-09 14:32:24 -08:00

68 lines
2.1 KiB
Python

from typing import AsyncGenerator, Union, Optional
import time
from .database import ComponentFactory
from .datamodel import TeamResult, TaskResult, ComponentConfigInput
from autogen_agentchat.messages import InnerMessage, ChatMessage
from autogen_core.base import CancellationToken
class TeamManager:
def __init__(self) -> None:
self.component_factory = ComponentFactory()
async def run_stream(
self,
task: str,
team_config: ComponentConfigInput,
cancellation_token: Optional[CancellationToken] = None
) -> AsyncGenerator[Union[InnerMessage, ChatMessage, TaskResult], None]:
"""Stream the team's execution results"""
start_time = time.time()
try:
# Let factory handle all config processing
team = await self.component_factory.load(team_config)
stream = team.run_stream(
task=task,
cancellation_token=cancellation_token
)
async for message in stream:
if cancellation_token and cancellation_token.is_cancelled():
break
if isinstance(message, TaskResult):
yield TeamResult(
task_result=message,
usage="",
duration=time.time() - start_time
)
else:
yield message
except Exception as e:
raise e
async def run(
self,
task: str,
team_config: ComponentConfigInput,
cancellation_token: Optional[CancellationToken] = None
) -> TeamResult:
"""Original non-streaming run method with optional cancellation"""
start_time = time.time()
# Let factory handle all config processing
team = await self.component_factory.load(team_config)
result = await team.run(
task=task,
cancellation_token=cancellation_token
)
return TeamResult(
task_result=result,
usage="",
duration=time.time() - start_time
)