Workflow agent
This guide covers building multi-agent workflows using BridgeBaseWorkflowAgent and LangGraph.
When to use BridgeBaseWorkflowAgent
Use BridgeBaseWorkflowAgent when your agent:
- Has multiple processing steps with conditional routing
- Needs multiple specialized sub-agents
- Has complex state management requirements
- Needs retry loops or fan-out patterns
Base Class Overview
from bridge_agent_sdk import BridgeBaseWorkflowAgent, WorkflowConfig
from langgraph.graph import StateGraph, END
class MyWorkflow(BridgeBaseWorkflowAgent):
"""Your custom workflow agent."""
CONFIG = WorkflowConfig(
name="my-workflow",
version="1.0.0",
description="My workflow agent",
)
def define_state(self) -> type:
"""Return the TypedDict class defining workflow state."""
return MyState
def setup_nodes(self, graph: StateGraph):
"""Add nodes to the workflow graph."""
pass
def setup_edges(self, graph: StateGraph):
"""Define edges and routing between nodes."""
passLifecycle Methods
Method | Purpose | When Called |
|---|---|---|
define_state() | Define state schema (TypedDict) | On graph construction |
setup_nodes() | Add processing nodes | On graph construction |
setup_edges() | Define flow routing | On graph construction |
prepare_initial_state() | Build initial state from AgentInput | Before graph execution |
extract_output() | Extract output from final state | After graph execution |
execute() | Orchestrates the full lifecycle | On each invocation |
There is no build() or run() method. The execute() method handles graph compilation and invocation automatically. Call await workflow.execute(agent_input).
Add workflow agent
Step 1: Define Workflow State
Create state in src/agents/state.py:
"""Workflow state definitions."""
from typing import TypedDict, Optional, List, Dict, Any, Annotated
from operator import add
class WorkflowStatus:
"""Canonical workflow status values."""
PENDING = "pending"
IN_PROGRESS = "in_progress"
SUCCESS = "success"
ERROR = "error"
class WorkflowState(TypedDict, total=False):
"""State shared across all workflow nodes.
All nodes can read and write to this state.
LangGraph manages state persistence and transitions.
"""
# Input
raw_input: dict
platform_context: dict
# Processing
status: str
current_node: str
# Results
output: Optional[str]
error: Optional[str]
error_message: Optional[str]
# Tracking (use Annotated for append-only lists)
messages: Annotated[List[Dict[str, Any]], add]Step 2: Create the Workflow Agent
Create agent in src/agents/orchestrator.py:
"""Orchestrator workflow agent using LangGraph."""
import logging
import os
from typing import Dict, Any, Literal
from bridge_agent_sdk import (
BridgeBaseWorkflowAgent,
WorkflowConfig,
AgentInput,
MCPClient,
)
from langgraph.graph import StateGraph, END
from src.agents.state import WorkflowState, WorkflowStatus
logger = logging.getLogger(__name__)
class OrchestratorWorkflow(BridgeBaseWorkflowAgent):
"""Multi-agent workflow that routes queries to specialized handlers.
Flow:
1. Classifier → Determines query type (SQL, RAG, Action)
2. Route to appropriate handler
3. Responder → Formats final response
"""
CONFIG = WorkflowConfig(
name="orchestrator",
version="1.0.0",
description="Multi-agent query orchestrator",
)
CLASSIFIER_PROMPT = """Classify the user query into one of these categories:
- SQL: Data queries, metrics, reports, statistics
- RAG: Knowledge questions, how-to, documentation
Respond with JSON: {"type": "sql|rag", "confidence": 0.0-1.0}"""# ─────────────────────────────────────────────────────────────
# Required Overrides
# ─────────────────────────────────────────────────────────────
def define_state(self) -> type:
"""Define the state schema for this workflow."""
return WorkflowState
def setup_nodes(self, graph: StateGraph):
"""Add processing nodes to the workflow graph."""
graph.add_node("classifier", self.classifier_node)
graph.add_node("sql_agent", self.sql_agent_node)
graph.add_node("rag_agent", self.rag_agent_node)
graph.add_node("responder", self.responder_node)
def setup_edges(self, graph: StateGraph):
"""Define edges and routing between nodes."""
graph.set_entry_point("classifier")
graph.add_conditional_edges(
"classifier",
self.route_query,
{
"sql": "sql_agent",
"rag": "rag_agent",
}
)
graph.add_edge("sql_agent", "responder")
graph.add_edge("rag_agent", "responder")
graph.add_edge("responder", END)
# ─────────────────────────────────────────────────────────────
# Node Implementations
# ─────────────────────────────────────────────────────────────
def classifier_node(self, state: WorkflowState) -> Dict[str, Any]:
"""Classify the user query to determine routing."""
logger.info("Classifying query...")
# self.llm is auto-created by the SDK's execute() method.
# Use self._create_llm(state) only if you need a fresh client
# with different settings (e.g., different temperature).
response = self.llm.invoke([
{"role": "system", "content": self.CLASSIFIER_PROMPT},
{"role": "user", "content": str(state.get("raw_input", {}).get("query", ""))},
])
classification = self._parse_classification(response.content)
return {
"status": WorkflowStatus.IN_PROGRESS,
"current_node": "classifier",
"messages": [{"role": "classifier", "content": response.content}],
}
async def sql_agent_node(self, state: WorkflowState) -> Dict[str, Any]:
"""Process SQL/data queries via MCP."""
logger.info("Processing SQL query...")
platform_context = state.get("platform_context", {})
mcp = MCPClient(
dev_mode=(os.getenv("KAIF_MODE") != "production"),
agent_payload=platform_context if os.getenv("KAIF_MODE") != "local_dev" else None,
)
result = await mcp.call_tool_parsed(
"bridge_execute_query",
{"query": 'SELECT * FROM "ITSM".incident LIMIT 10'},
)
success, parsed = result
return {
"output": str(parsed.get("data", []) if success else []),
"status": WorkflowStatus.IN_PROGRESS,
"current_node": "sql_agent",
}
def rag_agent_node(self, state: WorkflowState) -> Dict[str, Any]:
"""Process knowledge/RAG queries."""
logger.info("Processing RAG query...")
query = str(state.get("raw_input", {}).get("query", ""))
response = self.llm.invoke([
{"role": "system", "content": "Answer the question based on your knowledge."},
{"role": "user", "content": query},
])
return {
"output": response.content,
"status": WorkflowStatus.IN_PROGRESS,
"current_node": "rag_agent",
}
def responder_node(self, state: WorkflowState) -> Dict[str, Any]:
"""Format the final response."""
logger.info("Generating final response...")
return {
"status": WorkflowStatus.SUCCESS,
"current_node": "responder",
}
# ─────────────────────────────────────────────────────────────
# Routing Function
# ─────────────────────────────────────────────────────────────
def route_query(self, state: WorkflowState) -> Literal["sql", "rag"]:
"""Route query to appropriate handler based on classification."""
messages = state.get("messages", [])
if messages:
last = messages[-1].get("content", "")
if '"sql"' in last.lower():
return "sql"
return "rag"
# ─────────────────────────────────────────────────────────────
# Helpers
# ─────────────────────────────────────────────────────────────
def _parse_classification(self, response: str) -> Dict[str, Any]:
import json
try:
return json.loads(response)
except json.JSONDecodeError:
return {"type": "rag", "confidence": 0.0}Step 3: Create the Entrypoint
Create main.py:
"""Application entrypoint for workflow agent."""
from dotenv import load_dotenv
load_dotenv(override=False)
import logging
from bridge_agent_sdk import run_agent
from src.agents.orchestrator import OrchestratorWorkflow
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
WORKFLOWS = {
"orchestrator": OrchestratorWorkflow,
}
AGENT_NAME_MAP = {
"orchestrator": "orchestrator",
}
if __name__ == "__main__":
run_agent(
WORKFLOWS,
agent_name_map=AGENT_NAME_MAP,
description="Multi-Agent Query Orchestrator",
)How Execution Works
When you call await workflow.execute(agent_input), the base class:
- Calls define_state() to get the TypedDict schema
- Creates a StateGraph and calls setup_nodes() + setup_edges()
- Compiles the graph
- Calls prepare_initial_state(agent_input) to build the initial state dict
- Creates the LLM client and assigns it to self.llm (via self._create_llm(state)) now that platform_context is available
- Activates automatic CADF audit hooks (if enable_audit=True, the default)
- Invokes the compiled graph with Langfuse callbacks (if enable_langfuse=True, the default)
- Calls extract_output(final_state) to build the return value
- Deactivates audit hooks in the finally block
prepare_initial_state() — Default State Keys
The SDK's default prepare_initial_state() auto-decomposes ExecutionContext into:
State Key | Source | Description |
|---|---|---|
input | agent_input.content | Primary input content |
metadata | agent_input.metadata | Execution metadata (thread_id, etc.) |
context | agent_input.context | Additional context |
previous_outputs | agent_input.previous_outputs | Outputs from previous agents |
platform_context | Extracted from ExecutionContext | Auth, account_id, workflow_id, etc. |
agent_payload | Same as platform_context | Deprecated alias |
raw_input | Extracted from ExecutionContext | content, input, parameters, runtime_data |
connection_details | From platform_context | List of connection dicts for external services |
Override prepare_initial_state() to add custom keys while keeping the SDK's auto-decomposition. Always call super() to get the base state:
# AgentBaseWorkflow is your project's base class from base.py (extends BridgeBaseWorkflowAgent)
class MyWorkflow(AgentBaseWorkflow):
def prepare_initial_state(self, agent_input: AgentInput) -> Dict[str, Any]:
"""Extend the SDK's default state with custom keys."""
base_state = super().prepare_initial_state(agent_input)
# base_state already has: platform_context, raw_input, connection_details, etc.
base_state.update({
"status": WorkflowStatus.PENDING,
"messages": [],
})
return base_state
def extract_output(self, final_state: Dict[str, Any]) -> Any:
"""Customize how final state maps to output."""
return {
"content": final_state.get("output", ""),
"status": final_state.get("status", WorkflowStatus.ERROR),
}Do not build platform_context manually. The SDK's super().prepare_initial_state() extracts it from ExecutionContext and validates all required fields. If you skip super(), auth will fail in bridge_dev and production modes.
Advanced patterns
Retry Loops
def setup_edges(self, graph: StateGraph):
graph.set_entry_point("attempt")
graph.add_conditional_edges(
"attempt",
self.check_result,
{
"success": "finalize",
"retry": "attempt",
"fail": END,
}
)
graph.add_edge("finalize", END)
def check_result(self, state) -> str:
if state.get("success"):
return "success"
elif state.get("retry_count", 0) < 3:
return "retry"
else:
return "fail"Sub-Workflow Orchestration (Pipeline Pattern)
Use run_sub_workflow_sync to call sub-workflows from LangGraph nodes (which are synchronous). Create sub-workflow instances once in __init__, not per-invocation:
from bridge_agent_sdk import run_sub_workflow_sync, AgentInput
class PipelineWorkflow(BridgeBaseWorkflowAgent):
def __init__(self):
super().__init__()
self._debug_workflow = DebugWorkflow()
self._remediation_workflow = RemediationWorkflow()
def node_run_debug(self, state):
agent_input = AgentInput(
content="debug task",
metadata={
"thread_id": state.get("metadata", {}).get("thread_id", "sub-debug"),
"execution_id": state.get("metadata", {}).get("execution_id", "unknown"),
},
context={
"platform_context": state["platform_context"],
"raw_input": state["raw_input"],
},
)
result = run_sub_workflow_sync(self._debug_workflow, agent_input)
return {
"debug_report": result.get("content", ""),
"debug_status": result.get("status", "unknown"),
}WorkflowConfig reference
Key WorkflowConfig fields and their defaults:
Field | Type | Default | Description |
|---|---|---|---|
name | str | required | Workflow name |
version | str | "1.0.0" | Workflow version |
description | str | required | Workflow description |
llm_provider | str | os.getenv("LLM_PROVIDER", "openai") | LLM backend: openai, azure_openai, hosted_llm, anthropic |
llm_model | str | os.getenv("LLM_MODEL", "gpt-4.1-mini") | Model name (reads from env at class load) |
llm_temperature | float | 0.7 | Sampling temperature |
max_iterations | int | 20 | Max workflow iterations |
recursion_limit | int | 200 | Max LangGraph recursion depth |
enable_langfuse | bool | True | Enable Langfuse tracing for all LLM calls |
enable_audit | bool | True | Enable automatic CADF audit hooks (start/stop events) |
enable_tenant_db | bool | False | Enable multi-tenant database isolation |
account_code | str | None | Tenant code for DB isolation (required when enable_tenant_db=True) |
Automatic Audit Hooks (v1.0.16+)
When enable_audit=True (the default), the SDK automatically wraps execute() with CADF audit events. You do NOT need to manually call _post_audit() for start/stop — the SDK emits:
- Start event: when execute() begins
- SDK-level events: on LLM calls, auth operations, connection lookups
- Stop event: when execute() completes (or fails)
To opt out, set enable_audit=False in WorkflowConfig.
GuardrailsService (v1.0.16+)
The SDK includes GuardrailsService for PII/DLP redaction in audit logs and Langfuse traces. It is initialized automatically on BridgeBaseWorkflowAgent when the guardrails endpoint is configured:
GUARDRAILS_ENDPOINT=https://your-bridge-host/kaif/v2/guardrails
GUARDRAILS_PII_PROVIDER=azure # or "google"
GUARDRAILS_DEFAULT_CONFIG_NAME=genaiassist-default
The _guardrails attribute is automatically passed to AuditService and Langfuse callbacks to redact PII from logged data.
Tenant Database (v1.0.16+)
For agents that need per-account data storage with schema isolation:
CONFIG = WorkflowConfig(
name="my-agent",
description="...",
enable_tenant_db=True,
account_code="ACME",
)
# In workflow nodes, access via self.db:
async def store_results_node(self, state):
await self.db.insert("execution_log", {"status": "completed", "result": "..."})
rows = await self.db.select("execution_log", where={"status": "completed"})
return {"stored": True}
Requires ACCOUNT_POSTGRES_* environment variables. See TenantDatabaseManager docs for schema isolation, RLS, and restricted roles.Testing Workflows
"""Tests for OrchestratorWorkflow."""
import pytest
from unittest.mock import patch, MagicMock
from bridge_agent_sdk import AgentInput
from src.agents.orchestrator import OrchestratorWorkflow
from src.agents.state import WorkflowStatus
@pytest.fixture
def workflow():
"""Create workflow instance."""
return OrchestratorWorkflow()
def test_classifier_routes_correctly(workflow):
"""Test that classifier produces expected state updates."""
state = {
"raw_input": {"query": "Show me incidents from last week"},
"platform_context": {},
"status": WorkflowStatus.PENDING,
"messages": [],
}
mock_response = MagicMock()
mock_response.content = '{"type": "sql", "confidence": 0.95}'
# self.llm is set by execute() at runtime; in unit tests, set it directly
workflow.llm = MagicMock()
workflow.llm.invoke.return_value = mock_response
result = workflow.classifier_node(state)
assert result["status"] == WorkflowStatus.IN_PROGRESS
assert len(result["messages"]) == 1
@pytest.mark.asyncio
async def test_full_workflow_execution(workflow):
"""Test complete workflow execution via execute()."""
agent_input = AgentInput(
content="What is our SLA policy?",
metadata={"thread_id": "test-456"},
context={"agent_payload": {"account_id": "test-account"}},
)
mock_response = MagicMock()
mock_response.content = '{"type": "rag", "confidence": 0.9}'
with patch.object(workflow, '_create_llm') as mock_create_llm:
mock_llm = MagicMock()
mock_llm.invoke.return_value = mock_response
mock_create_llm.return_value = mock_llm
result = await workflow.execute(agent_input)
assert result is not None