197 lines
7.6 KiB
Python
197 lines
7.6 KiB
Python
"""ADK Multi-Agent Architecture for GCP Solution Architecture Agent.
|
|
|
|
Uses google.adk.agents primitives:
|
|
- SourceDiscoveryAgent (LlmAgent)
|
|
- DiscoveryAgent (LlmAgent)
|
|
- ArchitectureDesignAgent (LlmAgent)
|
|
- ValidationReviewAgent (LlmAgent)
|
|
- PackagingAgent (LlmAgent)
|
|
- OrchestratorLoopAgent (LoopAgent)
|
|
"""
|
|
|
|
import logging
|
|
import uuid
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from app.adk.compat import BaseAgent, LlmAgent, LoopAgent, SequentialAgent
|
|
from app.adk.tools import ADK_TOOLS
|
|
from app.config import get_settings
|
|
from app.database import get_db_manager
|
|
from app.nodes import design_node, discover_node, package_node, source_discover_node, validate_node
|
|
from app.skills.loader import SkillLoader
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class SourceDiscoveryAgent(LlmAgent):
|
|
"""ADK Agent responsible for Phase 0a Pre-emptive Source Environment Discovery."""
|
|
|
|
def __init__(self, skill_loader: SkillLoader) -> None:
|
|
super().__init__(
|
|
name="SourceDiscoveryAgent",
|
|
description="Audits and documents the pre-existing source environment (As-Is Architecture) before target migration.",
|
|
instruction=skill_loader.format_skills_for_prompt("source_discover"),
|
|
tools=[],
|
|
)
|
|
self.skill_loader = skill_loader
|
|
|
|
def execute(self, state: Dict[str, Any]) -> Dict[str, Any]:
|
|
logger.info("Executing SourceDiscoveryAgent")
|
|
res = source_discover_node(state, self.skill_loader)
|
|
state_copy = dict(state)
|
|
state_copy.update(res)
|
|
return state_copy
|
|
|
|
|
|
class DiscoveryAgent(LlmAgent):
|
|
"""ADK Agent responsible for Phase 0 Requirements Discovery."""
|
|
|
|
def __init__(self, skill_loader: SkillLoader) -> None:
|
|
super().__init__(
|
|
name="DiscoveryAgent",
|
|
description="Extracts functional and non-functional requirements with product selection deferred.",
|
|
instruction=skill_loader.format_skills_for_prompt("discover"),
|
|
tools=[],
|
|
)
|
|
self.skill_loader = skill_loader
|
|
|
|
def execute(self, state: Dict[str, Any]) -> Dict[str, Any]:
|
|
logger.info("Executing DiscoveryAgent")
|
|
res = discover_node(state, self.skill_loader)
|
|
state_copy = dict(state)
|
|
state_copy.update(res)
|
|
return state_copy
|
|
|
|
|
|
class ArchitectureDesignAgent(LlmAgent):
|
|
"""ADK Agent responsible for Phase 1 Product Selection, Diagrams, & Terraform IaC."""
|
|
|
|
def __init__(self, skill_loader: SkillLoader) -> None:
|
|
super().__init__(
|
|
name="ArchitectureDesignAgent",
|
|
description="Selects GCP products, generates Mermaid diagram, and produces Terraform IaC grounded by Developer Knowledge MCP.",
|
|
instruction=skill_loader.format_skills_for_prompt("design"),
|
|
tools=ADK_TOOLS,
|
|
)
|
|
self.skill_loader = skill_loader
|
|
|
|
def execute(self, state: Dict[str, Any]) -> Dict[str, Any]:
|
|
logger.info("Executing ArchitectureDesignAgent")
|
|
res = design_node(state, self.skill_loader)
|
|
state_copy = dict(state)
|
|
state_copy.update(res)
|
|
return state_copy
|
|
|
|
|
|
class ValidationReviewAgent(LlmAgent):
|
|
"""ADK Agent responsible for Phase 2 Artifact Review & Quality Validation."""
|
|
|
|
def __init__(self, skill_loader: SkillLoader) -> None:
|
|
super().__init__(
|
|
name="ValidationReviewAgent",
|
|
description="Validates Mermaid syntax, Terraform configuration, and required sections.",
|
|
instruction=skill_loader.format_skills_for_prompt("validate"),
|
|
tools=ADK_TOOLS,
|
|
)
|
|
self.skill_loader = skill_loader
|
|
|
|
def execute(self, state: Dict[str, Any]) -> Dict[str, Any]:
|
|
logger.info("Executing ValidationReviewAgent")
|
|
res = validate_node(state, self.skill_loader)
|
|
state_copy = dict(state)
|
|
state_copy.update(res)
|
|
return state_copy
|
|
|
|
|
|
class PackagingAgent(LlmAgent):
|
|
"""ADK Agent responsible for Phase 3 Solution Guide Packaging."""
|
|
|
|
def __init__(self, skill_loader: SkillLoader) -> None:
|
|
super().__init__(
|
|
name="PackagingAgent",
|
|
description="Packages final solution-architecture-guide.md.",
|
|
instruction=skill_loader.format_skills_for_prompt("package"),
|
|
tools=[],
|
|
)
|
|
self.skill_loader = skill_loader
|
|
|
|
def execute(self, state: Dict[str, Any]) -> Dict[str, Any]:
|
|
logger.info("Executing PackagingAgent")
|
|
res = package_node(state, self.skill_loader)
|
|
state_copy = dict(state)
|
|
state_copy.update(res)
|
|
return state_copy
|
|
|
|
|
|
class OrchestratorLoopAgent(LoopAgent):
|
|
"""ADK Orchestrator LoopAgent that coordinates multi-agent execution & iterative quality review.
|
|
|
|
Uses google.adk.agents.LoopAgent to execute discovery -> design -> validation review -> packaging
|
|
in a loop until validation passes 100% or max_iterations is reached. Writes every review iteration
|
|
event to PostgreSQL database (`orchestrator_review_logs`).
|
|
"""
|
|
|
|
def __init__(self, skill_loader: SkillLoader, max_iterations: Optional[int] = None) -> None:
|
|
settings = get_settings()
|
|
max_iters = max_iterations or settings.ADK_MAX_LOOP_ITERATIONS
|
|
self.db_manager = get_db_manager()
|
|
|
|
self.source_discovery_agent = SourceDiscoveryAgent(skill_loader)
|
|
self.discovery_agent = DiscoveryAgent(skill_loader)
|
|
self.design_agent = ArchitectureDesignAgent(skill_loader)
|
|
self.validation_agent = ValidationReviewAgent(skill_loader)
|
|
self.packaging_agent = PackagingAgent(skill_loader)
|
|
|
|
sub_pipeline = SequentialAgent(
|
|
name="MultiAgentGCPPipeline",
|
|
sub_agents=[
|
|
self.source_discovery_agent,
|
|
self.discovery_agent,
|
|
self.design_agent,
|
|
self.validation_agent,
|
|
self.packaging_agent,
|
|
],
|
|
description="Sequential pipeline of GCP architecture multi-agents.",
|
|
)
|
|
|
|
def review_validator(state: Dict[str, Any]) -> bool:
|
|
"""Check if solution meets production quality validation standards."""
|
|
is_valid = state.get("validation_passed", False)
|
|
errors = state.get("errors", [])
|
|
iteration = state.get("loop_count", 1)
|
|
execution_id = state.get("execution_id", str(uuid.uuid4()))
|
|
|
|
status_str = "APPROVED" if is_valid else "NEEDS_REVISION"
|
|
feedback_str = "All architecture validation rules passed." if is_valid else f"Validation errors: {', '.join(errors)}"
|
|
|
|
# Record review iteration in PostgreSQL database
|
|
self.db_manager.record_orchestrator_log(
|
|
log_id=str(uuid.uuid4()),
|
|
execution_id=execution_id,
|
|
iteration=iteration,
|
|
review_status=status_str,
|
|
reviewer_agent="ValidationReviewAgent",
|
|
feedback=feedback_str,
|
|
)
|
|
logger.info(
|
|
"Orchestrator Review Loop #%d: status=%s, valid=%s",
|
|
iteration,
|
|
status_str,
|
|
is_valid,
|
|
)
|
|
|
|
return is_valid
|
|
|
|
super().__init__(
|
|
name="OrchestratorLoopAgent",
|
|
sub_agent=sub_pipeline,
|
|
max_iterations=max_iters,
|
|
description="Production-ready multi-agent orchestrator loop agent.",
|
|
validator_fn=review_validator,
|
|
)
|
|
|
|
|
|
def build_adk_multi_agent_system(skill_loader: SkillLoader, max_iterations: Optional[int] = None) -> OrchestratorLoopAgent:
|
|
"""Factory function for building the complete ADK Orchestrator LoopAgent system."""
|
|
return OrchestratorLoopAgent(skill_loader, max_iterations=max_iterations)
|