المشغّل

بيئة التنفيذ من adk-runner التي تنسّق تنفيذ الوكيل.

نظرة عامة

يتولى Runner إدارة دورة الحياة الكاملة لتنفيذ الوكيل:

  • إدارة الجلسات (إنشاء/استرجاع الجلسات)
  • حقن الذاكرة (البحث عن الذاكرات ذات الصلة وحقنها)
  • التعامل مع الأصول (وصول مقيّد بنطاق إلى الأصول)
  • بث الأحداث (معالجة الأحداث وتمريرها)
  • نقل الوكلاء (التعامل مع عمليات التسليم بين عدة وكلاء)
Rendering architecture…

التثبيت

[dependencies]
adk-runner = "2.0.0"

RunnerConfig

قم بتهيئة المشغّل بالخدمات المطلوبة:

use adk_runner::{Runner, RunnerConfig};
use adk_session::InMemorySessionService;
use adk_artifact::InMemoryArtifactService;
use std::sync::Arc;

let config = RunnerConfig {
    app_name: "my_app".to_string(),
    agent: Arc::new(my_agent),
    session_service: Arc::new(InMemorySessionService::new()),
    artifact_service: Some(Arc::new(InMemoryArtifactService::new())),
    memory_service: None,
    plugin_manager: None,
    run_config: None,
    compaction_config: None,
    context_cache_config: None,
    cache_capable: None,
    request_context: None,
    cancellation_token: None,
};

let runner = Runner::new(config)?;

استخدم منشئ typestate لبناء Runner. يفرض المنشئ الحقول المطلوبة في وقت الترجمة ويعيّن قيماً افتراضية لجميع الحقول الاختيارية، لذلك فإن إضافة حقول جديدة في الإصدارات المستقبلية لن تكسر شيفرتك:

use adk_runner::Runner;

let runner = Runner::builder()
    .app_name("my_app")
    .agent(Arc::new(my_agent))
    .session_service(Arc::new(InMemorySessionService::new()))
    // Optional fields — only set what you need
    .artifact_service(Arc::new(InMemoryArtifactService::new()))
    .build()?;

يتطلب المنشئ ثلاثة حقول: app_name، وagent، وsession_service. كل ما عدا ذلك اختياري وله قيم افتراضية مناسبة. تكون طريقة build() متاحة فقط بعد تعيين الحقول المطلوبة الثلاثة كلها — وغياب أحدها يُعد خطأً في وقت الترجمة، وليس في وقت التشغيل.

حقول التهيئة

الحقلالنوعمطلوبالوصف
app_nameStringنعممعرف التطبيق
agentArc<dyn Agent>نعمالوكيل الجذري المراد تنفيذه
session_serviceArc<dyn SessionService>نعمالواجهة الخلفية لتخزين الجلسات
artifact_serviceOption<Arc<dyn ArtifactService>>لاتخزين العناصر
memory_serviceOption<Arc<dyn Memory>>لاالذاكرة طويلة الأمد
plugin_managerOption<Arc<PluginManager>>لاخطافات دورة حياة المكونات الإضافية
compaction_configOption<EventsCompactionConfig>لاإعدادات ضغط السياق
run_configOption<RunConfig>لاخيارات التنفيذ
context_cache_configOption<ContextCacheConfig>لادورة حياة ذاكرة التخزين المؤقت للسياق على مستوى المشغّل (تجريبي — انظر أدناه)
cache_capableOption<Arc<dyn CacheCapable>>لامرجع نموذج يدعم ذاكرة التخزين المؤقت (تجريبي — انظر أدناه)
request_contextOption<RequestContext>لاسياق وسيط المصادقة
cancellation_tokenOption<CancellationToken>لاالإلغاء التعاوني

التخزين المؤقت للمطالبة

التخزين المؤقت هو مسألة على مستوى مزوّد الخدمة ولا يحتاج إلى أي إعداد لـ Runner. تتولى كل تكاملات المزوّد التعامل معه في الموضع الذي يتم فيه تجميع الطلب:

المزوّدالآليةالافتراضي
Anthropic / Bedrockنقاط توقف cache_controlمفعّل (AnthropicConfig::prompt_caching، مع إمكانية إلغاء الاشتراك عبر with_prompt_caching(false))
OpenAIالتخزين المؤقت للمطالبة على جانب الخادم، PromptCacheRetention للاحتفاظتلقائي
Geminiالتخزين المؤقت الضمني على 2.5/3.x — تحصل البادئة المشتركة على خصم دون أي تغيير في الشيفرةتلقائي

تكون إصابات ذاكرة التخزين المؤقت قابلة للرصد دون أي توصيل إضافي: يسجّل تكامل Gemini cachedContentTokenCount على كل استجابة.

context_cache_config وcache_capable تجريبيان ويجب أن يُتركا غير مضبوطين. إنهما يوجهان cachedContents API الصريح الخاص بـ Gemini من خلال Runner. يتطلب ذلك API أن تستبدل ذاكرة التخزين المؤقت system_instruction وtools وtool_config — وإرسال ذاكرة تخزين مؤقت مع أيٍّ منها يُرفض مع INVALID_ARGUMENT. يختار Runner ذاكرة تخزين مؤقت قبل أن يحلّ الوكيل أدواته، لذلك لا يمكنه تجميع ذلك الطلب، ولا يؤدي تمكين هذه الحقول حاليًا إلى إصابات في ذاكرة التخزين المؤقت. إن التخزين المؤقت المضمون (بدلًا من بذل أفضل جهد) لـ Gemini يخص تكامل النموذج، إلى جانب كيفية قيام مزودي الخدمة الآخرين بذلك.

تشغيل الوكلاء

نفّذ وكيلًا مع إدخال المستخدم:

use adk_core::{Content, SessionId, UserId};
use futures::StreamExt;

let user_content = Content::new("user").with_text("Hello!");

let mut stream = runner.run(
    UserId::new("user-123")?,
    SessionId::new("session-456")?,
    user_content,
).await?;

while let Some(event) = stream.next().await {
    match event {
        Ok(e) => {
            if let Some(content) = e.content() {
                for part in &content.parts {
                    if let Some(text) = part.text() {
                        print!("{}", text);
                    }
                }
            }
        }
        Err(e) => eprintln!("Error: {}", e),
    }
}

طريقة ملائمة للسلاسل النصية

للحالات البسيطة، يقبل run_str() وسائط &str نصية عادية ويتولى التحويل إلى النوع الجديد داخليًا:

let mut stream = runner.run_str(
    "user-123",
    "session-456",
    Content::new("user").with_text("Hello!"),
).await?;

إذا فشل التحقق من صحة السلسلة النصية (فارغة، تحتوي على بايتات null، أو تتجاوز حد الطول)، فإن run_str() يعيد خطأ قبل بدء حلقة الوكيل. تبقى طريقة run() الحالية ذات UserId/SessionId المطبّعة دون تغيير.

المقاطعة وعزل التشغيل

يُسجَّل التشغيل بمجرد أن يعيد run() تدفّقه، ويُلغى تسجيله عندما يُسقط ذلك التدفق — بما في ذلك عندما يُسقط دون أن يُستقصى ولو مرة واحدة.

الطريقةالنطاق
interrupt(session_id)يلغي كل تشغيل جارٍ لذلك المعرّف الخاص بالجلسة، عبر التطبيقات والمستخدمين
interrupt_identity(app_name, user_id, session_id)يلغي عمليات التشغيل لهوية واحدة مطابقة تمامًا
active_runs()هوية كل تشغيل جارٍ؛ تعني الهوية المكررة وجود تشغيلات متزامنة
active_session_ids()معرّفات الجلسات غير المكررة للتشغيلات الجارية
// Cancel one tenant's run without touching another that shares the session ID
let cancelled = runner.interrupt_identity("my-app", "user-1", "session-1");

معرّف الجلسة يكون فريدًا فقط ضمن التطبيق والمستخدم، لذا فإن interrupt(session_id) هو الشكل العام وinterrupt_identity هو الشكل الدقيق. يُفضَّل interrupt_identity عندما تخدم Runner واحدة أكثر من تطبيق واحد أو مستخدم واحد.

تُتتبَّع عمليات التشغيل بواسطة معرّف تشغيل فريد بدلًا من معرّف الجلسة، لذا تُتتبَّع عمليتان للتشغيل لنفس الهوية بشكل منفصل، وكل واحدة تُلغِي تسجيل نفسها فقط.

الاستمرارية مرتبطة بالهوية

كل حدث يحتفظ به Runner — دورات المستخدم، استجابات النموذج، أحداث النقل، أحداث الإضافة، وأحداث الدمج — يُكتب عبر SessionService::append_event_for_identity مع (app_name, user_id, session_id) الثلاثي الكامل. وبذلك يمكن لـ SessionService الذي يكون مفتاحه الطبيعي مركبًا أن يربط كل حدث بمستأجره، ويمكنه رفض أو تجاهل مسار معرّف الجلسة الخام append_event بالكامل.

تدفق التنفيذ

┌─────────────────────────────────────────────────────────────┐
│                     Runner.run()                            │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                  1. Session Retrieval                       │
│                                                             │
│   SessionService.get(app_name, user_id, session_id)        │
│   → Creates new session if not exists                       │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                  2. Agent Selection                         │
│                                                             │
│   Check session state for active agent                      │
│   → Use root agent or transferred agent                     │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                3. Context Creation                          │
│                                                             │
│   InvocationContext with:                                   │
│   - Session (mutable)                                       │
│   - Artifacts (scoped to session)                          │
│   - Memory (if configured)                                  │
│   - Run config                                              │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                  4. Agent Execution                         │
│                                                             │
│   agent.run(ctx) → EventStream                             │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                 5. Event Processing                         │
│                                                             │
│   For each event:                                           │
│   - Update session state                                    │
│   - Handle transfers                                        │
│   - Forward to caller                                       │
└─────────────────────────────────────────────────────────────┘
                              │
                              ▼
┌─────────────────────────────────────────────────────────────┐
│                  6. Session Save                            │
│                                                             │
│   SessionService.append_event(session, events)             │
└─────────────────────────────────────────────────────────────┘

InvocationContext

السياق المزوّد للوكلاء أثناء التنفيذ:

pub trait InvocationContext: CallbackContext {
    /// The agent being executed
    fn agent(&self) -> Arc<dyn Agent>;
    
    /// Memory service (if configured)
    fn memory(&self) -> Option<Arc<dyn Memory>>;
    
    /// Current session
    fn session(&self) -> &dyn Session;
    
    /// Execution configuration
    fn run_config(&self) -> &RunConfig;
    
    /// Signal end of invocation
    fn end_invocation(&self);
    
    /// Check if invocation has ended
    fn ended(&self) -> bool;
}

RunConfig

خيارات التنفيذ:

pub struct RunConfig {
    /// Streaming mode for responses
    pub streaming_mode: StreamingMode,
    // ... other fields (tool_confirmation_decisions, cached_content, etc.)
}

ToolExecutionStrategy

يتحكم في كيفية إرسال عدة استدعاءات للأدوات من رد واحد لـ LLM:

الاستراتيجيةالسلوك
Sequential (الافتراضي)تنفيذ الأدوات واحدة تلو الأخرى وفق الترتيب الذي أرجعته LLM
Parallelتنفيذ جميع الأدوات بالتوازي؛ المستدعي مسؤول عن السلامة
Autoنفّذ المجموعة الفرعية الآمنة للقراءة فقط بالتوازي، ثم نفّذ جميع الاستدعاءات المتبقية بالتتابع

يتم تحديده لكل وكيل عبر LlmAgentBuilder:

use adk_core::ToolExecutionStrategy;

let agent = LlmAgentBuilder::new("fast_agent")
    .model(model)
    .tool_execution_strategy(ToolExecutionStrategy::Auto)
    .tool(Arc::new(
        search_tool
            .with_read_only(true)
            .with_concurrency_safe(true),
    ))
    .tool(Arc::new(save_tool)) // runs after the concurrent safe subset
    .build()?;

في وضع Auto، تستعلم حلقة التوزيع كلًّا من is_read_only() وis_concurrency_safe(). تُشغَّل أولًا بشكل متزامن الاستدعاءات التي تعيد أدواتها المختارة true لكلا الطريقتين؛ ثم تُنفَّذ جميع الاستدعاءات المتبقية بالتتابع. تتجاوز Parallel هذه عمليات فحص البيانات الوصفية كتجاوز صريح من المستدعي. تُعاد تجميع النتائج دائمًا بالترتيب الأصلي الذي أعاده LLM بغضّ النظر عن الاستراتيجية. تُنتِج الأدوات الفاشلة استجابة خطأ JSON دون إيقاف الدفعة.

pub enum StreamingMode {
    /// No streaming, return complete response
    None,
    /// Server-Sent Events (default)
    SSE,
    /// Bidirectional streaming (realtime)
    Bidi,
}

انتقالات الوكلاء

يتولى Runner انتقالات عدة وكلاء تلقائيًا:

// In an agent's tool or callback
if should_transfer {
    // Set transfer in event actions
    ctx.set_actions(EventActions {
        transfer_to_agent: Some("specialist_agent".to_string()),
        ..Default::default()
    });
}

سيقوم Runner بـ:

  1. اكتشاف طلب النقل في الحدث
  2. العثور على الوكيل الهدف في sub_agents
  3. تحديث حالة الجلسة بالوكيل النشط الجديد
  4. متابعة التنفيذ مع الوكيل الجديد

تكثيف السياق

للجلسات الطويلة، فعّل تكثيف السياق التلقائي لإبقاء نافذة سياق LLM محدودة:

use adk_runner::{Runner, RunnerConfig, EventsCompactionConfig};
use adk_agent::LlmEventSummarizer;
use std::sync::Arc;

let summarizer = LlmEventSummarizer::new(model.clone());

let config = RunnerConfig {
    // ... other fields ...
    compaction_config: Some(EventsCompactionConfig {
        compaction_interval: 3,  // Compact every 3 invocations
        overlap_size: 1,         // Keep 1 event overlap for continuity
        summarizer: Arc::new(summarizer),
    }),
    // ...
};

عند بدء التكثيف، تُستبدل الأحداث الأقدم بحدث ملخّص. يستخدم conversation_history() الملخّص تلقائيًا بدلًا من الأحداث الأصلية.

راجع تكثيف السياق للحصول على الوثائق الكاملة.

التكامل مع Launcher

يستخدم Launcher Runner داخليًا:

// Launcher creates Runner with default services
Launcher::new(agent)
    .app_name("my_app")
    .run()
    .await?;

// Equivalent to using the builder:
let runner = Runner::builder()
    .app_name("my_app")
    .agent(agent)
    .session_service(Arc::new(InMemorySessionService::new()))
    .build()?;

استخدام Runner المخصص

للحالات المتقدمة، استخدم Runner مباشرةً:

use adk_runner::Runner;

// Production configuration using the builder
let runner = Runner::builder()
    .app_name("production_app")
    .agent(my_agent)
    .session_service(Arc::new(SqliteSessionService::new(db_pool)))
    .artifact_service(Arc::new(S3ArtifactService::new(s3_client)))
    .memory_service(Arc::new(QdrantMemoryService::new(qdrant_client)))
    .build()?;

// Use in HTTP handler with run_str() for convenience
async fn chat_handler(runner: &Runner, request: ChatRequest) -> Response {
    let stream = runner.run_str(
        &request.user_id,
        &request.session_id,
        request.content,
    ).await?;
    
    // Stream events to client
    Response::sse(stream)
}

السابق: ← الأنواع الأساسية | التالي: Launcher →

المشغّل - وثائق ADK-Rust | ADK-Rust