Skip to content

Commit 6a61ddf

Browse files
authored
feat: introduce agent orchestrator to bridge agents and backend server (#20)
1 parent fdb7a6b commit 6a61ddf

31 files changed

Lines changed: 3324 additions & 124 deletions

python/third_party/TradingAgents/adapter/__main__.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@
1010
from langgraph.graph import StateGraph, MessagesState, START, END
1111
from pydantic import BaseModel, Field, field_validator
1212
from valuecell.core.agent.decorator import create_wrapped_agent
13-
from valuecell.core.agent.types import BaseAgent
13+
from valuecell.core.types import BaseAgent
1414

1515
from tradingagents.graph.trading_graph import TradingAgentsGraph
1616
from tradingagents.default_config import DEFAULT_CONFIG

python/third_party/ai-hedge-fund/adapter/__main__.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99
from langchain_core.messages import HumanMessage
1010
from pydantic import BaseModel, Field, field_validator
1111
from valuecell.core.agent.decorator import create_wrapped_agent
12-
from valuecell.core.agent.types import BaseAgent
12+
from valuecell.core.types import BaseAgent
1313

1414
from src.main import create_workflow
1515
from src.utils.analysts import ANALYST_ORDER

python/valuecell/agents/__init__.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99
from pathlib import Path
1010
from typing import List
1111

12-
from valuecell.core.agent.types import BaseAgent
12+
from valuecell.core.types import BaseAgent
1313

1414

1515
def _discover_and_import_agents() -> List[str]:
Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
from valuecell.core.agent.decorator import serve
2-
from valuecell.core.agent.types import BaseAgent
2+
from valuecell.core.types import BaseAgent
33

44

55
@serve()
@@ -9,7 +9,7 @@ class HelloWorldAgent(BaseAgent):
99
"""
1010

1111
async def stream(self, query, session_id, task_id):
12-
return {
12+
yield {
1313
"content": f"Hello! You said: {query}",
1414
"is_task_complete": True,
1515
}

python/valuecell/agents/sec_13F_agent.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
from edgar import Company, set_identity
88
from pydantic import BaseModel, Field, field_validator
99

10-
from valuecell.core.agent.types import BaseAgent
10+
from valuecell.core.types import BaseAgent
1111
from valuecell.core.agent.decorator import create_wrapped_agent
1212

1313
# Configure logging
Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import pytest
2+
from a2a.types import AgentCard
23
from valuecell.core.agent.connect import RemoteConnections
34

45

@@ -10,12 +11,15 @@ async def test_run_hello_world():
1011
available = connections.list_available_agents()
1112
assert name in available
1213

13-
url = await connections.start_agent("HelloWorldAgent")
14-
assert isinstance(url, str) and url
14+
agent_card = await connections.start_agent("HelloWorldAgent")
15+
assert isinstance(agent_card, AgentCard) and agent_card
1516

1617
client = await connections.get_client("HelloWorldAgent")
17-
task, event = await client.send_message("Hi there!")
18-
assert task is not None
19-
assert event is None
18+
turns = 0
19+
async for task, event in await client.send_message("Hi there!"):
20+
assert task is not None
21+
assert event is None
22+
turns += 1
23+
assert turns == 1
2024
finally:
2125
await connections.stop_all()

python/valuecell/core/__init__.py

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
# Session management
2+
from .session import (
3+
InMemorySessionStore,
4+
Message,
5+
Role,
6+
Session,
7+
SessionManager,
8+
SessionStore,
9+
)
10+
11+
# Task management
12+
from .task import (
13+
InMemoryTaskStore,
14+
Task,
15+
TaskManager,
16+
TaskStatus,
17+
TaskStore,
18+
)
19+
20+
# Type system
21+
from .types import (
22+
UserInput,
23+
UserInputMetadata,
24+
BaseAgent,
25+
StreamResponse,
26+
RemoteAgentResponse,
27+
)
28+
29+
__all__ = [
30+
# Session exports
31+
"Message",
32+
"Role",
33+
"Session",
34+
"SessionManager",
35+
"SessionStore",
36+
"InMemorySessionStore",
37+
# Task exports
38+
"Task",
39+
"TaskStatus",
40+
"TaskManager",
41+
"TaskStore",
42+
"InMemoryTaskStore",
43+
# Type system exports
44+
"UserInput",
45+
"UserInputMetadata",
46+
"BaseAgent",
47+
"StreamResponse",
48+
"RemoteAgentResponse",
49+
]
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
"""Agent module initialization"""
2+
3+
# Core agent functionality
4+
from .client import AgentClient
5+
from .connect import RemoteConnections
6+
from .decorator import serve
7+
from .registry import AgentRegistry
8+
9+
# Import types from the unified types module
10+
from ..types import BaseAgent, RemoteAgentResponse, StreamResponse
11+
12+
13+
__all__ = [
14+
# Core agent exports
15+
"AgentClient",
16+
"RemoteConnections",
17+
"serve",
18+
"AgentRegistry",
19+
"BaseAgent",
20+
"RemoteAgentResponse",
21+
"StreamResponse",
22+
]

python/valuecell/core/agent/client.py

Lines changed: 25 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@
55
from a2a.types import Message, Part, PushNotificationConfig, Role, TextPart
66
from valuecell.utils import generate_uuid
77

8-
from .types import MessageResponse
8+
from ..types import RemoteAgentResponse
99

1010

1111
class AgentClient:
@@ -48,8 +48,12 @@ async def _setup_client(self):
4848
self._client = client_factory.create(card)
4949

5050
async def send_message(
51-
self, text: str, context_id: str = None, streaming: bool = False
52-
) -> MessageResponse | AsyncIterator[MessageResponse]:
51+
self,
52+
query: str,
53+
context_id: str = None,
54+
metadata: dict = None,
55+
streaming: bool = False,
56+
) -> AsyncIterator[RemoteAgentResponse]:
5357
"""Send message to Agent.
5458
5559
If `streaming` is True, return an async iterator producing (task, event) pairs.
@@ -59,18 +63,28 @@ async def send_message(
5963

6064
message = Message(
6165
role=Role.user,
62-
parts=[Part(root=TextPart(text=text))],
66+
parts=[Part(root=TextPart(text=query))],
6367
message_id=generate_uuid("msg"),
6468
context_id=context_id or generate_uuid("ctx"),
69+
metadata=metadata if metadata else None,
6570
)
6671

67-
generator = self._client.send_message(message)
68-
if streaming:
69-
return generator
70-
71-
task, event = await generator.__anext__()
72-
await generator.aclose()
73-
return task, event
72+
source_gen = self._client.send_message(message)
73+
74+
async def wrapper() -> AsyncIterator[RemoteAgentResponse]:
75+
try:
76+
if streaming:
77+
async for item in source_gen:
78+
yield item
79+
else:
80+
# yield only the first item
81+
item = await source_gen.__anext__()
82+
yield item
83+
finally:
84+
# ensure underlying generator is closed
85+
await source_gen.aclose()
86+
87+
return wrapper()
7488

7589
async def get_agent_card(self):
7690
await self._ensure_initialized()

0 commit comments

Comments
 (0)