mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-02 02:58:10 +02:00
Add multi-pattern orchestrator with plan-then-execute and supervisor (#739)
Introduce an agent orchestrator service that supports three execution patterns (ReAct, plan-then-execute, supervisor) with LLM-based meta-routing to select the appropriate pattern and task type per request. Update the agent schema to support orchestration fields (correlation, sub-agents, plan steps) and remove legacy response fields (answer, thought, observation).
This commit is contained in:
parent
7af1d60db8
commit
849987f0e6
21 changed files with 3006 additions and 172 deletions
|
|
@ -16,6 +16,14 @@ class AgentRequestTranslator(MessageTranslator):
|
|||
collection=data.get("collection", "default"),
|
||||
streaming=data.get("streaming", False),
|
||||
session_id=data.get("session_id", ""),
|
||||
conversation_id=data.get("conversation_id", ""),
|
||||
pattern=data.get("pattern", ""),
|
||||
task_type=data.get("task_type", ""),
|
||||
framing=data.get("framing", ""),
|
||||
correlation_id=data.get("correlation_id", ""),
|
||||
parent_session_id=data.get("parent_session_id", ""),
|
||||
subagent_goal=data.get("subagent_goal", ""),
|
||||
expected_siblings=data.get("expected_siblings", 0),
|
||||
)
|
||||
|
||||
def from_pulsar(self, obj: AgentRequest) -> Dict[str, Any]:
|
||||
|
|
@ -28,6 +36,14 @@ class AgentRequestTranslator(MessageTranslator):
|
|||
"collection": getattr(obj, "collection", "default"),
|
||||
"streaming": getattr(obj, "streaming", False),
|
||||
"session_id": getattr(obj, "session_id", ""),
|
||||
"conversation_id": getattr(obj, "conversation_id", ""),
|
||||
"pattern": getattr(obj, "pattern", ""),
|
||||
"task_type": getattr(obj, "task_type", ""),
|
||||
"framing": getattr(obj, "framing", ""),
|
||||
"correlation_id": getattr(obj, "correlation_id", ""),
|
||||
"parent_session_id": getattr(obj, "parent_session_id", ""),
|
||||
"subagent_goal": getattr(obj, "subagent_goal", ""),
|
||||
"expected_siblings": getattr(obj, "expected_siblings", 0),
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -40,24 +56,15 @@ class AgentResponseTranslator(MessageTranslator):
|
|||
def from_pulsar(self, obj: AgentResponse) -> Dict[str, Any]:
|
||||
result = {}
|
||||
|
||||
# Check if this is a streaming response (has chunk_type)
|
||||
if hasattr(obj, 'chunk_type') and obj.chunk_type:
|
||||
if obj.chunk_type:
|
||||
result["chunk_type"] = obj.chunk_type
|
||||
if obj.content:
|
||||
result["content"] = obj.content
|
||||
result["end_of_message"] = getattr(obj, "end_of_message", False)
|
||||
result["end_of_dialog"] = getattr(obj, "end_of_dialog", False)
|
||||
else:
|
||||
# Legacy format (non-streaming)
|
||||
if obj.answer:
|
||||
result["answer"] = obj.answer
|
||||
if obj.thought:
|
||||
result["thought"] = obj.thought
|
||||
if obj.observation:
|
||||
result["observation"] = obj.observation
|
||||
# Include completion flags for legacy format too
|
||||
result["end_of_message"] = getattr(obj, "end_of_message", False)
|
||||
result["end_of_dialog"] = getattr(obj, "end_of_dialog", False)
|
||||
if obj.content:
|
||||
result["content"] = obj.content
|
||||
result["end_of_message"] = getattr(obj, "end_of_message", False)
|
||||
result["end_of_dialog"] = getattr(obj, "end_of_dialog", False)
|
||||
|
||||
if getattr(obj, "message_id", ""):
|
||||
result["message_id"] = obj.message_id
|
||||
|
||||
# Include explainability fields if present
|
||||
explain_id = getattr(obj, "explain_id", None)
|
||||
|
|
@ -76,11 +83,5 @@ class AgentResponseTranslator(MessageTranslator):
|
|||
|
||||
def from_response_with_completion(self, obj: AgentResponse) -> Tuple[Dict[str, Any], bool]:
|
||||
"""Returns (response_dict, is_final)"""
|
||||
# For streaming responses, check end_of_dialog
|
||||
if hasattr(obj, 'chunk_type') and obj.chunk_type:
|
||||
is_final = getattr(obj, 'end_of_dialog', False)
|
||||
else:
|
||||
# For legacy responses, check if answer is present
|
||||
is_final = (obj.answer is not None)
|
||||
|
||||
is_final = getattr(obj, 'end_of_dialog', False)
|
||||
return self.from_pulsar(obj), is_final
|
||||
Loading…
Add table
Add a link
Reference in a new issue