وكلاء الرسوم البيانية
أنشئ سير عمل معقدة وذات حالة باستخدام تنسيق LangGraph مع تكامل أصيل مع ADK-Rust.
نظرة عامة
يتيح لك GraphAgent تعريف سير العمل على شكل رسوم بيانية موجهة بعُقد وروابط، مع دعم:
- AgentNode: تغليف وكلاء LLM كعُقد في الرسم البياني مع موائمات إدخال/إخراج مخصصة
- سير العمل الدائري: دعم أصيل للحلقات والاستدلال التكراري (نمط ReAct)
- التوجيه الشرطي: توجيه ديناميكي للروابط بناءً على الحالة
- إدارة الحالة: حالة ذات أنواع مع مخففات (الاستبدال، الإلحاق، الجمع، مخصص)
- التحقق المرحلي: حالة دائمة لتحمل الأعطال ولإدخال الإنسان في الحلقة
- البث المتدفق: أوضاع تدفق متعددة (القيم، التحديثات، الرسائل، التصحيح)
توفر الحزمة adk-graph تنسيق سير عمل بنمط LangGraph لبناء سير عمل وكلاء معقدة وذات حالة. وهي تجلب قدرات سير العمل المعتمدة على الرسوم البيانية إلى منظومة ADK-Rust مع الحفاظ على التوافق الكامل مع نظام الوكلاء في ADK.
الفوائد الرئيسية:
- تصميم مرئي لسير العمل: عرّف منطقًا معقدًا على شكل رسوم بيانية بديهية من العُقد والروابط
- التنفيذ المتوازي: يمكن تشغيل عدة عُقد في الوقت نفسه لتحسين الأداء
- استمرارية الحالة: التحقق المرحلي المدمج لتحمل الأعطال ولإدخال الإنسان في الحلقة
- تكامل LLM: دعم أصيل لتغليف وكلاء ADK كعُقد في الرسم البياني
- توجيه مرن: روابط ثابتة، وتوجيه شرطي، واتخاذ قرار ديناميكي
ما الذي ستبنيه
في هذا الدليل، ستنشئ خط أنابيب لمعالجة النصوص يشغّل الترجمة والتلخيص بالتوازي:
┌─────────────────────┐
User Input │ │
────────────────▶ │ START │
│ │
└──────────┬──────────┘
│
┌───────────────┴───────────────┐
│ │
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ TRANSLATOR │ │ SUMMARIZER │
│ │ │ │
│ 🇫🇷 French │ │ 📝 One sentence │
│ Translation │ │ Summary │
└─────────┬────────┘ └─────────┬────────┘
│ │
└───────────────┬───────────────┘
│
▼
┌─────────────────────┐
│ COMBINE │
│ │
│ 📋 Merge Results │
└──────────┬──────────┘
│
▼
┌─────────────────────┐
│ END │
│ │
│ ✅ Complete │
└─────────────────────┘
المفاهيم الأساسية:
- العُقد - وحدات معالجة تنفذ العمل (وكلاء LLM، أو الدوال، أو منطق مخصص)
- الروابط - تدفق التحكم بين العُقد (اتصالات ثابتة أو توجيه شرطي)
- الحالة - بيانات مشتركة تتدفق عبر الرسم البياني وتبقى بين العُقد
- التنفيذ المتوازي - يمكن تشغيل عدة عُقد في الوقت نفسه لتحسين الأداء
فهم المكوّنات الأساسية
🔧 العُقد: العاملون العُقد هي المكان الذي يحدث فيه العمل الفعلي. يمكن لكل عقدة أن:
- AgentNode: تغلف وكيل LLM لمعالجة اللغة الطبيعية
- عقدة دالة: تنفذ كود Rust مخصص لمعالجة البيانات
- العُقد المدمجة: تستخدم منطقًا مُسبق التعريف مثل العدادات أو أدوات التحقق
تخيل العُقد كعمال متخصصين في خط تجميع - لكل منها مهمة وخبرة محددة.
🔀 الروابط: التحكم في التدفق تحدد الروابط كيف ينتقل التنفيذ عبر الرسم البياني:
- الروابط الثابتة: اتصالات مباشرة (
A → B → C) - الروابط الشرطية: توجيه ديناميكي بناءً على الحالة (
if sentiment == "positive" → positive_handler) - الروابط المتوازية: مسارات متعددة من عقدة واحدة (
START → [translator, summarizer])
الروابط تشبه إشارات المرور واللافتات التي توجه تدفق العمل.
💾 الحالة: الذاكرة المشتركة الحالة هي مخزن أزواج مفتاح/قيمة يمكن لجميع العُقد القراءة منه والكتابة إليه:
- بيانات الإدخال: المعلومات الأولية المُدخلة إلى الرسم البياني
- النتائج الوسيطة: مخرجات عقدة تصبح مدخلات لأخرى
- المخرج النهائي: النتيجة المكتملة بعد انتهاء كل المعالجة
تعمل الحالة مثل لوحة بيضاء مشتركة يمكن للعُقد أن تترك عليها معلومات ليستخدمها الآخرون.
⚡ التنفيذ المتوازي: دفعة السرعة عندما تغادر عدة روابط عقدةً ما، تعمل العقد الهدف في الوقت نفسه:
- معالجة أسرع: المهام المستقلة تعمل في الوقت ذاته
- كفاءة الموارد: استغلال أفضل لوحدة المعالجة المركزية وعمليات الإدخال/الإخراج
- قابلية التوسع: التعامل مع سير عمل أكثر تعقيدًا دون تباطؤ خطي
هذا يشبه وجود عدة عمال يعالجون أجزاء مختلفة من المهمة في الوقت نفسه بدلًا من الانتظار في الطابور.
البداية السريعة
1. أنشئ مشروعك
cargo new graph_demo
cd graph_demo
أضف الاعتمادات إلى Cargo.toml:
[dependencies]
adk-graph = { version = "2.0.0", features = ["sqlite"] }
adk-agent = "2.0.0"
adk-model = "2.0.0"
adk-core = "2.0.0"
tokio = { version = "1", features = ["full"] }
dotenvy = "0.15"
serde_json = "1.0"
أنشئ .env باستخدام مفتاح API الخاص بك:
echo 'GOOGLE_API_KEY=your-api-key' > .env
2. مثال على المعالجة المتوازية
إليك مثالًا كاملًا يعمل ويعالج النص بالتوازي:
use adk_agent::LlmAgentBuilder;
use adk_graph::{
agent::GraphAgent,
edge::{END, START},
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create specialized LLM agents
let translator_agent = Arc::new(
LlmAgentBuilder::new("translator")
.description("Translates text to French")
.model(model.clone())
.instruction("Translate the input text to French. Only output the translation.")
.build()?,
);
let summarizer_agent = Arc::new(
LlmAgentBuilder::new("summarizer")
.description("Summarizes text")
.model(model.clone())
.instruction("Summarize the input text in one sentence.")
.build()?,
);
// Wrap agents as graph nodes with input/output mappers
let translator_node = AgentNode::new(translator_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("translation".to_string(), json!(text));
}
}
}
updates
});
let summarizer_node = AgentNode::new(summarizer_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("summary".to_string(), json!(text));
}
}
}
updates
});
// Build the graph with parallel execution
let agent = GraphAgent::builder("text_processor")
.description("Processes text with translation and summarization in parallel")
.channels(&["input", "translation", "summary", "result"])
.node(translator_node)
.node(summarizer_node)
.node_fn("combine", |ctx| async move {
let translation = ctx.get("translation").and_then(|v| v.as_str()).unwrap_or("N/A");
let summary = ctx.get("summary").and_then(|v| v.as_str()).unwrap_or("N/A");
let result = format!(
"=== Processing Complete ===\n\n\
French Translation:\n{}\n\n\
Summary:\n{}",
translation, summary
);
Ok(NodeOutput::new().with_update("result", json!(result)))
})
// Parallel execution: both nodes start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge("translator", "combine")
.edge("summarizer", "combine")
.edge("combine", END)
.build()?;
// Execute the graph
let mut input = State::new();
input.insert("input".to_string(), json!("AI is transforming how we work and live."));
let result = agent.invoke(input, ExecutionConfig::new("thread-1")).await?;
println!("{}", result.get("result").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
مخرجات المثال:
=== Processing Complete ===
French Translation:
L'IA transforme notre façon de travailler et de vivre.
Summary:
AI is revolutionizing work and daily life through technological transformation.
كيف يعمل تنفيذ الرسم البياني
الصورة الكبيرة
ينفذ وكلاء الرسوم البيانية على شكل خطوات فائقة - تعمل جميع العُقد الجاهزة بالتوازي، ثم ينتظر الرسم البياني اكتمالها كلها قبل الخطوة التالية:
Step 1: START ──┬──▶ translator (running)
└──▶ summarizer (running)
⏳ Wait for both to complete...
Step 2: translator ──┬──▶ combine (running)
summarizer ──┘
⏳ Wait for combine to complete...
Step 3: combine ──▶ END ✅
تدفق الحالة عبر العُقد
يمكن لكل عقدة أن تقرأ من الحالة المشتركة وتكتب إليها:
┌─────────────────────────────────────────────────────────────────────┐
│ STEP 1: Initial state │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ State: { "input": "AI is transforming how we work" } │
│ │
│ ↓ │
│ │
│ ┌──────────────────┐ ┌──────────────────┐ │
│ │ translator │ │ summarizer │ │
│ │ reads "input" │ │ reads "input" │ │
│ └──────────────────┘ └──────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────────────────┐
│ STEP 2: After parallel execution │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ State: { │
│ "input": "AI is transforming how we work", │
│ "translation": "L'IA transforme notre façon de travailler", │
│ "summary": "AI is revolutionizing work through technology" │
│ } │
│ │
│ ↓ │
│ │
│ ┌──────────────────────────────────────┐ │
│ │ combine │ │
│ │ reads "translation" + "summary" │ │
│ │ writes "result" │ │
│ └──────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────────────────┐
│ STEP 3: Final state │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ State: { │
│ "input": "AI is transforming how we work", │
│ "translation": "L'IA transforme notre façon de travailler", │
│ "summary": "AI is revolutionizing work through technology", │
│ "result": "=== Processing Complete ===\n\nFrench..." │
│ } │
│ │
└─────────────────────────────────────────────────────────────────────┘
ما الذي يجعل هذا يعمل
| المكوّن | الدور |
|---|---|
AgentNode | يغلّف LLM agents مع مُحوِّلات الإدخال/الإخراج |
input_mapper | يحوّل الحالة → مدخل agent Content |
output_mapper | يحوّل أحداث الوكيل → تحديثات الحالة |
channels | يعلن حقول الحالة التي سيستخدمها الرسم البياني |
edge() | يعرّف تدفق التنفيذ بين العقد |
ExecutionConfig | يوفّر معرّف الخيط لالتقاط الحالة |
التوجيه الشرطي مع تصنيف LLM
ابنِ أنظمة توجيه ذكية حيث يقرر LLMs مسار التنفيذ:
مرئي: التوجيه القائم على المشاعر
┌─────────────────────┐
User Feedback │ │
────────────────▶ │ CLASSIFIER │
│ 🧠 Analyze tone │
└──────────┬──────────┘
│
┌───────────────┼───────────────┐
│ │ │
▼ ▼ ▼
┌──────────────────┐ ┌──────────────────┐ ┌──────────────────┐
│ POSITIVE │ │ NEGATIVE │ │ NEUTRAL │
│ │ │ │ │ │
│ 😊 Thank you! │ │ 😔 Apologize │ │ 😐 Ask more │
│ Celebrate │ │ Help fix │ │ questions │
└──────────────────┘ └──────────────────┘ └──────────────────┘
مثال الشيفرة الكامل
use adk_agent::LlmAgentBuilder;
use adk_graph::{
edge::{END, Router, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create classifier agent
let classifier_agent = Arc::new(
LlmAgentBuilder::new("classifier")
.description("Classifies text sentiment")
.model(model.clone())
.instruction(
"You are a sentiment classifier. Analyze the input text and respond with \
ONLY one word: 'positive', 'negative', or 'neutral'. Nothing else.",
)
.build()?,
);
// Create response agents for each sentiment
let positive_agent = Arc::new(
LlmAgentBuilder::new("positive")
.description("Handles positive feedback")
.model(model.clone())
.instruction(
"You are a customer success specialist. The customer has positive feedback. \
Express gratitude, reinforce the positive experience, and suggest ways to \
share their experience. Be warm and appreciative. Keep response under 3 sentences.",
)
.build()?,
);
let negative_agent = Arc::new(
LlmAgentBuilder::new("negative")
.description("Handles negative feedback")
.model(model.clone())
.instruction(
"You are a customer support specialist. The customer has a complaint. \
Acknowledge their frustration, apologize sincerely, and offer help. \
Be empathetic. Keep response under 3 sentences.",
)
.build()?,
);
let neutral_agent = Arc::new(
LlmAgentBuilder::new("neutral")
.description("Handles neutral feedback")
.model(model.clone())
.instruction(
"You are a customer service representative. The customer has neutral feedback. \
Ask clarifying questions to better understand their needs. Be helpful and curious. \
Keep response under 3 sentences.",
)
.build()?,
);
// Create AgentNodes with mappers
let classifier_node = AgentNode::new(classifier_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase()
.trim()
.to_string();
let sentiment = if text.contains("positive") { "positive" }
else if text.contains("negative") { "negative" }
else { "neutral" };
updates.insert("sentiment".to_string(), json!(sentiment));
}
}
updates
});
// Response nodes (similar pattern for each)
let positive_node = AgentNode::new(positive_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("response".to_string(), json!(text));
}
}
updates
});
// Build graph with conditional routing
let graph = StateGraph::with_channels(&["feedback", "sentiment", "response"])
.add_node(classifier_node)
.add_node(positive_node)
// ... add negative_node and neutral_node similarly
.add_edge(START, "classifier")
.add_conditional_edges(
"classifier",
Router::by_field("sentiment"), // Route based on sentiment field
[
("positive", "positive"),
("negative", "negative"),
("neutral", "neutral"),
],
)
.add_edge("positive", END)
.add_edge("negative", END)
.add_edge("neutral", END)
.compile()?;
// Test with different feedback
let mut input = State::new();
input.insert("feedback".to_string(), json!("Your product is amazing! I love it!"));
let result = graph.invoke(input, ExecutionConfig::new("feedback-1")).await?;
println!("Sentiment: {}", result.get("sentiment").and_then(|v| v.as_str()).unwrap_or(""));
println!("Response: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
تدفق المثال:
Input: "Your product is amazing! I love it!"
↓
Classifier: "positive"
↓
Positive Agent: "Thank you so much for the wonderful feedback!
We're thrilled you love our product.
Would you consider leaving a review to help others?"
نمط ReAct: الاستدلال + التنفيذ
ابنِ وكلاء يمكنهم استخدام الأدوات تكراريًا لحل المشكلات المعقدة:
مرئي: دورة ReAct
┌─────────────────────┐
User Question │ │
────────────────▶ │ REASONER │
│ 🧠 Think + Act │
└──────────┬──────────┘
│
▼
┌─────────────────────┐
│ Has tool calls? │
│ │
└──────────┬──────────┘
│
┌───────────────┴───────────────┐
│ │
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ YES │ │ NO │
│ │ │ │
│ 🔄 Loop back │ │ ✅ Final answer │
│ to reasoner │ │ END │
└─────────┬────────┘ └──────────────────┘
│
└─────────────────┐
│
▼
┌─────────────────────┐
│ REASONER │
│ 🧠 Think + Act │
│ (next iteration) │
└─────────────────────┘
مثال ReAct كامل
use adk_agent::LlmAgentBuilder;
use adk_core::{Part, Tool};
use adk_graph::{
edge::{END, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use adk_tool::FunctionTool;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create tools
let weather_tool = Arc::new(FunctionTool::new(
"get_weather",
"Get the current weather for a location. Takes a 'location' parameter (city name).",
|_ctx, args| async move {
let location = args.get("location").and_then(|v| v.as_str()).unwrap_or("unknown");
Ok(json!({
"location": location,
"temperature": "72°F",
"condition": "Sunny",
"humidity": "45%"
}))
},
)) as Arc<dyn Tool>;
let calculator_tool = Arc::new(FunctionTool::new(
"calculator",
"Perform mathematical calculations. Takes an 'expression' parameter (string).",
|_ctx, args| async move {
let expr = args.get("expression").and_then(|v| v.as_str()).unwrap_or("0");
let result = match expr {
"2 + 2" => "4",
"10 * 5" => "50",
"100 / 4" => "25",
"15 - 7" => "8",
_ => "Unable to evaluate",
};
Ok(json!({ "result": result, "expression": expr }))
},
)) as Arc<dyn Tool>;
// Create reasoner agent with tools
let reasoner_agent = Arc::new(
LlmAgentBuilder::new("reasoner")
.description("Reasoning agent with tools")
.model(model.clone())
.instruction(
"You are a helpful assistant with access to tools. Use tools when needed to answer questions. \
When you have enough information, provide a final answer without using more tools.",
)
.tool(weather_tool)
.tool(calculator_tool)
.build()?,
);
// Create reasoner node that detects tool usage
let reasoner_node = AgentNode::new(reasoner_agent)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build ReAct graph with cycle
let graph = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back for more reasoning
} else {
END.to_string() // Done - final answer
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
// Test the ReAct agent
let mut input = State::new();
input.insert("question".to_string(), json!("What's the weather in Paris and what's 15 + 25?"));
let result = graph.invoke(input, ExecutionConfig::new("react-1")).await?;
println!("Final answer: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
println!("Iterations: {}", result.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0));
Ok(())
}
تدفق المثال:
Question: "What's the weather in Paris and what's 15 + 25?"
Iteration 1:
Reasoner: "I need to get weather info and do math"
→ Calls get_weather(location="Paris") and calculator(expression="15 + 25")
→ has_tool_calls = true → Loop back
Iteration 2:
Reasoner: "Based on the results: Paris is 72°F and sunny, 15 + 25 = 40"
→ No tool calls → has_tool_calls = false → END
Final Answer: "The weather in Paris is 72°F and sunny with 45% humidity.
And 15 + 25 equals 40."
AgentNode
يغلّف أي ADK Agent (عادةً LlmAgent) كعقدة في الرسم البياني:
ما الذي يراه الوكيل
يعمل الوكيل داخل الرسم البياني ضمن سياق مشتق من الاستدعاء الذي
بدأ الرسم البياني، لذلك يتصرف كما يفعل خارجه. عندما يستدعي Runner
GraphAgent، تُنقل هوية المستدعي وخدماته تلقائيًا:
| مُرحَّل عبر | ملاحظة |
|---|---|
app_name، user_id، session_id | الخاص بالمتصل، وليس اصطناعيًا |
| النطاقات وبيانات تعريف الطلب | بحيث ترى عمليات فحص النطاقات المنح الممنوحة للمتصل |
| الخدمة السرية، الذاكرة، العناصر، الحالة المشتركة | متاح تمامًا كما هو خارج الرسم البياني |
| الإلغاء | Runner::interrupt يصل إلى عامل يعمل كعقدة |
RunConfig | موروث من المستدعي |
branch | مُشتق، مثل {caller_branch}.{agent_name}، بحيث يمكن إسناد أحداث العقدة إلى مصدرها |
الرسم البياني المُستدعى مباشرة — graph.invoke(state, ExecutionConfig::new("thread")) —
ليس لديه استدعاء ليرثه. هذا هو الوضع المستقل: تحصل العقدة على
user_id = "graph_user"، app_name = "graph_app"، الفرع main، بدون أسرار، وبدون
ذاكرة. إنه وضع مقصود لتشغيل رسم بياني خارج Runner، وليس
خيارًا احتياطيًا يُلجأ إليه في الإنتاج.
للوصل يدويًا — مثلًا عند تشغيل رسم بياني من المنفذ الخاص بك — مرِّر الاستدعاء صراحةً:
let config = ExecutionConfig::new(ctx.session_id()).with_parent_context(ctx.clone());
ملاحظة: ما تزال العقدة تعمل على جلسة الرسم البياني في الذاكرة الخاصة بها، لذا فإن سجل محادثة الوكيل داخل العقدة يكون ضمن نطاق العقدة نفسها بدلًا من إضافته إلى جلسة المُستدعي.
let node = AgentNode::new(llm_agent)
.with_input_mapper(|state| {
// Transform graph state to agent input Content
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
// Transform agent events to state updates
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("output".to_string(), json!(text));
}
}
updates
});
عقد الدوال
دوال async بسيطة تعالج الحالة:
.node_fn("process", |ctx| async move {
let input = ctx.state.get("input").unwrap();
let output = process_data(input).await?;
Ok(NodeOutput::new().with_update("output", output))
})
أنواع الحواف
الحواف الثابتة
اتصالات مباشرة بين العقد:
.edge(START, "first_node")
.edge("first_node", "second_node")
.edge("second_node", END)
الحواف الشرطية
توجيه ديناميكي بناءً على الحالة:
.conditional_edge(
"router",
|state| {
match state.get("next").and_then(|v| v.as_str()) {
Some("research") => "research_node".to_string(),
Some("write") => "write_node".to_string(),
_ => END.to_string(),
}
},
[
("research_node", "research_node"),
("write_node", "write_node"),
(END, END),
],
)
مساعدو الموجّه
استخدم الموجّهات المدمجة للأنماط الشائعة:
use adk_graph::edge::Router;
// Route based on a state field value
.conditional_edge("classifier", Router::by_field("sentiment"), [
("positive", "positive_handler"),
("negative", "negative_handler"),
("neutral", "neutral_handler"),
])
// Route based on boolean field
.conditional_edge("check", Router::by_bool("approved"), [
("true", "execute"),
("false", "reject"),
])
// Limit iterations
.conditional_edge("loop", Router::max_iterations("count", 5), [
("continue", "process"),
("done", END),
])
التنفيذ المتوازي
تُنفَّذ عدة حواف من عقدة واحدة بالتوازي:
let agent = GraphAgent::builder("parallel_processor")
.channels(&["input", "translation", "summary", "analysis"])
.node(translator_node)
.node(summarizer_node)
.node(analyzer_node)
.node(combiner_node)
// All three start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge(START, "analyzer")
// Wait for all to complete before combining
.edge("translator", "combiner")
.edge("summarizer", "combiner")
.edge("analyzer", "combiner")
.edge("combiner", END)
.build()?;
الرسوم البيانية الدورية (نمط ReAct)
ابنِ وكلاء استدلال تكراريين باستخدام الدورات:
use adk_core::Part;
// Create agent with tools
let reasoner = Arc::new(
LlmAgentBuilder::new("reasoner")
.model(model)
.instruction("Use tools to answer questions. Provide final answer when done.")
.tool(search_tool)
.tool(calculator_tool)
.build()?
);
let reasoner_node = AgentNode::new(reasoner)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { name, .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build graph with cycle
let react_agent = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back
} else {
END.to_string() // Done
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
المشرف متعدد الوكلاء
وجّه المهام إلى وكلاء متخصصين:
// Create supervisor agent
let supervisor = Arc::new(
LlmAgentBuilder::new("supervisor")
.model(model.clone())
.instruction("Route tasks to: researcher, writer, or coder. Reply with agent name only.")
.build()?
);
let supervisor_node = AgentNode::new(supervisor)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase();
let next = if text.contains("researcher") { "researcher" }
else if text.contains("writer") { "writer" }
else if text.contains("coder") { "coder" }
else { "done" };
updates.insert("next_agent".to_string(), json!(next));
}
}
updates
});
// Build supervisor graph
let graph = StateGraph::with_channels(&["task", "next_agent", "research", "content", "code"])
.add_node(supervisor_node)
.add_node(researcher_node)
.add_node(writer_node)
.add_node(coder_node)
.add_edge(START, "supervisor")
.add_conditional_edges(
"supervisor",
Router::by_field("next_agent"),
[
("researcher", "researcher"),
("writer", "writer"),
("coder", "coder"),
("done", END),
],
)
// Agents report back to supervisor
.add_edge("researcher", "supervisor")
.add_edge("writer", "supervisor")
.add_edge("coder", "supervisor")
.compile()?;
إدارة الحالة
مخطط الحالة مع المختزِلات
تحكّم في كيفية دمج تحديثات الحالة:
let schema = StateSchema::builder()
.channel("current_step") // Overwrite (default)
.list_channel("messages") // Append to list
.channel_with_reducer("count", Reducer::Sum) // Sum values
.channel_with_reducer("data", Reducer::Custom(Arc::new(|old, new| {
// Custom merge logic
merge_json(old, new)
})))
.build();
let agent = GraphAgent::builder("stateful")
.state_schema(schema)
// ... nodes and edges
.build()?;
أنواع المختزِلات
| المخفِّض | السلوك |
|---|---|
Overwrite | استبدال القيمة القديمة بالجديدة (الافتراضي) |
Append | الإلحاق بالقائمة |
Sum | إضافة قيم رقمية |
Custom | دالة دمج مخصصة |
إنشاء نقاط التحقق
فعّل الحالة المستمرة من أجل تحمّل الأعطال وعمليات الإنسان في الحلقة:
في الذاكرة (التطوير)
use adk_graph::checkpoint::MemoryCheckpointer;
let checkpointer = Arc::new(MemoryCheckpointer::new());
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer_arc(checkpointer.clone());
SQLite (الإنتاج)
use adk_graph::checkpoint::SqliteCheckpointer;
let checkpointer = SqliteCheckpointer::new("checkpoints.db").await?;
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer(checkpointer);
ما الذي تسجله نقطة التحقق
تخزّن نقطة التحقق الحالة المتراكمة، ورقم الخطوة، والحدّ الأمامي — أي العقد التي ما زال يتعين تشغيلها. تُكتب بعد تقدّم الحدّ الأمامي، لذلك لا تؤدي عملية الاستئناف إلى إعادة تنفيذ عقدة اكتملت بالفعل، ولا إلى تطبيق تحديثاتها مرتين. التشغيل الذي يكتمل يدوّن نقطة تحقق بحدّ أمامي فارغ، لذا فإن استئناف سلسلة مكتملة يعيد الحالة النهائية بدلاً من إعادة بدء الرسم البياني.
تسجّل حالتان عن قصد الحدّ الأمامي الذي كان قيد التنفيذ بدلاً من التالي، لأن العقدة المتقطعة لم تُنتج تحديثاتها بعد ويجب أن تُشغَّل مرة أخرى عند الاستئناف:
| الحالة | الجبهة المحفوظة |
|---|---|
| اكتملت خطوة فائقة | العُقد التالية للتنفيذ |
| اكتمل التشغيل | فارغ |
| تم رفع المقاطعة (حظر أو تدفق) | العقد التي كانت قيد التنفيذ |
تُحفظ نقاط التحقق للتشغيل المتدفق على الجدول الزمني نفسه للتشغيل الحاجب، بما في ذلك عندما ينهي مقاطِع البثّ، لذا يمكن استئناف التوقف المؤقت مع وجود إنسان في الحلقة في أي نمط تنفيذ.
سجل نقاط التحقق (السفر عبر الزمن)
للقراءة فقط. تُرجع
TimeTravelHandle::state_history(from, to)الحالة التي كانت مخزنة عند كل خطوة تم إنشاء نقطة تحقق لها. وهي لا تنفذ أي شيء — لا يتم تشغيل أي عقدة، ولا يُعاد توليد أي حدث، ولا تتكرر أي آثار جانبية. لإعادة التشغيل من نقطة في السجل، استخدمfork_atلتفرّع تلك نقطة التحقق واستدعِ الرسم البياني على الخيط المتشعب. كان اسم الطريقة سابقًاreplayووُثِّقت على أنها تعيد تنفيذ الرسم البياني، وهو ما لم تفعله قط.
تُمكّن نقاط التحقق أيضًا من الاستئناف الدائم — إذا تعطّل تنفيذ الرسم البياني أو أُعيد تشغيل العملية، فسيُستأنف التنفيذ من آخر نقطة تحقق محفوظة بدلًا من البدء من جديد. استخدم SqliteCheckpointer أو PostgresCheckpointer من أجل حفظ آمن ضد الأعطال.
// List all checkpoints for a thread
let checkpoints = checkpointer.list("thread-id").await?;
for cp in checkpoints {
println!("Step {}: {:?}", cp.step, cp.state.get("status"));
}
// Load a specific checkpoint
if let Some(checkpoint) = checkpointer.load_by_id(&checkpoint_id).await? {
println!("State at step {}: {:?}", checkpoint.step, checkpoint.state);
}
الإنسان في الحلقة
أوقف التنفيذ مؤقتًا من أجل موافقة بشرية باستخدام المقاطعات الديناميكية:
use adk_graph::{error::GraphError, node::NodeOutput};
// Planner agent assesses risk
let planner_node = AgentNode::new(planner_agent)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
// Extract risk level from LLM response
let risk = if text.to_lowercase().contains("risk: high") { "high" }
else if text.to_lowercase().contains("risk: medium") { "medium" }
else { "low" };
updates.insert("plan".to_string(), json!(text));
updates.insert("risk_level".to_string(), json!(risk));
}
}
updates
});
// Review node with dynamic interrupt
let graph = StateGraph::with_channels(&["task", "plan", "risk_level", "approved", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_node_fn("review", |ctx| async move {
let risk = ctx.get("risk_level").and_then(|v| v.as_str()).unwrap_or("low");
let approved = ctx.get("approved").and_then(|v| v.as_bool());
// Already approved - continue
if approved == Some(true) {
return Ok(NodeOutput::new());
}
// High/medium risk - interrupt for approval
if risk == "high" || risk == "medium" {
return Ok(NodeOutput::interrupt_with_data(
&format!("{} RISK: Human approval required", risk.to_uppercase()),
json!({
"plan": ctx.get("plan"),
"risk_level": risk,
"action": "Set 'approved' to true to continue"
})
));
}
// Low risk - auto-approve
Ok(NodeOutput::new().with_update("approved", json!(true)))
})
.add_edge(START, "planner")
.add_edge("planner", "review")
.add_edge("review", "executor")
.add_edge("executor", END)
.compile()?
.with_checkpointer_arc(checkpointer.clone());
// Execute - may pause for approval
let thread_id = "task-001";
let result = graph.invoke(input, ExecutionConfig::new(thread_id)).await;
match result {
Err(GraphError::Interrupted(interrupt)) => {
println!("*** EXECUTION PAUSED ***");
println!("Reason: {}", interrupt.interrupt);
println!("Plan awaiting approval: {:?}", interrupt.state.get("plan"));
// Human reviews and approves...
// Update state with approval
graph.update_state(thread_id, [("approved".to_string(), json!(true))]).await?;
// Resume execution
let final_result = graph.invoke(State::new(), ExecutionConfig::new(thread_id)).await?;
println!("Final result: {:?}", final_result.get("result"));
}
Ok(result) => {
println!("Completed without interrupt: {:?}", result);
}
Err(e) => {
println!("Error: {}", e);
}
}
المقاطعات الثابتة
استخدم interrupt_before أو interrupt_after لنقاط التوقف الإلزامية:
let graph = StateGraph::with_channels(&["task", "plan", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_edge(START, "planner")
.add_edge("planner", "executor")
.add_edge("executor", END)
.compile()?
.with_interrupt_before(&["executor"]); // Always pause before execution
التنفيذ المتدفق
بث الأحداث أثناء تنفيذ الرسم البياني:
use futures::StreamExt;
use adk_graph::stream::StreamMode;
let stream = agent.stream(input, config, StreamMode::Updates);
while let Some(event) = stream.next().await {
match event? {
StreamEvent::NodeStart(name) => println!("Starting: {}", name),
StreamEvent::Updates { node, updates } => {
println!("{} updated state: {:?}", node, updates);
}
StreamEvent::NodeEnd(name) => println!("Completed: {}", name),
StreamEvent::Done(state) => println!("Final state: {:?}", state),
_ => {}
}
}
أنماط البث
| الوضع | الوصف |
|---|---|
Values | بث الحالة الكاملة بعد كل عقدة |
Updates | بث تغييرات الحالة فقط |
Messages | بث تحديثات نوع الرسائل |
Debug | بث جميع الأحداث الداخلية |
يقرأ وضع Messages الرموز من Node::execute_stream كما يتم إنتاجها.
تعمل كل عقدة مرة واحدة لكل super-step في هذا الوضع: تُبلغ العقدة عن تحديثات حالتها على التدفق كحدث StreamEvent::Updates، ويطبّق المنفذ تلك التحديثات بدلًا من تنفيذ العقدة مرة ثانية لجمعها. هذا يهمّ أكثر في AgentNode، حيث إن التنفيذ الثاني يعني استدعاءً ثانيًا مُفوترةً للنموذج لكل عقدة.
مهم: يجب على
Nodeمخصص يَتجاوزexecute_streamأن يُصدر حدثStreamEvent::Updatesيحمل تحديثات حالته. من دونه تبث العقدة الأحداث لكنها لا تُسهم بأي حالة في وضعMessages. أماexecute_streamالافتراضي، الذي يغلّفexecute، فيقوم بذلك نيابةً عنك.
تنطبق سياسات المهلة الزمنية على التنفيذ المتدفق نفسه. بالنسبة إلى التدفق، يعني idle_timeout أنه لم يتم إنتاج أي حدث ضمن المهلة.
تكامل ADK
يُنفّذ GraphAgent الخاصية ADK Agent، لذا يعمل مع:
- Runner: استخدمه مع
adk-runnerللتنفيذ القياسي - Callbacks: دعم كامل لنداءات ما قبل/بعد
- Sessions: يعمل مع
adk-sessionلسجل المحادثة - Streaming: يُرجع ADK
EventStream
use adk_runner::Runner;
let graph_agent = GraphAgent::builder("workflow")
.before_agent_callback(|ctx| async {
println!("Starting graph execution for session: {}", ctx.session_id());
Ok(())
})
.after_agent_callback(|ctx, event| async {
if let Some(content) = event.content() {
println!("Graph completed with content");
}
Ok(())
})
// ... graph definition
.build()?;
// GraphAgent implements Agent trait - use with Launcher or Runner
// See adk-runner README for Runner configuration
أمثلة
أمثلة رسومية مُتحقَّق منها في هذا المستودع:
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 11-standard-graph
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 12-standard-sequential
cargo run --manifest-path examples/competitive_graph_resume/Cargo.toml
تتوفر مجموعة الرسوم البيانية الكاملة، مع تكاملات LLM حقيقية، في ADK-Rust Playground المضمّن في هذا الموقع.
مقارنة مع LangGraph
| الميزة | LangGraph | adk-graph |
|---|---|---|
| إدارة الحالة | TypedDict + Reducers | StateSchema + Reducers |
| نموذج التنفيذ | Pregel super-steps | Pregel super-steps |
| حفظ نقاط التحقق | الذاكرة، SQLite، Postgres | الذاكرة، SQLite |
| إنسان في الحلقة | interrupt_before/after | interrupt_before/after + ديناميكي |
| التدفق | 5 أوضاع | 5 أوضاع |
| الحلقات | دعم أصلي | دعم أصلي |
| أمان الأنواع | الكتابة في Python | نظام الأنواع في Rust |
| LLM تكامل | LangChain | AgentNode + ADK وكلاء |
السابق: ← أنظمة متعددة الوكلاء | التالي: وكلاء الوقت الحقيقي →